Skip to content
Merged
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
18 changes: 18 additions & 0 deletions .changeset/sdk-reconnect-on-clean-close.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
---
"@langchain/langgraph-sdk": patch
---

fix(sdk): recover from a server-side thread stream drop instead of freezing

The protocol SSE transport now reconnects when the server closes the event
stream cleanly. The thread stream is open-ended, so a clean close only happens
when the server's own upstream consumer died or it is restarting; before, the
client treated it as the end of the thread and the UI froze mid-run with no
error. A connection that delivered events also resets the reconnect budget, so
long-lived pages survive repeated deploys. `maxReconnectAttempts: 0` keeps the
old end-on-close behavior.

Unsolicited server error frames (no command id) and a shared stream that gives
up reconnecting now reach `stream.error`: `ThreadStream.onError` exposes them,
`useStream` sets `error`, clears `isLoading`, and settles the in-flight
`submit()` as failed.
28 changes: 28 additions & 0 deletions libs/sdk/src/client/stream/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -625,6 +625,7 @@ export class ThreadStream<
// to consume discovery and interrupt events without opening extra
// server subscriptions.
readonly #onEventListeners = new Set<(event: Event) => void>();
readonly #onErrorListeners = new Set<(error: Error) => void>();

#messagesIterable?: AsyncIterable<StreamingMessageHandle>;
#valuesProjection?: AsyncIterable<unknown> & PromiseLike<unknown>;
Expand Down Expand Up @@ -1410,6 +1411,18 @@ export class ThreadStream<
};
}

/**
* Register a listener for stream-level failures: an unsolicited server
* error frame, or the transport giving up on reconnecting. Returns an
* unsubscribe function.
*/
onError(listener: (error: Error) => void): () => void {
this.#onErrorListeners.add(listener);
return () => {
this.#onErrorListeners.delete(listener);
};
}

