Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 27 additions & 3 deletions open-sse/services/autoCombo/chaosEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,22 @@ function dispatchOnePanelModel(opts: {
log?.info?.(
`CHAOS panel ${index} (${model}) ok=${res.ok} status=${res.status} textLen=${text.length}`
);
const part: ChaosPart = { model, index, ok: true, text };
// G5b: honor the upstream response status — a 4xx/5xx is a panel FAILURE,
// not a success (previously ok:true was hardcoded, so an all-error panel
// never reached the all-failed branch and the error text was streamed as
// if it were a successful answer).
if (res.ok) {
const part: ChaosPart = { model, index, ok: true, text };
await onResult?.(part);
return part;
}
const part: ChaosPart = {
model,
index,
ok: false,
text: "",
error: `upstream ${res.status}: ${text.slice(0, 200) || res.statusText || "error"}`,
};
await onResult?.(part);
return part;
} catch (err) {
Expand Down Expand Up @@ -466,8 +481,17 @@ export async function handleChaosChat(opts: {
}

if (successes.length === 0) {
const errText = "All chaos panel models failed";
await safeEnqueue(chatChunk(chunkId, panelToDispatch[0] ?? "", errText));
// G5 (silent-stop fix): make an all-panel failure visible server-side.
// The status stays 200 (SSE envelope must stay well-formed), but the
// failure is now logged with the per-model errors so operators can see
// why the chaos panel produced nothing.
const modelErrors = allParts.map((p) => `${p.model}: ${p.error ?? "unknown"}`).join(" | ");
log?.warn?.(
"CHAOS",
`All chaos panel models failed for ${comboName ?? "panel"}: ${modelErrors}`
);
const errText = `All chaos panel models failed — ${modelErrors}`;
await safeEnqueue(chatChunk(chunkId, panelToDispatch[0] ?? panel[0] ?? "", errText));
await safeEnqueue(SSE_DONE);
await enqueueChain;
closed = true;
Expand Down
11 changes: 11 additions & 0 deletions open-sse/services/autoCombo/pipelineRouter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -343,6 +343,17 @@ export async function handlePipelineCombo({
}
}

// G6 (silent-stop fix): if the reflection loop burned its retry budget and the
// verdict is still "fail", the fall-through below returns a FAILED result
// indistinguishable from a first-attempt failure. Surface it loudly so the
// caller (and operator logs) can tell "retries exhausted" apart.
if (result.reflectVerdict === "fail" && reflectionCount > 0) {
log.warn(
"PIPELINE",
`Reflection retries exhausted (${reflectionCount}/${maxReflectionLoops}) — pipeline verdict still "fail", returning the original failed result`
);
}

// ── Return result ─────────────────────────────────────────────────────────
// Check if the last stage has a streaming Response
const lastStage = result.stages[result.stages.length - 1];
Expand Down
21 changes: 17 additions & 4 deletions open-sse/services/autoRefreshDaemon.ts
Original file line number Diff line number Diff line change
Expand Up @@ -125,8 +125,13 @@ class AutoRefreshDaemon {
`[AutoRefreshDaemon] Credential expired for "${providerId}" (${config.displayName})`
);
}
} catch {
// Network errors are non-fatal — retry next cycle
} catch (err) {
// Network errors are non-fatal — retry next cycle. G8: log which
// provider failed so credential problems are not silently masked.
console.warn(
`[AutoRefreshDaemon] Network error validating credential for "${providerId}" — retry next cycle`,
err instanceof Error ? err.message : err
);
}
}

Expand Down Expand Up @@ -165,8 +170,16 @@ class AutoRefreshDaemon {
}

return true;
} catch {
// Network errors (timeout, DNS failure) don't mean the credential is bad
} catch (err) {
// Network errors (timeout, DNS failure) don't mean the credential is bad.
// G8 (silent-stop fix): the previous bare `catch { return true; }` swallowed
// the error entirely — operators could never tell a credential was failing
// to validate due to network trouble. Log it (provider + reason) before
// returning the fail-open result.
console.warn(
`[AutoRefreshDaemon] Network error validating credential for "${providerId}" — treated as valid (fail-open), will retry next cycle`,
err instanceof Error ? err.message : err
);
return true;
} finally {
clearTimeout(timeout);
Expand Down
44 changes: 38 additions & 6 deletions open-sse/services/batchProcessor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -506,14 +506,46 @@ async function processSingleItemWithRetry(item: BatchRequestItem, apiKey: string
}
}

