mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-07-31 04:12:10 +03:00
* 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>
372 lines
15 KiB
TypeScript
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})`
|
|
);
|
|
}
|
|
);
|