diff --git a/open-sse/services/accountFallback.ts b/open-sse/services/accountFallback.ts index c6465321f94..88ea9cc52ed 100644 --- a/open-sse/services/accountFallback.ts +++ b/open-sse/services/accountFallback.ts @@ -32,7 +32,7 @@ import { } from "../../src/lib/resilience/settings"; import { resolveModelLockoutSettings } from "../../src/lib/resilience/modelLockoutSettings"; import { - getAllCircuitBreakerStatuses, + getBlockedCircuitBreakerSnapshots, getCircuitBreaker, } from "../../src/shared/utils/circuitBreaker"; import { @@ -1302,6 +1302,10 @@ export function clearProviderFailure(provider: string | null | undefined): void /** * Get all providers currently blocked by the shared breaker. + * + * Read-only: the blocked set comes from the `peek*` API, so asking which providers are in + * cooldown never probes an elapsed OPEN breaker into HALF_OPEN and never persists such a + * transition. Looking at health must not be what lets a request through to a broken provider. */ export function getProvidersInCooldown(): Array<{ provider: string; @@ -1309,17 +1313,12 @@ export function getProvidersInCooldown(): Array<{ cooldownRemainingMs: number | null; lastFailureAt: number | null; }> { - return getAllCircuitBreakerStatuses() - .filter((status) => { - const breaker = getProviderBreaker(status.name); - return Boolean(breaker && !breaker.canExecute()); - }) - .map((status) => ({ - provider: status.name, - failureCount: status.failureCount, - cooldownRemainingMs: status.retryAfterMs || null, - lastFailureAt: status.lastFailureTime, - })); + return getBlockedCircuitBreakerSnapshots().map((status) => ({ + provider: status.name, + failureCount: status.failureCount, + cooldownRemainingMs: status.retryAfterMs || null, + lastFailureAt: status.lastFailureTime, + })); } /** diff --git a/open-sse/services/autoCombo/routingDecision.ts b/open-sse/services/autoCombo/routingDecision.ts index 187334cdc89..1138f892bcf 100644 --- a/open-sse/services/autoCombo/routingDecision.ts +++ b/open-sse/services/autoCombo/routingDecision.ts @@ -147,11 +147,26 @@ export interface StrategySelection { connectionId?: string; } +/** + * How much of the decision the caller keeps. The live recording path passes the decision store's + * own bounds so the decision is built in the shape it is retained in, instead of materialising + * every candidate with every factor and letting the store throw the rest away. Omit it (preview, + * diagnostics) to build the full, uncompacted decision. + */ +export interface RoutingDecisionRetention { + /** Candidates kept, in decision order. The selected candidate is always kept. */ + maxCandidates: number; + /** Leading candidates that keep their factor breakdown, besides the selected one. */ + maxCandidatesWithFactors: number; +} + export interface BuildRoutingDecisionInput { request: RoutingRequest; config: AutoComboConfig; /** Every candidate the router considered, including the ones cut before scoring. */ candidates: DecisionCandidateInput[]; + /** Bounds to build the decision to. Unset builds every candidate with every factor. */ + retention?: RoutingDecisionRetention | null; /** Engine run over the routable candidates; null when none was routable. */ outcome: { selection: SelectionResult; trace: AutoSelectionTrace } | null; /** True when a strict budget cap refused every routable candidate. */ @@ -173,19 +188,24 @@ function engineExclusionReasons( return [input.outcome.trace.excludedProviders.get(candidate.provider) ?? "self_healing_excluded"]; } +/** + * A candidate without its factor breakdown. Factors are the bulk of a candidate's allocation, and + * a live decision retains them for only a handful of candidates, so they are filled in afterwards + * for the candidates that keep them (see `retainedRoutingCandidates`). Nothing before that point — + * ordering, selection, exclusion — reads `factors`. + */ function toRoutingCandidate( candidate: DecisionCandidateInput, input: BuildRoutingDecisionInput, - scoredByKey: Map, + scored: ScoredProvider | undefined, eligibleKeys: Set ): RoutingCandidate { - const scored = scoredByKey.get(candidateKey(candidate)); const exclusionReasons = engineExclusionReasons(candidate, input, eligibleKeys); return { providerId: candidate.provider, modelId: candidate.model, score: roundScore(scored?.score ?? 0), - factors: toRoutingFactors(scored, input.outcome?.trace.weights), + factors: [], eligible: exclusionReasons.length === 0 && input.outcome !== null, exclusionReasons, quota: quotaStateOf(candidate), @@ -206,6 +226,7 @@ function selectionModeOf(input: BuildRoutingDecisionInput): RoutingDecision["sel interface DecisionEntry { key: string; candidate: DecisionCandidateInput; + scored: ScoredProvider | undefined; routing: RoutingCandidate; } @@ -247,6 +268,63 @@ function selectedEntry( return entries.find((entry) => entry.key === candidateKey(chosen) && entry.routing.eligible); } +function sameRoutingCandidate(a: RoutingCandidate, b: RoutingCandidate | undefined): boolean { + return b !== undefined && a.providerId === b.providerId && a.modelId === b.modelId; +} + +/** + * The candidates kept, in decision order, bounded by `retention`. Mirrors the decision store's + * retention rule: the leading `maxCandidates` entries, with the selected candidate taking the last + * slot when it fell outside them, so a bounded build and a full build compacted afterwards keep + * exactly the same candidates in exactly the same order. + */ +function retainedEntries( + entries: DecisionEntry[], + selected: DecisionEntry | undefined, + retention: RoutingDecisionRetention | null | undefined +): DecisionEntry[] { + if (!retention) return entries; + const kept = entries.slice(0, retention.maxCandidates); + if ( + selected && + kept.length > 0 && + !kept.some((entry) => sameRoutingCandidate(entry.routing, selected.routing)) + ) { + kept[kept.length - 1] = selected; + } + return kept; +} + +/** + * The retained candidates with their factor breakdown filled in where the decision keeps it: the + * leading `maxCandidatesWithFactors` entries and the selected candidate. Without `retention` every + * candidate keeps its factors, which is what a preview returns. + */ +function retainedRoutingCandidates( + kept: DecisionEntry[], + selected: DecisionEntry | undefined, + selectedRouting: RoutingCandidate | undefined, + input: BuildRoutingDecisionInput +): RoutingCandidate[] { + const weights = input.outcome?.trace.weights; + const withFactorsCount = input.retention?.maxCandidatesWithFactors ?? kept.length; + return kept.map((entry, index) => { + if (entry === selected && selectedRouting) return selectedRouting; + if (index >= withFactorsCount && !sameRoutingCandidate(entry.routing, selectedRouting)) { + return entry.routing; + } + return withRoutingFactors(entry, weights); + }); +} + +function withRoutingFactors( + entry: DecisionEntry, + weights: ScoringWeights | undefined +): RoutingCandidate { + const factors = toRoutingFactors(entry.scored, weights); + return factors.length === 0 ? entry.routing : { ...entry.routing, factors }; +} + /** Turn a selection into the shared decision contract. Pure apart from the clock. */ export function buildRoutingDecision( input: BuildRoutingDecisionInput, @@ -257,26 +335,37 @@ export function buildRoutingDecision( if (!scoredByKey.has(candidateKey(scored))) scoredByKey.set(candidateKey(scored), scored); } const eligibleKeys = new Set((input.outcome?.trace.eligible ?? []).map(candidateKey)); - const entries: DecisionEntry[] = input.candidates.map((candidate) => ({ - key: candidateKey(candidate), - candidate, - routing: toRoutingCandidate(candidate, input, scoredByKey, eligibleKeys), - })); + const entries: DecisionEntry[] = input.candidates.map((candidate) => { + const key = candidateKey(candidate); + const scored = scoredByKey.get(key); + return { + key, + candidate, + scored, + routing: toRoutingCandidate(candidate, input, scored, eligibleKeys), + }; + }); entries.sort( (a, b) => Number(b.routing.eligible) - Number(a.routing.eligible) || b.routing.score - a.routing.score ); const selected = selectedEntry(entries, input); + const selectedRouting = selected + ? withRoutingFactors(selected, input.outcome?.trace.weights) + : undefined; + const kept = retainedEntries(entries, selected, input.retention); + const omittedCandidates = entries.length - kept.length; return { decisionId: clock.newDecisionId(), requestId: input.request.requestId, - ...(selected ? { selected: selected.routing } : {}), - candidates: entries.map((entry) => entry.routing), + ...(selectedRouting ? { selected: selectedRouting } : {}), + candidates: retainedRoutingCandidates(kept, selected, selectedRouting, input), policyVersion: computeRoutingPolicyVersion(input.config), generatedAt: new Date(clock.now()).toISOString(), liveRequestExecuted: input.liveRequestExecuted, selectionMode: selectionModeOf(input), strategy: input.strategySelection?.strategy ?? "rules", + ...(omittedCandidates > 0 ? { omittedCandidates } : {}), }; } diff --git a/open-sse/services/combo/autoRoutingDecision.ts b/open-sse/services/combo/autoRoutingDecision.ts index a5db71cb396..fbb442d95e3 100644 --- a/open-sse/services/combo/autoRoutingDecision.ts +++ b/open-sse/services/combo/autoRoutingDecision.ts @@ -24,7 +24,11 @@ import { type StrategySelection, summarizeRoutingDecision, } from "../autoCombo/routingDecision.ts"; -import { recordRoutingDecision } from "../routing/decisionStore.ts"; +import { + MAX_CANDIDATES_WITH_FACTORS, + MAX_STORED_CANDIDATES, + recordRoutingDecision, +} from "../routing/decisionStore.ts"; import type { ComboLogger } from "./types.ts"; export interface AutoDecisionContext { @@ -56,11 +60,24 @@ function requestProtocol(body: Record): string { return body.input !== undefined ? "responses" : "unknown"; } +/** + * Live recording builds the decision already in the shape the store retains, instead of + * materialising every candidate with every factor and letting the store drop the rest: an auto + * combo over the whole catalog considers hundreds of candidates, and the factor breakdown of the + * ones that are never retained is pure waste on every routed request. The bounds are the store's + * own, so the stored decision is identical either way. + */ +const LIVE_DECISION_RETENTION = { + maxCandidates: MAX_STORED_CANDIDATES, + maxCandidatesWithFactors: MAX_CANDIDATES_WITH_FACTORS, +} as const; + function recordDecision( context: AutoDecisionContext, selection: Pick ): RoutingDecision { const decision = buildRoutingDecision({ + retention: LIVE_DECISION_RETENTION, request: { requestId: getRequestId() ?? "", model: context.config.name, diff --git a/open-sse/services/webSessionPoolHealth.ts b/open-sse/services/webSessionPoolHealth.ts index 4ad8544b34c..cb636b51ef0 100644 --- a/open-sse/services/webSessionPoolHealth.ts +++ b/open-sse/services/webSessionPoolHealth.ts @@ -12,10 +12,9 @@ import { PoolRegistry } from "./sessionPool/poolRegistry.ts"; import { - isProviderInCooldown, - getProviderCooldownRemainingMs, -} from "./accountFallback.ts"; -import { getAllCircuitBreakerStatuses } from "../../src/shared/utils/circuitBreaker.ts"; + getAllCircuitBreakerSnapshots, + peekCircuitBreaker, +} from "../../src/shared/utils/circuitBreaker.ts"; import type { PoolStats, PoolSessionDetail } from "./sessionPool/types.ts"; // ─── Dependency Injection (for testability) ───────────────────────────────── @@ -34,14 +33,25 @@ export interface WebSessionPoolHealthDeps { } | null; } +/** + * Every breaker fact this report reads comes from the `peek*` API: building a health report must + * not transition an elapsed OPEN breaker into HALF_OPEN, persist that, or consume its half-open + * probe. Looking at a pool's health may not be what lets the next request reach a broken provider. + * Live routing keeps the probing accessors — there, the transition is the point. + */ const defaultDeps: WebSessionPoolHealthDeps = { listProviders: () => PoolRegistry.listProviders(), getStats: (p) => PoolRegistry.getStats(p), getSessionDetails: (p) => PoolRegistry.getSessionDetails(p), - isProviderInCooldown: (p) => isProviderInCooldown(p), - getProviderCooldownRemainingMs: (p) => getProviderCooldownRemainingMs(p), + isProviderInCooldown: (p) => peekCircuitBreaker(p)?.peekCanExecute() === false, + getProviderCooldownRemainingMs: (p) => { + const breaker = peekCircuitBreaker(p); + if (!breaker || breaker.peekCanExecute()) return null; + const remaining = breaker.peekStatus().retryAfterMs; + return remaining > 0 ? remaining : null; + }, getProviderBreakerState: (p) => { - const statuses = getAllCircuitBreakerStatuses(); + const statuses = getAllCircuitBreakerSnapshots(); const match = statuses.find((s) => s.name === p); if (!match) return null; return { diff --git a/src/app/api/monitoring/health/route.ts b/src/app/api/monitoring/health/route.ts index 68eaf44038e..089d2e6c929 100644 --- a/src/app/api/monitoring/health/route.ts +++ b/src/app/api/monitoring/health/route.ts @@ -161,7 +161,7 @@ async function rebuildHealthPayload(): Promise { circuitBreakerModule.status === "fulfilled" ? readHealthValue( "circuit breakers", - () => circuitBreakerModule.value.getAllCircuitBreakerStatuses(), + () => circuitBreakerModule.value.getAllCircuitBreakerSnapshots(), [] ) : []; diff --git a/src/app/api/resilience/connections/route.ts b/src/app/api/resilience/connections/route.ts index 55e4714de76..cf87a7f5f77 100644 --- a/src/app/api/resilience/connections/route.ts +++ b/src/app/api/resilience/connections/route.ts @@ -2,7 +2,7 @@ import { NextRequest, NextResponse } from "next/server"; import { z } from "zod"; import { getRawProviderConnections, getProviderConnectionsCount } from "@/lib/db/providers"; -import { getAllCircuitBreakerStatuses } from "@/shared/utils/circuitBreaker"; +import { getAllCircuitBreakerSnapshots } from "@/shared/utils/circuitBreaker"; import { resolveProviderId } from "@/shared/constants/providers"; import { TERMINAL_CONNECTION_STATUSES } from "@/lib/quota/connectionRecovery"; import { sanitizeErrorMessage, buildErrorBody } from "@omniroute/open-sse/utils/error"; @@ -143,14 +143,16 @@ export async function GET(req: NextRequest) { console.error("[API] resilience/connections database error:", err); } - // NOTE: getAllCircuitBreakerStatuses() calls getStatus() internally. If a single - // getStatus() throws (e.g., onStateChange callback error), the entire function - // throws before reaching our loop. This is an accepted limitation - per-item - // fault tolerance is not possible with the current getAllCircuitBreakerStatuses() - // API. The outer try/catch handles this case. + // Read-only: getAllCircuitBreakerSnapshots() calls peekStatus(), so listing connections + // reports an elapsed OPEN breaker's effective state without transitioning it, persisting + // that transition or consuming its half-open probe. Reading health must not let a request + // through to a provider that is still broken. + // NOTE: if a single peekStatus() throws, the whole call throws before reaching our loop. + // This is an accepted limitation - per-item fault tolerance is not possible with the + // current API. The outer try/catch handles this case. let breakers: BreakerWithHistory[] = []; try { - const allStatuses = getAllCircuitBreakerStatuses(); + const allStatuses = getAllCircuitBreakerSnapshots(); breakers = allStatuses.map((status) => ({ name: status.name, state: status.state, diff --git a/src/lib/a2a/skills/providerDiscovery.ts b/src/lib/a2a/skills/providerDiscovery.ts index 75c22dc7c2e..e3979126dbf 100644 --- a/src/lib/a2a/skills/providerDiscovery.ts +++ b/src/lib/a2a/skills/providerDiscovery.ts @@ -5,6 +5,7 @@ */ import type { A2ATask, TaskArtifact } from "../taskManager"; +import type { CircuitBreakerStatus } from "@/shared/utils/circuitBreaker"; import { AI_PROVIDERS, AUDIO_ONLY_PROVIDERS, @@ -119,7 +120,7 @@ export interface ProviderDiscoveryResult { } export async function executeProviderDiscovery(task: A2ATask): Promise { - const [{ getProviderConnections }, { getAllCircuitBreakerStatuses }] = await Promise.all([ + const [{ getProviderConnections }, { getAllCircuitBreakerSnapshots }] = await Promise.all([ import("@/lib/db/providers"), import("@/shared/utils/circuitBreaker"), ]); @@ -137,11 +138,11 @@ export async function executeProviderDiscovery(task: A2ATask): Promise connection.provider) .map((connection) => connection.provider as string) ); - const breakers = new Map( - getAllCircuitBreakerStatuses().map((breaker: CircuitBreakerLike) => [ - breaker.name || "", - breaker, - ]) + // Read-only: discovery only reports each provider's health, so it reads breakers through + // `peekStatus()`. Listing providers must not transition an elapsed OPEN breaker into HALF_OPEN, + // persist that, or hand the next request a probe into a provider that is still broken. + const breakers = new Map( + getAllCircuitBreakerSnapshots().map((breaker) => [breaker.name, breaker]) ); const candidates = Object.entries(AI_PROVIDERS) diff --git a/src/lib/monitoring/providerHealthAutopilot.ts b/src/lib/monitoring/providerHealthAutopilot.ts index 894553ab260..db6a513af6d 100644 --- a/src/lib/monitoring/providerHealthAutopilot.ts +++ b/src/lib/monitoring/providerHealthAutopilot.ts @@ -258,7 +258,7 @@ export async function buildProviderHealthAutopilotReport( const includeActions = options.includeActions !== false; const providerFilter = canonicalProviderId(options.provider); - const [{ getAllCircuitBreakerStatuses }, { getAllModelLockouts }, quotaMonitor] = + const [{ getAllCircuitBreakerSnapshots }, { getAllModelLockouts }, quotaMonitor] = await Promise.all([ import("@/shared/utils/circuitBreaker"), import("@omniroute/open-sse/services/accountFallback"), @@ -272,7 +272,7 @@ export async function buildProviderHealthAutopilotReport( const provider = canonicalProviderId(connection.provider); return provider && (!providerFilter || provider === providerFilter); }); - const breakers = getAllCircuitBreakerStatuses().filter((breaker) => { + const breakers = getAllCircuitBreakerSnapshots().filter((breaker) => { const name = toString(breaker.name); const provider = canonicalProviderId(name); if (!name || !provider || name.startsWith("test-") || name.startsWith("test_")) return false; diff --git a/src/lib/monitoring/providerHealthMatrix.ts b/src/lib/monitoring/providerHealthMatrix.ts index f9af936417a..32e06a6dbf5 100644 --- a/src/lib/monitoring/providerHealthMatrix.ts +++ b/src/lib/monitoring/providerHealthMatrix.ts @@ -1,7 +1,7 @@ import { getSyncedAvailableModelsByConnection } from "@/lib/db/models"; import { getProviderConnections } from "@/lib/db/providers"; import { getDbInstance } from "@/lib/db/core"; -import { getAllCircuitBreakerStatuses } from "@/shared/utils/circuitBreaker"; +import { getAllCircuitBreakerSnapshots } from "@/shared/utils/circuitBreaker"; import { getAllModelLockouts } from "@omniroute/open-sse/services/accountFallback"; import { resolveProviderAlias } from "@omniroute/open-sse/services/model"; import { getWebSessionPoolHealth } from "@omniroute/open-sse/services/webSessionPoolHealth"; @@ -360,7 +360,7 @@ export async function buildProviderHealthMatrix( // so a provider has one health row with every related signal attached. const [connections, breakers, lockouts, rawStats] = await Promise.all([ getProviderConnections({}), - getAllCircuitBreakerStatuses(), + getAllCircuitBreakerSnapshots(), getAllModelLockouts(), Promise.resolve(queryCallLogTargetStats(cutoff, null)), ]); diff --git a/src/lib/omnirouteStatus.ts b/src/lib/omnirouteStatus.ts index 74b8afdbaa5..8d00ea61729 100644 --- a/src/lib/omnirouteStatus.ts +++ b/src/lib/omnirouteStatus.ts @@ -25,7 +25,7 @@ export async function buildOmniRouteStatus() { import("../../open-sse/services/quotaMonitor").catch(() => null), ]); const pools = listPools().items; - const circuitStatuses = circuitModule?.getAllCircuitBreakerStatuses() ?? null; + const circuitStatuses = circuitModule?.getAllCircuitBreakerSnapshots() ?? null; const quotaSummary = quotaMonitorModule?.getQuotaMonitorSummary() ?? null; const active = connections.filter( (connection) => connection.is_active !== 0 && connection.is_active !== false diff --git a/src/shared/utils/circuitBreaker.ts b/src/shared/utils/circuitBreaker.ts index 95226e686a6..9ef8d502307 100644 --- a/src/shared/utils/circuitBreaker.ts +++ b/src/shared/utils/circuitBreaker.ts @@ -412,6 +412,18 @@ export class CircuitBreaker { return this.state; } + /** + * `canExecute()` without the probing transition: whether live routing would be let through if it + * asked right now. An OPEN breaker whose cooldown has elapsed would be granted a fresh probe + * budget by the next live call, so it reads as executable — but nothing transitions, persists or + * consumes a probe here. For health reads, which must not decide a breaker's fate by looking. + */ + peekCanExecute(): boolean { + if (this.state === STATE.CLOSED || this.state === STATE.DEGRADED) return true; + if (this.state === STATE.HALF_OPEN) return this.halfOpenAllowed > 0; + return this._shouldAttemptReset() && this.halfOpenRequests > 0; + } + getStatus(): CircuitBreakerStatus { this._refreshOpenState(); return this._statusAs(this.state); @@ -784,6 +796,44 @@ export function getAllCircuitBreakerSnapshots(): CircuitBreakerStatus[] { return snapshots; } +/** + * Read-only view of one breaker: the registered instance, or an unregistered copy built from its + * persisted state. Never registers a breaker, evicts one, or transitions one. Callers must only + * use the `peek*` accessors on the result — the mutating ones would write through the copy. + */ +export function peekCircuitBreaker(name: string): CircuitBreaker | null { + if (!name) return null; + const registered = registry.get(name); + if (registered) return registered; + try { + return new CircuitBreaker(name); + } catch { + return null; + } +} + +/** + * Snapshots of every breaker that would refuse a request right now — the read-only counterpart of + * filtering breakers by `!canExecute()`, which probes an elapsed OPEN breaker into HALF_OPEN and + * persists that. Covers registered breakers plus persisted ones not loaded in this process. + */ +export function getBlockedCircuitBreakerSnapshots(): CircuitBreakerStatus[] { + const blocked: CircuitBreakerStatus[] = []; + for (const cb of registry.values()) { + if (!cb.peekCanExecute()) blocked.push(cb.peekStatus()); + } + try { + for (const persisted of loadAllCircuitBreakerStates()) { + if (registry.has(persisted.name)) continue; + const breaker = new CircuitBreaker(persisted.name); + if (!breaker.peekCanExecute()) blocked.push(breaker.peekStatus()); + } + } catch { + // Use registry only + } + return blocked; +} + export function resetAllCircuitBreakers() { for (const cb of registry.values()) { cb.reset(); diff --git a/stryker.conf.json b/stryker.conf.json index 735b564ba25..bb116ab13ac 100644 --- a/stryker.conf.json +++ b/stryker.conf.json @@ -270,6 +270,7 @@ "tests/unit/guardrails/visionBridge-responses-9597.test.ts", "tests/unit/headroom-codex-quota-snapshot-6379.test.ts", "tests/unit/headroom-proxy-lifecycle.test.ts", + "tests/unit/health-reads-do-not-probe-breakers.test.ts", "tests/unit/idempotency-fusion-collision.test.ts", "tests/unit/isLocalStreamLifecycleError-abort-shape.test.ts", "tests/unit/issue-6343-v0-web-alias-collision.test.ts", diff --git a/tests/unit/health-reads-do-not-probe-breakers.test.ts b/tests/unit/health-reads-do-not-probe-breakers.test.ts new file mode 100644 index 00000000000..13cd3b25de0 --- /dev/null +++ b/tests/unit/health-reads-do-not-probe-breakers.test.ts @@ -0,0 +1,199 @@ +/** + * Reading health must not decide a circuit breaker's fate. + * + * `getAllCircuitBreakerStatuses()` calls `getStatus()`, which transitions an OPEN breaker whose + * cooldown has elapsed to HALF_OPEN and persists that, so merely *looking* at health flipped the + * breaker and handed the next request a probe through to a provider that is still broken. PR #38 + * fixed `/api/metrics` and the SLO timer with the read-only `getAllCircuitBreakerSnapshots()` / + * `peekStatus()` API; these are the remaining read paths, now converted. + * + * "Unchanged" here means the breaker's own state: it stays OPEN, records no transition, keeps no + * granted half-open probe, and its persisted row is untouched. The *reported* state is the + * effective one — HALF_OPEN for an elapsed breaker — which is what PR #38 established and what + * dashboards need to show; reporting it is not the same as causing it. + */ +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import type { A2ATask } from "../../src/lib/a2a/taskManager.ts"; + +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-breaker-reads-")); +process.env.DATA_DIR = TEST_DATA_DIR; + +const core = await import("../../src/lib/db/core.ts"); +const domainState = await import("../../src/lib/db/domainState.ts"); +const breakers = await import("../../src/shared/utils/circuitBreaker.ts"); +const accountFallback = await import("../../open-sse/services/accountFallback.ts"); +const webSessionPoolHealth = await import("../../open-sse/services/webSessionPoolHealth.ts"); +const omnirouteStatus = await import("../../src/lib/omnirouteStatus.ts"); +const healthMatrix = await import("../../src/lib/monitoring/providerHealthMatrix.ts"); +const autopilot = await import("../../src/lib/monitoring/providerHealthAutopilot.ts"); +const connectionsRoute = await import("../../src/app/api/resilience/connections/route.ts"); +const monitoringHealthRoute = await import("../../src/app/api/monitoring/health/route.ts"); +const providerDiscovery = await import("../../src/lib/a2a/skills/providerDiscovery.ts"); + +const PROVIDER = "breaker-read-only-provider"; + +/** The minimal A2A task `executeProviderDiscovery` needs to pick a capability. */ +function discoveryTask(): A2ATask { + const now = new Date().toISOString(); + return { + id: "task-breaker-read", + skill: "provider-discovery", + state: "working", + input: { skill: "provider-discovery", messages: [{ role: "user", content: "chat" }] }, + artifacts: [], + events: [], + metadata: {}, + createdAt: now, + updatedAt: now, + expiresAt: now, + }; +} + +test.after(() => { + breakers.resetAllCircuitBreakers(); + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 }); +}); + +let clock = Date.now(); + +/** An OPEN breaker whose cooldown has elapsed: the state a read used to probe out of. */ +async function openBreakerWithElapsedCooldown() { + breakers.resetAllCircuitBreakers(); + clock = Date.now(); + const breaker = breakers.getCircuitBreaker(PROVIDER, { + failureThreshold: 1, + resetTimeout: 1000, + now: () => clock, + }); + await assert.rejects( + breaker.execute(async () => { + throw new Error("upstream down"); + }) + ); + assert.equal(breaker.state, "OPEN"); + clock += 60_000; + return breaker; +} + +function persistedState(): string { + return JSON.stringify(domainState.loadCircuitBreakerState(PROVIDER)); +} + +/** Every read path converted away from the mutating `getAllCircuitBreakerStatuses()`. */ +const readPaths: Array<[string, () => Promise]> = [ + ["accountFallback.getProvidersInCooldown", async () => accountFallback.getProvidersInCooldown()], + [ + "webSessionPoolHealth", + async () => { + const report = webSessionPoolHealth.getWebSessionPoolHealth(PROVIDER); + assert.ok(report.providers[0]?.breaker, "the report really read the breaker"); + return report; + }, + ], + ["omnirouteStatus", () => omnirouteStatus.buildOmniRouteStatus()], + ["providerHealthMatrix", () => healthMatrix.buildProviderHealthMatrix({})], + ["providerHealthAutopilot", () => autopilot.buildProviderHealthAutopilotReport({})], + [ + "/api/resilience/connections", + async () => { + const response = await connectionsRoute.GET( + new Request("http://localhost/api/resilience/connections") + ); + assert.equal(response.status, 200, "the route answered, so it really read the breakers"); + return response.json(); + }, + ], + [ + "a2a providerDiscovery", + async () => { + const result = await providerDiscovery.executeProviderDiscovery(discoveryTask()); + assert.ok(result.metadata.totalCandidates > 0, "discovery really listed candidates"); + return result; + }, + ], + [ + "/api/monitoring/health", + async () => { + monitoringHealthRoute.__test_resetMonitoringHealthPayloadCache(); + const response = await monitoringHealthRoute.GET( + new Request("http://localhost/api/monitoring/health") + ); + assert.equal(response.status, 200, "the route answered, so it really read the breakers"); + return response.json(); + }, + ], +]; + +for (const [name, read] of readPaths) { + test(`${name} does not probe an OPEN breaker whose cooldown elapsed`, async () => { + const breaker = await openBreakerWithElapsedCooldown(); + const historyBefore = breaker.transitionHistory.length; + const persistedBefore = persistedState(); + assert.match(persistedBefore, /"OPEN"/, "the breaker was persisted OPEN"); + + await read(); + await read(); + + assert.equal(breaker.state, "OPEN", "the read did not move the breaker to HALF_OPEN"); + assert.equal(breaker.transitionHistory.length, historyBefore, "no transition was recorded"); + assert.equal(breaker.halfOpenAllowed, 0, "no half-open probe was granted by the read"); + assert.equal(persistedState(), persistedBefore, "the persisted state is unchanged"); + assert.equal(breaker.canExecute(), true, "live routing still gets its probe afterwards"); + assert.equal(breaker.state, "HALF_OPEN", "and live routing is what transitions the breaker"); + }); +} + +test("the read paths still report the breaker, with its effective state", async () => { + await openBreakerWithElapsedCooldown(); + + const connections = (await ( + await connectionsRoute.GET(new Request("http://localhost/api/resilience/connections")) + ).json()) as { breakers: Array<{ name: string; state: string }> }; + const reported = connections.breakers.find((entry) => entry.name === PROVIDER); + assert.equal(reported?.state, "HALF_OPEN", "the effective state is reported, not caused"); + + const status = (await omnirouteStatus.buildOmniRouteStatus()) as { + circuits: { halfOpen?: number }; + }; + assert.ok((status.circuits.halfOpen ?? 0) >= 1, "omnirouteStatus counted the breaker"); + + const pools = webSessionPoolHealth.getWebSessionPoolHealth(PROVIDER); + assert.equal(pools.providers[0]?.breaker?.state, "HALF_OPEN"); + assert.equal(pools.providers[0]?.breaker?.inCooldown, false, "an elapsed breaker is not blocked"); +}); + +test("a breaker still inside its cooldown is reported in cooldown, without mutating it", async () => { + breakers.resetAllCircuitBreakers(); + clock = Date.now(); + const breaker = breakers.getCircuitBreaker(PROVIDER, { + failureThreshold: 1, + resetTimeout: 600_000, + now: () => clock, + }); + await assert.rejects( + breaker.execute(async () => { + throw new Error("upstream down"); + }) + ); + const historyBefore = breaker.transitionHistory.length; + const persistedBefore = persistedState(); + + assert.ok( + accountFallback.getProvidersInCooldown().some((entry) => entry.provider === PROVIDER), + "an open breaker inside its cooldown is still listed" + ); + const blocked = breakers.getBlockedCircuitBreakerSnapshots(); + assert.ok(blocked.some((status) => status.name === PROVIDER)); + assert.ok((blocked.find((status) => status.name === PROVIDER)?.retryAfterMs ?? 0) > 0); + assert.equal(breakers.peekCircuitBreaker(PROVIDER)?.peekCanExecute(), false); + + assert.equal(breaker.state, "OPEN"); + assert.equal(breaker.transitionHistory.length, historyBefore); + assert.equal(persistedState(), persistedBefore); +}); diff --git a/tests/unit/routing-decision-compact-build.test.ts b/tests/unit/routing-decision-compact-build.test.ts new file mode 100644 index 00000000000..43df141b3f7 --- /dev/null +++ b/tests/unit/routing-decision-compact-build.test.ts @@ -0,0 +1,159 @@ +/** + * Live recording builds the routing decision already in the shape the decision store retains. + * + * `buildRoutingDecision` used to materialise every candidate with its full factor breakdown and + * leave it to `recordRoutingDecision` to throw all but 40 candidates / 10 factor breakdowns away, + * so an auto combo over the whole catalog paid hundreds of KiB and double-digit milliseconds of + * pure waste on every routed request. Passing the store's own bounds as `retention` builds the + * decision compact from the start. + * + * What must not change: + * - what `GET /api/omniroute/route/decisions/{id}` returns, byte for byte, for the same input; + * - the preview path, which must keep returning the full, uncompacted candidate list. + */ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + buildRoutingDecision, + previewRoutingDecision, + type BuildRoutingDecisionInput, + type DecisionCandidateInput, +} from "../../open-sse/services/autoCombo/routingDecision.ts"; +import { + previewSelectionDeps, + selectProviderWithTrace, + type AutoComboConfig, +} from "../../open-sse/services/autoCombo/engine.ts"; +import { DEFAULT_WEIGHTS } from "../../open-sse/services/autoCombo/scoring.ts"; +import { + getRoutingDecision, + MAX_CANDIDATES_WITH_FACTORS, + MAX_STORED_CANDIDATES, + recordRoutingDecision, + resetRoutingDecisionStore, +} from "../../open-sse/services/routing/decisionStore.ts"; + +const clock = { + now: () => Date.parse("2026-09-18T12:00:00.000Z"), + newDecisionId: () => "rd_compact_test", +}; + +const config = { + id: "compact-auto", + name: "compact-auto", + type: "auto", + candidatePool: [], + weights: DEFAULT_WEIGHTS, + explorationRate: 0, +} as unknown as AutoComboConfig; + +const request = { + requestId: "req-compact", + model: "compact-auto", + protocol: "messages", + stream: false, +}; + +const RETENTION = { + maxCandidates: MAX_STORED_CANDIDATES, + maxCandidatesWithFactors: MAX_CANDIDATES_WITH_FACTORS, +}; + +/** A pool wide enough that the store's bounds bite, with a mix of routable and excluded entries. */ +function candidatePool(count: number): DecisionCandidateInput[] { + const pool: DecisionCandidateInput[] = []; + for (let index = 0; index < count; index += 1) { + pool.push({ + provider: `provider-${index % 37}`, + model: `catalog-model-${index}`, + connectionId: `conn-${index}`, + quotaRemaining: 95 - (index % 90), + quotaTotal: 100, + circuitBreakerState: index % 23 === 0 ? "OPEN" : "CLOSED", + costPer1MTokens: 1 + (index % 19), + p95LatencyMs: 150 + (index % 1200), + latencyStdDev: 25 + (index % 65), + errorRate: (index % 13) / 200, + modelAvailable: index % 31 !== 0, + } as unknown as DecisionCandidateInput); + } + return pool; +} + +function buildInput(pool: DecisionCandidateInput[]): BuildRoutingDecisionInput { + const routable = pool.filter((candidate) => candidate.modelAvailable !== false); + const outcome = selectProviderWithTrace( + config, + routable, + "default", + undefined, + previewSelectionDeps() + ); + return { request, config, candidates: pool, outcome, liveRequestExecuted: true }; +} + +/** The decision as `GET /api/omniroute/route/decisions/{id}` serialises it. */ +function storedJson(input: BuildRoutingDecisionInput): string { + resetRoutingDecisionStore(); + recordRoutingDecision(buildRoutingDecision(input, clock), 1000); + const stored = getRoutingDecision("rd_compact_test", 1000); + assert.ok(stored, "the decision was recorded"); + return JSON.stringify(stored); +} + +test.afterEach(() => resetRoutingDecisionStore()); + +test("a decision built compact stores exactly what a full build stored", () => { + for (const size of [300, 50, MAX_STORED_CANDIDATES, 7]) { + const input = buildInput(candidatePool(size)); + const full = storedJson(input); + const compact = storedJson({ ...input, retention: RETENTION }); + assert.equal(compact, full, `stored decision differs for a pool of ${size} candidates`); + } +}); + +test("the compact build carries the store's omittedCandidates count", () => { + const input = buildInput(candidatePool(300)); + const compact = buildRoutingDecision({ ...input, retention: RETENTION }, clock); + + assert.equal(compact.candidates.length, MAX_STORED_CANDIDATES); + assert.equal(compact.omittedCandidates, 300 - MAX_STORED_CANDIDATES); + assert.equal( + compact.candidates.filter((candidate) => candidate.factors.length > 0).length, + MAX_CANDIDATES_WITH_FACTORS + ); + assert.ok(compact.selected, "a candidate was selected"); + assert.ok(compact.selected.factors.length > 0, "the selected candidate keeps its factors"); + assert.ok( + compact.candidates.some( + (candidate) => + candidate.providerId === compact.selected?.providerId && + candidate.modelId === compact.selected?.modelId + ), + "the selected candidate is among the retained ones" + ); +}); + +test("the retained candidates keep the order a full build produced", () => { + const input = buildInput(candidatePool(300)); + const full = buildRoutingDecision(input, clock); + const compact = buildRoutingDecision({ ...input, retention: RETENTION }, clock); + + assert.deepEqual( + compact.candidates.map((candidate) => [candidate.providerId, candidate.modelId]), + full.candidates + .slice(0, MAX_STORED_CANDIDATES) + .map((candidate) => [candidate.providerId, candidate.modelId]) + ); +}); + +test("preview still returns every candidate with its factor breakdown", () => { + const pool = candidatePool(300); + const decision = previewRoutingDecision({ request, config, candidates: pool }, clock); + + assert.equal(decision.candidates.length, 300); + assert.equal(decision.omittedCandidates, undefined); + const scored = decision.candidates.filter((candidate) => candidate.factors.length > 0); + assert.ok(scored.length > MAX_CANDIDATES_WITH_FACTORS, "preview is not compacted"); +});