// G10 (silent-stop fix): individual batch-item dispatches can hang indefinitely
// if the upstream route stalls (no signal/timeout plumbed through). Bound each
// item with a wall-clock timeout so a stuck item fails fast (recorded as an item
// error) instead of freezing the whole batch loop. The orphaned dispatch keeps
// running in the background but can no longer block the batch.
export const BATCH_ITEM_DISPATCH_TIMEOUT_MS = 120_000;

/**
* G10: race a promise against a wall-clock deadline. Exported for unit testing
* (batch dispatch is a module-internal import, so the timeout mechanism itself
* is verified directly here).
*/
export function withItemDispatchTimeout<T>(
promise: Promise<T>,
timeoutMs: number,
label: string
): Promise<T> {
let timer: ReturnType<typeof setTimeout> | undefined;
const timeoutPromise = new Promise<never>((_, reject) => {
timer = setTimeout(
() => reject(new Error(`${label} timed out after ${timeoutMs}ms`)),
timeoutMs
);
});
return Promise.race([promise, timeoutPromise]).finally(() => {
if (timer) clearTimeout(timer);
});
}

async function processSingleItem(item: BatchRequestItem, apiKey: string) {
const body = buildRequestBody(item);

return await dispatch.dispatchBatchApiRequest({
endpoint: item.url,
body,
apiKey,
});
return withItemDispatchTimeout(
dispatch.dispatchBatchApiRequest({
endpoint: item.url,
body,
apiKey,
}),
BATCH_ITEM_DISPATCH_TIMEOUT_MS,
`Batch item dispatch (${item.url})`
);
}

