diff --git a/src/app/api/conversations/route.test.ts b/src/app/api/conversations/route.test.ts index bf225a7..4796404 100644 --- a/src/app/api/conversations/route.test.ts +++ b/src/app/api/conversations/route.test.ts @@ -120,4 +120,37 @@ describe("conversations list route", () => { const body = (await response.json()) as { threads: unknown[] }; expect(body.threads).toEqual([]); }); + + it("forwards a valid status filter to the thread search", async () => { + const deployment = scopedDeployment([]); + mocks.getActiveDeployment.mockResolvedValue(deployment); + + const response = await routes.GET( + new NextRequest("http://localhost/api/conversations?status=busy", { + method: "GET", + }) + ); + + expect(response.status).toBe(200); + expect(deployment.threadClient.threads.search).toHaveBeenCalledWith( + expect.objectContaining({ status: "busy" }) + ); + }); + + it("rejects an unknown status filter instead of silently ignoring it", async () => { + const deployment = scopedDeployment([]); + mocks.getActiveDeployment.mockResolvedValue(deployment); + + const response = await routes.GET( + new NextRequest("http://localhost/api/conversations?status=flying", { + method: "GET", + }) + ); + + expect(response.status).toBe(400); + await expect(response.json()).resolves.toEqual({ + error: "Unknown thread status filter.", + }); + expect(deployment.threadClient.threads.search).not.toHaveBeenCalled(); + }); }); diff --git a/src/app/api/conversations/route.ts b/src/app/api/conversations/route.ts index 4a18286..8de43ba 100644 --- a/src/app/api/conversations/route.ts +++ b/src/app/api/conversations/route.ts @@ -36,12 +36,28 @@ export async function GET(request: NextRequest) { Number(request.nextUrl.searchParams.get("offset") ?? 0) || 0, 0 ); + const THREAD_STATUSES = new Set([ + "idle", + "busy", + "interrupted", + "error", + ] as const); + const statusParam = request.nextUrl.searchParams.get("status"); + if (statusParam !== null && !THREAD_STATUSES.has(statusParam as never)) { + return NextResponse.json( + { error: "Unknown thread status filter." }, + { status: 400 } + ); + } const threads = await deployment.threadClient.threads.search({ limit, offset, sortBy: "updated_at", sortOrder: "desc", metadata: { graph_id: deployment.assistantId }, + ...(statusParam + ? { status: statusParam as "idle" | "busy" | "interrupted" | "error" } + : {}), }); const visible = threads.filter( (thread) => diff --git a/src/app/components/ThreadList.tsx b/src/app/components/ThreadList.tsx index 7470131..44b8495 100644 --- a/src/app/components/ThreadList.tsx +++ b/src/app/components/ThreadList.tsx @@ -35,6 +35,8 @@ import { formatTime, formatFullTime } from "@/lib/time"; import type { ThreadItem } from "@/app/hooks/useThreads"; import { useThreads, + useBusyThreadIds, + useLocalRunningThreadIds, deleteThread, renameThread, pinThread, @@ -175,6 +177,10 @@ export function ThreadList({ status: statusFilter === "all" ? undefined : statusFilter, limit: 20, }); + // Running indicators: heartbeat for out-of-band runs, local store for + // runs this client started itself (zero-latency). + const heartbeatBusyIds = useBusyThreadIds(); + const localRunningIds = useLocalRunningThreadIds(); // Dedupe by id, keeping the first occurrence — page 0 wins, so the freshest // `updated_at` survives. Without this, a thread whose `updated_at` advances @@ -435,6 +441,8 @@ export function ThreadList({ const renderThreadCard = (thread: ThreadItem) => { const pinBusy = pinBusyIds.has(thread.id); const exportBusy = exportBusyIds.has(thread.id); + const isRunning = + heartbeatBusyIds.has(thread.id) || localRunningIds.has(thread.id); return (
- + {isRunning ? ( + + + ) : ( + + )}
diff --git a/src/app/hooks/useChat.ts b/src/app/hooks/useChat.ts index 9e577bc..afd071d 100644 --- a/src/app/hooks/useChat.ts +++ b/src/app/hooks/useChat.ts @@ -9,7 +9,7 @@ import { parseSummarizationEvent } from "@/lib/summarization"; import { findActiveTurnId } from "@/lib/usageTurn"; import { toast } from "sonner"; import type { ModelRef, ThreadModelSelection } from "@/lib/modelRegistry"; -import { setThreadModelSelection } from "@/app/hooks/useThreads"; +import { setThreadModelSelection, setThreadRunning } from "@/app/hooks/useThreads"; import { clearRunCursor, getRunStreamCheckpoint, @@ -308,6 +308,15 @@ export function useChat({ const isRunLoading = isSubmitting || runConnection !== null; const isReconnecting = runConnection?.phase === "reconnecting"; + // Zero-latency running indicator for the thread list: this client knows its + // own run lifecycle, no heartbeat wait. The 10s busy heartbeat reconciles + // any drift (e.g. SSE ending before the run actually pauses server-side). + useEffect(() => { + if (!threadId) return; + setThreadRunning(threadId, isRunLoading); + return () => setThreadRunning(threadId, false); + }, [threadId, isRunLoading]); + // Do not attach message modes to the SDK hook's synchronous external store. // The background subscriber below batches messages-tuple chunks before // publishing them, while this hook remains responsible for history/interrupts. diff --git a/src/app/hooks/useThreads.running.test.ts b/src/app/hooks/useThreads.running.test.ts new file mode 100644 index 0000000..e10398f --- /dev/null +++ b/src/app/hooks/useThreads.running.test.ts @@ -0,0 +1,38 @@ +import { describe, expect, it, vi } from "vitest"; + +import { + getRunningThreadIds, + setThreadRunning, + subscribeRunningThreads, +} from "./useThreads"; + +describe("running thread store", () => { + it("tracks running threads and notifies subscribers", () => { + const listener = vi.fn(); + const unsubscribe = subscribeRunningThreads(listener); + + setThreadRunning("thread-a", true); + expect(getRunningThreadIds().has("thread-a")).toBe(true); + expect(listener).toHaveBeenCalledTimes(1); + + setThreadRunning("thread-a", false); + expect(getRunningThreadIds().has("thread-a")).toBe(false); + expect(listener).toHaveBeenCalledTimes(2); + + unsubscribe(); + setThreadRunning("thread-b", true); + expect(listener).toHaveBeenCalledTimes(2); + setThreadRunning("thread-b", false); + }); + + it("keeps snapshot identity stable for useSyncExternalStore", () => { + const before = getRunningThreadIds(); + // Re-asserting the same state must not swap the snapshot (no-op writes). + setThreadRunning("thread-c", true); + const during = getRunningThreadIds(); + expect(during).not.toBe(before); + setThreadRunning("thread-c", true); + expect(getRunningThreadIds()).toBe(during); + setThreadRunning("thread-c", false); + }); +}); diff --git a/src/app/hooks/useThreads.ts b/src/app/hooks/useThreads.ts index 87aad73..139e9b6 100644 --- a/src/app/hooks/useThreads.ts +++ b/src/app/hooks/useThreads.ts @@ -1,3 +1,5 @@ +import { useMemo, useSyncExternalStore } from "react"; +import useSWR from "swr"; import useSWRInfinite from "swr/infinite"; import { patchConversation } from "@/lib/conversationApi"; import type { ConversationThreadSummary } from "@/lib/conversationSummary"; @@ -66,8 +68,64 @@ export function useThreads(props: { status?: string; limit?: number }) { ); } -export async function deleteThread(id: string): Promise { - const response = await fetch(`/api/conversations/${encodeURIComponent(id)}`, { +// --- Running-thread indicators ---------------------------------------------- +// Two sources, merged by the list UI: a 10s heartbeat for runs started +// out-of-band (scheduled tasks, other clients) and a local store fed by this +// client's own run lifecycle for zero-latency updates. + +const BUSY_HEARTBEAT_INTERVAL_MS = 10_000; + +export function useBusyThreadIds(): ReadonlySet { + const { data } = useSWR( + "conversations:busy", + async () => { + const response = await fetch("/api/conversations?status=busy&limit=100", { + cache: "no-store", + }); + const payload = (await response.json().catch(() => null)) as { + threads?: ConversationThreadSummary[]; + error?: string; + } | null; + if (!response.ok) { + throw new Error(payload?.error || "Failed to load running threads."); + } + return (payload?.threads ?? []).map((thread) => thread.thread_id); + }, + { refreshInterval: BUSY_HEARTBEAT_INTERVAL_MS } + ); + return useMemo(() => new Set(data ?? []), [data]); +} + +const runningThreadIds = new Set(); +const runningListeners = new Set<() => void>(); +let runningSnapshot: ReadonlySet = runningThreadIds; + +export function setThreadRunning(threadId: string, running: boolean): void { + const changed = running + ? !runningThreadIds.has(threadId) + : runningThreadIds.delete(threadId); + if (!changed) return; + if (running) runningThreadIds.add(threadId); + runningSnapshot = new Set(runningThreadIds); + for (const listener of runningListeners) listener(); +} + +export function subscribeRunningThreads(listener: () => void): () => void { + runningListeners.add(listener); + return () => { + runningListeners.delete(listener); + }; +} + +export function getRunningThreadIds(): ReadonlySet { + return runningSnapshot; +} + +export function useLocalRunningThreadIds(): ReadonlySet { + return useSyncExternalStore(subscribeRunningThreads, getRunningThreadIds); +} + +export async function deleteThread(id: string): Promise { const response = await fetch(`/api/conversations/${encodeURIComponent(id)}`, { method: "DELETE", }); if (!response.ok) {