Files
OmniRoute/tests/integration/live-gemini-agentic-loop.test.ts
Markus Hartung cc17b304ab fix(sse): Gemini TPM/RPD quota classification + combo cooldown-wait resilience (#8213)
* fix(sse): Gemini TPM classification, combo-cooldown-wait for auto/quota-share, and target-timeout floor

Gemini TPM/RPM 429s were misclassified as QUOTA_EXHAUSTED because
sanitizeErrorMessage() truncates to the first line, hiding Google's
metric name and retry hint on lines 2-3. Added a rawMessage field
(internal-only, never reaches the client) and classifyGeminiQuotaMetricFromText()
to classify from the untruncated text, reordered ahead of the generic
credits/daily-quota checks.

Widened comboCooldownWaitEnabled (wait out a short transient cooldown
instead of crystallizing a 429/503) from quota-share-only to also cover
auto-strategy combos, and raised the wait ceiling to 65s/130s-budget/90s-cap
to match Gemini's ~60s TPM/RPM windows.

The per-target timeout (DEFAULT_COMBO_TARGET_TIMEOUT_MS, 120s) was shorter
than the new 130s cooldown-wait budget, so a target could get cut off
mid-wait with a synthetic 524 instead of completing the retry. Added
resolveComboTargetTimeoutMsForCombo()/isComboCooldownWaitEligible() in
comboConfig.ts to raise the per-target floor to budgetMs+buffer only for
wait-eligible strategies (auto/quota-share), verified live: a 12-request
concurrent burst against a TPM-exhausted combo went from 2/12 succeeding
(10 x 524) to 12/12 succeeding with zero 503/524.

Also: liveGeminiShared.ts's sendAndValidate now fails fast on a 503
instead of retrying past it, and the health dashboard + request logger
surface TPM stats alongside RPM/RPD.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): combo-exhausted rejection logs now capture request body + attempted models

recordRejectedRequestUsage() (the fast path for combo requests that never
reach handleChatCore, e.g. all targets locked by resilience cooldown)
hardcoded provider: "-" and never passed a request body to saveCallLog(),
so /dashboard/logs entries for these failures were nearly useless for
debugging: no way to see the client's request or which models were tried.