export function buildRequestBody(item: BatchRequestItem) {
Expand Down
169 changes: 159 additions & 10 deletions open-sse/services/combo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,8 @@ import {
TRANSIENT_FOR_SEMAPHORE,
MAX_FALLBACK_WAIT_MS,
MAX_GLOBAL_ATTEMPTS,
COMBO_LOOP_SAFETY_TIMEOUT_MS,
COMBO_SAFETY_DRAIN_MS,
isAllAccountsRateLimitedResponse,
clampComboDepth,
shouldSkipForPredictedTtft,
Expand Down Expand Up @@ -1046,11 +1048,54 @@ async function handleComboChatInner({
const globalPromise = new Promise<Response>((res) => {
globalResolve = res;
});

// G1 (silent-stop fix): the speculative loop's `Promise.race` waits on
// `globalPromise`, which is ONLY resolved from inside a task (success or
// fatal error). If a target hangs — e.g. the operator disabled the per-model
// timeout (`targetTimeoutMs: 0`) and the upstream never settles — the race
// never resolves and the request hangs forever with no response. This safety
// promise force-resolves after the combo budget (comboTimeoutMs when set,
// otherwise a hard ceiling) so the request ALWAYS terminates with an
// actionable 504 instead of dying silently. `comboExpired` is flipped so the
// target loop stops launching new work; the existing comboExpired branch
// returns the aggregated 504.
const loopSafetyMs =
comboTimeoutMs > 0 ? comboTimeoutMs : COMBO_LOOP_SAFETY_TIMEOUT_MS;
let loopSafetyFired = false;
let loopSafetyTimer: ReturnType<typeof setTimeout> | null = null;
const loopSafetyPromise = new Promise<Response>((resolve) => {
loopSafetyTimer = setTimeout(() => {
loopSafetyFired = true;
log.warn(
"COMBO",
`Combo loop safety timeout (${loopSafetyMs}ms) reached without a terminal response — force-terminating`
);
resolve(
errorResponseWithComboDiagnostics(
504,
`Combo global timeout (${loopSafetyMs}ms) without a terminal response`,
buildComboDiag("combo_timeout"),
{ code: "COMBO_TIMEOUT", type: "server_error" }
)
);
}, loopSafetyMs);
loopSafetyTimer.unref?.();
});
const runningTasks = new Set<Promise<void>>();
let anySuccess = false;
// #10681: steps already recorded as dispatched (so per-target retries do not
// duplicate the decision).
const dispatchedTargets = new Set<string>();
// G1: flip comboExpired as soon as the safety timer fires so the next loop
// iteration breaks instead of launching more targets after the budget, and
// abort every in-flight target so a hung upstream actually gets cancelled
// (not just "response stops").
const markLoopExpiredIfSafetyFired = () => {
if (loopSafetyFired) {
comboExpired = true;
for (const [, ac] of abortControllers.entries()) ac.abort();
}
};
const abortControllers = new Map<number, AbortController>();
const zeroLatencyOptimizationsEnabled = config.zeroLatencyOptimizationsEnabled === true;
const hasProtectedPriorityTarget =
Expand Down Expand Up @@ -2363,6 +2408,17 @@ async function handleComboChatInner({
})().catch((err) => {
const logError = log.error ?? log.warn;
logError("COMBO", `Speculative task error for target ${i}`, err);
// G2 (silent-stop fix): never leave the speculative loop waiting on an
// unresolved globalPromise. If a task throws unexpectedly (outside
// executeTarget's error handling) and no other task succeeds, the post-loop
// `Promise.race([globalPromise, ...])` would hang forever. Resolve with a
// 502 so the request terminates with an actionable error.
if (!anySuccess && globalResolve) {
anySuccess = true;
globalResolve(
errorResponse(502, `Combo target ${i} failed with an unexpected error`)
);
}
});

runningTasks.add(task);
Expand All @@ -2380,10 +2436,11 @@ async function handleComboChatInner({
timeoutResolve = r;
setTimeout(r, hedgeDelay);
});
await Promise.race([task, globalPromise, timeoutPromise]);
await Promise.race([task, globalPromise, timeoutPromise, loopSafetyPromise]);
} else {
await Promise.race([task, globalPromise]);
await Promise.race([task, globalPromise, loopSafetyPromise]);
}
markLoopExpiredIfSafetyFired();

// Global combo timeout check: after each target completes, stop trying
// further targets if the total elapsed time exceeds comboTimeoutMs.
Expand All @@ -2398,13 +2455,51 @@ async function handleComboChatInner({
}

if (!anySuccess && runningTasks.size > 0) {
await Promise.race([globalPromise, Promise.all([...runningTasks])]);
// G1: include loopSafetyPromise so a hung last task (per-model timeout
// disabled) cannot freeze this post-loop race forever.
await Promise.race([globalPromise, Promise.all([...runningTasks]), loopSafetyPromise]);
markLoopExpiredIfSafetyFired();
}

// G1: if the safety timer won the race (request would otherwise hang), give
// in-flight tasks a short drain window to land their per-model errors into
// comboErrors so the 504 carries the same "tried: a (500)" summary the
// regular comboExpired branch produces — then return the safety 504.
if (loopSafetyFired && !anySuccess) {
if (runningTasks.size > 0) {
await Promise.race([
Promise.allSettled([...runningTasks]),
new Promise((resolve) => setTimeout(resolve, COMBO_SAFETY_DRAIN_MS)),
]);
}
const summary = comboErrors
.slice(0, 5)
.map((e) => `${e.model} (${e.status})`)
.join(", ");
const msg =
`Combo global timeout (${loopSafetyMs}ms) after ${recordedAttempts}/${orderedTargets.length} targets` +
(comboErrors.length > 0
? ` | tried: ${summary}${comboErrors.length > 5 ? `... (+${comboErrors.length - 5})` : ""}`
: "") +
" without a terminal response";
return errorResponseWithComboDiagnostics(
504,
msg,
buildComboDiag("combo_timeout"),
{ code: "COMBO_TIMEOUT", type: "server_error" }
);
}

// #10681: finalize the decision trace (success).
finalizeComboTrace(traceInvocationId, orderedTargets);
finishComboTrace(traceInvocationId, { status: 200 });
if (anySuccess) {
// G1: clear the safety timer on the happy path so a successful combo does
// not leave a 10-minute timer alive per request.
if (loopSafetyTimer) {
clearTimeout(loopSafetyTimer);
loopSafetyTimer = null;
}
return await globalPromise;
}

Expand Down Expand Up @@ -2923,6 +3018,33 @@ async function handleRoundRobinCombo({
// and the "Done with this model" path below), mirroring handleComboChat.
const rrOutcomes: Array<ComboErrorEntry> = [];

// G4 (silent-stop fix): round-robin has NO global timeout — a hung model
// (per-model timeout disabled via targetTimeoutMs: 0) would freeze the request
// forever with no response. Safety promise + timer bound the whole loop; when
// it fires, rrExpired flips and every subsequent model attempt short-circuits
// to the 504. Cleaned up in the loop's finally.
const rrConfiguredTimeoutMs =
(config as { comboTimeoutMs?: number }).comboTimeoutMs ?? 0;
const rrLoopSafetyMs =
rrConfiguredTimeoutMs > 0 ? rrConfiguredTimeoutMs : COMBO_LOOP_SAFETY_TIMEOUT_MS;
let rrExpired = false;
let rrLoopSafetyTimer: ReturnType<typeof setTimeout> | null = null;
let rrResolveSafety: ((res: Response) => void) | null = null;
const rrSafetyPromise = new Promise<Response>((resolve) => {
rrResolveSafety = resolve;
});
rrLoopSafetyTimer = setTimeout(() => {
rrExpired = true;
log.warn(
"COMBO-RR",
`Round-robin loop exceeded ${rrLoopSafetyMs}ms without a terminal response — force-terminating`
);
rrResolveSafety?.(
errorResponse(504, `Round-robin combo exceeded ${rrLoopSafetyMs}ms without a terminal response`)
);
}, rrLoopSafetyMs);
rrLoopSafetyTimer.unref?.();

// #1731: Per-request in-memory set of providers whose quota is fully exhausted.
// When a target returns a quota-exhausted 429, remaining targets from the same
// provider are skipped to avoid the cascade through N same-provider targets.
Expand All @@ -2931,8 +3053,11 @@ async function handleRoundRobinCombo({
const transientRateLimitedProviders = new Set<string>();

// Try each model starting from the round-robin target
for (let offset = 0; offset < modelCount; offset++) {
const modelIndex = (rrStartIndex + offset) % modelCount;
try {
for (let offset = 0; offset < modelCount; offset++) {
// G4: stop launching new work once the safety timer fired.
if (rrExpired) break;
const modelIndex = (rrStartIndex + offset) % modelCount;
const target = filteredTargets[modelIndex];
const modelStr = target.modelStr;
const provider = target.provider;
Expand Down Expand Up @@ -3077,11 +3202,15 @@ async function handleRoundRobinCombo({
fingerprint: resolveTargetFingerprint(target) ?? "",
});

const result = await handleSingleModel(attemptBody, modelStr, {
...targetForAttempt,
effectiveComboStrategy: "round-robin",
failoverBeforeRetry: config.failoverBeforeRetry,
});
const result = await Promise.race([
handleSingleModel(attemptBody, modelStr, {
...targetForAttempt,
effectiveComboStrategy: "round-robin",
failoverBeforeRetry: config.failoverBeforeRetry,
}),
rrSafetyPromise,
]);
if (rrExpired) return result; // G4: safety timer won — stop everything

// Quota-aware scheduling: reserve the estimated budget for this
// dispatch (opt-in, same env gate as the pre-request check). Best-effort
Expand Down Expand Up @@ -3519,6 +3648,26 @@ async function handleRoundRobinCombo({
release();
}
}
} catch (err) {
// G4: unexpected exception in the round-robin loop must never crash the
// request silently — surface a 500 instead of hanging the client.
log.error?.("COMBO-RR", "Unexpected error in round-robin loop", err);
return errorResponse(500, "Unexpected error in round-robin combo");
} finally {
if (rrLoopSafetyTimer) {
clearTimeout(rrLoopSafetyTimer);
rrLoopSafetyTimer = null;
}
}

// G4: if the safety timer fired between iterations (no race captured it),
// terminate with the actionable 504 instead of the generic exhaustion path.
if (rrExpired) {
return errorResponse(
504,
`Round-robin combo exceeded ${rrLoopSafetyMs}ms without a terminal response`
);
}

// All models exhausted
const latencyMs = Date.now() - startTime;
Expand Down
Loading
Loading