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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 11 additions & 12 deletions open-sse/services/accountFallback.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -1302,24 +1302,23 @@ 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;
failureCount: number;
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,
}));
}

/**
Expand Down
109 changes: 99 additions & 10 deletions open-sse/services/autoCombo/routingDecision.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand All @@ -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<string, ScoredProvider>,
scored: ScoredProvider | undefined,
eligibleKeys: Set<string>
): 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),
Expand All @@ -206,6 +226,7 @@ function selectionModeOf(input: BuildRoutingDecisionInput): RoutingDecision["sel
interface DecisionEntry {
key: string;
candidate: DecisionCandidateInput;
scored: ScoredProvider | undefined;
routing: RoutingCandidate;
}

Expand Down Expand Up @@ -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,
Expand All @@ -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 } : {}),
};
}

Expand Down
19 changes: 18 additions & 1 deletion open-sse/services/combo/autoRoutingDecision.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -56,11 +60,24 @@ function requestProtocol(body: Record<string, unknown>): 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<BuildRoutingDecisionInput, "outcome" | "budgetExceeded" | "strategySelection">
): RoutingDecision {
const decision = buildRoutingDecision({
retention: LIVE_DECISION_RETENTION,
request: {
requestId: getRequestId() ?? "",
model: context.config.name,
Expand Down
24 changes: 17 additions & 7 deletions open-sse/services/webSessionPoolHealth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) ─────────────────────────────────
Expand All @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion src/app/api/monitoring/health/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ async function rebuildHealthPayload(): Promise<unknown> {
circuitBreakerModule.status === "fulfilled"
? readHealthValue(
"circuit breakers",
() => circuitBreakerModule.value.getAllCircuitBreakerStatuses(),
() => circuitBreakerModule.value.getAllCircuitBreakerSnapshots(),
[]
)
: [];
Expand Down
16 changes: 9 additions & 7 deletions src/app/api/resilience/connections/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
Expand Down
13 changes: 7 additions & 6 deletions src/lib/a2a/skills/providerDiscovery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
*/

import type { A2ATask, TaskArtifact } from "../taskManager";
import type { CircuitBreakerStatus } from "@/shared/utils/circuitBreaker";
import {
AI_PROVIDERS,
AUDIO_ONLY_PROVIDERS,
Expand Down Expand Up @@ -119,7 +120,7 @@ export interface ProviderDiscoveryResult {
}

export async function executeProviderDiscovery(task: A2ATask): Promise<ProviderDiscoveryResult> {
const [{ getProviderConnections }, { getAllCircuitBreakerStatuses }] = await Promise.all([
const [{ getProviderConnections }, { getAllCircuitBreakerSnapshots }] = await Promise.all([
import("@/lib/db/providers"),
import("@/shared/utils/circuitBreaker"),
]);
Expand All @@ -137,11 +138,11 @@ export async function executeProviderDiscovery(task: A2ATask): Promise<ProviderD
.filter((connection) => 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<string, CircuitBreakerStatus>(
getAllCircuitBreakerSnapshots().map((breaker) => [breaker.name, breaker])
);

const candidates = Object.entries(AI_PROVIDERS)
Expand Down
4 changes: 2 additions & 2 deletions src/lib/monitoring/providerHealthAutopilot.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand All @@ -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;
Expand Down
Loading
Loading