- recordRejectedRequestUsage() now accepts requestBody and persists it
  through the existing saveCallLog() artifact mechanism (same path
  handleChatCore's own logging uses).
- Added summarizeComboAttemptedModels(), which reads the combo's own model
  list (always available, unlike the response's combo-diagnostics headers —
  a model-level resilience-lockout skip never touches the
  exhaustedProviders/exhaustedConnections sets those headers are built
  from) to populate a real "provider" value instead of "-".
- Wired both into the call site in src/sse/handlers/chat.ts.

NOTE: unrelated to the Gemini TPM/combo-cooldown-wait fix on this branch —
landed here per operator request, to be split into its own branch/PR.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* feat(sse): synthetic streaming keep-alive event + 5-minute Gemini cooldown-wait ceiling

Many clients enforce a first-SSE-byte timeout, which made it unsafe to wait
out a longer upstream rate-limit cooldown on a streaming request — the
client would abandon the connection before any bytes arrived. This landed
in two parts:

1. Synthetic startup "thinking" event (OpenAI chat/completions format):
   the already-existing withEarlyStreamKeepalive wrapper (open-sse/utils/
   earlyStreamKeepalive.ts, wired into /v1/chat/completions, /v1/messages,
   /v1/responses since #2544) opens the SSE stream immediately once a
   request runs past its threshold, but only ever sent empty/no-op
   keepalive frames. Added a `startupFrame` option (defaults to
   `keepaliveFrame` — zero behavior change unless a route opts in) so the
   very first frame can carry real content instead. Wired
   OPENAI_STARTUP_THINKING_FRAME (a reasoning_content delta: "OmniRoute:
   got request, sending to provider") into /v1/chat/completions only —
   Claude Messages and Responses API formats both require a preceding
   envelope event (message_start / response.created) that a synthetic
   pre-dispatch frame can't safely fabricate without risking a duplicate
   envelope once the real stream arrives, so those two routes keep their
   existing (safe, proven) keepalive frames unchanged.

2. Raised the "wait out a known cooldown, then retry" ceiling to 5 minutes
   for both retry mechanisms, now that a client-side first-byte timeout is
   no longer a risk on the (opted-in) route:
   - comboCooldownWait (auto/quota-share combos, open-sse/services/combo.ts):
     maxWaitMs hard clamp raised 90s -> 300s (src/lib/resilience/settings/
     normalize.ts); defaults raised to maxWaitMs:90s/maxAttempts:5/
     budgetMs:300s. comboConfig.ts's resolveComboTargetTimeoutMsForCombo
     already derives the per-target timeout floor from budgetMs, so it
     tracks the new ceiling with no further changes.
   - waitForCooldown (direct, non-combo model requests, src/sse/handlers/
     chat.ts): this mechanism had NO cumulative cap before — only a
     per-wait cap (maxRetryWaitMs) and a retry count (maxRetries), so
     maxRetries x maxRetryWaitMs could exceed 5 minutes with no ceiling.
     Added a budgetMs field (mirrors comboCooldownWait) to
     WaitForCooldownSettings/CooldownAwareRetrySettings, threaded a
     requestRetryBudgetLeftMs tracker through chat.ts's requestAttemptLoop
     (mirrors combo.ts's comboCooldownBudgetLeftMs), and made
     getCooldownAwareRetryDecision refuse to wait once the cumulative
     budget is exhausted even if the single wait is under maxRetryWaitMs.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): extend the synthetic keep-alive thinking event to /v1/responses

Live incident (OpenClaw, log id 1784407081908-cbc24f): a /v1/responses
request to gemini/gemma-4-31b-it took 56s to produce a first byte and the
client disconnected (499 request_signal_aborted) — the same client-first-byte-
timeout problem the previous commit fixed for /v1/chat/completions, but
/v1/responses only had the generic bare-comment keepalive (no content), so it
wasn't covered.

Added RESPONSES_STARTUP_THINKING_FRAME: a self-contained synthetic reasoning
item (response.output_item.added -> reasoning_summary_part.added ->
reasoning_summary_text.delta -> reasoning_summary_part.done), opened AND
closed within this one frame rather than left dangling — it never carries a
response_id, so it can't collide with the real upstream response's own
independent response.created lifecycle that follows. Mirrors the abbreviated
delta+part.done close pattern open-sse/utils/stream.ts's own
emitSyntheticResponsesReasoningSummary already uses for real mid-stream
reasoning content.

Wired into src/app/api/v1/responses/route.ts via the startupFrame option
added in the previous commit.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): combo cooldown-wait vars reset every setTry, crystallizing a bogus 503 instead of waiting

Live incident (log id 1784416706646-51): a request to the "default" combo
(strategy=auto, maxSetRetries=3) hit a real Gemini TPM 429 on both gemma-4
targets, correctly classified as a short 40s rate_limit lockout — then
crystallized a 503 "all upstream accounts are inactive" in 6.9s instead of
ever reaching the cooldown-aware wait.

Root cause: `lastError`/`earliestRetryAfter`/`lastStatus` were declared with
`let` INSIDE the `for (setTry...)` loop body, so they reset to null at the
start of every set-try. When both targets lock out on setTry 0, every
subsequent setTry (1..maxSetRetries) pre-skips both targets via the
isModelLocked check with no real dispatch — so on the FINAL setTry (the only
one whose values the post-loop decision reads, since it's gated behind
`if (setTry < maxSetRetries) continue`), lastStatus was null, hitting the
"!lastStatus" branch (ALL_ACCOUNTS_INACTIVE 503) and completely bypassing the
comboCooldownWaitEnabled / earliestRetryAfter wait logic — even though a
real 429 with a known ~40s retry-after WAS observed on setTry 0.

This bug predates today's Gemini TPM work (any combo with maxSetRetries > 0
whose targets all lock out on the first pass was affected) but was masked in
existing tests: the "auto strategy (2 models...)" regression test uses
maxSetRetries: 0, so it only ever runs ONE setTry iteration and never
exercises the reset-on-retry path. It also explains why the dedicated
12-concurrent-request burst test passed cleanly — with concurrent requests,
timing variance meant some request's FINAL setTry iteration still had a live
target to dispatch to, giving lastStatus/earliestRetryAfter fresh data. A
single isolated request has no such luck.

Fix: hoist lastError/earliestRetryAfter/lastStatus to just inside
dispatchWithCooldownRetry, before the setTry loop, so they persist across
set-tries (still reset fresh on each recursive dispatchWithCooldownRetry()
call after a wait, which is correct). recordedAttempts/fallbackCount/
exhaustedProviders etc. are intentionally left per-iteration (unrelated to
this bug).

New regression test in tests/unit/combo-quota-share-cooldown-wait.test.ts
reproduces the exact live scenario (2 targets, both lock out on setTry 0,
maxSetRetries: 3) — confirmed red (503) against the pre-fix code, green
(200, waits and retries) against the fix.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* test(sse): extend live Gemini workload to Responses API + add large-context TPM test

Two additions to the live Gemini test suite, both live-verified against the
dev instance:

1. sendAndValidate() (tests/integration/liveGeminiShared.ts) now accepts an
   apiFormat: "chat" | "responses" parameter, building the Responses-API
   request shape (input array, max_output_tokens) and parsing its SSE events
   (response.output_text.delta / response.reasoning_summary_text.delta /
   response.completed) via the new readResponsesSSEStream(). Wired into two
   new tests in live-gemini-workload.test.ts ([30]/[31]), mirroring the
   existing Chat Completions streaming coverage. Verified live: 24/25 + 5/5
   payloads succeeded end-to-end through the new code path (the one failure
   was a ~300s test-client fetch timeout unrelated to the Responses API code
   itself — a separate, not-yet-addressed test-harness limitation).

2. genHugeContextMessage() builds a single message large enough (~4
   chars/token estimate) to approach or exceed Gemini's free-tier TPM ceiling
   (16000 input tokens/min for gemma-4) by itself. Every other prompt
   generator in this file tops out around 1-2k tokens — nowhere near that
   ceiling — so none of the existing workload tests ever exercised a REAL TPM
   429, only RPM-style rate limiting. tests/integration/gemini-large-context-tpm.test.ts
   sends two ~12-13k-token requests back-to-back (comfortably exceeding
   16000/min together) to exercise the full path against production Gemini:
   TPM classification, the comboCooldownWait retry, and the synthetic
   keep-alive frame on a genuinely slow request.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): abandoned combo target dispatch now observes its own per-target timeout, fixing a permanent "pending" dashboard leak

