Skip to content

Commit 5259f9f

Browse files
committed
fix(core): abort stalled provider streams with header and chunk timeouts
The custom fetch on the OpenAI client only set the keep-alive agent. The SDK's own timeout covers the wait for response headers only, so a stream that went silent after the headers arrived - a gateway hiccup or a black-holed connection - never aborted and the CLI spun until the user interrupted it. timeoutFetch() wraps the client's fetch with two bounds: - A 300s response-header timeout, cleared once the headers arrive. - A 300s per-chunk idle timeout for the response body; every received chunk resets it, so a healthy stream of any length is never cut off. Both timers abort with a descriptive reason ("No response headers after 300s..." / "No data received for 300s while streaming..."), which surfaces as the request's error message. The caller's own signal is combined, and user interrupts keep working.
1 parent c6abd6a commit 5259f9f

2 files changed

Lines changed: 195 additions & 1 deletion

File tree

‎packages/core/src/common/openai-client.ts‎

Lines changed: 112 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,117 @@ export { resolveOpenAIConnection, DEEPCODE_PLUS_BASE_URL } from "./plus-subscrip
1515
// keep connections reusable for three minutes after the last request.
1616
const keepAliveAgent = new Agent({ keepAliveTimeout: 180_000 });
1717

18+
// A stalled provider connection used to hang a request forever: the OpenAI
19+
// SDK's own timeout only covers the wait for response headers, so a stream
20+
// that goes silent after the headers arrived (gateway hiccup, black-holed
21+
// connection) never aborts and the CLI spins until the user interrupts.
22+
// These defaults bound both waits; they mirror what other coding agents use.
23+
const STREAM_HEADER_TIMEOUT_MS = 300_000;
24+
const STREAM_CHUNK_TIMEOUT_MS = 300_000;
25+
26+
type TimeoutFetchOptions = {
27+
headerTimeoutMs?: number;
28+
chunkTimeoutMs?: number;
29+
};
30+
31+
// eslint-disable-next-line @typescript-eslint/no-explicit-any
32+
type FetchLike = (input: any, init?: any) => Promise<Response>;
33+
34+
// Applies a response-header timeout and a per-chunk idle timeout at the fetch
35+
// layer. The header timer is cleared once headers arrive; after that every
36+
// received chunk resets the idle timer, so a healthy stream of any length is
37+
// never cut off — only a silent connection is. The caller's own signal (the
38+
// SDK's or an abort request) is combined, and our timers abort with a
39+
// descriptive reason that surfaces as the request's error message.
40+
export function timeoutFetch(fetchImpl: FetchLike, options: TimeoutFetchOptions = {}): FetchLike {
41+
const headerTimeoutMs = options.headerTimeoutMs ?? STREAM_HEADER_TIMEOUT_MS;
42+
const chunkTimeoutMs = options.chunkTimeoutMs ?? STREAM_CHUNK_TIMEOUT_MS;
43+
44+
// eslint-disable-next-line @typescript-eslint/no-explicit-any
45+
return async (input: any, init?: any) => {
46+
const opts = { ...(init ?? {}) };
47+
const headerController = new AbortController();
48+
const chunkController = new AbortController();
49+
const signals = [...(opts.signal ? [opts.signal] : []), headerController.signal, chunkController.signal];
50+
if (signals.length > 1) {
51+
opts.signal = AbortSignal.any(signals);
52+
} else if (signals.length === 1) {
53+
opts.signal = signals[0];
54+
}
55+
56+
const headerTimer = setTimeout(
57+
() =>
58+
headerController.abort(
59+
new Error(`No response headers after ${Math.round(headerTimeoutMs / 1000)}s from the provider.`)
60+
),
61+
headerTimeoutMs
62+
);
63+
64+
let response: Response;
65+
try {
66+
response = await fetchImpl(input, opts);
67+
} finally {
68+
clearTimeout(headerTimer);
69+
}
70+
if (!response.body) {
71+
return response;
72+
}
73+
return wrapBodyWithIdleTimeout(response, response.body, chunkController, chunkTimeoutMs);
74+
};
75+
}
76+
77+
function wrapBodyWithIdleTimeout(
78+
response: Response,
79+
upstream: ReadableStream<Uint8Array>,
80+
controller: AbortController,
81+
chunkTimeoutMs: number
82+
): Response {
83+
let chunkTimer: ReturnType<typeof setTimeout> | undefined;
84+
const clearChunkTimer = () => {
85+
if (chunkTimer) {
86+
clearTimeout(chunkTimer);
87+
chunkTimer = undefined;
88+
}
89+
};
90+
const scheduleChunkAbort = () => {
91+
chunkTimer = setTimeout(
92+
() =>
93+
controller.abort(
94+
new Error(`No data received for ${Math.round(chunkTimeoutMs / 1000)}s while streaming from the provider.`)
95+
),
96+
chunkTimeoutMs
97+
);
98+
};
99+
const upstreamReader = upstream.getReader();
100+
const body = new ReadableStream<Uint8Array>({
101+
async pull(downstream) {
102+
try {
103+
const { done, value } = await upstreamReader.read();
104+
if (done) {
105+
clearChunkTimer();
106+
downstream.close();
107+
return;
108+
}
109+
clearChunkTimer();
110+
scheduleChunkAbort();
111+
downstream.enqueue(value);
112+
} catch (error) {
113+
clearChunkTimer();
114+
downstream.error(error);
115+
}
116+
},
117+
cancel(reason) {
118+
clearChunkTimer();
119+
return upstreamReader.cancel(reason);
120+
},
121+
});
122+
return new Response(body, {
123+
status: response.status,
124+
statusText: response.statusText,
125+
headers: response.headers,
126+
});
127+
}
128+
18129
// Module-level cache for the OpenAI client instance. The client itself is
19130
// a stateless fetch wrapper, so it is safe to share across calls as long as
20131
// the apiKey + baseURL stay the same. Model, thinking-mode and other
@@ -102,7 +213,7 @@ export function createOpenAIClient(
102213
apiKey: connection.apiKey,
103214
baseURL: connection.baseURL || undefined,
104215
// eslint-disable-next-line @typescript-eslint/no-explicit-any
105-
fetch: (url: any, init: any) => undiciFetch(url, { ...init, dispatcher: keepAliveAgent }),
216+
fetch: timeoutFetch((url: any, init: any) => undiciFetch(url, { ...init, dispatcher: keepAliveAgent })),
106217
});
107218
cachedOpenAIKey = cacheKey;
108219

Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,83 @@
1+
import { test } from "node:test";
2+
import assert from "node:assert/strict";
3+
import { timeoutFetch } from "../common/openai-client";
4+
5+
function hangingFetch(onSignal?: (signal: AbortSignal) => void) {
6+
return (_input: unknown, init?: { signal?: AbortSignal }) =>
7+
new Promise<Response>((_resolve, reject) => {
8+
const signal = init?.signal;
9+
if (!signal) return;
10+
onSignal?.(signal);
11+
signal.addEventListener("abort", () => reject(signal.reason));
12+
});
13+
}
14+
15+
function stalledStreamFetch(firstChunk: string, onSignal?: (signal: AbortSignal) => void) {
16+
return (_input: unknown, init?: { signal?: AbortSignal }) => {
17+
const encoder = new TextEncoder();
18+
const signal = init?.signal;
19+
onSignal?.(signal!);
20+
const body = new ReadableStream<Uint8Array>({
21+
start(controller) {
22+
controller.enqueue(encoder.encode(firstChunk));
23+
signal?.addEventListener("abort", () => controller.error(signal.reason));
24+
},
25+
});
26+
return Promise.resolve(new Response(body, { status: 200, headers: { "content-type": "text/event-stream" } }));
27+
};
28+
}
29+
30+
test("timeoutFetch passes a fast response through unchanged", async () => {
31+
const fetchImpl = timeoutFetch(async () => new Response('{"ok":true}', { status: 200 }));
32+
const res = await fetchImpl("https://api.example.com/v1/chat", { method: "POST" });
33+
assert.equal(res.status, 200);
34+
assert.equal(await res.text(), '{"ok":true}');
35+
});
36+
37+
test("timeoutFetch keeps the caller's abort working", async () => {
38+
let observed: AbortSignal | undefined;
39+
const fetchImpl = timeoutFetch(hangingFetch((signal) => (observed = signal)));
40+
const caller = new AbortController();
41+
const pending = fetchImpl("https://api.example.com/v1/chat", { signal: caller.signal });
42+
caller.abort(new Error("user interrupt"));
43+
await assert.rejects(pending, /user interrupt/);
44+
assert.ok(observed?.aborted);
45+
});
46+
47+
test("timeoutFetch aborts when response headers never arrive", async () => {
48+
const fetchImpl = timeoutFetch(hangingFetch(), { headerTimeoutMs: 20 });
49+
await assert.rejects(fetchImpl("https://api.example.com/v1/chat", { method: "POST" }), /No response headers/);
50+
});
51+
52+
test("timeoutFetch delivers chunks that keep arriving", async () => {
53+
const encoder = new TextEncoder();
54+
const fetchImpl = timeoutFetch(
55+
async () => {
56+
let sent = 0;
57+
const body = new ReadableStream<Uint8Array>({
58+
pull(controller) {
59+
if (sent >= 5) {
60+
controller.close();
61+
return;
62+
}
63+
sent += 1;
64+
controller.enqueue(encoder.encode(`chunk-${sent}\n\n`));
65+
},
66+
});
67+
return new Response(body, { status: 200 });
68+
},
69+
{ chunkTimeoutMs: 200 }
70+
);
71+
const res = await fetchImpl("https://api.example.com/v1/chat", { method: "POST" });
72+
const text = await res.text();
73+
assert.equal(text.split("\n\n").filter(Boolean).length, 5);
74+
});
75+
76+
test("timeoutFetch aborts a stream that stalls between chunks", async () => {
77+
const fetchImpl = timeoutFetch(stalledStreamFetch("data: hello\n\n"), { chunkTimeoutMs: 20 });
78+
const res = await fetchImpl("https://api.example.com/v1/chat", { method: "POST" });
79+
const reader = res.body!.getReader();
80+
const first = await reader.read();
81+
assert.equal(new TextDecoder().decode(first.value), "data: hello\n\n");
82+
await assert.rejects(reader.read(), /No data received/);
83+
});

0 commit comments

Comments
 (0)