diff --git a/async/unstable_circuit_breaker.ts b/async/unstable_circuit_breaker.ts index 71a9ca1cad7a..5ba51a18aaf9 100644 --- a/async/unstable_circuit_breaker.ts +++ b/async/unstable_circuit_breaker.ts @@ -770,45 +770,32 @@ export class CircuitBreaker { previousState: CircuitState, ): void { this.#failures.increment(); - const totalRequests = this.#requests.total; - const failureCount = this.#failures.total; - - const shouldOpen = previousState === "half_open" || - (totalRequests >= this.#minimumThroughput && - failureCount / totalRequests >= this.#failureRateThreshold); - - const existingOpenedAt = this.#state.state === "open" - ? this.#state.openedAt - : undefined; - - if (shouldOpen) { - this.#state = { - ...this.#state, - state: "open", - openedAt: existingOpenedAt ?? Date.now(), - consecutiveSuccesses: 0, - }; - } else { - this.#state = { - ...this.#state, - consecutiveSuccesses: 0, - }; - } - + this.#state = { ...this.#state, consecutiveSuccesses: 0 }; this.#onFailure?.( error ?? new Error("Result classified as failure"), - failureCount, - totalRequests, + this.#failures.total, + this.#requests.total, ); - if (shouldOpen && existingOpenedAt === undefined) { - this.#onStateChange?.(previousState, "open"); - this.#onOpen?.(failureCount, totalRequests); + if ( + this.#state.state !== "open" && + (previousState === "half_open" || this.#exceedsFailureRate()) + ) { + this.#open(); } } - /** Records a success and potentially closes the circuit from half-open. */ + /** + * Records a success. Closes the circuit from half-open once enough + * successes accrue. In closed state, a success can still be the request + * that lifts the window over `minimumThroughput`, so the rate is evaluated. + */ #handleSuccess(previousState: CircuitState, generation: number): void { - if (previousState === "closed") return; + if (previousState === "closed") { + if (this.#state.state === "closed" && this.#exceedsFailureRate()) { + this.#open(); + } + return; + } if (this.#state.state !== "half_open") return; if (generation !== this.#halfOpenGeneration) return; @@ -823,4 +810,24 @@ export class CircuitBreaker { this.#state = { ...this.#state, consecutiveSuccesses: newSuccessCount }; } + + /** Whether the window has enough requests and the rate meets the threshold. */ + #exceedsFailureRate(): boolean { + const totalRequests = this.#requests.total; + return totalRequests >= this.#minimumThroughput && + this.#failures.total / totalRequests >= this.#failureRateThreshold; + } + + /** Transitions the circuit to open and notifies observers. */ + #open(): void { + const from = this.#state.state; + this.#state = { + ...this.#state, + state: "open", + openedAt: Date.now(), + consecutiveSuccesses: 0, + }; + this.#onStateChange?.(from, "open"); + this.#onOpen?.(this.#failures.total, this.#requests.total); + } } diff --git a/async/unstable_circuit_breaker_test.ts b/async/unstable_circuit_breaker_test.ts index 53ee5be958ce..b58f723bf7b6 100644 --- a/async/unstable_circuit_breaker_test.ts +++ b/async/unstable_circuit_breaker_test.ts @@ -490,6 +490,56 @@ Deno.test("CircuitBreaker.execute() evicts expired history before evaluating a d assertEquals(breaker.state, "closed"); }); +Deno.test("CircuitBreaker.execute() opens when a success brings the window to minimumThroughput", async () => { + const opens: [number, number][] = []; + const changes: string[] = []; + const breaker = new CircuitBreaker({ + failureRateThreshold: 0.5, + minimumThroughput: 3, + onOpen: (failures, requests) => opens.push([failures, requests]), + onStateChange: (from, to) => changes.push(`${from}->${to}`), + }); + + await failN(breaker, 2); + assertEquals(breaker.state, "closed"); + + await succeedN(breaker, 1); + assertEquals(breaker.state, "open"); + assertEquals(opens, [[2, 3]]); + assertEquals(changes, ["closed->open"]); +}); + +Deno.test("CircuitBreaker.execute() reports a tripping failure before opening", async () => { + const events: string[] = []; + const breaker = new CircuitBreaker({ + failureRateThreshold: 0.5, + minimumThroughput: 2, + onFailure: (_error, failures, requests) => + events.push(`failure ${failures}/${requests}`), + onStateChange: (from, to) => events.push(`${from}->${to}`), + onOpen: (failures, requests) => events.push(`open ${failures}/${requests}`), + }); + + await failN(breaker, 2); + assertEquals(events, [ + "failure 1/1", + "failure 2/2", + "closed->open", + "open 2/2", + ]); +}); + +Deno.test("CircuitBreaker.execute() stays closed when a success keeps the rate below the threshold", async () => { + const breaker = new CircuitBreaker({ + failureRateThreshold: 0.5, + minimumThroughput: 3, + }); + + await failN(breaker, 1); + await succeedN(breaker, 2); + assertEquals(breaker.state, "closed"); +}); + Deno.test("CircuitBreaker.execute() prevents stale half_open success from closing after concurrent failure", async () => { using time = new FakeTime();