/**
* Lazily open the wildcard discovery watcher stream.
*
Expand Down Expand Up @@ -1620,6 +1633,16 @@ export class ThreadStream<
}
}

#fireOnError(error: Error): void {
for (const listener of this.#onErrorListeners) {
try {
listener(error);
} catch {
// A throwing listener must not block delivery to the others.
}
}
}

async close(): Promise<void> {
if (this.#closed) {
return;
Expand Down Expand Up @@ -1657,6 +1680,7 @@ export class ThreadStream<
const lifecycleWatcherStartPromise = this.#lifecycleWatcherStartPromise;
this.#lifecycleWatcherStartPromise = undefined;
this.#onEventListeners.clear();
this.#onErrorListeners.clear();
for (const subscription of this.#subscriptions.values()) {
subscription.close();
}
Expand Down Expand Up @@ -2094,6 +2118,7 @@ export class ThreadStream<
for (const subscription of this.#subscriptions.values()) {
subscription.close();
}
this.#fireOnError(normalized);
}

/**
Expand Down Expand Up @@ -2412,6 +2437,9 @@ export class ThreadStream<
const pending =
messageId === undefined ? undefined : this.#pending.get(messageId);
if (!pending) {
if (message.type === "error" && message.id == null) {
this.#fireOnError(new ProtocolError(message));
}
return;
}
if (messageId !== undefined) {
Expand Down
39 changes: 39 additions & 0 deletions libs/sdk/src/client/stream/shared-stream.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { describe, expect, it } from "vitest";

import { ProtocolError } from "./error.js";
import { ThreadStream } from "./index.js";
import { MockSseTransport, eventOf, nextValue } from "./test/utils.js";

Expand Down Expand Up @@ -460,6 +461,44 @@ describe("ThreadStream (SSE shared stream)", () => {
const narrowed = await nextValue(narrow);
expect(narrowed).toMatchObject({ event_id: "evt_inside" });

await thread.close();
});
it("reports an unsolicited server error frame to onError listeners", async () => {
const transport = new MockSseTransport();
const thread = new ThreadStream(transport, { assistantId: "a" });
const errors: Error[] = [];
thread.onError((err) => errors.push(err));
await thread.subscribe({ channels: ["values"] });

transport.pushMessage({
type: "error",
id: null,
error: "unknown_error",
message: "Thread event stream failed: redis gone",
});
await flush();

expect(errors).toHaveLength(1);
expect(errors[0]).toBeInstanceOf(ProtocolError);
expect((errors[0] as ProtocolError).code).toBe("unknown_error");
expect(errors[0].message).toBe("Thread event stream failed: redis gone");
expect(transport.activeStreamCount).toBe(1);

await thread.close();
});

it("reports a failed shared stream to onError listeners", async () => {
const transport = new MockSseTransport();
const thread = new ThreadStream(transport, { assistantId: "a" });
const errors: Error[] = [];
thread.onError((err) => errors.push(err));
await thread.subscribe({ channels: ["values"] });

transport.failStream(0, new Error("gave up reconnecting"));
await flush();

expect(errors.map((e) => e.message)).toEqual(["gave up reconnecting"]);

await thread.close();
});
});
29 changes: 24 additions & 5 deletions libs/sdk/src/client/stream/test/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,7 +139,7 @@ export class MockSseTransport implements TransportAdapter {
private readonly buffer: Event[] = [];
private readonly streams = new Set<{
record: MockSseStreamRecord;
push: (event: Event) => void;
push: (message: Message) => void;
pushError: (err: unknown) => void;
}>();

Expand Down Expand Up @@ -261,7 +261,7 @@ export class MockSseTransport implements TransportAdapter {

const stream = {
record,
push: (event: Event) => {
push: (event: Message) => {
if (record.closed) return;
const waiter = waiters.shift();
if (waiter) {
Expand Down Expand Up @@ -309,9 +309,17 @@ export class MockSseTransport implements TransportAdapter {
if (record.closed) {
return { done: true, value: undefined };
}
return await new Promise<IteratorResult<Message>>((resolve) => {
waiters.push(resolve);
});
const result = await new Promise<IteratorResult<Message>>(
(resolve) => {
waiters.push(resolve);
}
);
if (rejectedWith !== undefined) {
const err = rejectedWith;
rejectedWith = undefined;
throw err;
}
return result;
},
return: async () => {
close();
Expand All @@ -337,6 +345,17 @@ export class MockSseTransport implements TransportAdapter {
}
}

/**
* Deliver a non-event protocol message (e.g. an unsolicited error frame)
* to every open stream, bypassing the subscription filter.
*/
pushMessage(message: Message): void {
for (const stream of this.streams) {
if (stream.record.closed) continue;
stream.push(message);
}
}

async close(): Promise<void> {
this.closed = true;
for (const stream of Array.from(this.streams)) {
Expand Down
112 changes: 112 additions & 0 deletions libs/sdk/src/client/stream/transport/http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -597,6 +597,118 @@ describe("ProtocolSseTransportAdapter SSE reconnect with custom fetch", () => {
await transport.close();
});

function cleanCloseFetch(options: { keepOpenFrom: number }) {
let streamOpens = 0;
const encoder = new TextEncoder();
const fetchImpl = vi.fn((input: URL | RequestInfo) => {
if (!String(input).includes("/stream/events")) {
return Promise.resolve(protocolSuccessResponse());
}
streamOpens += 1;
const seq = streamOpens;
return Promise.resolve(
new Response(
new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(
encoder.encode(
`event: values\ndata: {"type":"event","method":"values","seq":${seq},"event_id":"e${seq}"}\n\n`
)
);
if (seq < options.keepOpenFrom) controller.close();
},
}),
{ status: 200, headers: { "content-type": "text/event-stream" } }
)
);
}) as MockFetch;
return { fetchImpl, streamOpens: () => streamOpens };
}

