Skip to content
Open
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
71 changes: 39 additions & 32 deletions async/unstable_circuit_breaker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -770,45 +770,32 @@ export class CircuitBreaker<T = unknown> {
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;

Expand All @@ -823,4 +810,24 @@ export class CircuitBreaker<T = unknown> {

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);
}
}
50 changes: 50 additions & 0 deletions async/unstable_circuit_breaker_test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down
Loading