Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions apps/presentation/dashboard/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
116 changes: 116 additions & 0 deletions apps/presentation/dashboard/smoke/chat-stream-stall-smoke.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
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<string, unknown>) {
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<Uint8Array>({
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 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");
65 changes: 62 additions & 3 deletions apps/presentation/dashboard/src/data/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -866,23 +866,58 @@ 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) {
return new Promise<void>((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<typeof globalThis.setTimeout> | 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 });
Expand All @@ -891,6 +926,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");
Expand All @@ -901,6 +937,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");
Expand All @@ -911,8 +948,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) {
Expand Down
Loading