diff --git a/apps/presentation/dashboard/package.json b/apps/presentation/dashboard/package.json index f4c7348c9f..7a9d0d4382 100644 --- a/apps/presentation/dashboard/package.json +++ b/apps/presentation/dashboard/package.json @@ -20,6 +20,7 @@ "smoke:benchmark-study-browser": "node ../../../examples/dashboard-benchmark-study-browser-smoke.mjs", "smoke:capability-configuration": "rm -rf node_modules/.cache/loopx-capability-configuration-smoke && tsc --ignoreConfig --target ES2022 --module NodeNext --moduleResolution NodeNext --types node --skipLibCheck --strict --outDir node_modules/.cache/loopx-capability-configuration-smoke smoke/capability-configuration-smoke.ts src/data/capability-configuration.ts && node node_modules/.cache/loopx-capability-configuration-smoke/smoke/capability-configuration-smoke.js", "smoke:chat-route": "tsc --ignoreConfig --target ES2022 --module ES2022 --moduleResolution Bundler --ignoreDeprecations 6.0 --skipLibCheck --strict --outDir node_modules/.cache/loopx-chat-route-smoke smoke/chat-route-smoke.ts src/data/chat-model.ts src/vite-env.d.ts && node node_modules/.cache/loopx-chat-route-smoke/smoke/chat-route-smoke.js", + "smoke:chat-stream-stall": "vite build --ssr smoke/chat-stream-stall-smoke.ts --outDir node_modules/.cache/loopx-chat-stream-stall --emptyOutDir && node node_modules/.cache/loopx-chat-stream-stall/chat-stream-stall-smoke.js", "smoke:chat-turn-acceptance-retry": "vite build --ssr smoke/chat-turn-acceptance-retry-smoke.ts --outDir node_modules/.cache/loopx-chat-turn-acceptance-retry --emptyOutDir && node node_modules/.cache/loopx-chat-turn-acceptance-retry/chat-turn-acceptance-retry-smoke.js", "smoke:demo-readiness": "bash ../../../scripts/loopx-python.sh --exec ../../../examples/dashboard-demo-readiness-smoke.py", "smoke:frontstage-browser": "node ../../../examples/dashboard-frontstage-browser-smoke.mjs", diff --git a/apps/presentation/dashboard/smoke/chat-stream-stall-smoke.ts b/apps/presentation/dashboard/smoke/chat-stream-stall-smoke.ts new file mode 100644 index 0000000000..1e944a9962 --- /dev/null +++ b/apps/presentation/dashboard/smoke/chat-stream-stall-smoke.ts @@ -0,0 +1,131 @@ +import assert from "node:assert/strict"; + +import { + CHAT_STREAM_STALL_TIMEOUT_MS, + ChatApiError, + type ChatStreamEvent, + resumeChatTurnStreaming, + streamChatTurn, +} from "../src/data/chat.ts"; + +// A stuck SSE connection delivers headers and a first event, then no bytes at +// all, not even the service heartbeat. The reader must abandon it, resume from +// its cursor, and tell the pending reply it is reconnecting. A caller abort is +// different: it ends transport activity at any point and is never retried. +const encoder = new TextEncoder(); +const originalFetch = globalThis.fetch; +const originalSetTimeout = globalThis.setTimeout; +const eventsUrl = "/api/chat/sessions/s/turns/t/events"; +const requests: string[] = []; +let serve: (url: URL, signal: AbortSignal | null | undefined) => Response; + +function sseBlock(eventId: string, kind: string, payload: Record) { + return encoder.encode(`id: ${eventId}\nevent: ${kind}\ndata: ${JSON.stringify({ created_at: "2026-09-28T00:00:00Z", event_id: eventId, kind, payload, sequence: 1 })}\n\n`); +} + +function silentAfter(first: Uint8Array | null, signal: AbortSignal | null | undefined) { + return new ReadableStream({ + start(controller) { + if (first) controller.enqueue(first); + signal?.addEventListener("abort", () => controller.error(new DOMException("aborted", "AbortError")), { once: true }); + }, + }); +} + +function isAbort(error: unknown) { + return error instanceof DOMException && error.name === "AbortError"; +} + +globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => { + const url = new URL(typeof input === "string" || input instanceof URL ? input : input.url); + requests.push(url.search); + return serve(url, init?.signal); +}) as typeof fetch; + +try { + // A stalled stream resumes from its cursor. + serve = (url, signal) => url.searchParams.has("after") + ? new Response(sseBlock("event-2", "turn.completed", { response: { message: "done" } }), { status: 200 }) + : new Response(silentAfter(sseBlock("event-1", "turn.started", {}), signal), { status: 200 }); + const events: ChatStreamEvent[] = []; + const started = Date.now(); + await streamChatTurn(eventsUrl, (event) => events.push(event), undefined, { stallTimeoutMs: 100 }); + assert.ok(Date.now() - started < 5_000, "a stalled stream is abandoned after the stall timeout"); + assert.deepEqual(requests, ["", "?after=event-1"], "the reader resumes from the last delivered event"); + assert.deepEqual(events.map((event) => event.kind), ["turn.started", "agent.phase", "turn.completed"]); + assert.equal(events[1].payload.method, "client/reconnect", "the pending reply learns it is reconnecting"); + assert.equal(events[1].event_id, "", "the local reconnect phase never moves the cursor"); + + // A caller abort from inside a delivered event ends the stream immediately. + requests.length = 0; + const inEvent = new AbortController(); + await assert.rejects(streamChatTurn(eventsUrl, () => inEvent.abort(), inEvent.signal, { stallTimeoutMs: 10_000 }), isAbort); + assert.deepEqual(requests, [""], "a caller abort is not retried"); + + // An abort that happened before the call opens no connection. + requests.length = 0; + const before = new AbortController(); + before.abort(); + const beforeEvents: ChatStreamEvent[] = []; + await assert.rejects(streamChatTurn(eventsUrl, (event) => beforeEvents.push(event), before.signal), isAbort); + assert.deepEqual(requests, [], "an already aborted caller opens no connection"); + assert.deepEqual(beforeEvents, []); + + // An abort during the retry backoff ends the wait and opens no new connection. + requests.length = 0; + serve = () => new Response("unavailable", { status: 503 }); + const backoff = new AbortController(); + const backoffEvents: ChatStreamEvent[] = []; + const backoffStarted = Date.now(); + const backoffRun = streamChatTurn(eventsUrl, (event) => backoffEvents.push(event), backoff.signal); + originalSetTimeout(() => backoff.abort(), 30); + await assert.rejects(backoffRun, isAbort); + assert.ok(Date.now() - backoffStarted < 250, "the backoff wait ends when the caller aborts"); + assert.deepEqual(requests, [""], "no connection opens after the caller aborts during backoff"); + assert.deepEqual(backoffEvents.map((event) => event.payload.method), ["client/reconnect"]); + + // An abort from the reconnect phase's own callback ends the stream without + // sitting out the backoff. + requests.length = 0; + serve = () => new Response("unavailable", { status: 503 }); + const onPhase = new AbortController(); + let phaseAt = 0; + const onPhaseRun = streamChatTurn(eventsUrl, (event) => { + if (event.payload.method !== "client/reconnect") return; + phaseAt = Date.now(); + onPhase.abort(); + }, onPhase.signal); + await assert.rejects(onPhaseRun, isAbort); + assert.ok(phaseAt > 0 && Date.now() - phaseAt < 200, "an abort from the reconnect callback skips the backoff"); + assert.deepEqual(requests, [""], "no connection opens after the reconnect callback aborts"); + + // An abort while a read waits ends that read without reconnecting. + requests.length = 0; + serve = (_url, signal) => new Response(silentAfter(null, signal), { status: 200 }); + const reading = new AbortController(); + const readingRun = streamChatTurn(eventsUrl, () => {}, reading.signal, { stallTimeoutMs: 10_000 }); + originalSetTimeout(() => reading.abort(), 30); + await assert.rejects(readingRun, isAbort); + assert.deepEqual(requests, [""], "an abort during a read is not retried"); + + // Four stalled connections end with the resume wrapper's typed recovery + // error, like any other exhausted reconnect. Only the watchdog clock is + // shortened; the retry backoff keeps its production delays. + requests.length = 0; + globalThis.setTimeout = ((handler: TimerHandler, timeout?: number, ...args: unknown[]) => + originalSetTimeout(handler, timeout === CHAT_STREAM_STALL_TIMEOUT_MS ? 50 : timeout, ...args)) as typeof setTimeout; + const exhausted = await resumeChatTurnStreaming("s", "t").then(() => null, (error: unknown) => error); + globalThis.setTimeout = originalSetTimeout; + assert.equal(requests.length, 4, "the watchdog reconnects a bounded number of times"); + assert.ok(exhausted instanceof ChatApiError, "stall exhaustion is a typed Chat error, not a raw abort"); + assert.equal(exhausted.payload.reconnectable, true); + assert.equal(exhausted.payload.session_id, "s"); + assert.equal(exhausted.payload.turn_id, "t"); + assert.equal(exhausted.payload.events_url, eventsUrl); + assert.equal(exhausted.payload.reconnect_attempts, 4); +} finally { + globalThis.fetch = originalFetch; + globalThis.setTimeout = originalSetTimeout; +} + +console.log("chat-stream-stall-smoke: ok"); diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index 66b81c6c16..1a77682671 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -866,23 +866,61 @@ function parseSseBlock(block: string): ChatStreamEvent | null { } } +// The Chat service sends an SSE heartbeat every 15 seconds while a Turn runs. +// Three missed heartbeats mean the connection is stuck rather than slow, so the +// reader reconnects from its cursor instead of waiting on a silent socket. +export const CHAT_STREAM_STALL_TIMEOUT_MS = 45_000; + +// Resolves early when the caller aborts, so the next attempt sees the abort +// instead of opening a connection the caller no longer wants. +function waitForRetry(ms: number, signal?: AbortSignal) { + // A callback may abort while handling the reconnect phase, before this wait + // starts listening; that abort must not sit out the backoff. + if (signal?.aborted) return Promise.resolve(); + return new Promise((resolve) => { + const done = () => { + globalThis.clearTimeout(timer); + signal?.removeEventListener("abort", done); + resolve(); + }; + const timer = globalThis.setTimeout(done, ms); + signal?.addEventListener("abort", done, { once: true }); + }); +} + export async function streamChatTurn( eventsUrl: string, onEvent: (event: ChatStreamEvent) => void, signal?: AbortSignal, + options: { stallTimeoutMs?: number } = {}, ) { + const stallTimeoutMs = options.stallTimeoutMs ?? CHAT_STREAM_STALL_TIMEOUT_MS; let cursor = ""; let attempts = 0; let terminal = false; while (!terminal && attempts < 4) { + signal?.throwIfAborted(); const origin = typeof window === "undefined" ? "http://127.0.0.1" : window.location.origin; const url = new URL(chatApiUrl(eventsUrl), origin); if (cursor) url.searchParams.set("after", cursor); + const attempt = new AbortController(); + const abortAttempt = () => attempt.abort(); + signal?.addEventListener("abort", abortAttempt, { once: true }); + let stalled = false; + let stallTimer: ReturnType | undefined; + const armStallTimer = () => { + if (stallTimer !== undefined) globalThis.clearTimeout(stallTimer); + stallTimer = globalThis.setTimeout(() => { + stalled = true; + attempt.abort(); + }, stallTimeoutMs); + }; try { + armStallTimer(); const response = await fetch(url, { cache: "no-store", headers: { Accept: "text/event-stream" }, - signal, + signal: attempt.signal, }); if (!response.ok || !response.body) { throw new ChatApiError(`SSE HTTP ${response.status}`, { status: response.status }); @@ -891,6 +929,7 @@ export async function streamChatTurn( const decoder = new TextDecoder(); let buffer = ""; while (true) { + armStallTimer(); const { done, value } = await reader.read(); buffer += decoder.decode(value, { stream: !done }).replaceAll("\r\n", "\n"); let boundary = buffer.indexOf("\n\n"); @@ -901,6 +940,7 @@ export async function streamChatTurn( if (event) { if (event.event_id) cursor = event.event_id; onEvent(event); + signal?.throwIfAborted(); terminal = ["turn.completed", "turn.interrupted", "turn.failed"].includes(event.kind); } boundary = buffer.indexOf("\n\n"); @@ -911,8 +951,30 @@ export async function streamChatTurn( } catch (error) { if (signal?.aborted) throw error; attempts += 1; - if (attempts >= 4) throw error; - await new Promise((resolve) => globalThis.setTimeout(resolve, 250 * 2 ** (attempts - 1))); + if (attempts >= 4) { + // A stalled connection is a transport failure, not a caller abort, so + // it ends with the same typed error as any other exhausted reconnect. + if (stalled) { + throw new ChatApiError("Agent 事件流连接已断开。", { + reconnect_attempts: attempts, + stall_timeout_ms: stallTimeoutMs, + }); + } + throw error; + } + // A local phase keeps the pending reply honest while the reader resumes + // from its cursor. It carries no event id, so the cursor is unchanged. + onEvent({ + created_at: new Date().toISOString(), + event_id: "", + kind: "agent.phase", + payload: { label: "连接中断,正在重连…", method: "client/reconnect" }, + sequence: 0, + }); + await waitForRetry(250 * 2 ** (attempts - 1), signal); + } finally { + if (stallTimer !== undefined) globalThis.clearTimeout(stallTimer); + signal?.removeEventListener("abort", abortAttempt); } } if (!terminal) {