mirror of
https://github.com/openclaw/openclaw.git
synced 2026-07-21 10:16:44 +00:00
fix(agents): preserve sanitized stream cancellation (#110427)
* fix(agents): add .catch() to reader.cancel() to prevent unhandled rejection * fix(agents): preserve sanitized stream cancellation Co-authored-by: zenglingbiao <zeng.lingbiao@xydigit.com> --------- Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
co-authored by
Peter Steinberger
parent
74919102a8
commit
490ab265ba
@@ -1323,6 +1323,58 @@ describe("buildGuardedModelFetch", () => {
|
||||
expect(items).toEqual([{ ok: true }]);
|
||||
});
|
||||
|
||||
it.each([
|
||||
{
|
||||
name: "JSON-to-SSE synthesis",
|
||||
contentType: "application/json",
|
||||
body: '{"ok": true}',
|
||||
},
|
||||
{
|
||||
name: "SSE sanitization",
|
||||
contentType: "text/event-stream",
|
||||
body: 'data: {"ok": true}\n\n',
|
||||
},
|
||||
])("ignores source cancellation failures during $name", async ({ contentType, body }) => {
|
||||
const cancel = vi.fn(async () => {
|
||||
throw new Error("upstream cancellation failed");
|
||||
});
|
||||
const release = vi.fn(async () => undefined);
|
||||
const encoder = new TextEncoder();
|
||||
fetchWithSsrFGuardMock.mockResolvedValue({
|
||||
response: new Response(
|
||||
new ReadableStream({
|
||||
start(controller) {
|
||||
controller.enqueue(encoder.encode(body));
|
||||
},
|
||||
cancel,
|
||||
}),
|
||||
{ headers: { "content-type": contentType } },
|
||||
),
|
||||
finalUrl: "https://openrouter.ai/api/v1/chat/completions",
|
||||
release,
|
||||
});
|
||||
const model = {
|
||||
id: "gpt-5.4",
|
||||
provider: "openrouter",
|
||||
api: "openai-completions",
|
||||
baseUrl: "https://openrouter.ai/api/v1",
|
||||
} as unknown as Model<"openai-completions">;
|
||||
|
||||
const response = await buildGuardedModelFetch(model)(
|
||||
"https://openrouter.ai/api/v1/chat/completions",
|
||||
{
|
||||
method: "POST",
|
||||
headers: { "content-type": "application/json" },
|
||||
body: JSON.stringify({ model: "gpt-5.4", stream: true }),
|
||||
},
|
||||
);
|
||||
|
||||
expect(response.body).not.toBeNull();
|
||||
await expect(response.body!.cancel("consumer stopped")).resolves.toBeUndefined();
|
||||
expect(cancel).toHaveBeenCalledTimes(1);
|
||||
expect(release).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("does not re-prefix SSE bodies mislabeled as JSON by streaming gateways", async () => {
|
||||
const source = openResponseStreamText(
|
||||
'data: {"id":"a","choices":[{"index":0,"delta":{"content":"Hi","role":"assistant"}}]}\n\n' +
|
||||
|
||||
@@ -96,6 +96,15 @@ function findSseEventBoundary(buffer: string): { index: number; length: number }
|
||||
return best;
|
||||
}
|
||||
|
||||
async function cancelReaderBestEffort(
|
||||
reader: ReadableStreamDefaultReader<Uint8Array> | undefined,
|
||||
reason?: unknown,
|
||||
): Promise<void> {
|
||||
// Reader cancellation is cleanup. An upstream cancel failure must not replace
|
||||
// the wrapper's authoritative stream error or downstream cancellation.
|
||||
await reader?.cancel(reason).catch(() => undefined);
|
||||
}
|
||||
|
||||
function capNonOkResponseBodyLazily(response: Response, maxBytes: number): Response {
|
||||
const source = response.body;
|
||||
if (!source) {
|
||||
@@ -123,18 +132,18 @@ function capNonOkResponseBodyLazily(response: Response, maxBytes: number): Respo
|
||||
}
|
||||
total = maxBytes;
|
||||
controller.close();
|
||||
void reader?.cancel().catch(() => undefined);
|
||||
void cancelReaderBestEffort(reader);
|
||||
return;
|
||||
}
|
||||
total += chunk.value.byteLength;
|
||||
controller.enqueue(chunk.value);
|
||||
} catch (error) {
|
||||
controller.error(error);
|
||||
void reader?.cancel(error).catch(() => undefined);
|
||||
void cancelReaderBestEffort(reader, error);
|
||||
}
|
||||
},
|
||||
async cancel(reason) {
|
||||
await reader?.cancel(reason).catch(() => undefined);
|
||||
await cancelReaderBestEffort(reader, reason);
|
||||
},
|
||||
});
|
||||
return new Response(capped, response);
|
||||
@@ -189,12 +198,12 @@ function sanitizeOpenAISdkSseResponse(
|
||||
buffer += decoder.decode(chunk.value, { stream: true });
|
||||
}
|
||||
} catch (error) {
|
||||
await reader?.cancel(error).catch(() => {});
|
||||
await cancelReaderBestEffort(reader, error);
|
||||
controller.error(error);
|
||||
}
|
||||
},
|
||||
async cancel(reason) {
|
||||
await reader?.cancel(reason);
|
||||
await cancelReaderBestEffort(reader, reason);
|
||||
},
|
||||
});
|
||||
const headers = new Headers(response.headers);
|
||||
@@ -277,12 +286,12 @@ function sanitizeOpenAISdkSseResponse(
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
await reader?.cancel(error).catch(() => {});
|
||||
await cancelReaderBestEffort(reader, error);
|
||||
controller.error(error);
|
||||
}
|
||||
},
|
||||
async cancel(reason) {
|
||||
await reader?.cancel(reason);
|
||||
await cancelReaderBestEffort(reader, reason);
|
||||
},
|
||||
});
|
||||
|
||||
@@ -361,7 +370,7 @@ async function classifyOpenAISdkStreamBody(response: Response): Promise<OpenAISd
|
||||
text += decoder.decode();
|
||||
return classifyOpenAISdkStreamBodyPrefix(text);
|
||||
} finally {
|
||||
void reader.cancel().catch(() => undefined);
|
||||
void cancelReaderBestEffort(reader);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user