Skip to content
Draft
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
2 changes: 1 addition & 1 deletion sdk/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -700,7 +700,7 @@ type WireData = Record<string, any>;
## 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
Expand Down
164 changes: 113 additions & 51 deletions sdk/src/acp/axon-stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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",
Expand All @@ -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 () => ({
Expand All @@ -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 () => {
Expand Down
62 changes: 46 additions & 16 deletions sdk/src/acp/axon-stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
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,
Expand Down Expand Up @@ -112,10 +135,13 @@ function createReadable(
initialAfterSequence?: number,
replayTargetSequence?: number,
): ReadableStream<AnyMessage> {
let cancelled = false;

return new ReadableStream<AnyMessage>({
async start(controller) {
let totalEvents = 0;
let attempt = 0;
let backoffMs = INITIAL_RECONNECT_BACKOFF_MS;
let lastSequence: number | undefined = initialAfterSequence;

const replaying = replayTargetSequence != null;
Expand All @@ -124,7 +150,7 @@ function createReadable(
// response is seen, the entry is deleted (resolved).
const replayBuffer = new Map<string, AnyMessage>();

while (!signal?.aborted) {
while (!signal?.aborted && !cancelled) {
attempt++;
let eventCount = 0;
try {
Expand All @@ -138,6 +164,7 @@ function createReadable(
eventCount++;
totalEvents++;
lastSequence = axonEvent.sequence;
backoffMs = INITIAL_RECONNECT_BACKOFF_MS;

onAxonEvent?.(axonEvent);

Expand Down Expand Up @@ -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,
Expand All @@ -272,6 +299,9 @@ function createReadable(
log?.("read", `SSE ended after ${totalEvents} total events`);
controller.close();
},
cancel() {
cancelled = true;
},
});
}

Expand Down
14 changes: 14 additions & 0 deletions sdk/src/claude/connection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
Loading
Loading