Repository navigation
fix: #2013 - correctly invoke controller.error during mid-stream closure and update tests accordingly - #2019
Conversation
…tream closure and update tests accordingly
🦋 Changeset detectedLatest commit: 5eb4c35 The changes in this PR will be included in the next version bump. This PR includes changesets to release 2 packages
Not sure what this means? Click here to learn what changesets are. Click here if you're a maintainer who wants to add another changeset to this PR |
agents
@cloudflare/ai-chat
@cloudflare/codemode
hono-agents
@cloudflare/shell
@cloudflare/think
@cloudflare/voice
@cloudflare/worker-bundler
commit: |
Co-authored-by: devin-ai-integration[bot] <158243242+devin-ai-integration[bot]@users.noreply.github.com>
|
I checked this against current Resume still works after this change. The concern was that erroring the stream would stop
With the change applied, the full To merge, this needs:
Rebased close handlersdiff --git a/packages/agents/src/chat/ws-chat-transport.ts b/packages/agents/src/chat/ws-chat-transport.ts
index a5442c45..caeb86c4 100644
--- a/packages/agents/src/chat/ws-chat-transport.ts
+++ b/packages/agents/src/chat/ws-chat-transport.ts
@@ -445,7 +445,11 @@ export class WebSocketChatTransport<
};
const onClose = () => {
- finish(() => controller.close(), false, false);
+ finish(
+ () => controller.error(new Error("WebSocket closed mid-stream")),
+ false,
+ false
+ );
};
agent.addEventListener("message", onMessage, {
@@ -812,7 +816,8 @@ export class WebSocketChatTransport<
const onClose = () =>
finish(() => {
batch.flush();
- controller.close();
+ if (requestId === null) controller.close();
+ else controller.error(new Error("WebSocket closed mid-stream"));
});
agent.addEventListener("message", onMessage, {
@@ -942,7 +947,7 @@ export class WebSocketChatTransport<
finish(
() => {
batch.flush();
- controller.close();
+ controller.error(new Error("WebSocket closed mid-stream"));
},
false,
falseReconnect-and-resume testdiff --git a/packages/ai-chat/src/react-tests/use-agent-chat.test.tsx b/packages/ai-chat/src/react-tests/use-agent-chat.test.tsx
index 55d6e36f..a6222c4a 100644
--- a/packages/ai-chat/src/react-tests/use-agent-chat.test.tsx
+++ b/packages/ai-chat/src/react-tests/use-agent-chat.test.tsx
@@ -7010,3 +7010,168 @@ describe("useAgentChat transparent reconnect re-probe (#1784)", () => {
consoleErrorSpy.mockRestore();
});
});
+
+describe("useAgentChat socket close mid-stream (#2013)", () => {
+ function createAgentWithTarget({ name, url }: { name: string; url: string }) {
+ const target = new EventTarget();
+ const sentMessages: string[] = [];
+ const agent = createAgent({
+ name,
+ url,
+ send: (data: string) => sentMessages.push(data)
+ });
+ (agent as unknown as Record<string, unknown>).addEventListener =
+ target.addEventListener.bind(target);
+ (agent as unknown as Record<string, unknown>).removeEventListener =
+ target.removeEventListener.bind(target);
+ return { agent, target, sentMessages };
+ }
+
+ function dispatch(target: EventTarget, data: Record<string, unknown>) {
+ target.dispatchEvent(
+ new MessageEvent("message", { data: JSON.stringify(data) })
+ );
+ }
+
+ function sentOfType(sentMessages: string[], type: string) {
+ return sentMessages
+ .map((message) => JSON.parse(message) as { type: string; id?: string })
+ .filter((message) => message.type === type);
+ }
+
+ it("resumes after the socket drops mid-response and ends ready", async () => {
+ const { agent, target, sentMessages } = createAgentWithTarget({
+ name: "close-mid-stream-resume",
+ url: "ws://localhost:3000/agents/chat/close-mid-stream-resume?_pk=abc"
+ });
+ const onError = vi.fn();
+ let chatInstance: ReturnType<typeof useAgentChat> | null = null;
+ const statuses: string[] = [];
+
+ const TestComponent = () => {
+ const chat = useAgentChat({
+ agent,
+ getInitialMessages: null,
+ messages: [] as UIMessage[],
+ onError
+ });
+ chatInstance = chat;
+ statuses.push(chat.status);
+ const assistant = chat.messages.filter((m) => m.role === "assistant");
+ const text = assistant
+ .flatMap((m) => m.parts)
+ .map((part) => (part.type === "text" ? part.text : ""))
+ .join("");
+ return (
+ <div>
+ <div data-testid="status">{chat.status}</div>
+ <div data-testid="text">{text}</div>
+ <div data-testid="assistants">{assistant.length}</div>
+ <div data-testid="error">{chat.error?.message ?? ""}</div>
+ </div>
+ );
+ };
+
+ const screen = await act(async () => {
+ const screen = render(<TestComponent />, {
+ wrapper: ({ children }) => (
+ <StrictMode>
+ <Suspense fallback="Loading...">{children}</Suspense>
+ </StrictMode>
+ )
+ });
+ await sleep(10);
+ return screen;
+ });
+
+ await act(async () => {
+ dispatch(target, { type: "cf_agent_stream_resume_none" });
+ await sleep(10);
+ });
+
+ await act(async () => {
+ void chatInstance!.sendMessage({ text: "hi" });
+ await sleep(20);
+ });
+ const requestId = sentOfType(sentMessages, "cf_agent_use_chat_request")[0]
+ ?.id;
+ expect(requestId).toBeDefined();
+
+ const chunk = (seq: number, body: string, replay = false) => ({
+ type: "cf_agent_use_chat_response",
+ id: requestId,
+ body,
+ done: false,
+ seq,
+ ...(replay && { replay: true })
+ });
+
+ await act(async () => {
+ dispatch(target, chunk(0, '{"type":"start","messageId":"a1"}'));
+ dispatch(target, chunk(1, '{"type":"text-start","id":"t1"}'));
+ dispatch(
+ target,
+ chunk(2, '{"type":"text-delta","id":"t1","delta":"Hel"}')
+ );
+ await sleep(10);
+ });
+ await expect.element(screen.getByTestId("text")).toHaveTextContent("Hel");
+
+ const resumeRequestsBefore = sentOfType(
+ sentMessages,
+ "cf_agent_stream_resume_request"
+ ).length;
+ await act(async () => {
+ target.dispatchEvent(new Event("close"));
+ await sleep(20);
+ });
+ const statusAfterClose = statuses.at(-1);
+
+ await act(async () => {
+ target.dispatchEvent(new Event("open"));
+ await sleep(20);
+ });
+ expect(
+ sentOfType(sentMessages, "cf_agent_stream_resume_request").length
+ ).toBeGreaterThan(resumeRequestsBefore);
+
+ await act(async () => {
+ dispatch(target, { type: "cf_agent_stream_resuming", id: requestId });
+ await sleep(10);
+ });
+ await act(async () => {
+ dispatch(target, chunk(0, '{"type":"start","messageId":"a1"}', true));
+ dispatch(target, chunk(1, '{"type":"text-start","id":"t1"}', true));
+ dispatch(
+ target,
+ chunk(2, '{"type":"text-delta","id":"t1","delta":"Hel"}', true)
+ );
+ dispatch(
+ target,
+ chunk(3, '{"type":"text-delta","id":"t1","delta":"lo"}')
+ );
+ dispatch(target, chunk(4, '{"type":"text-end","id":"t1"}'));
+ dispatch(target, {
+ type: "cf_agent_use_chat_response",
+ id: requestId,
+ body: "",
+ done: true
+ });
+ await sleep(20);
+ });
+
+ await expect
+ .element(screen.getByTestId("status"))
+ .toHaveTextContent("ready");
+ await expect
+ .element(screen.getByTestId("text"))
+ .toHaveTextContent(/^Hello$/);
+ await expect
+ .element(screen.getByTestId("assistants"))
+ .toHaveTextContent("1");
+ // The drop surfaces as an error once, and the resume clears it.
+ expect(statusAfterClose).toBe("error");
+ expect(onError).toHaveBeenCalledTimes(1);
+ await expect.element(screen.getByTestId("error")).toBeEmptyDOMElement();
+ });
+}); |
Co-authored-by: Cursor <cursoragent@cursor.com> # Conflicts: # packages/agents/src/chat/ws-chat-transport.ts
…terrupted turn (cloudflare#2013) Co-authored-by: Cursor <cursoragent@cursor.com>
| batch?.flush(); | ||
| controller.enqueue({ type: "error", errorText: SOCKET_CLOSED_MID_STREAM }); | ||
| controller.close(); |
There was a problem hiding this comment.
🔍 Socket closure changes the stream error contract
The description promises controller.error(), but interruptChatStream emits an error chunk and closes normally. Direct stream readers receive a fulfilled read, not a rejection; confirm that this public transport contract is intended.
Was this helpful? React with 👍 or 👎 to provide feedback.
There was a problem hiding this comment.
Intended. controller.error() discards chunks the consumer hasn't read yet, so it would drop content already received (e.g. a replayed resume burst); the error chunk keeps that content while still putting useChat into the error state. The PR description was stale and now documents this contract, including that direct readers see an error chunk rather than a rejected read.
|
✅ agents import sizes: no significant changes ( |
…port - appliedChunks ledger: a continuation replayed after a reconnect skips the frames this client already applied, decided as each frame arrives and applied to the projected chunks so the projector still sees the whole run (cloudflare#1951, upstream cloudflare#2348) - a socket close before the terminal frame marks the stream interrupted, and the AI SDK transport ends it with an error chunk behind the unread chunks (cloudflare#2013, upstream cloudflare#2019) - activeServerTurnId getter for the hook's held-turn tracking (cloudflare#2361) - onBuffered / onRequestBuffered when the request frame was buffered, and a close while the body is still being prepared no longer drops the request (cloudflare#1983, upstream cloudflare#2039) - the chunk stream keeps pulling past events that render nothing, which stalled a waiting reader on RUN_STARTED followed by STEP_STARTED The two stack tests that pinned a clean close mid-stream now expect the error chunk, as upstream changed its own in cloudflare#2019.
Fixes #2013. A WebSocket close before the terminal
done: trueframe was ending the chat stream cleanly, so the AI SDK'suseChatpresented a truncated answer as a completed turn.WebSocketChatTransportnow ends thesendMessages, resumed, and tool-continuation streams with an error chunk ("WebSocket closed mid-stream") followed by a normal close.useAgentChat/useChatenter theerrorstate and callonError; when the socket reconnects and the stream resumes, the error clears as before.An error chunk is used rather than
controller.error()because erroring the stream discards chunks the consumer has not read yet, which would drop content already received (including a replayed resume burst). Direct readers of the stream therefore see anerrorchunk and thendone, not a rejected read.Unchanged: a close after
done, and a close on a resumed stream before the resume handshake (the hook re-probes on the next open), still end cleanly.Tests: transport tests in ai-chat,
useAgentChatReact tests, and the #1913 replay-burst test now expects the interrupted-turn error while still keeping all replayed content.