async function collectEventIds(
handle: ReturnType<ProtocolSseTransportAdapter["openEventStream"]>,
until: string
): Promise<string[]> {
const received: string[] = [];
for await (const message of handle.events) {
received.push((message as { event_id: string }).event_id);
if (received.includes(until)) break;
}
return received;
}

it("reconnects when the server closes the stream cleanly", async () => {
const onReconnect = vi.fn();
const { fetchImpl, streamOpens } = cleanCloseFetch({ keepOpenFrom: 2 });

const transport = new ProtocolSseTransportAdapter({
apiUrl: "http://localhost:8123",
threadId: THREAD_ID,
fetch: fetchImpl,
maxReconnectAttempts: 3,
reconnectDelayMs: () => 0,
onReconnect,
idleReconnect: 0,
});

const handle = transport.openEventStream({ channels: ["values"] });
await handle.ready;

expect(await collectEventIds(handle, "e2")).toEqual(["e1", "e2"]);
expect(onReconnect).toHaveBeenCalledTimes(1);
expect(streamOpens()).toBe(2);

await transport.close();
});

it("resets the reconnect budget once a connection has delivered events", async () => {
const onReconnect = vi.fn();
const { fetchImpl, streamOpens } = cleanCloseFetch({ keepOpenFrom: 3 });

const transport = new ProtocolSseTransportAdapter({
apiUrl: "http://localhost:8123",
threadId: THREAD_ID,
fetch: fetchImpl,
maxReconnectAttempts: 1,
reconnectDelayMs: () => 0,
onReconnect,
idleReconnect: 0,
});

const handle = transport.openEventStream({ channels: ["values"] });
await handle.ready;

expect(await collectEventIds(handle, "e3")).toEqual(["e1", "e2", "e3"]);
expect(onReconnect.mock.calls.map((c) => c[0].attempt)).toEqual([1, 1]);
expect(streamOpens()).toBe(3);

await transport.close();
});

it("ends the stream on a clean close when maxReconnectAttempts is 0", async () => {
const { fetchImpl, streamOpens } = cleanCloseFetch({ keepOpenFrom: 2 });

const transport = new ProtocolSseTransportAdapter({
apiUrl: "http://localhost:8123",
threadId: THREAD_ID,
fetch: fetchImpl,
maxReconnectAttempts: 0,
idleReconnect: 0,
});

const handle = transport.openEventStream({ channels: ["values"] });
await handle.ready;

const received: string[] = [];
for await (const message of handle.events) {
received.push((message as { event_id: string }).event_id);
}
expect(received).toEqual(["e1"]);
expect(streamOpens()).toBe(1);

await transport.close();
});

it("keeps reconnect disabled when maxReconnectAttempts is explicitly 0", async () => {
const sentinel = new TypeError("net::ERR_QUIC_PROTOCOL_ERROR");
let streamOpens = 0;
Expand Down
20 changes: 18 additions & 2 deletions libs/sdk/src/client/stream/transport/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -268,6 +268,7 @@ export class ProtocolSseTransportAdapter implements TransportAdapter {

const startStream = async () => {
let attempt = 0;
let receivedEvent = false;

while (!ac.signal.aborted && !this.closed) {
try {
Expand Down Expand Up @@ -336,12 +337,27 @@ export class ProtocolSseTransportAdapter implements TransportAdapter {
break;
}
if (isRecord(event.data)) {
receivedEvent = true;
streamQueue.push(event.data as Message);
}
}
streamQueue.close();
return;
if (
ac.signal.aborted ||
this.closed ||
this.maxReconnectAttempts <= 0
) {
streamQueue.close();
return;
}
// The thread stream is open-ended: the server only ends it when
// its own upstream consumer died or it is shutting down, so a
// clean close is a disconnect, not the end of the thread.
throw new Error("Event stream closed by the server");
} catch (error) {
if (receivedEvent) {
attempt = 0;
receivedEvent = false;
}
if (ac.signal.aborted || this.closed) {
if (!readySettled) {
rejectReady(error);
Expand Down
Loading
Loading