diff --git a/.changeset/fair-panthers-repair.md b/.changeset/fair-panthers-repair.md new file mode 100644 index 00000000..56e53aab --- /dev/null +++ b/.changeset/fair-panthers-repair.md @@ -0,0 +1,5 @@ +--- +'@3loop/transaction-decoder': patch +--- + +Fix circuit breaker outcome tracking and atomic state transitions, and use a bounded request-pool outcome window for adaptive concurrency. diff --git a/packages/transaction-decoder/src/abi-loader.ts b/packages/transaction-decoder/src/abi-loader.ts index 1d814d48..d7aef3cd 100644 --- a/packages/transaction-decoder/src/abi-loader.ts +++ b/packages/transaction-decoder/src/abi-loader.ts @@ -202,7 +202,7 @@ export const AbiLoaderRequestResolver = RequestResolver.makeBatched((requests: A } } - const concurrency = Math.min(...[...concurrencyMap.values(), 50]) // Use minimum concurrency across all chains, capped at 25 + const concurrency = Math.min(...[...concurrencyMap.values(), 50]) // Use minimum concurrency across all chains, capped at 50 yield* Effect.logDebug(`Executing ${remaining.length} remaining requests with concurrency ${concurrency}`) diff --git a/packages/transaction-decoder/src/circuit-breaker/circuit-breaker.ts b/packages/transaction-decoder/src/circuit-breaker/circuit-breaker.ts index efdddb35..c7519d96 100644 --- a/packages/transaction-decoder/src/circuit-breaker/circuit-breaker.ts +++ b/packages/transaction-decoder/src/circuit-breaker/circuit-breaker.ts @@ -326,93 +326,146 @@ export const make = ( const getState = (strategyId: string): Effect.Effect => Ref.get(states).pipe(Effect.map((map) => map.get(strategyId) ?? defaultState)) - const updateState = (strategyId: string, newState: CircuitBreaker.CircuitBreakerState) => - Ref.update(states, (map) => new Map(map).set(strategyId, newState)) + // Create tripping strategy if provided, otherwise use default failure count + const trippingStrategy = finalConfig.strategy + ? yield* finalConfig.strategy + : yield* failureCount(finalConfig.maxFailures ?? 5) - const shouldAllowRequest = (state: CircuitBreaker.CircuitBreakerState): Effect.Effect => - Effect.gen(function* () { - const now = yield* Clock.currentTimeMillis + type Admission = 'Allowed' | 'Rejected' | 'HalfOpened' + + const tryAdmit = (strategyId: string, now: number): Effect.Effect => + Ref.modify(states, (map): [Admission, Map] => { + const state = map.get(strategyId) ?? defaultState switch (state.state) { case 'Closed': - return true - case 'Open': - return now - state.lastFailureTime >= Duration.toMillis(finalConfig.resetTimeout ?? Duration.seconds(60)) + return ['Allowed', map] case 'HalfOpen': - return state.halfOpenCalls < (finalConfig.halfOpenMaxCalls ?? 3) + if (state.halfOpenCalls >= (finalConfig.halfOpenMaxCalls ?? 3)) { + return ['Rejected', map] + } + + return [ + 'Allowed', + new Map(map).set(strategyId, { + ...state, + halfOpenCalls: state.halfOpenCalls + 1, + }), + ] + case 'Open': + if (now - state.lastFailureTime < Duration.toMillis(finalConfig.resetTimeout ?? Duration.seconds(60))) { + return ['Rejected', map] + } + + return [ + 'HalfOpened', + new Map(map).set(strategyId, { + ...state, + state: 'HalfOpen', + halfOpenCalls: 1, + }), + ] } }) - // Create tripping strategy if provided, otherwise use default failure count - const trippingStrategy = finalConfig.strategy - ? yield* finalConfig.strategy - : yield* failureCount(finalConfig.maxFailures ?? 5) - - const onSuccess = (strategyId: string, state: CircuitBreaker.CircuitBreakerState) => + const onSuccess = (strategyId: string) => Effect.gen(function* () { - if (state.state === 'HalfOpen') { + yield* trippingStrategy.shouldTrip(true) + + const closedCircuit = yield* Ref.modify( + states, + (map): [boolean, Map] => { + const state = map.get(strategyId) ?? defaultState + + if (state.state === 'HalfOpen') { + return [ + true, + new Map(map).set(strategyId, { + ...defaultState, + state: 'Closed', + }), + ] + } + + if (state.failures > 0) { + return [ + false, + new Map(map).set(strategyId, { + ...state, + failures: 0, + }), + ] + } + + return [false, map] + }, + ) + + if (closedCircuit) { // Reset to closed state after successful half-open calls yield* trippingStrategy.onReset - yield* notifyStateChange(state.state, 'Closed') - yield* updateState(strategyId, { - ...defaultState, - state: 'Closed', - }) + yield* notifyStateChange('HalfOpen', 'Closed') yield* withMetrics((metrics) => Metric.increment(metrics.stateChanges).pipe( Effect.zipRight(Metric.set(metrics.state, stateToCode('Closed'))), ), ) - } else if (state.failures > 0) { - // Reset failures on success - yield* updateState(strategyId, { - ...state, - failures: 0, - }) } }) - const onFailure = (strategyId: string, state: CircuitBreaker.CircuitBreakerState) => + const onFailure = (strategyId: string) => Effect.gen(function* () { const now = yield* Clock.currentTimeMillis const shouldTrip = yield* trippingStrategy.shouldTrip(false) - if (state.state === 'HalfOpen') { - // Failed during half-open, go back to open - yield* notifyStateChange(state.state, 'Open') - yield* updateState(strategyId, { - ...state, - state: 'Open', - failures: state.failures + 1, - lastFailureTime: now, - halfOpenCalls: 0, - }) + const openedFrom = yield* Ref.modify( + states, + (map): [CircuitBreaker.State | undefined, Map] => { + const state = map.get(strategyId) ?? defaultState + + if (state.state === 'HalfOpen') { + return [ + 'HalfOpen', + new Map(map).set(strategyId, { + ...state, + state: 'Open', + failures: state.failures + 1, + lastFailureTime: now, + halfOpenCalls: 0, + }), + ] + } + + if (shouldTrip && state.state === 'Closed') { + return [ + 'Closed', + new Map(map).set(strategyId, { + ...state, + state: 'Open', + failures: state.failures + 1, + lastFailureTime: now, + }), + ] + } + + return [ + undefined, + new Map(map).set(strategyId, { + ...state, + failures: state.failures + 1, + lastFailureTime: now, + }), + ] + }, + ) + + if (openedFrom !== undefined) { + yield* notifyStateChange(openedFrom, 'Open') yield* withMetrics((metrics) => Metric.increment(metrics.stateChanges).pipe( Effect.zipRight(Metric.set(metrics.state, stateToCode('Open'))), ), ) - } else if (shouldTrip && state.state === 'Closed') { - // Threshold reached, open the circuit - yield* notifyStateChange(state.state, 'Open') - yield* updateState(strategyId, { - ...state, - state: 'Open', - failures: state.failures + 1, - lastFailureTime: now, - }) - yield* withMetrics((metrics) => - Metric.increment(metrics.stateChanges).pipe( - Effect.zipRight(Metric.set(metrics.state, stateToCode('Open'))), - ), - ) - } else { - // Increment failures but keep closed - yield* updateState(strategyId, { - ...state, - failures: state.failures + 1, - lastFailureTime: now, - }) } }) @@ -421,38 +474,27 @@ export const make = ( effect: Effect.Effect, ): Effect.Effect => Effect.gen(function* () { - const state = yield* getState(strategyId) - const shouldAllow = yield* shouldAllowRequest(state) + const now = yield* Clock.currentTimeMillis + const admission = yield* tryAdmit(strategyId, now) - if (!shouldAllow) { + if (admission === 'Rejected') { yield* withMetrics((metrics) => Metric.increment(metrics.rejectedCalls)) return yield* Effect.fail(OpenError(strategyId)) } - // Transition to half-open if we're allowing a request from open state - if (state.state === 'Open') { - yield* notifyStateChange(state.state, 'HalfOpen') - yield* updateState(strategyId, { - ...state, - state: 'HalfOpen', - halfOpenCalls: 1, - }) + if (admission === 'HalfOpened') { + yield* notifyStateChange('Open', 'HalfOpen') yield* withMetrics((metrics) => Metric.increment(metrics.stateChanges).pipe( Effect.zipRight(Metric.set(metrics.state, stateToCode('HalfOpen'))), ), ) - } else if (state.state === 'HalfOpen') { - yield* updateState(strategyId, { - ...state, - halfOpenCalls: state.halfOpenCalls + 1, - }) } const result = yield* Effect.either(effect) if (Either.isRight(result)) { - yield* onSuccess(strategyId, yield* getState(strategyId)) + yield* onSuccess(strategyId) yield* withMetrics((metrics) => Metric.increment(metrics.successfulCalls)) return result.right } else { @@ -460,14 +502,13 @@ export const make = ( const shouldCountFailure = finalConfig.isFailure ? finalConfig.isFailure(result.left) : true if (shouldCountFailure) { - yield* onFailure(strategyId, yield* getState(strategyId)) + yield* onFailure(strategyId) yield* withMetrics((metrics) => Metric.increment(metrics.failedCalls)) } return yield* Effect.fail(result.left) } }) - const currentState = (strategyId: string): Effect.Effect => getState(strategyId).pipe(Effect.map((state) => state.state)) diff --git a/packages/transaction-decoder/src/circuit-breaker/request-pool.ts b/packages/transaction-decoder/src/circuit-breaker/request-pool.ts index 626a6b25..8dab8934 100644 --- a/packages/transaction-decoder/src/circuit-breaker/request-pool.ts +++ b/packages/transaction-decoder/src/circuit-breaker/request-pool.ts @@ -12,8 +12,7 @@ export interface RequestPoolState { readonly activeRequests: number readonly maxConcurrency: number readonly successRate: number - readonly totalRequests: number - readonly successfulRequests: number + readonly recentOutcomes: ReadonlyArray } export interface RequestPool { @@ -33,8 +32,7 @@ const defaultState: RequestPoolState = { activeRequests: 0, maxConcurrency: 10, // Start conservative successRate: 1.0, - totalRequests: 0, - successfulRequests: 0, + recentOutcomes: [], } export const make = (configParam: Partial = {}): Effect.Effect => @@ -49,12 +47,6 @@ export const make = (configParam: Partial = {}): Effect.Effec // Setup metrics if provided const metrics = Option.fromNullable(config.metricLabels).pipe(Option.map(makeRequestPoolMetrics)) - const getState = (chainID: number): Effect.Effect => - Ref.get(poolStates).pipe(Effect.map((map) => map.get(chainID) ?? defaultState)) - - const updateState = (chainID: number, newState: RequestPoolState) => - Ref.update(poolStates, (map) => new Map(map).set(chainID, newState)) - const incrementActive = (chainID: number) => Effect.gen(function* () { yield* Ref.update(activeCounters, (map) => { @@ -116,26 +108,38 @@ export const make = (configParam: Partial = {}): Effect.Effec const getOptimalConcurrency = (chainID: number): Effect.Effect => Effect.gen(function* () { - const state = yield* getState(chainID) - const optimalConcurrency = calculateOptimalConcurrency(state) - - // Update the max concurrency in state if it changed - if (optimalConcurrency !== state.maxConcurrency) { - yield* updateState(chainID, { - ...state, - maxConcurrency: optimalConcurrency, - }) + const adjustment = yield* Ref.modify( + poolStates, + ( + map, + ): [{ readonly optimalConcurrency: number; readonly changed: boolean }, Map] => { + const state = map.get(chainID) ?? defaultState + const optimalConcurrency = calculateOptimalConcurrency(state) + const changed = optimalConcurrency !== state.maxConcurrency + + return [ + { optimalConcurrency, changed }, + changed + ? new Map(map).set(chainID, { + ...state, + maxConcurrency: optimalConcurrency, + }) + : map, + ] + }, + ) + if (adjustment.changed) { // Track concurrency adjustment in metrics yield* withRequestPoolMetrics(metrics, (m) => Effect.all([ Metric.increment(m.concurrencyAdjustments), - Metric.set(m.maxConcurrency, optimalConcurrency), + Metric.set(m.maxConcurrency, adjustment.optimalConcurrency), ]).pipe(Effect.asVoid), ) } - return optimalConcurrency + return adjustment.optimalConcurrency }) const withPoolManagement = ( @@ -161,22 +165,26 @@ export const make = (configParam: Partial = {}): Effect.Effec const updateMetrics = (chainID: number, success: boolean): Effect.Effect => Effect.gen(function* () { - const state = yield* getState(chainID) - const newTotalRequests = state.totalRequests + 1 - const newSuccessfulRequests = state.successfulRequests + (success ? 1 : 0) - - // Use sliding window to prevent metrics from becoming stale const windowSize = 100 - const effectiveTotalRequests = Math.min(newTotalRequests, windowSize) - const effectiveSuccessfulRequests = Math.min(newSuccessfulRequests, windowSize) - const effectiveSuccessRate = effectiveSuccessfulRequests / effectiveTotalRequests - - yield* updateState(chainID, { - ...state, - totalRequests: effectiveTotalRequests, - successfulRequests: effectiveSuccessfulRequests, - successRate: effectiveSuccessRate, - }) + const updated = yield* Ref.modify( + poolStates, + (map): [{ readonly successRate: number; readonly maxConcurrency: number }, Map] => { + const state = map.get(chainID) ?? defaultState + const recentOutcomes = state.recentOutcomes.slice(-(windowSize - 1)) + recentOutcomes.push(success) + const successful = recentOutcomes.reduce((count, outcome) => count + (outcome ? 1 : 0), 0) + const successRate = recentOutcomes.length === 0 ? 1 : successful / recentOutcomes.length + + return [ + { successRate, maxConcurrency: state.maxConcurrency }, + new Map(map).set(chainID, { + ...state, + recentOutcomes, + successRate, + }), + ] + }, + ) // Update metrics yield* withRequestPoolMetrics(metrics, (m: RequestPoolMetrics) => @@ -186,11 +194,11 @@ export const make = (configParam: Partial = {}): Effect.Effec } else { yield* Metric.increment(m.failedRequests) } - yield* Metric.set(m.successRate, effectiveSuccessRate) + yield* Metric.set(m.successRate, updated.successRate) // Determine and update pool state const activeCount = yield* Ref.get(activeCounters).pipe(Effect.map((map) => map.get(chainID) ?? 0)) - const poolState = determinePoolState(effectiveSuccessRate, activeCount, state.maxConcurrency) + const poolState = determinePoolState(updated.successRate, activeCount, updated.maxConcurrency) yield* Metric.set(m.poolState, poolStateToCode(poolState)) }), ) diff --git a/packages/transaction-decoder/test/circuit-breaker-with-metrics.test.ts b/packages/transaction-decoder/test/circuit-breaker-with-metrics.test.ts index 4fc184af..05e8d250 100644 --- a/packages/transaction-decoder/test/circuit-breaker-with-metrics.test.ts +++ b/packages/transaction-decoder/test/circuit-breaker-with-metrics.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect } from 'vitest' -import { Effect, MetricLabel, Duration, Either, TestContext } from 'effect' +import { Deferred, Duration, Effect, Either, Fiber, MetricLabel, Ref, TestContext } from 'effect' import * as CircuitBreaker from '../src/circuit-breaker/circuit-breaker.js' describe('CircuitBreaker with Metrics', () => { @@ -180,4 +180,102 @@ describe('CircuitBreaker with Metrics', () => { const result = await Effect.runPromise(testEffect.pipe(Effect.provide(TestContext.TestContext))) expect(result).toBe(true) }) + + it('resets the consecutive failure count after a success', async () => { + const testEffect = Effect.gen(function* () { + const circuitBreaker = yield* CircuitBreaker.make({ + strategy: CircuitBreaker.failureCount(2), + }) + const strategyId = 'consecutive-failures' + + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + yield* circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(false)) + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + + expect(yield* circuitBreaker.currentState(strategyId)).toBe('Closed') + + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + expect(yield* circuitBreaker.currentState(strategyId)).toBe('Open') + }) + + await Effect.runPromise(testEffect.pipe(Effect.provide(TestContext.TestContext))) + }) + + it('includes successful calls in the failure-rate window', async () => { + const testEffect = Effect.gen(function* () { + const circuitBreaker = yield* CircuitBreaker.make({ + strategy: CircuitBreaker.failureRate({ + threshold: 0.75, + minCalls: 4, + windowSize: 4, + }), + }) + const strategyId = 'failure-rate' + + yield* circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(false)) + yield* circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(false)) + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + + expect(yield* circuitBreaker.currentState(strategyId)).toBe('Closed') + + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + expect(yield* circuitBreaker.currentState(strategyId)).toBe('Open') + }) + + await Effect.runPromise(testEffect.pipe(Effect.provide(TestContext.TestContext))) + }) + + it('atomically limits concurrent half-open probes', async () => { + const testEffect = Effect.gen(function* () { + const firstTransitionStarted = yield* Deferred.make() + const secondTransitionStarted = yield* Deferred.make() + const releaseTransition = yield* Deferred.make() + const secondCompleted = yield* Deferred.make() + const transitionCount = yield* Ref.make(0) + const executedProbes = yield* Ref.make(0) + + const circuitBreaker = yield* CircuitBreaker.make({ + maxFailures: 1, + resetTimeout: Duration.millis(0), + halfOpenMaxCalls: 1, + onStateChange: (change) => + change.from === 'Open' && change.to === 'HalfOpen' + ? Effect.gen(function* () { + const count = yield* Ref.updateAndGet(transitionCount, (current) => current + 1) + yield* Deferred.succeed(count === 1 ? firstTransitionStarted : secondTransitionStarted, undefined) + yield* Deferred.await(releaseTransition) + }) + : Effect.void, + }) + const strategyId = 'half-open-limit' + const probe = circuitBreaker.withCircuitBreaker( + strategyId, + Ref.update(executedProbes, (count) => count + 1), + ) + + yield* Effect.either(circuitBreaker.withCircuitBreaker(strategyId, mockApiCall(true))) + expect(yield* circuitBreaker.currentState(strategyId)).toBe('Open') + + const firstFiber = yield* Effect.fork(probe) + yield* Deferred.await(firstTransitionStarted) + + const secondFiber = yield* Effect.fork( + Effect.either(probe).pipe(Effect.tap(() => Deferred.succeed(secondCompleted, undefined))), + ) + yield* Effect.race(Deferred.await(secondCompleted), Deferred.await(secondTransitionStarted)) + yield* Deferred.succeed(releaseTransition, undefined) + + yield* Fiber.join(firstFiber) + const secondResult = yield* Fiber.join(secondFiber) + + expect(yield* Ref.get(executedProbes)).toBe(1) + expect(Either.isLeft(secondResult)).toBe(true) + if (Either.isLeft(secondResult)) { + expect(CircuitBreaker.isCircuitBreakerOpenError(secondResult.left)).toBe(true) + } + }) + + await Effect.runPromise(testEffect.pipe(Effect.provide(TestContext.TestContext))) + }) }) diff --git a/packages/transaction-decoder/test/request-pool.test.ts b/packages/transaction-decoder/test/request-pool.test.ts index 0ad9959a..8020a9ce 100644 --- a/packages/transaction-decoder/test/request-pool.test.ts +++ b/packages/transaction-decoder/test/request-pool.test.ts @@ -295,4 +295,31 @@ describe('RequestPool with Metrics', () => { const result = await Effect.runPromise(testEffect.pipe(Effect.provide(TestContext.TestContext))) expect(result).toBe(true) }) + + it('adapts to failures after the outcome window is saturated', async () => { + const testEffect = Effect.gen(function* () { + const pool = yield* makeRequestPool({ + maxConcurrentRequests: 20, + adaptiveConcurrency: true, + healthThreshold: 0.8, + concurrencyStep: 2, + }) + const chainId = 1 + + yield* pool.getOptimalConcurrency(chainId) + for (let i = 0; i < 100; i++) { + yield* pool.updateMetrics(chainId, true) + } + const concurrencyAfterSuccesses = yield* pool.getOptimalConcurrency(chainId) + + for (let i = 0; i < 100; i++) { + yield* pool.updateMetrics(chainId, false) + } + const concurrencyAfterFailures = yield* pool.getOptimalConcurrency(chainId) + + expect(concurrencyAfterFailures).toBeLessThan(concurrencyAfterSuccesses) + }) + + await Effect.runPromise(testEffect.pipe(Effect.provide(TestContext.TestContext))) + }) })