deer-flow/frontend/tests/unit/core/api/api-client.test.ts
Huixin615 1cd5dea336
fix(streaming): signal replay history gaps (#4426)
* fix(streaming): signal replay history gaps

* fix(streaming): guard initial Redis replay window

* fix(frontend): align inactive gap recovery

---------

Co-authored-by: Willem Jiang <willem.jiang@gmail.com>
2026-07-27 07:13:06 +08:00

676 lines
20 KiB
TypeScript

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<string, string>();
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");
});