Live incident (dashboard log id 1784418258231-14961a, reported as "an ongoing
request even though there's already a 200"): a combo target dispatch
abandoned by comboTargetTimeoutMs (open-sse/services/combo/targetTimeoutRunner.ts)
left a permanent phantom "pending" entry in the dashboard, even after the
overall combo request had already succeeded via a different retry.

Root cause: chatCore.ts's createStreamController — and everything downstream
that depends on it (withRateLimit's Promise.race against Bottleneck,
acquireAccountSemaphore) — only ever watches clientRawRequest.signal, which
is the ORIGINAL client's request signal (set once via buildClientRawRequest
and reused unchanged across every target dispatch in a combo). It has no
connection to targetTimeoutRunner.ts's OWN AbortController
(target.modelAbortSignal), which is what actually fires when
comboTargetTimeoutMs (300s) elapses. src/sse/handlers/chat.ts's
handleSingleModel bridge between combo.ts and handleSingleModelChat received
`target.modelAbortSignal` but silently dropped it — never forwarded it
anywhere. So when a target got abandoned (e.g. stuck inside a wedged
Bottleneck rate-limiter queue, see the WEDGED force-reset log line from the
same incident), its per-target timeout fired and let the COMBO move on and
retry successfully elsewhere — but the abandoned dispatch's own promise
chain never learned it had been superseded, so it hung forever waiting on a
signal that was never going to fire, and trackPendingRequest(false) (the
finalize call) never ran.

