Files
OmniRoute/src/shared/utils/circuitBreaker.ts
Diego Rodrigues de Sa e Souza 555b21d296 Release v3.8.37 (#5053)
* chore(release): open v3.8.37 development cycle

* chore(ci): harden release flow — ratchet decoupling, fast-path drift gates, build-scope guard, heap default (#5054)

Implements improvements 1-4 from the v3.8.36 release benchmark (_tasks/release-bench/v3.8.36/PLANO-MELHORIA.md):

1. Quality Ratchet decoupled from flaky coverage (ci.yml): the shard→coverage→ratchet
   chain meant a single flaky Coverage Shard SKIPPED the whole Quality Ratchet on the
   release PR (v3.8.36 #4854), so cycle drift only surfaced post-merge in #5029. The job
   now runs on !cancelled(); coverage download is continue-on-error and the ratchet runs
   --allow-missing, so the DETERMINISTIC gates (eslint/complexity/cognitive/duplication/
   codeql) stay blocking even when coverage is unavailable.

2. Fast-path drift gates (quality.yml PR→release): added check:complexity,
   check:cognitive-complexity, and a new lightweight check:pack-policy (pack-artifact
   unexpected-files check WITHOUT a build, via --policy-only) so drift + stray-tarball-file
   regressions are caught/rebaselined PER-PR instead of cascading onto the release PR.

3. Build heap default 4096→8192 MB (build-next-isolated.mjs): the clean graph peaks
   ~3.9 GB and brushed the old 4 GB ceiling; 8 GB gives headroom. Comment notes heap is
   NOT the fix for a poisoned scope (run check:build-scope instead).

4. check:build-scope gate (new): fails if .ts/.tsx/.js/.jsx files in the tsconfig scope
   exceed a threshold — catches worktrees/cruft leaking into the build scope (the v3.8.36
   OOM root cause: 355,215 vs 4,547 files) BEFORE it detonates next build. Wired into the
   fast-path.

* fix(auth): only trust forwarding headers from loopback TCP peers (#4689)

Integrated into release/v3.8.37 — loopback-gated forwarding headers (IP spoofing fix). Cherry-picked onto current release tip; ipUtils.test.ts 9/9 green.

* fix(codex): treat OAuth 401 as unrecoverable refresh failure (#4686)

Integrated into release/v3.8.37 — codex OAuth 401 treated as unrecoverable refresh. Cherry-picked onto release tip; token-refresh-service.test.ts 38/38 green.

* fix(translator): preserve reasoning_effort for non-Copilot Responses clients (#4688)

Integrated into release/v3.8.37 — preserve reasoning_effort for non-Copilot Responses clients. Cherry-picked onto release tip; tests 47/47 green.

* fix(translator): coerce tool descriptions to strings in OpenAI normalization (#4675)

Integrated into release/v3.8.37 — coerce tool descriptions to strings in OpenAI normalization. Cherry-picked onto release tip; tests 3/3 green.

* feat(sse): x-omniroute-strip-reasoning header to drop reasoning_content (#4678)

Integrated into release/v3.8.37 — x-omniroute-strip-reasoning header. Cherry-picked onto release tip (resolved chatCore.ts/headers.ts adjacency conflict, kept resolveCompressionHeader + isStripReasoningRequested); tests 8/8 green.

* fix(combo): flatten Anthropic tool messages + tool history to prevent upstream 503 (#4648)

Integrated into release/v3.8.37 — flattenToolHistory helper (combo anti-503). Cherry-picked onto release tip; tests 9/9 green.

* feat(headroom): proxy lifecycle management + dashboard UI (Docker sidecar supported) (#4649)

Integrated into release/v3.8.37 — headroom proxy lifecycle (status/start/stop, local-only + spawn-capable per Rules #15/#17). Cherry-picked onto release tip; lifecycle 7/7 + route-guard 43/43 + check:cycles green.

* feat(cli): multi-model support for Factory Droid CLI (#4682)

Integrated into release/v3.8.37 — Factory Droid multi-model support. Cherry-picked onto release tip (kept readJsoncConfig + droidCustomModels imports); droid-custom-models 11/11 green.

* fix(providers): require Default Model in compatible-provider API-key setup (#4641)

Integrated into release/v3.8.37 — require Default Model in compatible-provider API-key setup. Cherry-picked fix + test-move onto release tip (kept release providerSpecificData + QuotaScrapingFields; fixed moved-test import path; baseline rebaseline unneeded, 865<866); UI test 2/2 green.

* fix(dashboard): stop double-masking already-masked API key in list (E2E 3/9 regression) (#4671)

Integrated into release/v3.8.37 — render server-masked key verbatim (drop redundant maskKey call). Note: release's maskKey already guards '****' (since v3.8.34), so this is a safe simplification; added a contract test pinning the **** passthrough invariant (2/2 green, would fail against the pre-guard maskKey = the historical double-mask bug).

* chore(quality): rebaseline file-size for rc17 PR batch drift

Own growth from the merged rc17 PRs (#4678/#4686/#4688) at existing chokepoints —
cohesive, not extractable:
- open-sse/handlers/responseSanitizer.ts 1103->1122 (SanitizeOpenAIResponseOptions + stripReasoning, #4678)
- open-sse/services/tokenRefresh.ts 2070->2090 (codex 401 unrecoverable-refresh guard, #4686)
- tests/unit/token-refresh-service.test.ts 1322->1353 (401 regression case, #4686)
- tests/unit/translator-openai-responses-req.test.ts 1047->1050 (reasoning_effort assertion, #4688)

* docs(env): document HEADROOM_URL in .env.example + ENVIRONMENT.md

The headroom proxy lifecycle (#4649) reads HEADROOM_URL (src/lib/headroom/detect.ts,
default http://localhost:8787) but it was missing from the env contract, tripping
check:env-doc-sync. Adds the var to both .env.example (commented, has a default) and
the Proxy Health table in ENVIRONMENT.md.

* fix(sse): stream writer mock abort() returns a Promise (#4788)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(cli): fall back to default data dir when DATA_DIR is not writable (#4767)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(oauth): verify Cursor installation on Linux before auto-import (#4770)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): track Ollama streaming usage from raw NDJSON chunks (#4754)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): strip enumDescriptions from antigravity tool schema (#4740)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): include low-level cause details in formatProviderError (#4741)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(translator): strip x-anthropic-billing-header in claude-to-openai (#4728)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): gate Kiro image attachments behind a Claude-capability check (#4763)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): read Antigravity usage from the response.usageMetadata envelope (#4785)

Integrated into release/v3.8.37 — Antigravity response.usageMetadata envelope. Cherry-picked onto release tip (resolved test-tail adjacency with #4754 Ollama block); usage-extractor 23/23 green.

* fix(api): fall back to existing access token for any OAuth provider on refresh failure (#4786)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(cli): verify launchd registration + skip self-SIGTERM in macOS autostart (#4765)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(executors): anthropic-compatible-* gateways get Bearer alongside x-api-key (#4729)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): json_schema fallback for OpenAI-compatible providers (#4766)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): use workos auth token shape for cline (#4787)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* feat(sse): parse Gemini CLI 429 retryDelay from structured RetryInfo (#4738)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; tests green.

* fix(sse): finalize tool_calls finish_reason on early stream end in OpenAI Responses translator (#4764)

Integrated into release/v3.8.37 — computeFinishReason finalizes tool_calls on early stream end (Responses translator). Cherry-picked onto release tip; responses-translation-fixes 29/29 green.

* test(sse): golden-lock provider.ts translate-path across all providers (#4734)

Integrated into release/v3.8.37 — golden-lock for provider.ts translate-path. Cherry-picked onto release tip; snapshot regenerated against the current provider set (UPDATE_GOLDEN=1, 167 entries); golden test 3/3 deterministic.

* chore(quality): rebaseline file-size for rc17 leva2 PR batch drift

Own growth from the merged leva2 PRs (cohesive, not extractable):
- src/lib/usage/providerLimits.ts 950->955 (#4786)
- open-sse/executors/default.ts NEW frozen @828 (#4729 + #4766 + #4787 header branches)
- open-sse/translator/request/openai-to-kiro.ts 807->814 (#4763)
- open-sse/translator/response/openai-responses.ts 923->937 (#4764)
- tests/unit/executor-default-base.test.ts 1339->1440 (#4766)
- tests/unit/translator-openai-to-kiro.test.ts 918->980 (#4763)

* fix(dashboard): align Engine Combos editor engines with API schema (#4955) (#5062)

The named-combos pipeline dropdown offered four engines (headroom,
session-dedup, ccr, llmlingua) that stackedPipelineStepSchema rejects, so
selecting one made PUT /api/context/combos/[id] return HTTP 400 while
saveCombo swallowed the non-OK response (if (!res.ok) return). Editing the
default 'Standard Savings' combo and changing an engine reproduced the 400.

- Add canonical STACKED_PIPELINE_ENGINE_INTENSITIES next to the schema as the
  single source of truth; the client dropdown imports it so it can never drift
  from the discriminated union the API validates against.
- Surface save errors and empty-name/empty-pipeline validation in the editor
  instead of failing silently.
- Add a parity unit test asserting the UI engine map equals the schema union
  and that every (engine, intensity) the UI emits is accepted.

* fix(sse): filter nameless hosted tools when converting Responses API to Chat format (#4789)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(dashboard): keep desktop sidebar visible via explicit CSS class (#4812)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip enumDescriptions from Antigravity tool schemas (#4813)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(dashboard): resolve passthrough model aliases by providerId in ModelSelectModal (#4815)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(oauth): allow per-connection refresh lead-time override via providerSpecificData.refreshLeadMs (#4818)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip X-Stainless-* headers and normalize SDK User-Agent for OpenAI-compatible endpoints (#4820)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip Gemini built-in tools when functionDeclarations present in Antigravity envelope (#4821)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(api): surface a Docker-localhost hint on provider-node validation connection errors (#4822)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): resolve bare model names to connection defaultModel before upstream calls (#4825)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(build): trace-include sql.js sql-wasm.wasm in standalone bundle (#4839)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip Composer <|final|> sentinel markers leaking after Composer reasoning (#4842)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(config): sync full SiliconFlow model list into registry (#4844)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): close reasoning before message content in Responses stream (#4848)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): reject unsupported Kiro [1m] context suffix (#4816)

Integrated into release/v3.8.37 — cherry-picked onto release tip; test-tail conflict with #4763 resolved (kept both image + [1m] test blocks); CHANGELOG re-merged; 29/29 green.

* fix(db): validate HuggingFace tokens via whoami-v2 auth probe (#4819)

Integrated into release/v3.8.37 — defining commit re-homed onto the god-file-split validation module (validateHuggingFaceProvider in validation/openaiFormat.ts + map wiring); 115/115 green.

* fix(sse): make anthropic-version default-guard case-insensitive (#4823)

Integrated into release/v3.8.37 — conflict with #4729 Bearer-fallback resolved (kept both Bearer fallback + case-insensitive anthropic-version guard); 48/48 green.

* fix(sse): sanitize Kiro tool schemas to avoid 400 "Improperly formed request" (#4847)

Integrated into release/v3.8.37 — conflict in kiro-to-openai.ts resolved (kept release fallbackToolCallId + adopted #1375 toolNameMap remap); 7/7 green.

* feat(sse): add GPT-4 to the GitHub Copilot provider (#4798)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(sse): add GPT-4o mini to GitHub Copilot provider (#4797)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(api): add MiniMax-M3 pricing row (#4814)

Integrated into release/v3.8.37 — pricing row re-homed onto god-file-split pricing/regional.ts (pricing.ts is now a barrel); 4/4 green.

* fix(cli): save runtime deps with --save-exact so a sibling install can't prune them (#4841)

Integrated into release/v3.8.37 — trayRuntime conflict resolved (kept release SYSTRAY_SPEC + added --save-exact); 2/2 green.

* fix(sse): preserve required fields in antigravity tool schemas (#4843)

Integrated into release/v3.8.37 — conflict resolved (kept #4740/#4813 enumDescriptions strip + typed normalizeSchemaTypes, added required-preservation helpers; test-tail merged keeping both enumDescriptions + required tests); 7/7 green.

* chore(quality): rebaseline file-size for rc17b leva3 PR batch drift

* fix(sse): strip reasoning blobs from agentic context to prevent O(n^2) token growth (#4849)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): unwrap Qoder HTTP 200 SSE error envelope so fallback can trigger (#4850)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): strip temperature for Claude models with extended thinking (#4853)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): emit valid concatenable kiro tool_calls.arguments deltas (#4855)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(sse): add toggleable tool-source diagnostics (#4856)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): redact api key from the AUTH debug log in the chat handler (#4858)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): forward AI SDK image parts in Responses translator (#4859)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): resolve custom combos by id and case-insensitive name (#4446) (#4869)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): exclude WS bridge controller-closed error from provider breaker (#4602) (#4870)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(providers): add xAI Grok inbound translators and thinking patcher (#4910)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(embeddings): add dimensions override field to embedding combos (#4913)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* feat(oauth): Codex bulk-import endpoint — POST /api/oauth/codex/import (#4914)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(antigravity): retry transient upstream failures (#4941)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): surface malformed HTTP-200 upstream responses (#4942)

Integrated into release/v3.8.37 — cherry-picked defining commit onto release tip; CHANGELOG re-merged; tests green.

* fix(sse): normalize Codex custom tools (apply_patch) to { input: string } schema (#4862)

Integrated into release/v3.8.37 — conflict in request/openai-responses.ts resolved (kept #4789 nameless-tool skip + added #1007 custom-tool {input:string} normalization); 48/48 green incl. #4789/#4859 regression.

* fix(sse): dense, deterministic output ordering in Responses API response.completed (#4906)

Integrated into release/v3.8.37 — manual integration with #4862 in response/openai-responses.ts (custom-tool funcItem + dense recordCompletedItem). Fixed a latent #4848 interaction: the close-reasoning-before-message guard force-closed <think>-tag reasoning prematurely, which dense output (#4906) then snapshotted as a partial buffer ("plan" vs "planning") — scoped the guard to native reasoning_content (!inThinking) in BOTH transformer + translator paths. Full Responses suite 203/203 green incl. #4848/#4862 regression.

* feat(sse): auto-promote successful combo model to position #1 (#4852)

Integrated into release/v3.8.37 — dropped the stale file-size-baseline.json hunk (re-derived against the rc17b rebaseline); code+test applied clean; 13/13 green.

* feat(providers): add Pioneer AI (Fastino Labs) provider (#4909)

Integrated into release/v3.8.37 — providers.ts apikey block re-homed onto god-file-split src/shared/constants/providers/apikey/frontier-labs.ts (inline APIKEY_PROVIDERS no longer exists); registry/pioneer + providers/index.ts applied clean; 6/6 green.

* add DGrid AI gateway provider (#4931)

Integrated into release/v3.8.37 — rebased the contributor's commit onto the release tip; providers.ts god-file-split conflict resolved by relocating the dgrid APIKEY_PROVIDERS entry into apikey/gateways.ts; CHANGELOG added. 7/7 green. Thanks @dgridOP!

* chore(quality): rebaseline file-size for rc17b leva4 PR batch drift

* docs(routing): sync combo strategy docs for Fusion (17 strategies) (#5067)

Fusion (16th strategy, panel fan-out + judge synthesis) and headroom
shipped but the strategy-count docs were stale (14/15) and omitted both.
Update every combo-strategy reference to the canonical 17, add fusion +
headroom to all strategy tables, and add a dedicated Fusion section to
AUTO-COMBO.md documenting judgeModel / fusionTuning config + an example.

- CLAUDE.md, README.md, FEATURES.md, RESILIENCE_GUIDE.md,
  ARCHITECTURE.md, OPEN_SSE_ARCHITECTURE.md, OMNIROUTE_VS_ALTERNATIVES.md,
  docs/README.md, request-pipeline.mmd: 14/15 -> 17, list fusion + headroom
- docs/routing/AUTO-COMBO.md: strategy table + new Fusion strategy section
- docs/openapi.yaml: add reset-window, headroom, fusion to the strategy enum

* fix(oauth): classify /api/oauth/cursor/auto-import as local-only (route-guard) (#5070)

The Cursor auto-import route runs execFile("which", ["cursor"]) to verify a
local Cursor install before importing credentials — a child-process spawn. The
check:route-guard-membership gate (Hard Rules #15/#17) flagged it as an
unclassified spawn-capable route: reachable past the loopback gate, an
RCE-via-tunnel surface (a leaked JWT over a tunnel could trigger the spawn).

Classify the specific path in LOCAL_ONLY_API_PREFIXES so loopback enforcement
runs unconditionally before any auth check. Scoped to the exact path — the rest
of /api/oauth/ (browser redirect/callback flows) stays remote-reachable.

TDD: added a failing-then-passing assertion in route-guard-local-prefix.test.ts
(classification + an over-broadening guard proving sibling OAuth paths stay
remote). check:route-guard-membership now reports 0 new gaps.

* chore(release): v3.8.37 — 2026-06-26

---------

Co-authored-by: dgridOP <dgrid_op@outlook.com>
2026-06-26 02:51:06 -03:00

612 lines
20 KiB
TypeScript

/**
* Circuit Breaker — FASE-04 Observability & Resilience (v2.0)
*
* Implements the circuit breaker pattern with:
* - States: CLOSED → DEGRADED → OPEN → HALF_OPEN → CLOSED
* - Adaptive backoff: resetTimeout escalates on repeated open→probe→open cycles
* - Failure-kind-aware thresholds: different limits per failure type
* - Progressive degradation: high failure rate triggers warning before full open
* - Transition history tracking for diagnostics
* - DB persistence via domainState.js
*
* States:
* CLOSED — Normal operation, requests pass through
* DEGRADED — Failure rate elevated, requests pass through but warnings logged
* OPEN — Requests are short-circuited
* HALF_OPEN — Probing: limited requests allowed to test recovery
*/
import {
saveCircuitBreakerState,
loadCircuitBreakerState,
loadAllCircuitBreakerStates,
deleteCircuitBreakerState,
deleteAllCircuitBreakerStates,
} from "../../lib/db/domainState";
import type { FailureKind } from "./classify429";
/**
* #4602 — Detect a LOCAL stream-lifecycle error that must NOT count as a
* whole-provider failure. The Codex WebSocket→SSE bridge can throw a bare
* `Invalid state: Controller is already closed` (an enqueue-after-close on our
* own ReadableStream controller). It carries no `statusCode`, so it defaults to
* HTTP 502 and would otherwise trip the provider circuit breaker — blacklisting
* the entire Codex provider for a bug that lives in our bridge, not upstream.
* Use this with the breaker's `isFailure` option so the bridge error is ignored
* by the provider breaker while genuine upstream 5xx failures still count.
*/
export function isLocalStreamLifecycleError(error: unknown): boolean {
if (!error) return false;
const message =
typeof error === "string"
? error
: typeof (error as { message?: unknown }).message === "string"
? ((error as { message: string }).message as string)
: "";
if (!message) return false;
return /controller is already closed/i.test(message);
}
export const STATE = {
CLOSED: "CLOSED",
DEGRADED: "DEGRADED",
OPEN: "OPEN",
HALF_OPEN: "HALF_OPEN",
} as const;
type CircuitState = (typeof STATE)[keyof typeof STATE];
/** Per-failure-kind threshold overrides */
interface FailureKindThresholds {
/** Max failures of this kind before escalating to next state */
threshold: number;
/** Cooldown override for this failure kind */
cooldown?: number;
/** Whether this failure kind should trigger immediate OPEN (skip DEGRADED) */
immediateOpen?: boolean;
}
interface CircuitBreakerOptions {
failureThreshold?: number;
resetTimeout?: number;
halfOpenRequests?: number;
onStateChange?: ((name: string, oldState: string, newState: string) => void) | null;
isFailure?: (error: unknown) => boolean;
cooldownByKind?: Partial<Record<FailureKind, number>>;
classifyError?: (error: unknown) => FailureKind | undefined;
/**
* Per-failure-kind thresholds.
* When set, different failure types have different limits.
*/
kindThresholds?: Partial<Record<FailureKind, Partial<FailureKindThresholds>>>;
/**
* Degradation threshold — failure count at which state becomes DEGRADED.
* Default: 60% of failureThreshold.
*/
degradationThreshold?: number;
/**
* Max backoff multiplier (exponential). Default: 16x resetTimeout.
*/
maxBackoffMultiplier?: number;
/**
* How many open→half_open→open cycles before escalating backoff.
* Default: 3.
*/
backoffEscalationCount?: number;
}
interface TransitionRecord {
from: string;
to: string;
timestamp: number;
failureCount: number;
reason?: string;
}
export class CircuitBreaker {
name: string;
failureThreshold: number;
resetTimeout: number;
halfOpenRequests: number;
onStateChange: ((name: string, oldState: string, newState: string) => void) | null;
isFailure: (error: unknown) => boolean;
state: CircuitState;
failureCount: number;
successCount: number;
lastFailureTime: number | null;
halfOpenAllowed: number;
cooldownByKind: Partial<Record<FailureKind, number>>;
classifyError: ((error: unknown) => FailureKind | undefined) | null;
lastFailureKind: FailureKind | null;
kindThresholds: Partial<Record<FailureKind, Partial<FailureKindThresholds>>>;
degradationThreshold: number;
maxBackoffMultiplier: number;
backoffEscalationCount: number;
/** Track failure counts per kind separately */
kindFailureCounts: Record<string, number>;
/** How many times has the breaker gone from OPEN → HALF_OPEN → OPEN */
openCycleCount: number;
/** State transition history */
transitionHistory: TransitionRecord[];
/** Max transition history entries */
maxTransitionHistory: number;
constructor(name: string, options: CircuitBreakerOptions = {}) {
this.name = name;
this.failureThreshold = options.failureThreshold ?? 5;
this.resetTimeout = options.resetTimeout ?? 30000;
this.halfOpenRequests = options.halfOpenRequests ?? 1;
this.onStateChange = options.onStateChange || null;
this.isFailure = options.isFailure || (() => true);
this.state = STATE.CLOSED;
this.failureCount = 0;
this.successCount = 0;
this.lastFailureTime = null;
this.halfOpenAllowed = 0;
this.cooldownByKind = options.cooldownByKind ?? {};
this.classifyError = options.classifyError ?? null;
this.lastFailureKind = null;
this.kindThresholds = options.kindThresholds ?? {};
this.degradationThreshold =
options.degradationThreshold ?? Math.ceil((this.failureThreshold * 60) / 100);
this.maxBackoffMultiplier = options.maxBackoffMultiplier ?? 16;
this.backoffEscalationCount = options.backoffEscalationCount ?? 3;
this.kindFailureCounts = {};
this.openCycleCount = 0;
this.transitionHistory = [];
this.maxTransitionHistory = 20;
this._restoreFromDb();
}
_restoreFromDb() {
try {
const saved = loadCircuitBreakerState(this.name);
if (saved) {
if (
saved.state === STATE.CLOSED ||
saved.state === STATE.DEGRADED ||
saved.state === STATE.OPEN ||
saved.state === STATE.HALF_OPEN
) {
this.state = saved.state;
}
this.failureCount = saved.failureCount;
this.lastFailureTime = saved.lastFailureTime;
const savedKind = saved.options?.lastFailureKind;
if (
savedKind === "rate_limit" ||
savedKind === "quota_exhausted" ||
savedKind === "transient"
) {
this.lastFailureKind = savedKind;
}
this.openCycleCount = (saved.options?.openCycleCount as number) ?? 0;
this.kindFailureCounts = (saved.options?.kindFailureCounts as Record<string, number>) ?? {};
if (this.state === STATE.HALF_OPEN) {
this.halfOpenAllowed = this.halfOpenRequests;
}
}
} catch {
// DB may not be ready yet (build phase)
}
}
_persistToDb() {
try {
saveCircuitBreakerState(this.name, {
state: this.state,
failureCount: this.failureCount,
lastFailureTime: this.lastFailureTime,
options: {
failureThreshold: this.failureThreshold,
resetTimeout: this.resetTimeout,
halfOpenRequests: this.halfOpenRequests,
lastFailureKind: this.lastFailureKind,
openCycleCount: this.openCycleCount,
kindFailureCounts: this.kindFailureCounts,
},
});
} catch {
// Non-critical
}
}
/**
* Get the effective reset timeout, escalated by open cycle count.
* Each open→half_open→open cycle multiplies the timeout.
*/
_effectiveResetTimeout(): number {
if (this.openCycleCount <= this.backoffEscalationCount) {
return this.resetTimeout;
}
const escalationFactor = Math.pow(2, this.openCycleCount - this.backoffEscalationCount);
return Math.min(
this.resetTimeout * escalationFactor,
this.resetTimeout * this.maxBackoffMultiplier
);
}
async execute<T>(fn: () => Promise<T>): Promise<T> {
this._refreshOpenState();
if (this.state === STATE.OPEN) {
throw new CircuitBreakerOpenError(
`Circuit breaker "${this.name}" is OPEN. Try again later.`,
this.name,
this._timeUntilReset()
);
}
if (this.state === STATE.HALF_OPEN && this.halfOpenAllowed <= 0) {
throw new CircuitBreakerOpenError(
`Circuit breaker "${this.name}" is HALF_OPEN, no more probe requests allowed.`,
this.name,
this._timeUntilReset()
);
}
if (this.state === STATE.HALF_OPEN) {
this.halfOpenAllowed--;
}
try {
const result = await fn();
this._onSuccess();
return result;
} catch (error) {
if (this.isFailure(error)) {
let kind: FailureKind | undefined;
if (this.classifyError) {
try {
kind = this.classifyError(error);
} catch {
kind = undefined;
}
}
this._onFailure(kind);
}
throw error;
}
}
canExecute() {
this._refreshOpenState();
if (this.state === STATE.CLOSED || this.state === STATE.DEGRADED) return true;
if (this.state === STATE.OPEN) return false;
if (this.state === STATE.HALF_OPEN) return this.halfOpenAllowed > 0;
return false;
}
getStatus() {
this._refreshOpenState();
return {
name: this.name,
state: this.state,
failureCount: this.failureCount,
lastFailureTime: this.lastFailureTime,
retryAfterMs: this.getRetryAfterMs(),
lastFailureKind: this.lastFailureKind,
openCycleCount: this.openCycleCount,
kindFailureCounts: { ...this.kindFailureCounts },
degradationThreshold: this.degradationThreshold,
effectiveResetTimeout: this._effectiveResetTimeout(),
};
}
getRetryAfterMs() {
this._refreshOpenState();
if (this.state === STATE.CLOSED || this.state === STATE.DEGRADED) return 0;
return this._timeUntilReset();
}
reset() {
this._transition(STATE.CLOSED, "manual-reset");
this.failureCount = 0;
this.successCount = 0;
this.lastFailureTime = null;
this.lastFailureKind = null;
this.openCycleCount = 0;
this.kindFailureCounts = {};
this._persistToDb();
}
// ─── Internal ─────────────────────────────────
_onSuccess() {
if (this.state === STATE.OPEN) {
this._transition(STATE.CLOSED, "success-recovery");
this.failureCount = 0;
this.successCount = 0;
this.lastFailureTime = null;
this.lastFailureKind = null;
this.openCycleCount = 0;
this.kindFailureCounts = {};
} else if (this.state === STATE.HALF_OPEN) {
this.successCount++;
this._transition(STATE.CLOSED, "probe-success");
this.failureCount = 0;
this.lastFailureKind = null;
this.openCycleCount = 0;
this.kindFailureCounts = {};
} else {
// CLOSED or DEGRADED: reset counts
this.failureCount = Math.max(0, this.failureCount - 1); // gradual recovery
if (this.state === STATE.DEGRADED && this.failureCount <= this.degradationThreshold) {
this._transition(STATE.CLOSED, "recovery");
}
}
this._persistToDb();
}
_onFailure(kind?: FailureKind | null) {
const failureKind = kind ?? null;
this.failureCount++;
this.lastFailureTime = Date.now();
this.lastFailureKind = failureKind;
// Track per-kind failure counts
if (failureKind) {
this.kindFailureCounts[failureKind] = (this.kindFailureCounts[failureKind] || 0) + 1;
}
// Check kind-specific thresholds
if (failureKind) {
const kindConfig = this.kindThresholds[failureKind];
if (kindConfig) {
const kindCount = this.kindFailureCounts[failureKind] || 0;
// Immediate open for critical failure kinds
if (kindConfig.immediateOpen && kindCount >= (kindConfig.threshold || 1)) {
this._openCircuit(failureKind);
return;
}
// Kind-specific threshold reached
if (kindCount >= (kindConfig.threshold || this.failureThreshold)) {
this._openCircuit(failureKind);
return;
}
}
}
// State transitions based on total failure count
if (this.state === STATE.OPEN) {
// Already OPEN — just update persistence
} else if (this.state === STATE.HALF_OPEN) {
// Probe failed: OPEN with cycle count escalation
this.openCycleCount++;
this._transition(STATE.OPEN, `probe-failed (cycle ${this.openCycleCount})`);
} else if (this.state === STATE.DEGRADED) {
// Degraded → Open when threshold reached
if (this.failureCount >= this.failureThreshold) {
this._openCircuit(failureKind);
}
} else {
// CLOSED → DEGRADED or OPEN
if (this.failureCount >= this.failureThreshold) {
this._openCircuit(failureKind);
} else if (this.failureCount >= this.degradationThreshold) {
this._transition(
STATE.DEGRADED,
`elevated-failures (${this.failureCount}/${this.failureThreshold})`
);
}
}
this._persistToDb();
}
_openCircuit(kind: FailureKind | null) {
this._transition(STATE.OPEN, kind ? `kind:${kind}` : undefined);
}
_shouldAttemptReset() {
if (!this.lastFailureTime) return true;
const cooldown = this._effectiveCooldown();
return Date.now() - this.lastFailureTime >= cooldown;
}
_effectiveCooldown() {
const baseTimeout = this._effectiveResetTimeout();
if (this.lastFailureKind !== null) {
const override = this.cooldownByKind[this.lastFailureKind];
if (typeof override === "number" && Number.isFinite(override) && override >= 0) {
return override;
}
}
return baseTimeout;
}
_timeUntilReset() {
if (!this.lastFailureTime) return 0;
const cooldown = this._effectiveCooldown();
return Math.max(0, cooldown - (Date.now() - this.lastFailureTime));
}
_refreshOpenState() {
if (this.state === STATE.OPEN && this._shouldAttemptReset()) {
this._transition(STATE.HALF_OPEN, "timeout-elapsed");
this._persistToDb();
}
}
_transition(newState: CircuitState, reason?: string) {
const oldState = this.state;
this.state = newState;
if (newState === STATE.HALF_OPEN) {
this.halfOpenAllowed = this.halfOpenRequests;
}
// Record transition
this.transitionHistory.push({
from: oldState,
to: newState,
timestamp: Date.now(),
failureCount: this.failureCount,
reason,
});
if (this.transitionHistory.length > this.maxTransitionHistory) {
this.transitionHistory.shift();
}
if (this.onStateChange && oldState !== newState) {
this.onStateChange(this.name, oldState, newState);
}
}
}
export class CircuitBreakerOpenError extends Error {
circuitName: string;
retryAfterMs: number;
constructor(message: string, circuitName: string, retryAfterMs: number) {
super(message);
this.name = "CircuitBreakerOpenError";
this.circuitName = circuitName;
this.retryAfterMs = retryAfterMs;
}
}
// ─── Registry ─────────────────────────────────────
const MAX_REGISTRY_SIZE = 500;
const registry = new Map<string, CircuitBreaker>();
/** Test-only: current number of registered circuit breakers. */
export function __getCircuitRegistrySizeForTests(): number {
return registry.size;
}
const _registrySweep = setInterval(() => {
const now = Date.now();
for (const [name, breaker] of registry) {
const status = breaker.getStatus();
if (
status.state === STATE.CLOSED &&
status.failureCount === 0 &&
(!status.lastFailureTime || now - status.lastFailureTime > 30 * 60 * 1000)
) {
registry.delete(name);
try {
deleteCircuitBreakerState(name);
} catch {}
}
}
}, 5 * 60_000);
if (typeof _registrySweep === "object" && "unref" in _registrySweep) {
(_registrySweep as { unref?: () => void }).unref?.();
}
/**
* Enforce MAX_REGISTRY_SIZE before inserting a new breaker. The cap was previously declared
* but never used — the only bound was the 5-min sweep, which evicts a breaker only if it is
* CLOSED, has zero failures, AND has been idle for >30 min. With high-cardinality breaker
* names that cap could be exceeded for up to 30 min. Evict idle CLOSED breakers (oldest first)
* to make room; never evict OPEN/HALF_OPEN breakers, since those carry meaningful state. A
* CLOSED breaker with zero failures is behaviorally identical to a freshly-created one, so
* evicting and lazily recreating it later changes nothing.
*/
function evictColdBreakersIfNeeded(): void {
if (registry.size < MAX_REGISTRY_SIZE) return;
const candidates: { name: string; lastFailureTime: number }[] = [];
for (const [name, breaker] of registry) {
const status = breaker.getStatus();
if (status.state === STATE.CLOSED && status.failureCount === 0) {
candidates.push({ name, lastFailureTime: status.lastFailureTime || 0 });
}
}
candidates.sort((a, b) => a.lastFailureTime - b.lastFailureTime);
const target = registry.size - MAX_REGISTRY_SIZE + 1;
for (let i = 0; i < candidates.length && i < target; i++) {
registry.delete(candidates[i].name);
try {
deleteCircuitBreakerState(candidates[i].name);
} catch {}
}
}
export function getCircuitBreaker(name: string, options?: CircuitBreakerOptions): CircuitBreaker {
if (!registry.has(name)) {
evictColdBreakersIfNeeded();
registry.set(name, new CircuitBreaker(name, options));
}
const breaker = registry.get(name)!;
if (options) {
if (typeof options.failureThreshold === "number") {
breaker.failureThreshold = options.failureThreshold;
}
if (typeof options.resetTimeout === "number") {
breaker.resetTimeout = options.resetTimeout;
}
if (typeof options.halfOpenRequests === "number") {
breaker.halfOpenRequests = options.halfOpenRequests;
if (breaker.state === STATE.HALF_OPEN) {
breaker.halfOpenAllowed = Math.min(breaker.halfOpenAllowed, breaker.halfOpenRequests);
}
}
if (typeof options.onStateChange === "function") {
breaker.onStateChange = options.onStateChange;
}
if (typeof options.isFailure === "function") {
breaker.isFailure = options.isFailure;
}
if (options.cooldownByKind) {
breaker.cooldownByKind = {
...breaker.cooldownByKind,
...options.cooldownByKind,
};
}
if (typeof options.classifyError === "function") {
breaker.classifyError = options.classifyError;
}
if (options.kindThresholds) {
breaker.kindThresholds = {
...breaker.kindThresholds,
...options.kindThresholds,
};
}
if (typeof options.degradationThreshold === "number") {
breaker.degradationThreshold = options.degradationThreshold;
}
if (typeof options.maxBackoffMultiplier === "number") {
breaker.maxBackoffMultiplier = options.maxBackoffMultiplier;
}
if (typeof options.backoffEscalationCount === "number") {
breaker.backoffEscalationCount = options.backoffEscalationCount;
}
breaker._persistToDb();
}
return breaker;
}
export function getAllCircuitBreakerStatuses() {
try {
const persisted = loadAllCircuitBreakerStates();
for (const cb of persisted) {
if (!registry.has(cb.name)) {
getCircuitBreaker(cb.name);
}
}
} catch {
// Use registry only
}
return Array.from(registry.values()).map((cb) => cb.getStatus());
}
export function resetAllCircuitBreakers() {
for (const cb of registry.values()) {
cb.reset();
}
registry.clear();
try {
deleteAllCircuitBreakerStates();
} catch {
// Non-critical
}
}