mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-26 04:15:48 -06:00
fix(agents): bound stalled binary provider response reads (#109067)
This commit is contained in:
@@ -436,6 +436,28 @@ describe("provider error utils", () => {
|
||||
).rejects.toThrow("stalled-provider: response body stalled for 20ms");
|
||||
});
|
||||
|
||||
it("bounds stalled binary provider responses with the shared default idle timeout", async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const response = new Response(
|
||||
new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
controller.enqueue(new Uint8Array([1]));
|
||||
},
|
||||
}),
|
||||
{ headers: { "content-type": "audio/mpeg" } },
|
||||
);
|
||||
const assertion = expect(
|
||||
readProviderBinaryResponse(response, "stalled-provider", "audio"),
|
||||
).rejects.toThrow("stalled-provider: response body stalled for 30000ms");
|
||||
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
await assertion;
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it("rejects stalled non-2xx error body read after chunk idle timeout", async () => {
|
||||
vi.useFakeTimers();
|
||||
const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
|
||||
|
||||
@@ -19,9 +19,7 @@ export { asBoolean } from "../utils/boolean.js";
|
||||
export { normalizeOptionalString as trimToUndefined } from "../../packages/normalization-core/src/string-coerce.js";
|
||||
|
||||
const ERROR_BODY_METADATA_LIMIT = 500;
|
||||
const PROVIDER_BINARY_RESPONSE_MAX_BYTES = 16 * 1024 * 1024;
|
||||
const PROVIDER_JSON_RESPONSE_MAX_BYTES = 16 * 1024 * 1024;
|
||||
const PROVIDER_TEXT_RESPONSE_MAX_BYTES = 16 * 1024 * 1024;
|
||||
const PROVIDER_RESPONSE_MAX_BYTES = 16 * 1024 * 1024;
|
||||
const SHORT_BEARER_TOKEN_PATTERN =
|
||||
/\b(Bearer)\s+[-A-Za-z0-9._~+/=]{1,17}(?![-A-Za-z0-9._~+/=…])/giu;
|
||||
|
||||
@@ -92,6 +90,26 @@ type ProviderResponseReadOptions = ReadResponseTextPrefixOptions & {
|
||||
onOverflow?: (params: { size: number; maxBytes: number; res: Response }) => Error;
|
||||
};
|
||||
|
||||
function readProviderResponseBytes(
|
||||
response: Response,
|
||||
label: string,
|
||||
kind: string,
|
||||
opts?: ProviderResponseReadOptions,
|
||||
onOverflow?: ProviderResponseReadOptions["onOverflow"],
|
||||
): Promise<Uint8Array> {
|
||||
return readResponseWithLimit(response, opts?.maxBytes ?? PROVIDER_RESPONSE_MAX_BYTES, {
|
||||
...opts,
|
||||
chunkTimeoutMs: opts?.chunkTimeoutMs ?? 30_000,
|
||||
onIdleTimeout:
|
||||
opts?.onIdleTimeout ??
|
||||
(({ chunkTimeoutMs }) =>
|
||||
new Error(`${label}: response body stalled for ${chunkTimeoutMs}ms`)),
|
||||
onOverflow:
|
||||
onOverflow ??
|
||||
(({ maxBytes: limit }) => new Error(`${label}: ${kind} response exceeds ${limit} bytes`)),
|
||||
});
|
||||
}
|
||||
|
||||
/** Options for bounded provider error-body normalization. */
|
||||
type ProviderHttpErrorOptions = {
|
||||
statusPrefix?: string;
|
||||
@@ -148,18 +166,7 @@ export async function readProviderTextResponse(
|
||||
label: string,
|
||||
opts?: ProviderResponseReadOptions,
|
||||
): Promise<string> {
|
||||
const maxBytes = opts?.maxBytes ?? PROVIDER_TEXT_RESPONSE_MAX_BYTES;
|
||||
const bytes = await readResponseWithLimit(response, maxBytes, {
|
||||
chunkTimeoutMs: opts?.chunkTimeoutMs ?? 30_000,
|
||||
onIdleTimeout:
|
||||
opts?.onIdleTimeout ??
|
||||
(({ chunkTimeoutMs }) =>
|
||||
new Error(`${label}: response body stalled for ${chunkTimeoutMs}ms`)),
|
||||
timeoutMs: opts?.timeoutMs,
|
||||
onTimeout: opts?.onTimeout,
|
||||
onOverflow: ({ maxBytes: maxBytesLocal }) =>
|
||||
new Error(`${label}: text response exceeds ${maxBytesLocal} bytes`),
|
||||
});
|
||||
const bytes = await readProviderResponseBytes(response, label, "text", opts);
|
||||
return new TextDecoder().decode(bytes);
|
||||
}
|
||||
|
||||
@@ -408,18 +415,7 @@ export async function readProviderJsonResponse<T>(
|
||||
label: string,
|
||||
opts?: ProviderResponseReadOptions,
|
||||
): Promise<T> {
|
||||
const maxBytes = opts?.maxBytes ?? PROVIDER_JSON_RESPONSE_MAX_BYTES;
|
||||
const bytes = await readResponseWithLimit(response, maxBytes, {
|
||||
chunkTimeoutMs: opts?.chunkTimeoutMs ?? 30_000,
|
||||
onIdleTimeout:
|
||||
opts?.onIdleTimeout ??
|
||||
(({ chunkTimeoutMs }) =>
|
||||
new Error(`${label}: response body stalled for ${chunkTimeoutMs}ms`)),
|
||||
timeoutMs: opts?.timeoutMs,
|
||||
onTimeout: opts?.onTimeout,
|
||||
onOverflow: ({ maxBytes: maxBytesLocal }) =>
|
||||
new Error(`${label}: JSON response exceeds ${maxBytesLocal} bytes`),
|
||||
});
|
||||
const bytes = await readProviderResponseBytes(response, label, "JSON", opts);
|
||||
try {
|
||||
return JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(bytes)) as T;
|
||||
} catch (cause) {
|
||||
@@ -495,14 +491,7 @@ export async function readProviderBinaryResponse(
|
||||
void response.body?.cancel().catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
const maxBytes = opts?.maxBytes ?? PROVIDER_BINARY_RESPONSE_MAX_BYTES;
|
||||
const bytes = await readResponseWithLimit(response, maxBytes, {
|
||||
...opts,
|
||||
onOverflow:
|
||||
opts?.onOverflow ??
|
||||
(({ maxBytes: maxBytesLocal }) =>
|
||||
new Error(`${label}: ${kind} response exceeds ${maxBytesLocal} bytes`)),
|
||||
});
|
||||
const bytes = await readProviderResponseBytes(response, label, kind, opts, opts?.onOverflow);
|
||||
if (bytes.byteLength === 0) {
|
||||
throw new Error(`${label}: malformed ${kind} response`);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user