Fix: thread target.modelAbortSignal through as a new modelAbortSignal
runtimeOption, and merge it into clientRawRequest.signal (via the existing
mergeAbortSignals helper from open-sse/executors/base.ts) right before
dispatch, so an abandoned target's own promise chain now observes its abort
and can reach its cleanup path — new resolveDispatchClientRawRequest() makes
this mechanically testable in isolation.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* fix(sse): combo cooldown-wait state recording, rate-limit wedge recovery, OpenAI-format SSE error frames

Five related fixes surfaced by live incidents (dashboard log ids 1784457764961-73,
1784465227489-a2cbc0, 1784504040241-6f8b9a) while validating the Gemini TPM/cooldown-wait
work on this branch against real OpenClaw traffic:

- combo.ts: the model-lockout bail-out branches in dispatchWithCooldownRetry never
  recorded lastStatus, so once every target in a set hit an existing lockout the final
  check crystallized a bogus ALL_ACCOUNTS_INACTIVE 503 instead of reaching the
  cooldown-wait decision, even with a real 429 + short retry-after observed.
- combo.ts/combo/types.ts: the "all credentials cooling down" pre-dispatch rejection
  (buildModelCooldownBody) nests its retry hint as error.retry_after/reset_seconds, not
  the top-level retryAfter every other 429 shape uses — combo's extraction only read the
  latter, so earliestRetryAfter stayed null for this shape even after lastStatus was fixed.
- rateLimitManager.ts: the wedge-recovery watchdog used disconnect(), which releases the
  heartbeat timer but never rejects jobs already QUEUED on that instance — orphaned
  dispatches hung until the outer ~300s per-target timeout, well past real clients'
  patience. Switched to stop({ dropWaitingJobs: true }), safe because the wedge condition
  already requires RUNNING===0 && EXECUTING===0.
- earlyStreamKeepalive.ts: the in-band error frame emitted after committing to a 200 SSE
  stream was hardcoded to Anthropic's `event: error` convention for every route, including
  the OpenAI-format ones (/v1/chat/completions, /v1/responses) where that framing is
  either invisible or malformed to a plain data-line parser. Added per-route
  OPENAI_CHAT_ERROR_FRAME / OPENAI_RESPONSES_ERROR_FRAME and wired them in.
- chatCore.ts: persisted a synthetic clientResponse error body even when the client had
  already disconnected (AbortError) before that body was ever computed — misleading the
  dashboard into showing "what the client received" for a response that was never sent.

Also: RequestLoggerDetail.tsx — Provider/Client Event Stream panes lost their collapse
toggle when StreamSection replaced the collapsible PayloadSection (692d6be80, unifying
active/finished request views) without carrying the toggle over.

Each fix has a TDD regression test with a confirmed red-before-green cycle.

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>

* test(sse): free-tier model + gemma-4 TPM-ceiling benchmark harness

Adds a live benchmark comparing free models OmniRoute exposes across
configured providers plus previously-unexercised no-auth providers
(felo-web, aihorde, opencode, duckduckgo-web — none need a connection
row, they were just never tried). Reuses liveGeminiShared.ts's SSE
parsers and CASE_BUILDERS instead of duplicating them.

Also adds a targeted TPM-stress test firing back-to-back large-context
prompts at the gemma-4-31b model across its 3 free hosts (Gemini,
NVIDIA, AI Horde) to isolate whether the documented 16k-tokens/minute
free-tier ceiling is Gemini-specific enforcement or an inherent
model property.

FORCE_TOOL_CHOICE_REQUIRED is a test-only, default-off env flag added
to liveGeminiShared.ts and live-gemini-agentic-loop.test.ts for an
earlier live A/B comparison of tool_choice: required vs unset — kept
as a reusable knob for future runs.

Co-Authored-By: Markus Hartung <markus.hartung@gmail.com>

* test(sse): benchmark for the 2026-07-22 newly-enabled provider batch

