fix(gateway): abort retry backoff on disconnect

This commit is contained in:
ShiroKSH
2026-07-29 15:58:35 +03:00
parent 3b99fa239b
commit 546719003c
3 changed files with 87 additions and 998 deletions
+20 -2
View File
@@ -1,3 +1,21 @@
export function delay(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
export const DEFAULT_RETRY_AFTER_MS = 1000;
export function delay(ms: number, signal?: AbortSignal): Promise<void> {
if (signal?.aborted) {
return Promise.resolve();
}
return new Promise((resolve) => {
let timer: ReturnType<typeof setTimeout> | undefined;
const complete = () => {
if (timer !== undefined) {
clearTimeout(timer);
}
signal?.removeEventListener("abort", complete);
resolve();
};
signal?.addEventListener("abort", complete, { once: true });
timer = setTimeout(complete, ms);
});
}
File diff suppressed because one or more lines are too long
@@ -0,0 +1,66 @@
import assert from "node:assert/strict";
import test from "node:test";
import { fetchUpstreamWithFallback } from "@ccr/core/gateway/upstream/executor.ts";
const retryConfig = {
Providers: [],
Router: { fallback: { mode: "retry", models: [], retryCount: 1 }, rules: [] },
virtualModelProfiles: []
};
const retryFallback = { mode: "retry", models: [], retryCount: 1 };
async function assertRetryBackoffStopsAfterAbort(fetchImpl) {
const originalFetch = globalThis.fetch;
const originalSetTimeout = globalThis.setTimeout;
const controller = new AbortController();
let fetchCount = 0;
globalThis.fetch = async (...args) => {
fetchCount += 1;
return fetchImpl(...args);
};
globalThis.setTimeout = (_callback, delay, ..._args) => {
const timer = originalSetTimeout(() => {}, delay);
timer.unref?.();
queueMicrotask(() => controller.abort(new Error("client disconnected")));
return timer;
};
try {
const outcome = await Promise.race([
fetchUpstreamWithFallback({
body: Buffer.from('{"model":"test-model"}'),
config: retryConfig,
coreAuthToken: "core-token",
fallback: retryFallback,
headers: {},
method: "POST",
path: "/v1/messages",
routedModel: "test-model",
signal: controller.signal,
upstreamUrl: "http://127.0.0.1:3456/v1/messages"
}).then(
() => ({ kind: "resolved" }),
(error) => ({ error, kind: "rejected" })
),
new Promise((resolve) => setImmediate(() => resolve({ kind: "pending" })))
]);
assert.notEqual(outcome.kind, "pending");
assert.equal(outcome.kind, "rejected");
assert.match(outcome.error.message, /client disconnected/);
assert.equal(fetchCount, 1);
} finally {
globalThis.fetch = originalFetch;
globalThis.setTimeout = originalSetTimeout;
}
}
test("retry backoff stops after client aborts a retryable HTTP response", async () => {
await assertRetryBackoffStopsAfterAbort(async () => new Response(null, { status: 503 }));
});
test("retry backoff stops after client aborts a network error", async () => {
await assertRetryBackoffStopsAfterAbort(async () => {
throw new Error("upstream unavailable");
});
});