From cbc343963c2cce736c2daab401eb3621c179366e Mon Sep 17 00:00:00 2001 From: m4 Date: Fri, 31 Jul 2026 09:57:08 +0800 Subject: [PATCH] feat(webui): show a running indicator on busy threads in the list Three mechanisms, merged in ThreadList: a 10s heartbeat polling the new ?status=busy filter (out-of-band runs: scheduled tasks, other clients), a module-level running store fed by useChat's own run lifecycle (zero latency for this client's runs), and a spinner replacing the status dot while either source marks the thread busy. The BFF list route now validates and forwards the status filter, which also unbreaks the previously ignored status dropdown. Co-Authored-By: Claude Opus 4.7 --- src/app/api/conversations/route.test.ts | 33 +++++++++++++ src/app/api/conversations/route.ts | 16 ++++++ src/app/components/ThreadList.tsx | 35 +++++++++---- src/app/hooks/useChat.ts | 11 ++++- src/app/hooks/useThreads.running.test.ts | 38 +++++++++++++++ src/app/hooks/useThreads.ts | 62 +++++++++++++++++++++++- 6 files changed, 183 insertions(+), 12 deletions(-) create mode 100644 src/app/hooks/useThreads.running.test.ts 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) {