Adds NEWLY_ENABLED_MODELS to freeModelBenchmarkShared.ts (Mistral
Leanstral, OpenRouter's live "free"-tagged roster, OpenCode Zen's
current free models — refetched live from
https://opencode.ai/zen/v1/models since the static catalog had
drifted) and a dedicated workload benchmark test for them.

Co-Authored-By: Markus Hartung <markus.hartung@gmail.com>

* test(sse): sync geminiRateLimitTracker tests with e74a1722b's corrected Gemma 4 limits

e74a1722b updated geminiRateLimits.json's gemma-4-* entries from the stale
15/1500/-1 (rpm/rpd/tpm) to the real published free-tier values
16000/14400/16000, but never updated the tests asserting the old numbers.
Surfaced by running the full test:unit suite as a post-rebase sanity check.

Co-Authored-By: Markus Hartung <markus.hartream@gmail.com>

* chore(quality): file-size baseline for own-growth (#8213)

Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>

---------

Co-authored-by: Markus Hartung <markus.hartung@gmail.com>
Co-authored-by: Markus Hartung <markus.hartream@gmail.com>
Co-authored-by: diegosouzapw <8016841+diegosouzapw@users.noreply.github.com>
2026-07-23 05:17:35 -03:00

372 lines
15 KiB
TypeScript

/**
* tests/integration/live-gemini-agentic-loop.test.ts
*
* Live test: a REAL, streaming, 3-turn agentic tool-calling flow against the
* "default" gemini combo (strategy=auto, 2 gemma-4 targets), scripted to
* exercise the exact cross-model cooldown-wait sequence live incidents have
* shown:
*
* 1. Turn 1 (~10k-token dispatch): model A serves and replies with a
* write_file tool call.
* 2. Turn 2: the tool result + ~10k MORE filler is appended (cumulative
* conversation now ~20k tokens — past gemma-4's published 16000 TPM
* free-tier ceiling within the same rolling 60s window). Model A hits a
* real 429; OmniRoute transparently falls back to model B — the client
* must see a normal 200 with a tool call from B, NEVER the 429.
* 3. Turn 3: B's tool result + more filler. Now BOTH models are cooling
* down, but A's remaining cooldown (recorded a full turn earlier) is
* shorter than B's and well under the 5-minute comboCooldownWait budget
* — the request must STALL for A rather than give up, using the
* synthetic startup "thinking" keep-alive frame
* (open-sse/utils/earlyStreamKeepalive.ts OPENAI_STARTUP_THINKING_FRAME)
* to hold the SSE connection open during the wait, then resolve with A's
* response once its cooldown clears.
*
* Streaming is required (not just used for realism): earlyStreamKeepalive's
* "slow path" only exists for SSE routes, and it's the only way a client can
* observe that a wait happened without seeing an error — once the slow path
* commits to HTTP 200, a would-be error response is reframed as an in-band
* `event: error` SSE frame rather than changing the status code. So this test
* treats an `event: error` frame exactly like a leaked 429/503 status: both
* are the same regression (comboCooldownWait giving up instead of waiting).
*
* Env vars: same as liveGeminiShared.ts (OMNIROUTE_API_KEY required).
*/
import test from "node:test";
import assert from "node:assert/strict";
import { Agent, fetch } from "undici";
import {
skip,
API_KEY,
BASE_URL,
MODEL,
ensureTestEnvironment,
pick,
LONG_DOCUMENTS,
CODE_BLOCKS,
TOOL_DEFINITION,
} from "./liveGeminiShared.ts";
// comboCooldownWait budgetMs default (src/lib/resilience/settings.ts) is
// 300_000ms. A single client request can span a full combo SET retry
// (maxSetRetries: 3 in the "default" combo config), each set trying both
// targets at up to comboTargetTimeoutMs (300_000ms) apiece — give this two
// full target-timeouts of slack so a legitimate one-set-retry cycle doesn't
// get killed client-side before the server can resolve it.
const TURN_TIMEOUT_MS = 700_000;
const FILLER_TOKENS_PER_TURN = 10_000;
const MODEL_A = "gemma-4-31b-it";
const MODEL_B = "gemma-4-26b-a4b-it";
const SYNTHETIC_MODEL_MARKER = "omniroute";
// Must match STARTUP_THINKING_TEXT in open-sse/utils/earlyStreamKeepalive.ts.
const STARTUP_THINKING_SUBSTRING = "OmniRoute:";
// Node's global fetch (undici) has its own client-side headersTimeout that
// defaults to 300_000ms — the SAME order of magnitude as comboCooldownWait's
// budget, so a genuine full-budget server-side wait can race the client's own
// timeout and get killed with UND_ERR_HEADERS_TIMEOUT before the server ever
// gets to respond. Use an explicit dispatcher with headroom above
// TURN_TIMEOUT_MS so only the server's behavior (and our own AbortSignal) is
// under test, not undici's unrelated default. Irrelevant for the actual SSE
// body once bytes start flowing (the keepalive frames reset it), but matters
// for the initial connection.
const dispatcher = new Agent({
headersTimeout: TURN_TIMEOUT_MS + 30_000,
bodyTimeout: TURN_TIMEOUT_MS + 30_000,
});
function buildFillerText(approxTokens: number): string {
const CHARS_PER_TOKEN = 4;
const targetChars = approxTokens * CHARS_PER_TOKEN;
const chunks = [...LONG_DOCUMENTS, ...CODE_BLOCKS];
let content = "";
let i = 0;
while (content.length < targetChars) {
content += `\n\n--- Reference ${i + 1} ---\n\n${pick(chunks)}`;
i++;
}
return content;
}
type ToolCallMsg = { id: string; type: "function"; function: { name: string; arguments: string } };
type AgentMessage =
| { role: "system" | "user"; content: string }
| { role: "assistant"; content: string | null; tool_calls?: ToolCallMsg[] }
| { role: "tool"; tool_call_id: string; content: string };
type TurnResult = {
status: number;
servedModel: string | null;
sawSyntheticKeepalive: boolean;
sawErrorEvent: string | null;
toolCalls: ToolCallMsg[];
content: string;
finishReason: string;
timeToFirstByteMs: number;
timeToFirstRealChunkMs: number | null;
totalDurationMs: number;
correlationId: string;
};
async function runStreamingAgentTurn(messages: AgentMessage[]): Promise<TurnResult> {
const start = performance.now();
const res = await fetch(`${BASE_URL}/v1/chat/completions`, {
method: "POST",
headers: { "Content-Type": "application/json", Authorization: `Bearer ${API_KEY}` },
body: JSON.stringify({
model: MODEL,
messages,
tools: [TOOL_DEFINITION],
stream: true,
max_tokens: 4096,
temperature: 0.2,
...(process.env.FORCE_TOOL_CHOICE_REQUIRED === "1" ? { tool_choice: "required" } : {}),
}),
signal: AbortSignal.timeout(TURN_TIMEOUT_MS),
dispatcher,
});
const correlationId = res.headers.get("x-correlation-id") || "?";
// A leaked 429/503 status is the direct regression. The early-keepalive slow
// path (see file header) never changes the HTTP status once committed to
// 200, so this only catches a FAST failure (before the 2s keepalive
// threshold) — the `event: error` scan below catches the slow-path case.
if (res.status === 429 || res.status === 503) {
const body = await res.text().catch(() => "");
assert.fail(
`turn leaked HTTP ${res.status} to the client instead of waiting for target model ` +
`availability (cid=${correlationId}): ${body.slice(0, 300)}`
);
}
assert.equal(res.status, 200, `unexpected HTTP ${res.status} (cid=${correlationId})`);
const reader = res.body!.getReader();
const decoder = new TextDecoder();
let buffer = "";
let servedModel: string | null = null;
let sawSyntheticKeepalive = false;
let sawErrorEvent: string | null = null;
let pendingEventType: string | null = null;
let timeToFirstByteMs = -1;
let timeToFirstRealChunkMs: number | null = null;
let content = "";
let finishReason = "unknown";
const toolCallDeltas = new Map<string, { id: string; name: string; arguments: string }>();
while (true) {
const { done, value } = await reader.read();
if (done) break;
if (timeToFirstByteMs < 0) timeToFirstByteMs = performance.now() - start;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split("\n");
buffer = lines.pop() || "";
for (const line of lines) {
if (line.startsWith("event: ")) {
pendingEventType = line.slice(7).trim();
continue;
}
if (!line.startsWith("data: ")) continue;
const data = line.slice(6).trim();
const eventType = pendingEventType;
pendingEventType = null;
if (data === "[DONE]") continue;
if (eventType === "error") {
sawErrorEvent = data.slice(0, 300);
continue;
}
try {
const parsed = JSON.parse(data) as Record<string, unknown>;
const model = parsed.model as string | undefined;
const choice = ((parsed.choices ?? []) as Array<Record<string, unknown>>)[0];
const delta = choice?.delta as Record<string, unknown> | undefined;
if (model === SYNTHETIC_MODEL_MARKER) {
const reasoningDelta = delta?.reasoning_content as string | undefined;
if (reasoningDelta?.includes(STARTUP_THINKING_SUBSTRING)) {
sawSyntheticKeepalive = true;
}
continue; // synthetic frame — not real model output
}
if (model && !servedModel) servedModel = model;
if (model && timeToFirstRealChunkMs === null) {
timeToFirstRealChunkMs = performance.now() - start;
}
if (delta?.content) content += delta.content as string;
if (choice?.finish_reason) finishReason = choice.finish_reason as string;
const tcDeltas = delta?.tool_calls as Array<Record<string, unknown>> | undefined;
if (tcDeltas) {
for (const tcd of tcDeltas) {
const idx = String(tcd.index as number);
if (!toolCallDeltas.has(idx)) {
toolCallDeltas.set(idx, { id: (tcd.id as string) ?? "", name: "", arguments: "" });
}
const entry = toolCallDeltas.get(idx)!;
if (tcd.id) entry.id = tcd.id as string;
const fn = tcd.function as Record<string, unknown> | undefined;
if (fn?.name) entry.name = fn.name as string;
if (fn?.arguments) entry.arguments += fn.arguments as string;
}
}
} catch {
// skip malformed chunks
}
}
}
const toolCalls: ToolCallMsg[] = [...toolCallDeltas.values()].map((tc) => ({
id: tc.id,
type: "function" as const,
function: { name: tc.name, arguments: tc.arguments },
}));
if (sawErrorEvent) {
assert.fail(
`turn leaked an in-band SSE error event instead of waiting for target model ` +
`availability (cid=${correlationId}): ${sawErrorEvent}`
);
}
return {
status: res.status,
servedModel,
sawSyntheticKeepalive,
sawErrorEvent,
toolCalls,
content,
finishReason,
timeToFirstByteMs,
timeToFirstRealChunkMs,
totalDurationMs: performance.now() - start,
correlationId,
};
}
function logTurn(label: string, r: TurnResult) {
console.log(
` ${label.padEnd(10)} HTTP ${r.status} | model=${r.servedModel ?? "?"} | ` +
`finish=${r.finishReason} | tools=${r.toolCalls.length} | ` +
`keepalive=${r.sawSyntheticKeepalive ? "yes" : "no"} | ` +
`ttfb=${Math.round(r.timeToFirstByteMs)}ms | ` +
`ttfRealChunk=${r.timeToFirstRealChunkMs === null ? "?" : Math.round(r.timeToFirstRealChunkMs) + "ms"} | ` +
`total=${Math.round(r.totalDurationMs)}ms | cid=${r.correlationId}`
);
}
test.before(async () => {
await ensureTestEnvironment();
});
test(
"[32] agentic loop: transparent cross-model cooldown-wait across 3 real streaming turns",
{ skip, timeout: 3 * TURN_TIMEOUT_MS + 60_000 },
async () => {
const messages: AgentMessage[] = [
{
role: "system",
content:
"You are building a small TypeScript library across 3 steps, one file per step. " +
"For each step, call write_file exactly once for that step's file (using the reference " +
"material provided as context), then briefly confirm you're ready for the next step.",
},
];
// ── Turn 1: initial ~10k-token dispatch — expect model A to serve a tool call ──
messages.push({
role: "user",
content: `Step 1/3: write file step1.ts. Reference material:\n${buildFillerText(FILLER_TOKENS_PER_TURN)}`,
});
const turn1 = await runStreamingAgentTurn(messages);
logTurn("turn 1/3", turn1);
assert.equal(turn1.finishReason, "tool_calls", "turn 1 should finish with a tool call");
assert.ok(turn1.toolCalls.length > 0, "turn 1 should have called write_file");
assert.ok(turn1.servedModel, "turn 1 should report a served model");
messages.push({
role: "assistant",
content: turn1.content || null,
tool_calls: turn1.toolCalls,
});
for (const tc of turn1.toolCalls) {
messages.push({ role: "tool", tool_call_id: tc.id, content: JSON.stringify({ ok: true }) });
}
// ── Turn 2: tool result + ~10k MORE filler (cumulative ~20k, past the 16k ──
// TPM ceiling within the rolling 60s window) — model A should hit a real
// 429 and OmniRoute should transparently fail over to model B. The served
// model changing between turn 1 and turn 2, with NO leaked error, is the
// client-observable proof that (a) a 429 really happened and (b) the other
// model was transparently retried — there is no other way for a client to
// see this, since hiding the 429 from the client is the entire point of
// comboCooldownWait.
messages.push({
role: "user",
content: `Step 2/3: write file step2.ts. Reference material:\n${buildFillerText(FILLER_TOKENS_PER_TURN)}`,
});
const turn2 = await runStreamingAgentTurn(messages);
logTurn("turn 2/3", turn2);
assert.equal(turn2.finishReason, "tool_calls", "turn 2 should finish with a tool call");
assert.ok(turn2.toolCalls.length > 0, "turn 2 should have called write_file");
assert.notEqual(
turn2.servedModel,
turn1.servedModel,
`expected turn 2 to transparently fail over to the OTHER model after turn 1's model ` +
`(${turn1.servedModel}) hit TPM contention — got the SAME model again ` +
`(${turn2.servedModel}), meaning no 429/fallback was observed. Re-run if the account ` +
`wasn't actually under contention this time.`
);
messages.push({
role: "assistant",
content: turn2.content || null,
tool_calls: turn2.toolCalls,
});
for (const tc of turn2.toolCalls) {
messages.push({ role: "tool", tool_call_id: tc.id, content: JSON.stringify({ ok: true }) });
}
// ── Turn 3: tool result + more filler — B (just used) should ALSO hit 429, ──
// but A's cooldown (recorded a full turn earlier) is now the shorter of the
// two and well under the 5-minute comboCooldownWait budget — the request
// must STALL for A specifically rather than crystallizing a 503. The
// synthetic startup "thinking" keep-alive frame is the client-visible proof
// the server actually waited instead of just failing fast; servedModel
// flipping back to A (not staying on B, not erroring) is the proof it
// waited for the SHORTER cooldown specifically.
messages.push({
role: "user",
content: `Step 3/3: write file step3.ts. Reference material:\n${buildFillerText(FILLER_TOKENS_PER_TURN)}`,
});
const turn3 = await runStreamingAgentTurn(messages);
logTurn("turn 3/3", turn3);
assert.equal(turn3.finishReason, "tool_calls", "turn 3 should finish with a tool call");
assert.ok(turn3.toolCalls.length > 0, "turn 3 should have called write_file");
assert.equal(
turn3.servedModel,
turn1.servedModel,
`expected turn 3 to wait for and return the LOWER-cooldown model (${turn1.servedModel}, ` +
`same as turn 1) once B also hit contention — got ${turn3.servedModel}`
);
assert.ok(
turn3.sawSyntheticKeepalive,
"expected the synthetic startup keep-alive ('thinking') frame during turn 3's stall — " +
"its absence means the wait was fast enough to not need it, or the keepalive path didn't engage"
);
console.log(
`\n Summary: turn1=${turn1.servedModel} → turn2=${turn2.servedModel} (fallback) → ` +
`turn3=${turn3.servedModel} (waited for lower cooldown, keepalive=${turn3.sawSyntheticKeepalive})`
);
}
);