mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-07 15:52:52 +03:00
refactor(combo): extract buildTargetTimeoutRunner from handleComboChat (#6036)
Bloco J (hot-path decomposition), Task 1. Extract the per-target-timeout dispatch wrapper
(handleComboChat's handleSingleModelWithTimeout closure) verbatim into the leaf
combo/targetTimeoutRunner.ts as a factory buildTargetTimeoutRunner({handleSingleModel,
comboTargetTimeoutMs, log}). The per-model abort still comes from target.modelAbortSignal,
so the outer request signal is intentionally not a dependency. Host call-sites unchanged.
combo.ts shrinks ~60 LOC; leaf is 91 LOC (<800). Body byte-identical (verbatim), no cycle.
This is the first slice toward extracting the shared attempt-loop/success/error handlers
(Tasks 3-4) that de-duplicate handleComboChat and handleRoundRobinCombo. Adds a dedicated
test (5) so the failover path can be mutated independently. Consumer tests stay green
(combo-strategy-fallbacks 24, combo-499-abort 5, empty-content-failover 3, body-400-stop 1,
priority-quota-exhaustion 2, rr-streaming-lock 1, rr-session-stickiness 2).
Plan: _tasks/superpowers/plans/2026-07-03-blocoJ-combo-hotpath-decomposition.md
This commit is contained in:
committed by
GitHub
parent
4e16a491a1
commit
373dd17ccb
@@ -18,6 +18,7 @@ import {
|
||||
selectLockoutCooldownMs,
|
||||
} from "./accountFallback.ts";
|
||||
import { errorResponse, unavailableResponse } from "../utils/error.ts";
|
||||
import { buildTargetTimeoutRunner } from "./combo/targetTimeoutRunner.ts";
|
||||
import {
|
||||
recordComboIntent,
|
||||
recordComboRequest,
|
||||
@@ -45,11 +46,7 @@ import { extractSessionAffinityKey } from "@/sse/services/auth";
|
||||
import { getHiddenModelsByProvider } from "@/models";
|
||||
import { resolveModelLockoutSettings } from "../../src/lib/resilience/modelLockoutSettings";
|
||||
import { fetchCodexQuota } from "./codexQuotaFetcher.ts";
|
||||
import {
|
||||
evaluateQuotaCutoff,
|
||||
getQuotaFetcher,
|
||||
type QuotaInfo,
|
||||
} from "./quotaPreflight.ts";
|
||||
import { evaluateQuotaCutoff, getQuotaFetcher, type QuotaInfo } from "./quotaPreflight.ts";
|
||||
import * as semaphore from "./rateLimitSemaphore.ts";
|
||||
import { getCircuitBreaker } from "../../src/shared/utils/circuitBreaker";
|
||||
import { fisherYatesShuffle, getNextFromDeck } from "../../src/shared/utils/shuffleDeck";
|
||||
@@ -688,72 +685,11 @@ export async function handleComboChat({
|
||||
} = phaseComboSetup(comboCtx);
|
||||
body = comboCtx.body;
|
||||
|
||||
const handleSingleModelWithTimeout = async (
|
||||
b: Record<string, unknown>,
|
||||
modelStr: string,
|
||||
target?: SingleModelTarget
|
||||
): Promise<Response> => {
|
||||
if (comboTargetTimeoutMs <= 0) {
|
||||
return handleSingleModel(b, modelStr, target).catch((err) =>
|
||||
errorResponse(502, err?.message ?? "Upstream model error")
|
||||
);
|
||||
}
|
||||
|
||||
const timeoutController = new AbortController();
|
||||
let timeoutId: ReturnType<typeof setTimeout> | undefined;
|
||||
let timedOut = false;
|
||||
const timeoutPromise = new Promise<Response>((resolve) => {
|
||||
timeoutId = setTimeout(() => {
|
||||
timedOut = true;
|
||||
log.warn(
|
||||
"COMBO",
|
||||
`Model ${modelStr} exceeded ${comboTargetTimeoutMs}ms timeout — falling back`
|
||||
);
|
||||
timeoutController.abort(new Error("combo-per-model-timeout"));
|
||||
resolve(
|
||||
new Response(JSON.stringify({ error: { message: `Model ${modelStr} timed out` } }), {
|
||||
status: 524,
|
||||
headers: { "Content-Type": "application/json" },
|
||||
})
|
||||
);
|
||||
}, comboTargetTimeoutMs);
|
||||
});
|
||||
const targetWithSignal = {
|
||||
...(target ?? {}),
|
||||
modelAbortSignal: timeoutController.signal,
|
||||
};
|
||||
const parentHedgeSignal = target?.modelAbortSignal ?? null;
|
||||
let onParentHedgeAbort: (() => void) | null = null;
|
||||
if (parentHedgeSignal) {
|
||||
if (parentHedgeSignal.aborted) {
|
||||
timeoutController.abort(new Error("hedge-cancelled"));
|
||||
} else {
|
||||
onParentHedgeAbort = () => {
|
||||
timeoutController.abort(new Error("hedge-cancelled"));
|
||||
};
|
||||
parentHedgeSignal.addEventListener("abort", onParentHedgeAbort, { once: true });
|
||||
}
|
||||
}
|
||||
try {
|
||||
return await Promise.race([
|
||||
handleSingleModel(b, modelStr, targetWithSignal).catch((err) => {
|
||||
if (timedOut) {
|
||||
// Inner call rejected because we aborted it. The synthetic 524 from
|
||||
// timeoutPromise already wins the race; return an empty response so
|
||||
// the loser branch resolves cleanly without leaking err.message.
|
||||
return new Response(null, { status: 599 });
|
||||
}
|
||||
return errorResponse(502, err?.message ?? "Upstream model error");
|
||||
}),
|
||||
timeoutPromise,
|
||||
]);
|
||||
} finally {
|
||||
clearTimeout(timeoutId);
|
||||
if (parentHedgeSignal && onParentHedgeAbort) {
|
||||
parentHedgeSignal.removeEventListener("abort", onParentHedgeAbort);
|
||||
}
|
||||
}
|
||||
};
|
||||
const handleSingleModelWithTimeout = buildTargetTimeoutRunner({
|
||||
handleSingleModel,
|
||||
comboTargetTimeoutMs,
|
||||
log,
|
||||
});
|
||||
|
||||
// Route to pinned model if context caching specifies one (Fix #679)
|
||||
if (pinnedModel) {
|
||||
|
||||
91
open-sse/services/combo/targetTimeoutRunner.ts
Normal file
91
open-sse/services/combo/targetTimeoutRunner.ts
Normal file
@@ -0,0 +1,91 @@
|
||||
/**
|
||||
* Wrap a single-model dispatch with a per-target timeout that aborts and falls back.
|
||||
*
|
||||
* Verbatim extraction of handleComboChat's `handleSingleModelWithTimeout` closure
|
||||
* (combo.ts). Behavior is byte-identical; the only change is that the closed-over locals
|
||||
* (`handleSingleModel`, `comboTargetTimeoutMs`, `log`) became explicit factory params.
|
||||
* The per-model abort signal still comes from the target (`target.modelAbortSignal`), so
|
||||
* the outer request signal is intentionally NOT a dependency here.
|
||||
*
|
||||
* See _tasks/superpowers/plans/2026-07-03-blocoJ-combo-hotpath-decomposition.md (Task 1).
|
||||
*/
|
||||
import { errorResponse } from "../../utils/error.ts";
|
||||
import type { HandleSingleModel, SingleModelTarget, ComboLogger } from "./types.ts";
|
||||
|
||||
export function buildTargetTimeoutRunner(deps: {
|
||||
handleSingleModel: HandleSingleModel;
|
||||
comboTargetTimeoutMs: number;
|
||||
log: ComboLogger;
|
||||
}): (
|
||||
b: Record<string, unknown>,
|
||||
modelStr: string,
|
||||
target?: SingleModelTarget
|
||||
) => Promise<Response> {
|
||||
const { handleSingleModel, comboTargetTimeoutMs, log } = deps;
|
||||
return async (
|
||||
b: Record<string, unknown>,
|
||||
modelStr: string,
|
||||
target?: SingleModelTarget
|
||||
): Promise<Response> => {
|
||||
if (comboTargetTimeoutMs <= 0) {
|
||||
return handleSingleModel(b, modelStr, target).catch((err) =>
|
||||
errorResponse(502, err?.message ?? "Upstream model error")
|
||||
);
|
||||
}
|
||||
|
||||
const timeoutController = new AbortController();
|
||||
let timeoutId: ReturnType<typeof setTimeout> | undefined;
|
||||
let timedOut = false;
|
||||
const timeoutPromise = new Promise<Response>((resolve) => {
|
||||
timeoutId = setTimeout(() => {
|
||||
timedOut = true;
|
||||
log.warn(
|
||||
"COMBO",
|
||||
`Model ${modelStr} exceeded ${comboTargetTimeoutMs}ms timeout — falling back`
|
||||
);
|
||||
timeoutController.abort(new Error("combo-per-model-timeout"));
|
||||
resolve(
|
||||
new Response(JSON.stringify({ error: { message: `Model ${modelStr} timed out` } }), {
|
||||
status: 524,
|
||||
headers: { "Content-Type": "application/json" },
|
||||
})
|
||||
);
|
||||
}, comboTargetTimeoutMs);
|
||||
});
|
||||
const targetWithSignal = {
|
||||
...(target ?? {}),
|
||||
modelAbortSignal: timeoutController.signal,
|
||||
};
|
||||
const parentHedgeSignal = target?.modelAbortSignal ?? null;
|
||||
let onParentHedgeAbort: (() => void) | null = null;
|
||||
if (parentHedgeSignal) {
|
||||
if (parentHedgeSignal.aborted) {
|
||||
timeoutController.abort(new Error("hedge-cancelled"));
|
||||
} else {
|
||||
onParentHedgeAbort = () => {
|
||||
timeoutController.abort(new Error("hedge-cancelled"));
|
||||
};
|
||||
parentHedgeSignal.addEventListener("abort", onParentHedgeAbort, { once: true });
|
||||
}
|
||||
}
|
||||
try {
|
||||
return await Promise.race([
|
||||
handleSingleModel(b, modelStr, targetWithSignal).catch((err) => {
|
||||
if (timedOut) {
|
||||
// Inner call rejected because we aborted it. The synthetic 524 from
|
||||
// timeoutPromise already wins the race; return an empty response so
|
||||
// the loser branch resolves cleanly without leaking err.message.
|
||||
return new Response(null, { status: 599 });
|
||||
}
|
||||
return errorResponse(502, err?.message ?? "Upstream model error");
|
||||
}),
|
||||
timeoutPromise,
|
||||
]);
|
||||
} finally {
|
||||
clearTimeout(timeoutId);
|
||||
if (parentHedgeSignal && onParentHedgeAbort) {
|
||||
parentHedgeSignal.removeEventListener("abort", onParentHedgeAbort);
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
78
tests/unit/combo-target-timeout-runner.test.ts
Normal file
78
tests/unit/combo-target-timeout-runner.test.ts
Normal file
@@ -0,0 +1,78 @@
|
||||
import { test } from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import { buildTargetTimeoutRunner } from "../../open-sse/services/combo/targetTimeoutRunner.ts";
|
||||
|
||||
const noopLog = { warn() {}, info() {}, error() {}, debug() {} } as any;
|
||||
|
||||
test("timeout<=0: passthrough direto (sem timer)", async () => {
|
||||
let called = false;
|
||||
const runner = buildTargetTimeoutRunner({
|
||||
handleSingleModel: async () => {
|
||||
called = true;
|
||||
return new Response("ok");
|
||||
},
|
||||
comboTargetTimeoutMs: 0,
|
||||
log: noopLog,
|
||||
});
|
||||
const res = await runner({}, "m");
|
||||
assert.equal(called, true);
|
||||
assert.equal(await res.text(), "ok");
|
||||
});
|
||||
|
||||
test("timeout<=0: erro do upstream vira errorResponse 502", async () => {
|
||||
const runner = buildTargetTimeoutRunner({
|
||||
handleSingleModel: async () => {
|
||||
throw new Error("boom");
|
||||
},
|
||||
comboTargetTimeoutMs: 0,
|
||||
log: noopLog,
|
||||
});
|
||||
const res = await runner({}, "m");
|
||||
assert.equal(res.status, 502);
|
||||
});
|
||||
|
||||
test("excede o limite: aborta e retorna 524 timed out", async () => {
|
||||
const runner = buildTargetTimeoutRunner({
|
||||
handleSingleModel: (_b, _m, target) =>
|
||||
new Promise<Response>((resolve) => {
|
||||
// resolve só se abortado (simula um upstream que respeita o signal)
|
||||
const sig = (target as any)?.modelAbortSignal as AbortSignal | undefined;
|
||||
sig?.addEventListener("abort", () => resolve(new Response(null, { status: 599 })));
|
||||
}),
|
||||
comboTargetTimeoutMs: 20,
|
||||
log: noopLog,
|
||||
});
|
||||
const res = await runner({}, "slow-model");
|
||||
assert.equal(res.status, 524);
|
||||
const body = await res.json();
|
||||
assert.match(JSON.stringify(body), /timed out/i);
|
||||
});
|
||||
|
||||
test("sucesso rápido vence a corrida do timeout", async () => {
|
||||
const runner = buildTargetTimeoutRunner({
|
||||
handleSingleModel: async () => new Response("fast", { status: 200 }),
|
||||
comboTargetTimeoutMs: 1000,
|
||||
log: noopLog,
|
||||
});
|
||||
const res = await runner({}, "m");
|
||||
assert.equal(res.status, 200);
|
||||
assert.equal(await res.text(), "fast");
|
||||
});
|
||||
|
||||
test("hedge do parent já abortado propaga o abort ao filho", async () => {
|
||||
const parent = new AbortController();
|
||||
parent.abort(new Error("hedge-cancelled"));
|
||||
let sawAbort = false;
|
||||
const runner = buildTargetTimeoutRunner({
|
||||
handleSingleModel: (_b, _m, target) =>
|
||||
new Promise<Response>((resolve) => {
|
||||
const sig = (target as any)?.modelAbortSignal as AbortSignal | undefined;
|
||||
if (sig?.aborted) sawAbort = true;
|
||||
resolve(new Response("ok"));
|
||||
}),
|
||||
comboTargetTimeoutMs: 1000,
|
||||
log: noopLog,
|
||||
});
|
||||
await runner({}, "m", { modelAbortSignal: parent.signal } as any);
|
||||
assert.equal(sawAbort, true);
|
||||
});
|
||||
Reference in New Issue
Block a user