diff --git a/sdk/README.md b/sdk/README.md index af84aac..5d1d4f7 100644 --- a/sdk/README.md +++ b/sdk/README.md @@ -700,7 +700,7 @@ type WireData = Record; ## Known Limitations - **Explicit `connect()` required** (ACP & Claude): Both `ACPAxonConnection` and `ClaudeAxonConnection` require an explicit `await conn.connect()` call before `initialize()`. The constructor is lightweight and synchronous. -- **Automatic reconnection (single retry)**: If an SSE stream drops unexpectedly, the SDK re-subscribes once and logs a `console.warn`. If the retry also fails, the connection is terminal — create a new instance. +- **Automatic reconnection**: If a managed ACP or Claude connection's SSE stream drops unexpectedly, the SDK keeps re-subscribing from the last seen Axon sequence until the connection is intentionally aborted/disconnected or a fatal broker error is received. - **Permission handling** (Claude): The `ClaudeAxonConnection` auto-approves all tool use by default. Register a `"can_use_tool"` handler via `onControlRequest()` to customize. ### ACP: `prompt()` resolves before all session updates arrive diff --git a/sdk/src/acp/axon-stream.test.ts b/sdk/src/acp/axon-stream.test.ts index b2ee3b0..f1c51f7 100644 --- a/sdk/src/acp/axon-stream.test.ts +++ b/sdk/src/acp/axon-stream.test.ts @@ -568,10 +568,12 @@ describe("axonStream", () => { }); describe("auto-reconnect", () => { - it("re-subscribes once when the SSE stream ends unexpectedly", async () => { + it("keeps re-subscribing when the SSE stream ends unexpectedly", async () => { + vi.useFakeTimers(); const ctrl1 = createControllableStream(); const ctrl2 = createControllableStream(); const published: PublishCall[] = []; + const abortController = new AbortController(); const axon = { id: "test-axon", subscribeSse: vi @@ -585,30 +587,46 @@ describe("axonStream", () => { const errorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); - const { readable } = axonStream({ axon: axon as never }); - - ctrl1.push(makeAgentEvent("session/update", { msg: "first" })); - ctrl1.end(); - - ctrl2.push(makeAgentEvent("session/update", { msg: "second" })); - ctrl2.end(); - - const messages = await drain(readable); - expect(messages).toHaveLength(2); - expect(messages[0]).toMatchObject({ params: { msg: "first" } }); - expect(messages[1]).toMatchObject({ params: { msg: "second" } }); - expect(axon.subscribeSse).toHaveBeenCalledTimes(2); - expect(errorSpy).toHaveBeenCalledWith( - "[axonStream]", - expect.stringContaining("SSE stream ended"), - ); - - errorSpy.mockRestore(); + try { + const { readable } = axonStream({ + axon: axon as never, + signal: abortController.signal, + }); + const reader = readable.getReader(); + + ctrl1.push(makeAgentEvent("session/update", { msg: "first" })); + const first = await reader.read(); + expect(first.value).toMatchObject({ params: { msg: "first" } }); + + ctrl1.end(); + await vi.advanceTimersByTimeAsync(1_000); + await vi.waitFor(() => expect(axon.subscribeSse).toHaveBeenCalledTimes(2)); + + ctrl2.push(makeAgentEvent("session/update", { msg: "second" })); + const second = await reader.read(); + expect(second.value).toMatchObject({ params: { msg: "second" } }); + + abortController.abort(); + ctrl2.end(); + const done = await reader.read(); + expect(done.done).toBe(true); + reader.releaseLock(); + + expect(errorSpy).toHaveBeenCalledWith( + "[axonStream]", + expect.stringContaining("SSE stream ended"), + ); + } finally { + errorSpy.mockRestore(); + vi.useRealTimers(); + } }); it("passes after_sequence on re-subscribe using last seen sequence", async () => { + vi.useFakeTimers(); const ctrl1 = createControllableStream(); const ctrl2 = createControllableStream(); + const abortController = new AbortController(); const axon = { id: "test-axon", subscribeSse: vi @@ -620,26 +638,43 @@ describe("axonStream", () => { const errorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); - const { readable } = axonStream({ axon: axon as never }); + try { + const { readable } = axonStream({ + axon: axon as never, + signal: abortController.signal, + }); + const reader = readable.getReader(); - ctrl1.push(makeAgentEvent("session/update", { msg: "first" }, 10)); - ctrl1.push(makeAgentEvent("session/update", { msg: "second" }, 15)); - ctrl1.end(); + ctrl1.push(makeAgentEvent("session/update", { msg: "first" }, 10)); + ctrl1.push(makeAgentEvent("session/update", { msg: "second" }, 15)); + expect((await reader.read()).value).toMatchObject({ params: { msg: "first" } }); + expect((await reader.read()).value).toMatchObject({ params: { msg: "second" } }); - ctrl2.push(makeAgentEvent("session/update", { msg: "third" }, 16)); - ctrl2.end(); + ctrl1.end(); + await vi.advanceTimersByTimeAsync(1_000); + await vi.waitFor(() => expect(axon.subscribeSse).toHaveBeenCalledTimes(2)); - const messages = await drain(readable); - expect(messages).toHaveLength(3); - expect(axon.subscribeSse).toHaveBeenCalledTimes(2); - expect(axon.subscribeSse).toHaveBeenNthCalledWith(1, undefined); - expect(axon.subscribeSse).toHaveBeenNthCalledWith(2, { after_sequence: 15 }); + ctrl2.push(makeAgentEvent("session/update", { msg: "third" }, 16)); + expect((await reader.read()).value).toMatchObject({ params: { msg: "third" } }); + + abortController.abort(); + ctrl2.end(); + expect((await reader.read()).done).toBe(true); + reader.releaseLock(); - errorSpy.mockRestore(); + expect(axon.subscribeSse).toHaveBeenCalledTimes(2); + expect(axon.subscribeSse).toHaveBeenNthCalledWith(1, undefined); + expect(axon.subscribeSse).toHaveBeenNthCalledWith(2, { after_sequence: 15 }); + } finally { + errorSpy.mockRestore(); + vi.useRealTimers(); + } }); - it("re-subscribes once on SSE stream error and continues", async () => { + it("re-subscribes on SSE stream error and continues", async () => { + vi.useFakeTimers(); const ctrl2 = createControllableStream(); + const abortController = new AbortController(); let callCount = 0; const axon = { id: "test-axon", @@ -663,24 +698,38 @@ describe("axonStream", () => { const errorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); - const { readable } = axonStream({ axon: axon as never }); + try { + const { readable } = axonStream({ + axon: axon as never, + signal: abortController.signal, + }); + const reader = readable.getReader(); - ctrl2.push(makeAgentEvent("session/update", { msg: "recovered" })); - ctrl2.end(); + await vi.advanceTimersByTimeAsync(1_000); + await vi.waitFor(() => expect(axon.subscribeSse).toHaveBeenCalledTimes(2)); - const messages = await drain(readable); - expect(messages).toHaveLength(1); - expect(messages[0]).toMatchObject({ params: { msg: "recovered" } }); - expect(axon.subscribeSse).toHaveBeenCalledTimes(2); - expect(errorSpy).toHaveBeenCalledWith( - "[axonStream]", - expect.stringContaining("SSE stream error"), - ); + ctrl2.push(makeAgentEvent("session/update", { msg: "recovered" })); + const message = await reader.read(); + expect(message.value).toMatchObject({ params: { msg: "recovered" } }); + + abortController.abort(); + ctrl2.end(); + expect((await reader.read()).done).toBe(true); + reader.releaseLock(); - errorSpy.mockRestore(); + expect(errorSpy).toHaveBeenCalledWith( + "[axonStream]", + expect.stringContaining("SSE stream error"), + ); + } finally { + errorSpy.mockRestore(); + vi.useRealTimers(); + } }); - it("closes the stream if the second subscription also fails", async () => { + it("keeps retrying repeated SSE failures until aborted", async () => { + vi.useFakeTimers(); + const abortController = new AbortController(); const axon = { id: "test-axon", subscribeSse: vi.fn().mockImplementation(async () => ({ @@ -697,13 +746,26 @@ describe("axonStream", () => { const errorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); - const { readable } = axonStream({ axon: axon as never }); + try { + const { readable } = axonStream({ + axon: axon as never, + signal: abortController.signal, + }); + const reader = readable.getReader(); - const reader = readable.getReader(); - await expect(reader.read()).rejects.toThrow("permanent failure"); - expect(axon.subscribeSse).toHaveBeenCalledTimes(2); + await vi.advanceTimersByTimeAsync(1_000); + await vi.waitFor(() => expect(axon.subscribeSse).toHaveBeenCalledTimes(2)); + await vi.advanceTimersByTimeAsync(2_000); + await vi.waitFor(() => expect(axon.subscribeSse).toHaveBeenCalledTimes(3)); - errorSpy.mockRestore(); + abortController.abort(); + await vi.advanceTimersByTimeAsync(4_000); + expect((await reader.read()).done).toBe(true); + reader.releaseLock(); + } finally { + errorSpy.mockRestore(); + vi.useRealTimers(); + } }); it("does NOT re-subscribe when signal is aborted", async () => { diff --git a/sdk/src/acp/axon-stream.ts b/sdk/src/acp/axon-stream.ts index 5f0be00..10cf771 100644 --- a/sdk/src/acp/axon-stream.ts +++ b/sdk/src/acp/axon-stream.ts @@ -10,6 +10,29 @@ import { isTurnFailedAxonEvent, tryParseSystemEvent } from "../shared/timeline.j import type { LogFn } from "../shared/types.js"; import type { AxonStreamOptions } from "./types.js"; +const INITIAL_RECONNECT_BACKOFF_MS = 1_000; +const MAX_RECONNECT_BACKOFF_MS = 30_000; + +function delayWithAbort(ms: number, signal: AbortSignal | undefined): Promise { + return new Promise((resolve) => { + if (signal?.aborted) { + resolve(); + return; + } + + const timer = setTimeout(() => { + signal?.removeEventListener("abort", onAbort); + resolve(); + }, ms); + const onAbort = () => { + clearTimeout(timer); + signal?.removeEventListener("abort", onAbort); + resolve(); + }; + signal?.addEventListener("abort", onAbort, { once: true }); + }); +} + /** * Set of event_types that are notifications (no request ID correlation). * `session/update` is the main one -- the broker sends streaming chunks, @@ -112,10 +135,13 @@ function createReadable( initialAfterSequence?: number, replayTargetSequence?: number, ): ReadableStream { + let cancelled = false; + return new ReadableStream({ async start(controller) { let totalEvents = 0; let attempt = 0; + let backoffMs = INITIAL_RECONNECT_BACKOFF_MS; let lastSequence: number | undefined = initialAfterSequence; const replaying = replayTargetSequence != null; @@ -124,7 +150,7 @@ function createReadable( // response is seen, the entry is deleted (resolved). const replayBuffer = new Map(); - while (!signal?.aborted) { + while (!signal?.aborted && !cancelled) { attempt++; let eventCount = 0; try { @@ -138,6 +164,7 @@ function createReadable( eventCount++; totalEvents++; lastSequence = axonEvent.sequence; + backoffMs = INITIAL_RECONNECT_BACKOFF_MS; onAxonEvent?.(axonEvent); @@ -240,25 +267,25 @@ function createReadable( if (msg) controller.enqueue(msg); } } catch (err) { - if (signal?.aborted) break; - if (attempt === 1) { - onError( - `[axonStream] SSE stream error after ${eventCount} events, re-subscribing: ${err}`, - ); - continue; + if (signal?.aborted || cancelled) break; + if (!signal) { + controller.error(err); + return; } - log?.("read", `error on reconnect attempt after ${eventCount} events: ${err}`); - controller.error(err); - return; + onError( + `[axonStream] SSE stream error after ${eventCount} events, re-subscribing: ${err}`, + ); + await delayWithAbort(backoffMs, signal); + backoffMs = Math.min(backoffMs * 2, MAX_RECONNECT_BACKOFF_MS); + continue; } - if (signal?.aborted) break; + if (signal?.aborted || cancelled) break; + if (!signal) break; - if (attempt === 1) { - onError(`[axonStream] SSE stream ended after ${eventCount} events, re-subscribing`); - continue; - } - break; + onError(`[axonStream] SSE stream ended after ${eventCount} events, re-subscribing`); + await delayWithAbort(backoffMs, signal); + backoffMs = Math.min(backoffMs * 2, MAX_RECONNECT_BACKOFF_MS); } // If replay ended because the stream closed before reaching the target, @@ -272,6 +299,9 @@ function createReadable( log?.("read", `SSE ended after ${totalEvents} total events`); controller.close(); }, + cancel() { + cancelled = true; + }, }); } diff --git a/sdk/src/claude/connection.test.ts b/sdk/src/claude/connection.test.ts index 624614f..9fdba5d 100644 --- a/sdk/src/claude/connection.test.ts +++ b/sdk/src/claude/connection.test.ts @@ -1072,6 +1072,20 @@ describe("ClaudeAxonConnection", () => { expect(initCalls.length).toBe(1); }); + it("keeps re-subscribing after repeated unexpected stream ends", async () => { + await createConnectedClient(transport); + + transport._end(); + await vi.waitFor(() => { + expect(transport.reconnect).toHaveBeenCalledTimes(1); + }); + + transport._end(); + await vi.waitFor(() => { + expect(transport.reconnect).toHaveBeenCalledTimes(2); + }); + }); + it("does not reconnect when disconnect() was called", async () => { const conn = await createConnectedClient(transport); diff --git a/sdk/src/claude/connection.ts b/sdk/src/claude/connection.ts index b461bb2..91f17fd 100644 --- a/sdk/src/claude/connection.ts +++ b/sdk/src/claude/connection.ts @@ -536,8 +536,6 @@ export class ClaudeAxonConnection { return; } - let reconnected = false; - const consumeStream = async (): Promise<"ended" | "error"> => { try { for await (const message of transport.readMessages()) { @@ -560,23 +558,26 @@ export class ClaudeAxonConnection { } }; - const outcome = await consumeStream(); + while (!this.closed && !this.suppressTransportAutoReconnect && !this.streamAborted) { + const outcome = await consumeStream(); + if ( + this.closed || + this.suppressTransportAutoReconnect || + this.streamAborted || + this.fatal + ) { + break; + } - if ( - !this.closed && - !this.suppressTransportAutoReconnect && - !this.streamAborted && - !reconnected - ) { const label = outcome === "error" ? "error" : "ended unexpectedly"; this.log("readLoop", `SSE stream ${label}, reconnecting...`); - reconnected = true; try { await transport.reconnect(); this.log("readLoop", "reconnected successfully"); - await consumeStream(); } catch (reconnectErr) { this.log("readLoop", `reconnect failed: ${reconnectErr}`); + this.handleError(reconnectErr); + break; } }