import { afterEach, expect, test, rs } from "@rstest/core"; import { clearReconnectRun, getAPIClient, isInactiveRunStreamError, isRunNotCancellableError, StreamReplayGapError, } from "@/core/api/api-client"; function makeSessionStorage() { const values = new Map(); return { getItem: rs.fn((key: string) => values.get(key) ?? null), removeItem: rs.fn((key: string) => { values.delete(key); }), setItem: rs.fn((key: string, value: string) => { values.set(key, value); }), }; } function makeSSEResponse(body: string, headers?: HeadersInit) { return new Response(body, { status: 200, headers: { "Content-Type": "text/event-stream", ...headers, }, }); } afterEach(() => { rs.unstubAllGlobals(); }); test("identifies inactive run stream errors", () => { const error = Object.assign( new Error( 'HTTP 409: {"detail":"Run run-1 is not active on this worker and cannot be streamed"}', ), { status: 409 }, ); expect(isInactiveRunStreamError(error)).toBe(true); }); test("does not classify unrelated conflict errors as inactive streams", () => { const error = Object.assign(new Error("HTTP 409: run is still active"), { status: 409, }); expect(isInactiveRunStreamError(error)).toBe(false); }); test("clears matching reconnect metadata", () => { const sessionStorage = makeSessionStorage(); sessionStorage.setItem("lg:stream:thread-1", "run-1"); rs.stubGlobal("window", { sessionStorage }); clearReconnectRun("thread-1", "run-1"); expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); }); test("keeps newer reconnect metadata", () => { const sessionStorage = makeSessionStorage(); sessionStorage.setItem("lg:stream:thread-1", "newer-run"); rs.stubGlobal("window", { sessionStorage }); clearReconnectRun("thread-1", "stale-run"); expect(sessionStorage.removeItem).not.toHaveBeenCalled(); }); test("ignores reconnect metadata storage access failures", () => { rs.stubGlobal("window", { get sessionStorage() { throw new DOMException("Blocked", "SecurityError"); }, }); expect(() => clearReconnectRun("thread-1", "run-1")).not.toThrow(); }); test("clears stale reconnect metadata when join stream cannot be resumed", 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 streamed", }), { status: 409 }, ); }), ); await expect( getAPIClient(true).runs.joinStream("thread-1", "run-1").next(), ).resolves.toMatchObject({ done: true }); expect(sessionStorage.removeItem).toHaveBeenCalledWith("lg:stream:thread-1"); }); test("rethrows unrelated streaming errors", 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 is still active" }), { status: 409, }); }), ); await expect( getAPIClient(true).runs.joinStream("thread-1", "run-1").next(), ).rejects.toThrow("HTTP 409"); 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("recovers a join stream gap from durable state and resumes after the retained tail", async () => { const sessionStorage = makeSessionStorage(); sessionStorage.setItem("lg:stream:thread-1", "run-1"); const gap = { code: "stream_replay_gap", run_id: "run-1", requested_event_id: "1-0", earliest_available_event_id: "2-0", latest_available_event_id: "3-0", recovery: "reload_durable_state", }; const recoveryRequests: RequestInit[] = []; const fetchFn = rs.fn(async (url: string | URL, init?: RequestInit) => { const path = url.toString(); if (path.endsWith("/runs/run-1")) { return new Response(JSON.stringify({ status: "running" }), { status: 200, }); } if (path.includes("/runs/run-1/stream")) { recoveryRequests.push(init ?? {}); if (recoveryRequests.length === 1) { return makeSSEResponse(`event: gap\ndata: ${JSON.stringify(gap)}\n\n`); } return makeSSEResponse("event: end\ndata: null\n\n"); } if (path.includes("/threads/thread-1/state")) { return new Response( JSON.stringify({ values: { messages: [{ type: "ai", content: "durable" }] }, next: [], tasks: [], metadata: {}, created_at: null, checkpoint: {}, parent_checkpoint: null, }), { status: 200 }, ); } return new Response(JSON.stringify({ detail: "unexpected request" }), { status: 500, }); }); rs.stubGlobal("window", { location: { origin: "http://localhost:2026" }, sessionStorage, }); rs.stubGlobal("fetch", fetchFn); const received: Array<{ id?: string; event: string; data: unknown }> = []; for await (const entry of getAPIClient(true).runs.joinStream( "thread-1", "run-1", { lastEventId: "1-0" }, )) { received.push(entry); } expect(received).toEqual([ { event: "custom", data: { type: "stream_replay_gap", ...gap }, }, { event: "values", data: { messages: [{ type: "ai", content: "durable" }] }, }, { event: "end", data: null }, ]); expect(new Headers(recoveryRequests[1]?.headers).get("Last-Event-ID")).toBe( "3-0", ); }); test("recovers a gap emitted by the initial run stream", async () => { const sessionStorage = makeSessionStorage(); const gap = { code: "stream_replay_gap", run_id: "run-2", requested_event_id: null, earliest_available_event_id: "4-0", latest_available_event_id: "5-0", recovery: "reload_durable_state", }; const recoveryHeaders: Headers[] = []; const fetchFn = rs.fn(async (url: string | URL, init?: RequestInit) => { const path = url.toString(); if (path.endsWith("/threads/thread-2/runs/stream")) { return makeSSEResponse(`event: gap\ndata: ${JSON.stringify(gap)}\n\n`, { "Content-Location": "/threads/thread-2/runs/run-2", }); } if (path.includes("/threads/thread-2/state")) { return new Response( JSON.stringify({ values: { messages: [{ type: "ai", content: "checkpoint" }] }, }), { status: 200 }, ); } if (path.includes("/runs/run-2/stream")) { recoveryHeaders.push(new Headers(init?.headers)); return makeSSEResponse("event: end\ndata: null\n\n"); } return new Response(JSON.stringify({ detail: "unexpected request" }), { status: 500, }); }); rs.stubGlobal("window", { location: { origin: "http://localhost:2026" }, sessionStorage, }); rs.stubGlobal("fetch", fetchFn); const received: Array<{ id?: string; event: string; data: unknown }> = []; for await (const entry of getAPIClient(true).runs.stream( "thread-2", "lead_agent", { streamResumable: true }, )) { received.push(entry); } expect(received).toEqual([ { event: "custom", data: { type: "stream_replay_gap", ...gap }, }, { event: "values", data: { messages: [{ type: "ai", content: "checkpoint" }] }, }, { event: "end", data: null }, ]); expect(recoveryHeaders[0]?.get("Last-Event-ID")).toBe("5-0"); }); test("clears reconnect metadata when an initial stream gap resume is inactive", async () => { const sessionStorage = makeSessionStorage(); const gap = { code: "stream_replay_gap", run_id: "run-inactive-resume", requested_event_id: null, earliest_available_event_id: "4-0", latest_available_event_id: "5-0", recovery: "reload_durable_state", }; let recoveryRequests = 0; const fetchFn = rs.fn(async (url: string | URL) => { const path = url.toString(); if (path.endsWith("/threads/thread-inactive-resume/runs/stream")) { return makeSSEResponse(`event: gap\ndata: ${JSON.stringify(gap)}\n\n`, { "Content-Location": "/threads/thread-inactive-resume/runs/run-inactive-resume", }); } if (path.includes("/threads/thread-inactive-resume/state")) { return new Response( JSON.stringify({ values: { messages: [{ type: "ai", content: "checkpoint" }] }, }), { status: 200 }, ); } if (path.includes("/runs/run-inactive-resume/stream")) { recoveryRequests += 1; return new Response( JSON.stringify({ detail: "Run run-inactive-resume is not active on this worker and cannot be streamed", }), { status: 409 }, ); } return new Response(JSON.stringify({ detail: "unexpected request" }), { status: 500, }); }); rs.stubGlobal("window", { location: { origin: "http://localhost:2026" }, sessionStorage, }); rs.stubGlobal("fetch", fetchFn); const stream = getAPIClient(true).runs.stream( "thread-inactive-resume", "lead_agent", { streamResumable: true }, ); await expect(stream.next()).resolves.toMatchObject({ done: false, value: { event: "custom", data: { type: "stream_replay_gap", ...gap }, }, }); await expect(stream.next()).resolves.toMatchObject({ done: false, value: { event: "values", data: { messages: [{ type: "ai", content: "checkpoint" }] }, }, }); await expect(stream.next()).resolves.toMatchObject({ done: true }); expect(recoveryRequests).toBe(1); expect(sessionStorage.getItem("lg:stream:thread-inactive-resume")).toBeNull(); }); test("stops after five consecutive stream gap recoveries", async () => { const sessionStorage = makeSessionStorage(); sessionStorage.setItem("lg:stream:thread-gap-loop", "run-gap-loop"); let streamCalls = 0; let stateCalls = 0; const fetchFn = rs.fn(async (url: string | URL) => { const path = url.toString(); if (path.endsWith("/runs/run-gap-loop")) { return new Response(JSON.stringify({ status: "running" }), { status: 200, }); } if (path.includes("/runs/run-gap-loop/stream")) { streamCalls += 1; const gap = { code: "stream_replay_gap", run_id: "run-gap-loop", requested_event_id: `${streamCalls}-0`, earliest_available_event_id: `${streamCalls + 1}-0`, latest_available_event_id: `${streamCalls + 2}-0`, recovery: "reload_durable_state", }; return makeSSEResponse(`event: gap\ndata: ${JSON.stringify(gap)}\n\n`); } if (path.includes("/threads/thread-gap-loop/state")) { stateCalls += 1; return new Response(JSON.stringify({ values: { messages: [] } }), { status: 200, }); } return new Response(JSON.stringify({ detail: "unexpected request" }), { status: 500, }); }); rs.stubGlobal("window", { location: { origin: "http://localhost:2026" }, sessionStorage, }); rs.stubGlobal("fetch", fetchFn); const consume = async () => { for await (const entry of getAPIClient(true).runs.joinStream( "thread-gap-loop", "run-gap-loop", { lastEventId: "1-0" }, )) { // Drain until the recovery budget is exhausted. void entry; } }; await expect(consume()).rejects.toBeInstanceOf(StreamReplayGapError); expect(streamCalls).toBe(6); expect(stateCalls).toBe(5); }); test("surfaces durable-state recovery failures as a structured gap error", async () => { const sessionStorage = makeSessionStorage(); sessionStorage.setItem("lg:stream:thread-gap-state", "run-gap-state"); sessionStorage.setItem.mockClear(); const gap = { code: "stream_replay_gap", run_id: "run-gap-state", requested_event_id: "1-0", earliest_available_event_id: "2-0", latest_available_event_id: "3-0", recovery: "reload_durable_state", }; const fetchFn = rs.fn(async (url: string | URL) => { const path = url.toString(); if (path.endsWith("/runs/run-gap-state")) { return new Response(JSON.stringify({ status: "running" }), { status: 200, }); } if (path.includes("/runs/run-gap-state/stream")) { return makeSSEResponse(`event: gap\ndata: ${JSON.stringify(gap)}\n\n`); } if (path.includes("/threads/thread-gap-state/state")) { return new Response(JSON.stringify({ detail: "state unavailable" }), { status: 404, }); } return new Response(JSON.stringify({ detail: "unexpected request" }), { status: 500, }); }); rs.stubGlobal("window", { location: { origin: "http://localhost:2026" }, sessionStorage, }); rs.stubGlobal("fetch", fetchFn); const consume = async () => { for await (const entry of getAPIClient(true).runs.joinStream( "thread-gap-state", "run-gap-state", { lastEventId: "1-0" }, )) { void entry; } }; const error = await consume().catch((cause: unknown) => cause); expect(error).toBeInstanceOf(StreamReplayGapError); expect(error).toMatchObject({ gap, recoveryAttempts: 1, recoveryCause: expect.any(Error), }); expect(sessionStorage.removeItem).toHaveBeenCalledWith( "lg:stream:thread-gap-state", ); expect(sessionStorage.setItem).not.toHaveBeenCalledWith( "lg:stream:thread-gap-state", "run-gap-state", ); }); 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"); });