From a817a0ed87856c3b18e23f83de90ca1ca5d2373b Mon Sep 17 00:00:00 2001 From: lllyfff <122260771+lllyfff@users.noreply.github.com> Date: Sat, 4 Jul 2026 09:15:36 +0800 Subject: [PATCH] Fix(frontend): stale run reconnect and cancel handling (#3908) * Fix stale run reconnect and cancel handling * fix the frontend * Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> * Fix prettier formatting and document error/timeout reconnect intent Wrap the over-long assertion in api-client.test.ts so the CI prettier --check job passes, and add a comment above TERMINAL_RUN_STATUSES noting that error/timeout short-circuit intentionally drops the transient onError toast on reload (the persisted error state still loads from the checkpoint via useThreadHistory). Co-Authored-By: Claude --------- Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> Co-authored-by: Claude --- frontend/src/core/api/api-client.ts | 127 ++++++++++- .../tests/unit/core/api/api-client.test.ts | 211 ++++++++++++++++++ 2 files changed, 332 insertions(+), 6 deletions(-) diff --git a/frontend/src/core/api/api-client.ts b/frontend/src/core/api/api-client.ts index ff03fac6f..14a2ae911 100644 --- a/frontend/src/core/api/api-client.ts +++ b/frontend/src/core/api/api-client.ts @@ -38,7 +38,50 @@ function injectCsrfHeader(_url: URL, init: RequestInit): RequestInit { return { ...init, headers }; } -export function isInactiveRunStreamError(error: unknown): boolean { +// Run statuses that have reached a terminal state where no further streaming +// is possible. Reconnecting (``joinStream``) to such a run either 409s or, once +// the backend's in-memory stream bridge is reaped (``worker.py`` calls +// ``publish_end`` unconditionally, including for interrupted runs, then reaps +// the bridge after 60s), blocks forever on a drained condition variable — +// pinning ``isLoading`` true so the submit button stays a stop button and the +// first message after a reload never sends. The ``joinStream`` wrapper below +// short-circuits these *before* joining. +// +// ``interrupted`` is included because in DeerFlow it is only ever written by +// ``RunManager.cancel()`` (a user-initiated stop); the resumable human-in-the- +// loop path uses ``Command(goto=END)`` (``ClarificationMiddleware``), which +// ends the run as ``success``, not ``interrupted``. So an interrupted run has +// nothing left to stream — its state lives in the checkpoint, fetched +// independently by ``useThreadHistory``, and resuming means a fresh ``submit``. +// +// ``error``/``timeout`` are terminal too, so a reload within the ~60s +// bridge-reap window no longer replays the buffered error event through +// ``onError`` — the transient error toast (``getStreamErrorMessage``) is +// dropped. The persisted error state still loads from the checkpoint via +// ``useThreadHistory``, so only the toast is lost; that is intentional, since +// surfacing a stale error toast on every reload is noise rather than signal. +const TERMINAL_RUN_STATUSES = new Set([ + "success", + "error", + "timeout", + "interrupted", +]); + +/** + * Shared matcher for the gateway's 409 conflict responses. The SDK surfaces + * non-2xx responses as ``HTTPError { status, message }`` where ``message`` is + * ``"HTTP 409: {\"detail\":\"...\"}"``, so a 409 may be detected either via the + * numeric ``status`` or a substring of ``message``. + * + * Every passed ``needles`` substring must be present; this AND semantics is what + * lets a caller distinguish sibling conflict branches by phrase (e.g. the + * terminal-state cancel branch from the still-active-on-another-worker branch). + * + * Match strings until the API exposes a structured error code; the source of + * truth is ``_cancel_conflict_detail`` / the store-only response in + * ``backend/app/gateway/routers/thread_runs.py``. + */ +function isRunConflictError(error: unknown, ...needles: string[]): boolean { const status = typeof error === "object" && error !== null ? Reflect.get(error, "status") @@ -52,16 +95,61 @@ export function isInactiveRunStreamError(error: unknown): boolean { ? String(Reflect.get(error, "message") ?? "") : ""; - // Match the gateway's store-only run response in - // backend/app/gateway/routers/thread_runs.py until the API exposes a - // structured error code for inactive run streams. return ( (status === 409 || message.includes("HTTP 409")) && - message.includes("not active on this worker") && - message.includes("cannot be streamed") + needles.every((needle) => message.includes(needle)) ); } +// Store-only run cannot be streamed (no in-memory stream bridge on this +// worker): reconnect has nothing to rejoin. +export function isInactiveRunStreamError(error: unknown): boolean { + return isRunConflictError( + error, + "not active on this worker", + "cannot be streamed", + ); +} + +/** + * Matches the gateway's terminal-state cancel conflict, raised by + * ``_cancel_conflict_detail`` in ``backend/app/gateway/routers/thread_runs.py`` + * as ``Run X is not cancellable (status: success|error|timeout)`` when + * ``RunManager.cancel`` refuses a run that already finished. + * + * The sibling ``"not active on this worker and cannot be cancelled"`` branch + * (run still pending/running on another worker in a multi-instance deploy) is + * intentionally NOT matched — that is a real cancel failure on a live run and + * must stay visible. Only the terminal-state branch is a true no-op. + */ +export function isRunNotCancellableError(error: unknown): boolean { + return isRunConflictError(error, "is not cancellable"); +} + +/** + * Preflight a reconnect: if the run already reached a terminal state, there is + * nothing to rejoin. Returns ``true`` when the caller should skip the + * underlying ``joinStream`` so the SDK's ``onSuccess`` path runs and + * ``isLoading`` flips back to false — instead of blocking forever on a drained + * stream bridge. + * + * Any error (404 for an evicted record, network blip, auth hiccup, …) falls + * back to the original join so a legitimately active reconnect is never + * silently suppressed. + */ +async function shouldSkipReconnect( + client: LangGraphClient, + threadId: string, + runId: string, +): Promise { + try { + const run = await client.runs.get(threadId, runId); + return TERMINAL_RUN_STATUSES.has(run.status); + } catch { + return false; + } +} + export function clearReconnectRun( threadId: string | null | undefined, runId: string, @@ -99,8 +187,35 @@ function createCompatibleClient(isMock?: boolean): LangGraphClient { sanitizeRunStreamOptions(payload), )) as typeof client.runs.stream; + const originalCancel = client.runs.cancel.bind(client.runs); + client.runs.cancel = (async (threadId, runId, wait, action, options) => { + try { + return await originalCancel(threadId, runId, wait, action, options); + } catch (error) { + if (isRunNotCancellableError(error)) { + // The run already reached a terminal state, so cancelling it is a + // no-op. Swallow the 409 so a stop click during the finish window + // (backend flipped to ``success`` but the SSE stream hasn't drained) + // doesn't surface as an unhandled rejection, and clear the now-stale + // reconnect key. clearReconnectRun only removes the key when it still + // matches this runId, so a newer run's key is never touched. + clearReconnectRun(threadId, runId); + return; + } + throw error; + } + }) as typeof client.runs.cancel; + const originalJoinStream = client.runs.joinStream.bind(client.runs); client.runs.joinStream = async function* (threadId, runId, options) { + // Short-circuit reconnects to runs that have already finished: otherwise a + // reload after the backend's stream bridge is reaped blocks forever on a + // drained condition variable, pinning ``isLoading`` true so the first + // post-reload message is routed to ``stop`` instead of ``submit``. + if (threadId && (await shouldSkipReconnect(client, threadId, runId))) { + clearReconnectRun(threadId, runId); + return; + } try { yield* originalJoinStream( threadId, diff --git a/frontend/tests/unit/core/api/api-client.test.ts b/frontend/tests/unit/core/api/api-client.test.ts index 45c5af394..0d6f39afd 100644 --- a/frontend/tests/unit/core/api/api-client.test.ts +++ b/frontend/tests/unit/core/api/api-client.test.ts @@ -4,6 +4,7 @@ import { clearReconnectRun, getAPIClient, isInactiveRunStreamError, + isRunNotCancellableError, } from "@/core/api/api-client"; function makeSessionStorage() { @@ -121,3 +122,213 @@ test("rethrows unrelated streaming errors", async () => { expect(sessionStorage.removeItem).not.toHaveBeenCalled(); }); + +test("identifies terminal-state cancel conflicts", () => { + const error = Object.assign( + new Error( + 'HTTP 409: {"detail":"Run run-1 is not cancellable (status: success)"}', + ), + { status: 409 }, + ); + + expect(isRunNotCancellableError(error)).toBe(true); +}); + +test("does not classify not-active-on-worker cancel as terminal", () => { + // A run still pending/running on another worker is a real cancel failure — + // it must stay visible and must NOT be swallowed. + const error = Object.assign( + new Error( + 'HTTP 409: {"detail":"Run run-1 is not active on this worker and cannot be cancelled"}', + ), + { status: 409 }, + ); + + expect(isRunNotCancellableError(error)).toBe(false); +}); + +test("swallows terminal-state cancel 409 and clears stale key", async () => { + const sessionStorage = makeSessionStorage(); + sessionStorage.setItem("lg:stream:thread-1", "run-1"); + rs.stubGlobal("window", { + location: { origin: "http://localhost:2026" }, + sessionStorage, + }); + rs.stubGlobal( + "fetch", + rs.fn(async () => { + return new Response( + JSON.stringify({ + detail: "Run run-1 is not cancellable (status: success)", + }), + { status: 409 }, + ); + }), + ); + + // Resolves (no throw) — cancelling an already-finished run is a no-op. + await expect( + getAPIClient(true).runs.cancel("thread-1", "run-1"), + ).resolves.toBeUndefined(); + + expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); +}); + +test("rethrows not-active-on-worker cancel 409", async () => { + const sessionStorage = makeSessionStorage(); + sessionStorage.setItem("lg:stream:thread-1", "run-1"); + rs.stubGlobal("window", { + location: { origin: "http://localhost:2026" }, + sessionStorage, + }); + rs.stubGlobal( + "fetch", + rs.fn(async () => { + return new Response( + JSON.stringify({ + detail: + "Run run-1 is not active on this worker and cannot be cancelled", + }), + { status: 409 }, + ); + }), + ); + + await expect( + getAPIClient(true).runs.cancel("thread-1", "run-1"), + ).rejects.toThrow("HTTP 409"); + + expect(sessionStorage.removeItem).not.toHaveBeenCalled(); +}); + +test("short-circuits reconnect to a terminal run", async () => { + const sessionStorage = makeSessionStorage(); + sessionStorage.setItem("lg:stream:thread-1", "run-1"); + const fetchFn = rs.fn(async (url: string | URL) => { + const path = url.toString(); + // Preflight GET /threads/{tid}/runs/{runId} reports a finished run. + if (path.endsWith("/runs/run-1")) { + return new Response(JSON.stringify({ status: "success" }), { + status: 200, + }); + } + // If join were attempted it must never run; fail loudly if it does. + return new Response(JSON.stringify({ detail: "unexpected join" }), { + status: 500, + }); + }); + rs.stubGlobal("window", { + location: { origin: "http://localhost:2026" }, + sessionStorage, + }); + rs.stubGlobal("fetch", fetchFn); + + const gen = getAPIClient(true).runs.joinStream("thread-1", "run-1"); + await expect(gen.next()).resolves.toMatchObject({ done: true }); + + // Preflight only — no stream/join request beyond the GET. + expect(fetchFn).toHaveBeenCalledTimes(1); + expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); +}); + +test("falls back to join when preflight cannot resolve the run", async () => { + const sessionStorage = makeSessionStorage(); + sessionStorage.setItem("lg:stream:thread-1", "run-1"); + const fetchFn = rs.fn(async (url: string | URL) => { + const path = url.toString(); + // Preflight GET 404s (record evicted) — must fall back to join. + if (path.endsWith("/runs/run-1")) { + return new Response(JSON.stringify({ detail: "Run run-1 not found" }), { + status: 404, + }); + } + // Join then surfaces the inactive-stream 409 and clears the key. + return new Response( + JSON.stringify({ + detail: "Run run-1 is not active on this worker and cannot be streamed", + }), + { status: 409 }, + ); + }); + rs.stubGlobal("window", { + location: { origin: "http://localhost:2026" }, + sessionStorage, + }); + rs.stubGlobal("fetch", fetchFn); + + await expect( + getAPIClient(true).runs.joinStream("thread-1", "run-1").next(), + ).resolves.toMatchObject({ done: true }); + + expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); +}); + +test("proceeds to join when the run is still active", async () => { + // Positive path: a running/pending run must NOT be short-circuited — the + // preflight must let the real join through so an in-flight stream can be + // rejoined. Proves the guard does not over-eagerly skip active runs. + const sessionStorage = makeSessionStorage(); + sessionStorage.setItem("lg:stream:thread-1", "run-1"); + const fetchFn = rs.fn(async (url: string | URL) => { + const path = url.toString(); + // Preflight GET reports an active run. + if (path.endsWith("/runs/run-1")) { + return new Response(JSON.stringify({ status: "running" }), { + status: 200, + }); + } + // The real join is attempted (here it surfaces the inactive-stream 409, + // which the wrapper catches and clears the key — same path as production). + return new Response( + JSON.stringify({ + detail: "Run run-1 is not active on this worker and cannot be streamed", + }), + { status: 409 }, + ); + }); + rs.stubGlobal("window", { + location: { origin: "http://localhost:2026" }, + sessionStorage, + }); + rs.stubGlobal("fetch", fetchFn); + + await expect( + getAPIClient(true).runs.joinStream("thread-1", "run-1").next(), + ).resolves.toMatchObject({ done: true }); + + // Two requests: preflight GET + the real join. A short-circuit would be one. + expect(fetchFn).toHaveBeenCalledTimes(2); + expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); +}); + +test("short-circuits reconnect to an interrupted (user-cancelled) run", async () => { + // Regression: interrupted is a persisted terminal status written by + // RunManager.cancel(). Reconnecting to it must short-circuit like other + // terminal states, otherwise — once the bridge is reaped — joinStream blocks + // forever and isLoading sticks. Keeps the frontend status set aligned with + // the backend RunStatus contract. + const sessionStorage = makeSessionStorage(); + sessionStorage.setItem("lg:stream:thread-1", "run-1"); + const fetchFn = rs.fn(async (url: string | URL) => { + const path = url.toString(); + if (path.endsWith("/runs/run-1")) { + return new Response(JSON.stringify({ status: "interrupted" }), { + status: 200, + }); + } + return new Response(JSON.stringify({ detail: "unexpected join" }), { + status: 500, + }); + }); + rs.stubGlobal("window", { + location: { origin: "http://localhost:2026" }, + sessionStorage, + }); + rs.stubGlobal("fetch", fetchFn); + + const gen = getAPIClient(true).runs.joinStream("thread-1", "run-1"); + await expect(gen.next()).resolves.toMatchObject({ done: true }); + + expect(fetchFn).toHaveBeenCalledTimes(1); + expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); +});