/** * @module @dotdo/postgres-shared/circuit-breaker-unified * * Unified Circuit Breaker API * * This module provides a unified circuit breaker interface that consolidates * features from multiple implementations into a single, flexible API. It supports * both single-instance mode (for services like CDC) and multi-instance mode * (for DO routing scenarios). * * ## When to Use This vs Base CircuitBreaker * * Use **UnifiedCircuitBreaker** when you need: * - Single-instance mode (no instance IDs) * - Built-in retry functionality with executeWithRetry() * - Active recovery strategy (timer-based state transitions) * - A simpler API that adapts to your use case * * Use **CircuitBreaker** (base) when you need: * - Direct DO stub wrapping with DOCircuitBreakerWrapper * - More explicit multi-instance control * - Direct compatibility with existing code * * ## Features * * - **Dual Mode**: Single-instance (like CDC) or multi-instance (like DO routing) * - **Recovery Strategies**: Passive (on getState check) or active (timer-based) * - **Retry Support**: Built-in executeWithRetry() with exponential backoff * - **Cleanup**: destroy() method for proper resource cleanup * - **Full Compatibility**: Wraps the base CircuitBreaker for consistency * * ## Quick Start * * @example Single-instance mode (for CDC, single service) * ```typescript * import { createUnifiedCircuitBreaker, CircuitBreakerError } from '@dotdo/postgres-shared' * * const cb = createUnifiedCircuitBreaker({ * mode: 'single', * failureThreshold: 5, * resetTimeoutMs: 10000, * recoveryStrategy: 'active', // Uses timers for automatic recovery * }) * * try { * // No instance ID needed in single mode * const result = await cb.execute(async () => { * return await fetchFromCDCService() * }) * } catch (error) { * if (error instanceof CircuitBreakerError) { * console.log(`Circuit open, retry after ${error.retryAfterMs}ms`) * } * } * * // Cleanup when done * cb.destroy() * ``` * * @example Multi-instance mode (for DO routing) * ```typescript * import { createUnifiedCircuitBreaker } from '@dotdo/postgres-shared' * * const cb = createUnifiedCircuitBreaker({ * mode: 'multi', * failureThreshold: 3, * resetTimeoutMs: 30000, * enableExponentialBackoff: true, * }) * * // Track each tenant/DO instance separately * const result = await cb.execute('tenant-123', async () => { * return await doStub.fetch(request) * }) * * // Check state per instance * console.log('Tenant 123 state:', cb.getState('tenant-123')) * console.log('Tenant 456 state:', cb.getState('tenant-456')) * ``` * * @example With retry support * ```typescript * const cb = createUnifiedCircuitBreaker({ mode: 'single' }) * * // Automatically retries with exponential backoff * const result = await cb.executeWithRetry( * async () => unreliableOperation(), * { maxRetries: 3, delayMs: 1000 } * ) * ``` * * @see CircuitBreaker for the base multi-instance implementation * @see DOCircuitBreakerWrapper for wrapping DO stubs directly */ import { CircuitBreaker, CircuitBreakerConfig, CircuitBreakerStats, CircuitBreakerEvent, CircuitBreakerMetricsCollector, CircuitBreakerMetricsSummary, DefaultMetricsCollector, DOCircuitBreakerWrapper, DOCircuitBreakerConfig, DOStub, } from './circuit-breaker.js' import { CircuitOpenError } from './errors.js' // Re-export types export type { CircuitBreakerEvent, CircuitBreakerMetricsCollector, CircuitBreakerMetricsSummary } /** * Circuit state type (string literal union for better TypeScript inference) */ export type CircuitState = 'CLOSED' | 'OPEN' | 'HALF_OPEN' /** * Unified circuit breaker configuration */ export interface UnifiedCircuitBreakerConfig extends Omit { /** Operating mode: 'single' for single-instance, 'multi' for multi-instance (default: 'multi') */ mode?: 'single' | 'multi' /** Recovery strategy: 'passive' checks on getState, 'active' uses timers (default: 'passive') */ recoveryStrategy?: 'passive' | 'active' /** Event handler for circuit breaker events */ onEvent?: (event: CircuitBreakerEvent) => void } /** * Unified statistics interface combining all implementations */ export interface UnifiedCircuitBreakerStats extends CircuitBreakerStats { /** Number of successes in half-open state (from CDC) */ halfOpenSuccesses: number /** Number of times circuit has opened (from CDC) */ openCount: number } /** * Unified Circuit Breaker class * * Wraps the base CircuitBreaker to provide a unified API that supports * both single-instance and multi-instance modes. */ export class UnifiedCircuitBreaker { private readonly cb: CircuitBreaker private readonly mode: 'single' | 'multi' private readonly recoveryStrategy: 'passive' | 'active' // @ts-expect-error - stored for future use (e.g., getConfig method) private readonly _config: UnifiedCircuitBreakerConfig private activeRecoveryTimers: Map> = new Map() private openCounts: Map = new Map() private static readonly DEFAULT_INSTANCE = 'default' constructor(config: UnifiedCircuitBreakerConfig = {}) { this.mode = config.mode ?? 'multi' this.recoveryStrategy = config.recoveryStrategy ?? 'passive' this._config = config // Create underlying circuit breaker const cbConfig: CircuitBreakerConfig = {} if (config.failureThreshold !== undefined) cbConfig.failureThreshold = config.failureThreshold if (config.resetTimeoutMs !== undefined) cbConfig.resetTimeoutMs = config.resetTimeoutMs if (config.halfOpenSuccessThreshold !== undefined) cbConfig.halfOpenSuccessThreshold = config.halfOpenSuccessThreshold if (config.failureWindowMs !== undefined) cbConfig.failureWindowMs = config.failureWindowMs if (config.enableExponentialBackoff !== undefined) cbConfig.enableExponentialBackoff = config.enableExponentialBackoff if (config.backoffMultiplier !== undefined) cbConfig.backoffMultiplier = config.backoffMultiplier if (config.maxResetTimeoutMs !== undefined) cbConfig.maxResetTimeoutMs = config.maxResetTimeoutMs if (config.metricsCollector) cbConfig.metricsCollector = config.metricsCollector if (config.fallback) cbConfig.fallback = config.fallback // Wrap the event handler to track open counts and schedule active recovery cbConfig.onEvent = (event) => { if (event.type === 'STATE_CHANGE') { const instanceId = event.instanceId if (event.to === 'OPEN') { // Track open count this.openCounts.set(instanceId, (this.openCounts.get(instanceId) ?? 0) + 1) // Schedule active recovery if enabled if (this.recoveryStrategy === 'active') { this.scheduleActiveRecovery(instanceId) } } else if (event.to === 'CLOSED') { // Clear recovery timer on close this.clearRecoveryTimer(instanceId) } } // Forward to user's event handler if (config.onEvent) { config.onEvent(event) } } this.cb = new CircuitBreaker(cbConfig) } /** * Schedule active recovery transition from OPEN to HALF_OPEN */ private scheduleActiveRecovery(instanceId: string): void { this.clearRecoveryTimer(instanceId) const timeout = this.cb.getBackoffTimeout(instanceId) const timer = setTimeout(() => { const currentState = this.cb.getState(instanceId) if (currentState === 'OPEN') { // Trigger transition to HALF_OPEN by calling getState // (which handles the transition based on time elapsed) this.cb.getState(instanceId) } this.activeRecoveryTimers.delete(instanceId) }, timeout) this.activeRecoveryTimers.set(instanceId, timer) } /** * Clear active recovery timer for an instance */ private clearRecoveryTimer(instanceId: string): void { const timer = this.activeRecoveryTimers.get(instanceId) if (timer) { clearTimeout(timer) this.activeRecoveryTimers.delete(instanceId) } } /** * Get the instance ID based on mode */ private resolveInstanceId(instanceIdOrUndefined?: string): string { if (this.mode === 'single') { return UnifiedCircuitBreaker.DEFAULT_INSTANCE } return instanceIdOrUndefined ?? UnifiedCircuitBreaker.DEFAULT_INSTANCE } /** * Get the current state of the circuit * * @param instanceId - Instance ID (only used in multi mode) */ getState(instanceId?: string): CircuitState { return this.cb.getState(this.resolveInstanceId(instanceId)) } /** * Check if a request is allowed to proceed * * @param instanceId - Instance ID (only used in multi mode) */ isAllowed(instanceId?: string): boolean { return this.cb.isAllowed(this.resolveInstanceId(instanceId)) } /** * Check if a request can be executed (alias for isAllowed) * * @param instanceId - Instance ID (only used in multi mode) */ canExecute(instanceId?: string): boolean { return this.cb.canExecute(this.resolveInstanceId(instanceId)) } /** * Record a successful operation * * @param instanceId - Instance ID (only used in multi mode) */ recordSuccess(instanceId?: string): void { this.cb.recordSuccess(this.resolveInstanceId(instanceId)) } /** * Record a failed operation * * @param instanceId - Instance ID (only used in multi mode) */ recordFailure(instanceId?: string): void { this.cb.recordFailure(this.resolveInstanceId(instanceId)) } /** * Execute a function with circuit breaker protection * * In single mode: execute(fn) * In multi mode: execute(instanceId, fn) or execute(fn) with default instance */ async execute(fnOrInstanceId: string | (() => Promise), fn?: () => Promise): Promise { if (this.mode === 'single') { // Single mode: first arg is the function const actualFn = fnOrInstanceId as () => Promise return this.cb.execute(UnifiedCircuitBreaker.DEFAULT_INSTANCE, actualFn) } // Multi mode if (typeof fnOrInstanceId === 'function') { // Called as execute(fn) - use default instance return this.cb.execute(UnifiedCircuitBreaker.DEFAULT_INSTANCE, fnOrInstanceId) } // Called as execute(instanceId, fn) if (!fn) { throw new Error('Function argument is required when instanceId is provided') } return this.cb.execute(fnOrInstanceId, fn) } /** * Execute a function with circuit breaker protection and automatic retry * * In single mode: executeWithRetry(fn, options) * In multi mode: executeWithRetry(instanceId, fn, options) or executeWithRetry(fn, options) */ async executeWithRetry( fnOrInstanceId: string | (() => Promise), fnOrOptions?: (() => Promise) | { maxRetries?: number; delayMs?: number }, options?: { maxRetries?: number; delayMs?: number } ): Promise { if (this.mode === 'single') { // Single mode: first arg is the function, second is options const actualFn = fnOrInstanceId as () => Promise const actualOptions = (fnOrOptions as { maxRetries?: number; delayMs?: number }) ?? {} return this.cb.executeWithRetry(UnifiedCircuitBreaker.DEFAULT_INSTANCE, actualFn, actualOptions) } // Multi mode if (typeof fnOrInstanceId === 'function') { // Called as executeWithRetry(fn, options) - use default instance const actualOptions = (fnOrOptions as { maxRetries?: number; delayMs?: number }) ?? {} return this.cb.executeWithRetry(UnifiedCircuitBreaker.DEFAULT_INSTANCE, fnOrInstanceId, actualOptions) } // Called as executeWithRetry(instanceId, fn, options) const actualFn = fnOrOptions as () => Promise if (!actualFn || typeof actualFn !== 'function') { throw new Error('Function argument is required when instanceId is provided') } return this.cb.executeWithRetry(fnOrInstanceId, actualFn, options ?? {}) } /** * Get statistics for an instance or the default instance * * @param instanceId - Instance ID (only used in multi mode) */ getStats(instanceId?: string): UnifiedCircuitBreakerStats { const resolvedId = this.resolveInstanceId(instanceId) const baseStats = this.cb.getStats(resolvedId) return { ...baseStats, halfOpenSuccesses: baseStats.successCount, // In half-open, successCount tracks half-open successes openCount: this.openCounts.get(resolvedId) ?? 0, } } /** * Get statistics for all tracked instances */ getAllStats(): Map { const baseStats = this.cb.getAllStats() const result = new Map() for (const [id, stats] of baseStats) { result.set(id, { ...stats, halfOpenSuccesses: stats.successCount, openCount: this.openCounts.get(id) ?? 0, }) } return result } /** * Get the current configuration */ getConfig(): Required> { const baseConfig = this.cb.getConfig() return { ...baseConfig, mode: this.mode, recoveryStrategy: this.recoveryStrategy, } } /** * Manually reset the circuit for an instance * * @param instanceId - Instance ID (only used in multi mode) */ reset(instanceId?: string): void { const resolvedId = this.resolveInstanceId(instanceId) this.clearRecoveryTimer(resolvedId) this.openCounts.delete(resolvedId) this.cb.reset(resolvedId) } /** * Manually force a circuit to open * * @param instanceId - Instance ID (only used in multi mode) */ forceOpen(instanceId?: string): void { this.cb.forceOpen(this.resolveInstanceId(instanceId)) } /** * Force the circuit to a specific state * * @param state - The state to force * @param instanceId - Instance ID (only used in multi mode) */ forceState(state: CircuitState, instanceId?: string): void forceState(instanceIdOrState: string, stateOrUndefined?: CircuitState): void { // Handle overloaded signatures let instanceId: string let state: CircuitState if (stateOrUndefined !== undefined) { // Called as forceState(instanceId, state) instanceId = instanceIdOrState state = stateOrUndefined } else if (['CLOSED', 'OPEN', 'HALF_OPEN'].includes(instanceIdOrState)) { // Called as forceState(state) in single mode instanceId = this.resolveInstanceId(undefined) state = instanceIdOrState as CircuitState } else { // Called as forceState(instanceId) - invalid, assume state is CLOSED instanceId = instanceIdOrState state = 'CLOSED' } this.cb.forceState(instanceId, state) } /** * Remove tracking for an instance * * @param instanceId - Instance ID to remove */ remove(instanceId: string): boolean { this.clearRecoveryTimer(instanceId) this.openCounts.delete(instanceId) return this.cb.remove(instanceId) } /** * Clear all tracked instances */ clear(): void { // Clear all recovery timers for (const timer of this.activeRecoveryTimers.values()) { clearTimeout(timer) } this.activeRecoveryTimers.clear() this.openCounts.clear() this.cb.clear() } /** * Cleanup all resources (timers, etc.) * Call this when the circuit breaker is no longer needed. * Note: This only clears timers, not the circuit state. */ destroy(): void { // Clear all recovery timers for (const timer of this.activeRecoveryTimers.values()) { clearTimeout(timer) } this.activeRecoveryTimers.clear() // Don't clear the circuit state - just the timers // This preserves the OPEN state after destroy } } /** * Create a unified circuit breaker instance * * @example * ```typescript * // Single-instance mode (like CDC) * const cb = createUnifiedCircuitBreaker({ mode: 'single' }) * * // Multi-instance mode (like DO routing) * const cb = createUnifiedCircuitBreaker({ mode: 'multi' }) * ``` */ export function createUnifiedCircuitBreaker(config?: UnifiedCircuitBreakerConfig): UnifiedCircuitBreaker { return new UnifiedCircuitBreaker(config) } // Re-export CircuitBreakerError as alias for CircuitOpenError for backwards compatibility export { CircuitOpenError as CircuitBreakerError } // Re-export other utilities export { CircuitBreaker, DefaultMetricsCollector, DOCircuitBreakerWrapper, CircuitOpenError } export type { CircuitBreakerConfig, CircuitBreakerStats, DOCircuitBreakerConfig, DOStub }