GitNexus/gitnexus/test/unit/integrations/resilient-fetch.test.ts
Gergő Magyar 152a0506c9
feat: shared resilient-fetch (retries + circuit breaker) (#1448)
* feat: shared resilient-fetch (retries + circuit breaker)

Add a small, runtime-agnostic resilience layer in gitnexus-shared and
migrate every backend HTTP outbound call (CLI, MCP, wiki LLM, web → backend)
through it.

Helpers (gitnexus-shared/src/integrations/):

- retry.ts            — withRetry(fn, opts) with caller-supplied
                        retryability classification and full-jitter
                        exponential backoff.
- circuit-breaker.ts  — closed/open/half-open per-process breaker with
                        injectable clock, plus a keyed registry so
                        callers targeting the same endpoint share state.
- resilient-fetch.ts  — composed wrapper: retries 5xx + 429 + retryable
                        network throws, treats AbortSignal.timeout()
                        and 4xx (other than 429) as terminal, honors
                        Retry-After (capped at 30s), throws
                        CircuitOpenError when the breaker opens.

Migrations (no behaviour regression — all existing tests pass):

- gitnexus/src/core/embeddings/http-client.ts (covers analyze + MCP
  query path) — replaces inline linear-backoff retry.
- gitnexus/src/core/wiki/llm-client.ts — preserves Azure content-filter
  branch; resilientFetch handles 5xx/429.
- gitnexus-web/src/services/backend-client.ts (fetchWithTimeout helper)
  — small retry budget (2 attempts, 250–1500 ms) so a dead local
  backend still fails fast for the user.
- gitnexus-web/src/core/llm/settings-service.ts (OpenRouter model list).

Deliberately not migrated:

- gitnexus-web/src/services/backend-client.ts streamJob() — Server-Sent
  Events stream; the existing reconnect-with-Last-Event-ID logic is
  not unary-fetch shaped.
- gitnexus-web/src/components/SettingsPanel.tsx checkOllamaStatus() —
  one-shot health probe; retrying delays the "Ollama not running"
  error rather than improving UX.

41 new helper tests cover backoff math, breaker state transitions,
Retry-After parsing (delta-seconds + HTTP-date), 401/422 terminal
classification, and breaker fail-fast on three exhausted retry batches.

* fix(review): apply autofix feedback

Address Claude's two MEDIUM blocking findings on PR #1448 plus the
CodeQL SSRF false-positive flag.

- backend-client `fetchWithTimeout` now uses `AbortSignal.timeout()`
  merged with the caller's signal via `AbortSignal.any()`. Timer-fired
  aborts surface as `DOMException(name='TimeoutError')` so
  resilientFetch routes them through the terminal-network branch
  (no retry, no breaker hit), instead of incrementing the breaker
  for user-side network slowness.
- Method-aware retry budget in `fetchWithTimeout`: idempotent verbs
  (GET/HEAD/OPTIONS) keep the 2-attempt budget; POST/PATCH/PUT/DELETE
  default to single-attempt so a 5xx on `startAnalyze` cannot start
  a duplicate job. New `forceRetry` parameter for callers that
  know-idempotent mutations (e.g. DELETE of a known-deleted resource).
- `resilient-fetch.ts` carries a documented suppression for CodeQL
  js/server-side-request-forgery on the inner fetch call. Every
  concrete caller passes a hardcoded URL constant or a value from
  configuration (env vars, saved settings); user request input never
  flows into the URL parameter.
- New test file `backend-client-retry.test.ts` covers all three
  paths: GET retries on 503, POST does not retry, timeout does not
  increment the breaker.

* fix(resilient-fetch): address Codex adversarial findings

Closes the three blocking issues from Codex's review on PR #1448.

U1 — Add `recordNeutral()` to CircuitBreaker.
  Third outcome path that's an explicit no-op for state and the
  consecutive-failure counter. Distinct from `recordSuccess` (closes
  the breaker) and `recordFailure` (may open it). Used for outcomes
  that are neither evidence of backend health nor evidence of
  backend failure.

U2 — Route terminal-client / terminal-network through `recordNeutral`.
  Previously a 401 or local timeout called `recordSuccess`, which
  reset `consecutiveFailures` to 0. A 5xx → 401 → 5xx → 401 → 5xx
  sequence would NEVER trip the breaker because each 4xx in between
  erased the running count. Also classify external `AbortError` as
  terminal-network (was retryable-network), so caller-driven
  cancellation no longer retries against an already-aborted signal
  or counts toward breaker failures on exhaustion.

U3 — Per-origin breaker key in web `fetchWithTimeout`.
  Was hardcoded to `'web-backend'` even though `_backendUrl` is
  mutable via `setBackendUrl`. Switching backend URLs after a
  circuit tripped on host-A would strand the user during the full
  cooldown. Key is now `web-backend:<origin>`, so each backend URL
  gets its own breaker state.

Tests: +5 recordNeutral, +4 resilient-fetch (interleaved 4xx/5xx,
external AbortError, prior-state preservation), +1 web switch-backend
regression. All 70 gitnexus integration tests + 15 web tests green.

* fix(resilient-fetch): tolerate header-less fetch mocks on 429

`classifyOutcome` called `resp.headers.get('Retry-After')` directly,
which crashed when a test stubs `fetch` with a plain object like
`{ ok: false, status: 429 }` (no `headers` field). Real `Response`
always has Headers, so this surfaces only in test setups, but the
helper has no business assuming caller-side correctness on this — the
defensive guard is cheap and a missing `Retry-After` falls through to
exponential-backoff retry like any 429 without the header.

Surfaced by `gitnexus/test/unit/http-embedder.test.ts > retries on
rate limit`, which the embeddings migration exercises against a
plain-object 429 stub. Locked in with a new
`classifies 429 from a header-less fetch mock without throwing` case.

* fix(review): apply autofix feedback

Closes findings from the third multi-agent review pass on PR #1448.

#1 (P1) callLLM had no per-attempt timeout
  Wiki LLM calls passed no `signal` to resilientFetch; each of three
  retry attempts could hang indefinitely on a frozen TCP connection.
  Add `signal: AbortSignal.timeout(60_000)` so the per-attempt budget
  matches what http-client.ts and backend-client.ts already provide.

#2 (P2) drop dead `lastRetryableResp` post-loop fallback
  Variable was set in one switch arm but only read in unreachable code
  after the loop. The retry loop always returns/throws on every
  iteration. Keep only the defensive `throw` so TypeScript's
  control-flow analysis still sees `Promise<Response>` as the return.

#5 (P2) gate test-only exports behind a subpath
  `__resetBreakerRegistry__` and `classifyOutcome` were reachable from
  the main `gitnexus-shared` barrel — production code calling
  `__resetBreakerRegistry__` from a tool implementation would silently
  nuke every circuit breaker process-wide. Move to a new
  `gitnexus-shared/test-helpers` subpath export. Production callers
  see the cleaner public API; tests import via the explicit
  `gitnexus-shared/test-helpers` path.

#6 (P2) exhaustiveness guard on Outcome switch
  Add a `default: const _: never = outcome` arm so a future sixth
  `Outcome.kind` won't compile silently — it'll surface at the switch
  site rather than fall through to a retry/no-retry default.

#9 (P3) document cumulative wall-clock budget
  Add a "Cumulative wall-clock budget" paragraph to resilientFetch's
  JSDoc explaining the worst-case total wait (`maxAttempts × (per-attempt
  timeout + capDelayMs)` ≈ 60s with defaults) and pointing callers at
  outer `AbortSignal.timeout()` when they want a tighter bound.

Deferred to follow-up PRs (per review's Auto-resolve recommendation):
  - #3 idempotency knob to shared API (forceRetry into ResilientFetchOptions)
  - #4 publish.ts migration to resilientFetch
  - #7 parseRetryAfter past-HTTP-date / negative-seconds asymmetry
  - #8 recordNeutral counter time-decay (documented breaker semantic)

* fix(circuit-breaker): gate half-open to a single in-flight probe

Closes the Codex adversarial-review finding on PR #1448 that flagged a
recovery-time thundering herd: when cooldown expired, every concurrent
caller transitioned the breaker to half-open and probed the still-
recovering dependency in lockstep, defeating the breaker's "fail fast"
promise.

U1 — probe-permit gate in CircuitBreaker.check()
  Added a `probeInFlight: boolean` field. After cooldown expires, the
  first `check()` admits the probe and consumes the permit; subsequent
  callers throw `CircuitOpenError` with a configurable
  `halfOpenRetryAfterMs` (default 1000ms) until the probe resolves.

  Critical design point: `recordNeutral` now RELEASES the permit but
  does NOT transition state. Without that split, a single `TimeoutError`
  from per-attempt `AbortSignal.timeout` (which routes through neutral
  classification) would permanently park the breaker in half-open. By
  separating permit-release from state-resolution, we keep the
  "neutral doesn't claim health" semantic without creating that wedge.

  Other changes:
  - `halfOpenRetryAfterMs` is now a constructor option for consumers
    with long-running protected ops (LLM streaming, large uploads).
  - `getState()` is documented as a pure read; the implicit
    Open -> Half-Open transition lives in `check()` only, so tests
    that inspect state never inadvertently consume a probe permit.
  - `isProbeInFlight()` test-only accessor for assertion clarity.
  - JSDoc on `check()` records the JS event-loop atomicity dependency
    and the load-bearing `try/finally` pairing invariant.

U2 — End-to-end concurrency regression through resilientFetch
  Three new scenarios in resilient-fetch.test.ts (26 -> 29):
  - 3 concurrent calls + probe gets 200 -> 1 hits fetch, 2 throw
    CircuitOpenError, breaker closes.
  - 3 concurrent calls + probe gets 503 -> ResilientFetchExhaustedError
    on probe; concurrent callers see halfOpenRetryAfterMs (1000ms);
    fresh caller after probe resolves sees the FULL new cooldown
    (10000ms), not the probe-in-flight default.
  - Probe cancelled mid-flight via AbortError -> permit released,
    state stays half-open, next caller becomes the new probe and
    succeeds.

Plus 9 new circuit-breaker unit tests (16 -> 25) covering the permit
gate, recordNeutral-releases-permit semantic, fresh-cooldown distinction,
default vs configurable halfOpenRetryAfterMs, getState() purity, and
the three-probes-via-neutrals chain.

Total integration test count: 70 -> 82. All 106 gitnexus + 15 web
tests pass; both packages typecheck.

Maintainer decisions (deferred per plan 003 Open Questions):
  - Plan 002's deferral judgement was reversed on Codex's argument
    without new measurement / incident data. The reversal is defensible
    on principle (Hystrix / Resilience4j alignment) but lacks workload-
    driven evidence.
  - Probe-blocked callers throw silently (no log / event hook). R4's
    "no new public API" prevents adding observability; loosen if a
    debug log on probe-blocked is wanted.

* refactor(embeddings): replace bespoke HF breaker with shared CircuitBreaker

Deleted the local `HfDownloadCircuitBreaker` class and the manual
retry loop in `withHfDownloadRetry`. Both are now backed by the
shared `gitnexus-shared` primitives:

- `hfDownloadCircuit` is `new CircuitBreaker({ failureThreshold,
  cooldownMs, key: 'hf-download' })` — same state machine as before
  PLUS the single-permit half-open gate that prevents recovery-time
  stampedes when CLI + MCP embedders concurrently re-load the model.
- `withHfDownloadRetry` delegates the loop to `withRetry` from the
  shared package. Per-attempt timeout (`withDownloadTimeout`),
  network-vs-non-network classification, circuit recording, and the
  `onRetry` callback wire through `withRetry`'s `isRetryable`
  callback.

Behaviour preserved:
- Pre-flight `CIRCUIT_OPEN_TAG` rejection when the breaker is open.
- Mid-loop `CIRCUIT_OPEN_TAG` "opened after N consecutive failures"
  when a network error trips the threshold.
- Non-network errors (e.g. CUDA unavailable) bypass retry and go
  through `recordNeutral` instead of resetting the breaker's
  failure-count progress.
- `onRetry(attempt+1, max, err)` fires only when there's a next
  attempt, matching the prior semantic.

Generic CircuitBreaker gained two inspection accessors:
- `getOpenedAt(): number | null`
- `getCooldownMs(): number`
Used by `withHfDownloadRetry` to compute `secsUntilReset` without
consuming a probe permit (which `check()` would do).

Test consolidation: the 7 bespoke `HfDownloadCircuitBreaker`
state-machine tests in hf-env.test.ts were 1:1 duplicates of
existing tests in `circuit-breaker.test.ts` and were deleted.
Remaining 42 hf-env tests all pass; full integration sweep (148
gitnexus + 15 web) green.
2026-05-09 15:18:09 +01:00

558 lines
21 KiB
TypeScript

import { describe, it, expect, beforeEach, vi } from 'vitest';
import {
CircuitBreaker,
CircuitOpenError,
parseRetryAfter,
resilientFetch,
ResilientFetchExhaustedError,
RETRY_AFTER_CAP_MS,
} from 'gitnexus-shared';
import { __resetBreakerRegistry__, classifyOutcome } from 'gitnexus-shared/test-helpers';
describe('parseRetryAfter', () => {
it('parses delta-seconds form', () => {
expect(parseRetryAfter('30')).toBe(30_000);
expect(parseRetryAfter('0')).toBe(0);
});
it('returns null on negative or non-numeric garbage', () => {
expect(parseRetryAfter(null)).toBeNull();
expect(parseRetryAfter('')).toBeNull();
expect(parseRetryAfter(' ')).toBeNull();
expect(parseRetryAfter('not-a-number')).toBeNull();
});
it('parses HTTP-date form against an injected clock', () => {
const now = () => Date.parse('Wed, 21 Oct 2025 07:28:00 GMT');
expect(parseRetryAfter('Wed, 21 Oct 2025 07:28:30 GMT', now)).toBe(30_000);
});
it('returns 0 (not negative) on past HTTP-date', () => {
const now = () => Date.parse('Wed, 21 Oct 2025 08:00:00 GMT');
expect(parseRetryAfter('Wed, 21 Oct 2025 07:28:00 GMT', now)).toBe(0);
});
});
describe('classifyOutcome', () => {
const now = () => 1_700_000_000_000;
it('classifies 2xx as success', () => {
const resp = new Response(null, { status: 204 });
const out = classifyOutcome({ kind: 'response', resp }, now);
expect(out.kind).toBe('success');
});
it('classifies 5xx as retryable-status without afterMs', () => {
const resp = new Response(null, { status: 503 });
const out = classifyOutcome({ kind: 'response', resp }, now);
expect(out.kind).toBe('retryable-status');
if (out.kind === 'retryable-status') expect(out.afterMs).toBeUndefined();
});
it('classifies 429 with Retry-After (capped) as retryable-status', () => {
const resp = new Response(null, { status: 429, headers: { 'Retry-After': '99999' } });
const out = classifyOutcome({ kind: 'response', resp }, now);
expect(out.kind).toBe('retryable-status');
if (out.kind === 'retryable-status') expect(out.afterMs).toBe(RETRY_AFTER_CAP_MS);
});
it('classifies 429 from a header-less fetch mock without throwing', () => {
// Tests sometimes stub `fetch` with a plain `{ ok, status }` object
// (e.g. http-embedder.test.ts). Real `Response` always carries
// `Headers`, but the helper must not crash when the stub does not.
// Falls through to exponential-backoff retry like a 429 with no
// Retry-After header.
const resp = { ok: false, status: 429 } as unknown as Response;
const out = classifyOutcome({ kind: 'response', resp }, now);
expect(out.kind).toBe('retryable-status');
if (out.kind === 'retryable-status') expect(out.afterMs).toBeUndefined();
});
it('classifies 401/403/404/422 as terminal-client', () => {
for (const status of [401, 403, 404, 422, 400]) {
const resp = new Response(null, { status });
const out = classifyOutcome({ kind: 'response', resp }, now);
expect(out.kind).toBe('terminal-client');
}
});
it('classifies TimeoutError as terminal-network', () => {
const err = new DOMException('aborted', 'TimeoutError');
const out = classifyOutcome({ kind: 'error', err }, now);
expect(out.kind).toBe('terminal-network');
});
it('classifies generic network throw as retryable-network', () => {
const err = new TypeError('fetch failed');
const out = classifyOutcome({ kind: 'error', err }, now);
expect(out.kind).toBe('retryable-network');
});
});
describe('resilientFetch', () => {
const URL_STR = 'https://example.test/api/dispatch';
beforeEach(() => __resetBreakerRegistry__());
function jsonResp(status: number, headers?: Record<string, string>): Response {
return new Response(null, { status, headers });
}
function makeBreaker(opts: Partial<ConstructorParameters<typeof CircuitBreaker>[0]> = {}) {
let t = 1_700_000_000_000;
const breaker = new CircuitBreaker({
failureThreshold: 3,
cooldownMs: 30_000,
key: 'test',
now: () => t,
...opts,
});
return { breaker, advance: (ms: number) => (t += ms) };
}
it('204 returns immediately, no retries, breaker stays closed', async () => {
const fetchImpl = vi.fn(async () => jsonResp(204));
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
const resp = await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep },
});
expect(resp.status).toBe(204);
expect(fetchImpl).toHaveBeenCalledTimes(1);
expect(sleep).not.toHaveBeenCalled();
expect(breaker.getState()).toBe('closed');
expect(breaker.getConsecutiveFailures()).toBe(0);
});
it('one 503 then 204 → retried once, returns 204, breaker stays closed', async () => {
let n = 0;
const fetchImpl = vi.fn(async () => {
n += 1;
return n === 1 ? jsonResp(503) : jsonResp(204);
});
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
const resp = await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep, random: () => 0.5, baseDelayMs: 100, capDelayMs: 1000 },
});
expect(resp.status).toBe(204);
expect(fetchImpl).toHaveBeenCalledTimes(2);
expect(sleep).toHaveBeenCalledTimes(1);
expect(breaker.getConsecutiveFailures()).toBe(0);
});
it('429 with Retry-After honored (capped at RETRY_AFTER_CAP_MS)', async () => {
let n = 0;
const fetchImpl = vi.fn(async () => {
n += 1;
return n === 1 ? jsonResp(429, { 'Retry-After': '1' }) : jsonResp(204);
});
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep },
});
expect(sleep).toHaveBeenCalledWith(1000); // 1s
});
it('429 with absurd Retry-After is capped to RETRY_AFTER_CAP_MS', async () => {
let n = 0;
const fetchImpl = vi.fn(async () => {
n += 1;
return n === 1 ? jsonResp(429, { 'Retry-After': '99999' }) : jsonResp(204);
});
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep, capDelayMs: 999_999 }, // ensure cap comes from RETRY_AFTER_CAP_MS, not retry config
});
expect(sleep).toHaveBeenCalledWith(RETRY_AFTER_CAP_MS);
});
it('429 without Retry-After falls back to exponential-backoff delay', async () => {
let n = 0;
const fetchImpl = vi.fn(async () => {
n += 1;
return n === 1 ? jsonResp(429) : jsonResp(204);
});
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep, baseDelayMs: 100, capDelayMs: 1000, random: () => 0.5 },
});
// attempt 0: full-jitter upper = min(1000, 100*1) = 100; floor(0.5*100) = 50
expect(sleep).toHaveBeenCalledWith(50);
});
it('401 returned as Response, no retry, breaker not incremented', async () => {
const fetchImpl = vi.fn(async () => jsonResp(401));
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
const resp = await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep },
});
expect(resp.status).toBe(401);
expect(fetchImpl).toHaveBeenCalledTimes(1);
expect(breaker.getConsecutiveFailures()).toBe(0);
});
it('422 returned as Response, no retry', async () => {
const fetchImpl = vi.fn(async () => jsonResp(422));
const { breaker } = makeBreaker();
const resp = await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {} },
});
expect(resp.status).toBe(422);
expect(fetchImpl).toHaveBeenCalledTimes(1);
});
it('TimeoutError rethrown immediately, no retry, breaker not incremented', async () => {
const fetchImpl = vi.fn(async () => {
throw new DOMException('aborted', 'TimeoutError');
});
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep },
}),
).rejects.toThrow(DOMException);
expect(fetchImpl).toHaveBeenCalledTimes(1);
expect(sleep).not.toHaveBeenCalled();
expect(breaker.getConsecutiveFailures()).toBe(0);
});
it('three consecutive 503 throws ResilientFetchExhaustedError; breaker increments by 1', async () => {
const fetchImpl = vi.fn(async () => jsonResp(503));
const sleep = vi.fn(async () => {});
const { breaker } = makeBreaker();
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep, maxAttempts: 3 },
}),
).rejects.toBeInstanceOf(ResilientFetchExhaustedError);
expect(fetchImpl).toHaveBeenCalledTimes(3);
expect(breaker.getConsecutiveFailures()).toBe(1);
});
it('after three exhausted 503 batches, breaker opens and fails fast', async () => {
const fetchImpl = vi.fn(async () => jsonResp(503));
const { breaker } = makeBreaker({ failureThreshold: 3, cooldownMs: 60_000 });
for (let i = 0; i < 3; i++) {
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 3 },
}),
).rejects.toBeInstanceOf(ResilientFetchExhaustedError);
}
expect(breaker.getState()).toBe('open');
// 4th call: breaker open, no fetch invoked.
const fetchCallsBefore = fetchImpl.mock.calls.length;
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 3 },
}),
).rejects.toBeInstanceOf(CircuitOpenError);
expect(fetchImpl.mock.calls.length).toBe(fetchCallsBefore);
});
it('retryable-network error retries, breaker counts only on exhaustion', async () => {
const fetchImpl = vi.fn(async () => {
throw new TypeError('fetch failed');
});
const { breaker } = makeBreaker();
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 3 },
}),
).rejects.toBeInstanceOf(TypeError);
expect(fetchImpl).toHaveBeenCalledTimes(3);
expect(breaker.getConsecutiveFailures()).toBe(1);
});
describe('U2: terminal outcomes route through recordNeutral', () => {
it('401 does not erase prior partial-failure progress on the breaker', async () => {
const { breaker } = makeBreaker();
// Pre-seed the breaker with 2 failures (still closed; threshold 3).
breaker.recordFailure();
breaker.recordFailure();
expect(breaker.getConsecutiveFailures()).toBe(2);
const fetchImpl = vi.fn(async () => jsonResp(401));
await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {} },
});
// Counter MUST stay at 2 — under the old behaviour recordSuccess
// would have reset to 0 and the next 5xx batch would have started
// from scratch instead of tipping over the threshold.
expect(breaker.getConsecutiveFailures()).toBe(2);
expect(breaker.getState()).toBe('closed');
});
it('TimeoutError does not erase prior partial-failure progress', async () => {
const { breaker } = makeBreaker();
breaker.recordFailure();
breaker.recordFailure();
const fetchImpl = vi.fn(async () => {
throw new DOMException('aborted by timeout', 'TimeoutError');
});
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {} },
}),
).rejects.toBeInstanceOf(DOMException);
expect(breaker.getConsecutiveFailures()).toBe(2);
});
it('external AbortError is terminal: no retry, breaker untouched', async () => {
const { breaker } = makeBreaker();
breaker.recordFailure();
const fetchImpl = vi.fn(async () => {
throw new DOMException('aborted by caller', 'AbortError');
});
const sleep = vi.fn(async () => {});
await expect(
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep, maxAttempts: 3 },
}),
).rejects.toMatchObject({ name: 'AbortError' });
expect(fetchImpl).toHaveBeenCalledTimes(1);
expect(sleep).not.toHaveBeenCalled();
// Counter unchanged — neither incremented (no failure) nor reset
// (no synthetic success).
expect(breaker.getConsecutiveFailures()).toBe(1);
});
it('interleaved 5xx + 401 + 5xx + 401 + 5xx opens breaker on third real failure', async () => {
const { breaker } = makeBreaker({ failureThreshold: 3 });
const sequence = [503, 401, 503, 401, 503];
let i = 0;
const fetchImpl = vi.fn(async () => jsonResp(sequence[i++]));
// Each call uses maxAttempts:1 so each surfaces a single response
// (5xx → ResilientFetchExhaustedError; 4xx → returned Response).
const driveOne = () =>
resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
await expect(driveOne()).rejects.toBeInstanceOf(ResilientFetchExhaustedError); // 5xx fail #1
await driveOne(); // 401 neutral
await expect(driveOne()).rejects.toBeInstanceOf(ResilientFetchExhaustedError); // 5xx fail #2
await driveOne(); // 401 neutral
await expect(driveOne()).rejects.toBeInstanceOf(ResilientFetchExhaustedError); // 5xx fail #3 → opens
expect(breaker.getState()).toBe('open');
expect(fetchImpl).toHaveBeenCalledTimes(5);
});
});
describe('half-open single-probe gating (U2)', () => {
/** Test helper: a fetch mock whose Response is controlled by the test. */
function deferredFetch(): {
promise: Promise<Response>;
resolve: (resp: Response) => void;
reject: (err: unknown) => void;
} {
let resolve!: (resp: Response) => void;
let reject!: (err: unknown) => void;
const promise = new Promise<Response>((res, rej) => {
resolve = res;
reject = rej;
});
return { promise, resolve, reject };
}
/** Builds a clock-injected breaker pre-opened with cooldown elapsed. */
function preOpenedBreaker(opts: { cooldownMs: number; halfOpenRetryAfterMs?: number }): {
breaker: CircuitBreaker;
advance: (ms: number) => void;
} {
let t = 1_700_000_000_000;
const breaker = new CircuitBreaker({
failureThreshold: 1,
cooldownMs: opts.cooldownMs,
halfOpenRetryAfterMs: opts.halfOpenRetryAfterMs ?? 1_000,
key: 'test',
now: () => t,
});
breaker.recordFailure();
t += opts.cooldownMs + 1; // cooldown elapsed
return { breaker, advance: (ms) => (t += ms) };
}
it('happy: 3 concurrent calls — exactly 1 hits fetch, others throw CircuitOpenError', async () => {
const { breaker } = preOpenedBreaker({ cooldownMs: 10 });
const deferred = deferredFetch();
const fetchImpl = vi.fn(() => deferred.promise);
// Synchronous portion of each `resilientFetch` runs eagerly up to
// the first await, so by the time r2/r3 are constructed the probe
// permit is already consumed by r1 and they reject synchronously.
const r1 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
const r2 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
const r3 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
// Resolve the probe with 200; r1 should now settle.
deferred.resolve(new Response(null, { status: 200 }));
const results = await Promise.allSettled([r1, r2, r3]);
expect(results[0].status).toBe('fulfilled');
if (results[0].status === 'fulfilled') {
expect(results[0].value.status).toBe(200);
}
expect(results[1].status).toBe('rejected');
if (results[1].status === 'rejected') {
expect(results[1].reason).toBeInstanceOf(CircuitOpenError);
}
expect(results[2].status).toBe('rejected');
if (results[2].status === 'rejected') {
expect(results[2].reason).toBeInstanceOf(CircuitOpenError);
}
// Only ONE underlying fetch was invoked.
expect(fetchImpl).toHaveBeenCalledTimes(1);
// Breaker closed after the probe's success.
expect(breaker.getState()).toBe('closed');
});
it('error: probe gets 503 — exhausted error; subsequent caller sees fresh full cooldown', async () => {
const { breaker } = preOpenedBreaker({ cooldownMs: 10_000, halfOpenRetryAfterMs: 1_000 });
const deferred = deferredFetch();
const fetchImpl = vi.fn(() => deferred.promise);
const r1 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
const r2 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
const r3 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
// Probe fails with 503 → exhausted (maxAttempts: 1) → recordFailure → reopen.
deferred.resolve(new Response(null, { status: 503 }));
const results = await Promise.allSettled([r1, r2, r3]);
expect(results[0].status).toBe('rejected');
if (results[0].status === 'rejected') {
expect(results[0].reason).toBeInstanceOf(ResilientFetchExhaustedError);
}
expect(results[1].status).toBe('rejected');
if (results[1].status === 'rejected') {
expect(results[1].reason).toBeInstanceOf(CircuitOpenError);
// Blocked-while-half-open used the halfOpenRetryAfterMs default.
expect((results[1].reason as CircuitOpenError).retryAfterMs).toBe(1_000);
}
// Breaker has re-opened with a fresh openedAt.
expect(breaker.getState()).toBe('open');
// r4: should see the fresh full cooldown, NOT the probe-in-flight 1000ms.
let r4Caught: CircuitOpenError | null = null;
try {
await resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
} catch (err) {
r4Caught = err as CircuitOpenError;
}
expect(r4Caught).toBeInstanceOf(CircuitOpenError);
expect(r4Caught?.retryAfterMs).toBe(10_000);
});
it('cancellation: probe AbortError releases permit; next caller becomes new probe', async () => {
const { breaker } = preOpenedBreaker({ cooldownMs: 10_000 });
const deferred1 = deferredFetch();
const deferred2 = deferredFetch();
let callIdx = 0;
const fetchImpl = vi.fn(() => (callIdx++ === 0 ? deferred1.promise : deferred2.promise));
// r1 admitted as the probe; r2 blocked while r1 still in flight.
const r1 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
const r2 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
await expect(r2).rejects.toBeInstanceOf(CircuitOpenError);
// Cancel the probe — `AbortError` routes through terminal-network →
// `recordNeutral` → permit released, state stays half-open.
deferred1.reject(new DOMException('aborted by caller', 'AbortError'));
await expect(r1).rejects.toMatchObject({ name: 'AbortError' });
expect(breaker.isProbeInFlight()).toBe(false);
expect(breaker.getState()).toBe('half-open');
// r3: now succeeds and becomes the new probe.
const r3 = resilientFetch(URL_STR, undefined, {
fetchImpl: fetchImpl as unknown as typeof fetch,
breaker,
retry: { sleep: async () => {}, maxAttempts: 1 },
});
deferred2.resolve(new Response(null, { status: 200 }));
const r3Resp = await r3;
expect(r3Resp.status).toBe(200);
expect(breaker.getState()).toBe('closed');
// Two fetches total: the cancelled probe + the recovery probe.
expect(fetchImpl).toHaveBeenCalledTimes(2);
});
});
});