From c9a4fc48886e387534a994c3f0dadb32057a7a8e Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Thu, 24 Sep 2026 09:15:43 -0400 Subject: [PATCH 01/10] feat: add a reusable retry state controller for RETRY-conformant backoff --- .../datasource/retry/RetryState.test.ts | 305 ++++++++++++++++++ .../datasource/retry/classification.test.ts | 27 ++ .../shared/common/src/datasource/index.ts | 22 ++ .../src/datasource/retry/ResetPolicy.ts | 73 +++++ .../common/src/datasource/retry/RetryState.ts | 267 +++++++++++++++ .../src/datasource/retry/classification.ts | 41 +++ .../common/src/datasource/retry/index.ts | 16 + 7 files changed, 751 insertions(+) create mode 100644 packages/shared/common/__tests__/datasource/retry/RetryState.test.ts create mode 100644 packages/shared/common/__tests__/datasource/retry/classification.test.ts create mode 100644 packages/shared/common/src/datasource/retry/ResetPolicy.ts create mode 100644 packages/shared/common/src/datasource/retry/RetryState.ts create mode 100644 packages/shared/common/src/datasource/retry/classification.ts create mode 100644 packages/shared/common/src/datasource/retry/index.ts diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts new file mode 100644 index 0000000000..93a32c4d9e --- /dev/null +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -0,0 +1,305 @@ +import { + AfterConsecutiveSuccesses, + AfterHealthyFor, +} from '../../../src/datasource/retry/ResetPolicy'; +import { + forPolling, + forStreaming, + RetryState, + RetryStateConfig, +} from '../../../src/datasource/retry/RetryState'; + +const MINUTE = 60 * 1000; +const HOUR = 60 * MINUTE; + +const noJitter = (): number => 0; + +let now = 0; +const clock = (): number => now; + +beforeEach(() => { + now = 0; +}); + +function streamingState(overrides: Partial = {}): RetryState { + return new RetryState({ + normalInitialDelayMs: 1000, + normalCeilingMs: 30 * 1000, + extendedInitialDelayMs: 5 * MINUTE, + extendedCeilingMs: HOUR, + resetPolicy: new AfterHealthyFor(MINUTE), + operatingCadenceMs: 0, + clock, + random: noJitter, + ...overrides, + }); +} + +function pollingState(intervalMs: number, overrides: Partial = {}): RetryState { + return new RetryState({ + normalInitialDelayMs: intervalMs, + normalCeilingMs: intervalMs, + extendedInitialDelayMs: Math.max(5 * MINUTE, intervalMs), + extendedCeilingMs: Math.max(HOUR, intervalMs), + resetPolicy: new AfterConsecutiveSuccesses(2), + operatingCadenceMs: intervalMs, + clock, + random: noJitter, + ...overrides, + }); +} + +it('reports the operating cadence before any outcome is recorded', () => { + expect(streamingState().nextDelay).toEqual(0); + expect(pollingState(30 * 1000).nextDelay).toEqual(30 * 1000); +}); + +it('doubles the delay on consecutive normal failures up to the normal ceiling', () => { + const state = streamingState(); + const delays: number[] = []; + for (let i = 0; i < 7; i += 1) { + state.recordFailure('normal'); + delays.push(state.nextDelay); + } + expect(delays).toEqual([1000, 2000, 4000, 8000, 16000, 30000, 30000]); +}); + +it('moves to the extended regime on the first unexpected failure', () => { + const state = streamingState(); + state.recordFailure('normal'); + state.recordFailure('normal'); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(5 * MINUTE); +}); + +it('keeps doubling through the extended regime up to the extended ceiling', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + const delays = [state.nextDelay]; + for (let i = 0; i < 5; i += 1) { + state.recordFailure('normal'); + delays.push(state.nextDelay); + } + expect(delays).toEqual([5 * MINUTE, 10 * MINUTE, 20 * MINUTE, 40 * MINUTE, HOUR, HOUR]); +}); + +it('does not re-pin the initial delay on a later unexpected failure', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(5 * MINUTE); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(10 * MINUTE); +}); + +it('does not lower the bounds when a normal failure follows an unexpected one', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(10 * MINUTE); +}); + +it('resets after the component has been healthy for the reset duration', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + state.recordFailure('normal'); + + now = 10 * MINUTE; + state.recordSuccess(); + now += MINUTE + 1; + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(1000); +}); + +it('does not reset when the healthy stretch is shorter than the reset duration', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + + now = 10 * MINUTE; + state.recordSuccess(); + now += MINUTE - 1; + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(10 * MINUTE); +}); + +it('anchors the healthy stretch at the first success rather than the latest', () => { + const state = streamingState(); + state.recordFailure('normal'); + state.recordFailure('normal'); + + now = 10 * MINUTE; + state.recordSuccess(); + now += 30 * 1000; + state.recordSuccess(); + now += 31 * 1000; + // 61 seconds since the first success; if repeated successes moved the + // anchor, only 31 seconds would have elapsed and no reset would occur. + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(1000); +}); + +it('continues the delay progression when a success is not enough to reset', () => { + const state = streamingState(); + state.recordFailure('normal'); + state.recordFailure('normal'); + state.recordSuccess(); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(4000); +}); + +it('returns to the operating cadence after a success even while the state is raised', () => { + const state = pollingState(30 * 1000); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(5 * MINUTE); + state.recordSuccess(); + expect(state.nextDelay).toEqual(30 * 1000); +}); + +it('stays in the extended regime until two consecutive polls succeed', () => { + const state = pollingState(30 * 1000); + state.recordFailure('unexpected'); + state.recordSuccess(); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(10 * MINUTE); +}); + +it('resets after two consecutive successful polls', () => { + const state = pollingState(30 * 1000); + state.recordFailure('unexpected'); + state.recordSuccess(); + state.recordSuccess(); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(30 * 1000); +}); + +it('re-arms the extended transition after a reset', () => { + const state = pollingState(30 * 1000); + state.recordFailure('unexpected'); + state.recordSuccess(); + state.recordSuccess(); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(5 * MINUTE); +}); + +it('clears the reset progress when a failure lands between successes', () => { + const state = pollingState(30 * 1000); + state.recordFailure('unexpected'); + state.recordSuccess(); + state.recordFailure('normal'); + state.recordSuccess(); + state.recordFailure('normal'); + // Still extended: the intervening failures kept the reset from occurring. + expect(state.nextDelay).toEqual(20 * MINUTE); +}); + +it('never waits less than the poll interval in the normal regime', () => { + const state = pollingState(30 * 1000, { random: Math.random }); + for (let i = 0; i < 100; i += 1) { + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(30 * 1000); + } +}); + +it('jitters the wait into the upper half of the target delay', () => { + for (let i = 0; i < 100; i += 1) { + const state = streamingState({ random: Math.random }); + state.recordFailure('normal'); + state.recordFailure('normal'); + // Target is 2000; the wait must be in (1000, 2000]. + expect(state.nextDelay).toBeGreaterThan(1000); + expect(state.nextDelay).toBeLessThanOrEqual(2000); + } +}); + +it('uses a server-directed retry time as the new base and restarts the doubling', () => { + const state = streamingState(); + state.recordFailure('normal'); + state.recordFailure('normal'); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(4000); + state.applyServerDirectedRetry(2500); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(2500); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(5000); +}); + +it('clamps a server-directed retry time to the regime ceiling', () => { + const state = streamingState(); + state.applyServerDirectedRetry(2 * HOUR); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(30 * 1000); +}); + +it('keeps a server-directed retry time across a reset', () => { + const state = streamingState(); + state.applyServerDirectedRetry(2500); + state.recordFailure('normal'); + + now = 10 * MINUTE; + state.recordSuccess(); + now += MINUTE + 1; + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(2500); +}); + +it.each([Number.NaN, -1, Number.POSITIVE_INFINITY])( + 'ignores the invalid server-directed retry time %p', + (value) => { + const state = streamingState(); + state.recordFailure('normal'); + state.recordFailure('normal'); + state.applyServerDirectedRetry(value); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(4000); + }, +); + +it('remains at the ceiling without overflowing after very many failures', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + for (let i = 0; i < 500; i += 1) { + state.recordFailure('normal'); + } + expect(state.nextDelay).toEqual(HOUR); + expect(Number.isFinite(state.nextDelay)).toEqual(true); +}); + +it('binds streaming defaults through the factory', () => { + const state = forStreaming(1000); + expect(state.nextDelay).toEqual(0); +}); + +it('binds polling defaults through the factory', () => { + const state = forPolling(30 * 1000); + expect(state.nextDelay).toEqual(30 * 1000); +}); + +it.each([0, -5, Number.NaN, Number.POSITIVE_INFINITY])( + 'warns and uses the default when the streaming initial delay is %p', + (value) => { + const warn = jest.fn(); + const logger = { error: jest.fn(), warn, info: jest.fn(), debug: jest.fn() }; + forStreaming(value, logger); + expect(warn).toHaveBeenCalledTimes(1); + expect(warn.mock.calls[0][0]).toMatch(/initialReconnectDelayMs/); + }, +); + +it.each([0, -5, Number.NaN, Number.POSITIVE_INFINITY])( + 'warns and uses the default when the poll interval is %p', + (value) => { + const warn = jest.fn(); + const logger = { error: jest.fn(), warn, info: jest.fn(), debug: jest.fn() }; + const state = forPolling(value, logger); + expect(warn).toHaveBeenCalledTimes(1); + expect(state.nextDelay).toEqual(30 * 1000); + }, +); + +it('does not warn for valid factory inputs', () => { + const warn = jest.fn(); + const logger = { error: jest.fn(), warn, info: jest.fn(), debug: jest.fn() }; + forStreaming(1000, logger); + forPolling(30 * 1000, logger); + expect(warn).not.toHaveBeenCalled(); +}); diff --git a/packages/shared/common/__tests__/datasource/retry/classification.test.ts b/packages/shared/common/__tests__/datasource/retry/classification.test.ts new file mode 100644 index 0000000000..ec7b307493 --- /dev/null +++ b/packages/shared/common/__tests__/datasource/retry/classification.test.ts @@ -0,0 +1,27 @@ +import { + classifyHttpStatus, + classifyTransportFailure, +} from '../../../src/datasource/retry/classification'; + +it.each([400, 408, 429])('classifies %i as a normal failure', (status) => { + expect(classifyHttpStatus(status)).toEqual('normal'); +}); + +it.each([401, 403, 404, 418, 451, 499])('classifies %i as an unexpected failure', (status) => { + expect(classifyHttpStatus(status)).toEqual('unexpected'); +}); + +it.each([500, 502, 503, 504, 599])('classifies server error %i as a normal failure', (status) => { + expect(classifyHttpStatus(status)).toEqual('normal'); +}); + +it.each([0, 100, 200, 301, 304, 399, 600])( + 'classifies non-error-range status %i as a normal failure', + (status) => { + expect(classifyHttpStatus(status)).toEqual('normal'); + }, +); + +it('classifies transport failures as normal', () => { + expect(classifyTransportFailure()).toEqual('normal'); +}); diff --git a/packages/shared/common/src/datasource/index.ts b/packages/shared/common/src/datasource/index.ts index 497a7588fe..db45c2778e 100644 --- a/packages/shared/common/src/datasource/index.ts +++ b/packages/shared/common/src/datasource/index.ts @@ -8,15 +8,37 @@ import { LDStreamingError, StreamingErrorHandler, } from './errors'; +import { + AfterConsecutiveSuccesses, + AfterHealthyFor, + classifyHttpStatus, + classifyTransportFailure, + FailureKind, + forPolling, + forStreaming, + ResetPolicy, + RetryState, + RetryStateConfig, +} from './retry'; export { + AfterConsecutiveSuccesses, + AfterHealthyFor, Backoff, + classifyHttpStatus, + classifyTransportFailure, CompositeDataSource, DefaultBackoff, DataSourceErrorKind, + FailureKind, + forPolling, + forStreaming, LDFileDataSourceError, LDFlagDeliveryFallbackError, LDPollingError, LDStreamingError, + ResetPolicy, + RetryState, + RetryStateConfig, StreamingErrorHandler, }; diff --git a/packages/shared/common/src/datasource/retry/ResetPolicy.ts b/packages/shared/common/src/datasource/retry/ResetPolicy.ts new file mode 100644 index 0000000000..ee1f2473d1 --- /dev/null +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -0,0 +1,73 @@ +/** + * Decides when a component has operated well enough, for long enough, that its + * retry state should reset. This is the only behavioral difference between the + * streaming and polling retry configurations. + */ +export interface ResetPolicy { + /** + * Records that the component is operating normally. + * + * @param nowMs The current time, from the clock of the RetryState that owns + * this policy. + */ + noteHealthy(nowMs: number): void; + + /** + * Records a failure, which ends any healthy stretch in progress. + */ + noteFailure(): void; + + /** + * Reports whether the reset condition is met. + * + * @param nowMs The current time, from the clock of the RetryState that owns + * this policy. + */ + isSatisfied(nowMs: number): boolean; +} + +/** + * Resets once the component has operated without failing for the given + * duration. Repeated healthy reports do not move the starting point; only a + * failure clears it. + */ +export class AfterHealthyFor implements ResetPolicy { + private _healthySinceMs?: number; + + constructor(private readonly _healthyForMs: number) {} + + noteHealthy(nowMs: number): void { + if (this._healthySinceMs === undefined) { + this._healthySinceMs = nowMs; + } + } + + noteFailure(): void { + this._healthySinceMs = undefined; + } + + isSatisfied(nowMs: number): boolean { + return this._healthySinceMs !== undefined && nowMs - this._healthySinceMs >= this._healthyForMs; + } +} + +/** + * Resets once the given number of operations in a row have succeeded. + */ +export class AfterConsecutiveSuccesses implements ResetPolicy { + private _successes = 0; + + constructor(private readonly _count: number) {} + + noteHealthy(_nowMs: number): void { + this._successes += 1; + } + + noteFailure(): void { + this._successes = 0; + } + + isSatisfied(_nowMs: number): boolean { + return this._successes >= this._count; + } +} diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts new file mode 100644 index 0000000000..d663e0ce6a --- /dev/null +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -0,0 +1,267 @@ +import { LDLogger } from '../../api/logging/LDLogger'; +import { FailureKind } from './classification'; +import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; + +// The delay bounds of the extended regime, in milliseconds. A component enters +// the extended regime after an unexpected failure and leaves it when its reset +// policy is satisfied. +const EXTENDED_INITIAL_DELAY_MS = 5 * 60 * 1000; +const EXTENDED_CEILING_MS = 60 * 60 * 1000; + +// The bounds, reset window, and default initial delay of streaming's normal regime. +const STREAMING_NORMAL_CEILING_MS = 30 * 1000; +const STREAMING_RESET_INTERVAL_MS = 60 * 1000; +const DEFAULT_STREAMING_INITIAL_DELAY_MS = 1000; + +// Polling's reset condition and default cadence. +const POLLING_RESET_SUCCESSES = 2; +const DEFAULT_POLL_INTERVAL_MS = 30 * 1000; + +// An upper bound on the backoff exponent so that a long outage cannot overflow +// the delay computation. Any real ceiling is reached long before this. +const MAX_BACKOFF_EXPONENT = 30; + +function defaultClock(): () => number { + const perf = (globalThis as any)?.performance; + if (perf && typeof perf.now === 'function') { + return () => perf.now(); + } + return Date.now; +} + +function positiveFiniteOrDefault( + value: number, + defaultValueMs: number, + name: string, + logger?: LDLogger, +): number { + if (typeof value === 'number' && Number.isFinite(value) && value > 0) { + return value; + } + logger?.warn( + `${name} must be a positive, finite number of milliseconds; using the default of ${defaultValueMs}`, + ); + return defaultValueMs; +} + +export interface RetryStateConfig { + /** The delay before the first retry in the normal regime, in milliseconds. */ + normalInitialDelayMs: number; + + /** The longest normal-regime delay, in milliseconds. */ + normalCeilingMs: number; + + /** The delay before the first retry in the extended regime, in milliseconds. */ + extendedInitialDelayMs: number; + + /** The longest extended-regime delay, in milliseconds. */ + extendedCeilingMs: number; + + /** Decides when the retry state resets. */ + resetPolicy: ResetPolicy; + + /** + * The wait between healthy operations, in milliseconds. No wait is ever + * shorter than this. Zero (the default) for a component that operates + * continuously. + */ + operatingCadenceMs?: number; + + /** + * The time source used for all time-sensitive decisions. Defaults to a + * monotonic clock. All state in an instance lives on this clock's timeline; + * timestamps from other sources must not be mixed in. + */ + clock?: () => number; + + /** The random source used for jitter. */ + random?: () => number; +} + +/** + * Tracks how long a data source, or another long-running component, should + * wait before its next attempt. + * + * Recording a failure advances the state and decides the next wait, which + * {@link RetryState.nextDelay} reports: the delay doubles with each + * consecutive failure, from the current regime's initial delay up to its + * ceiling, less a random jitter of up to half of it, and never less than the + * operating cadence. An unexpected failure moves the state to the extended + * regime, which raises both bounds; a normal failure that follows cannot lower + * them. The bounds stay raised until the reset policy is satisfied. + * + * Recording a success sets the next wait back to the operating cadence, even + * while the retry state is raised, because a backoff wait applies to a retry + * and not to every operation. + * + * This class never starts timers, performs I/O, or logs; callers schedule the + * waits it computes. Use one instance per component; instances are never + * shared between components. + */ +export class RetryState { + private _attempts = 0; + private _extended = false; + private _minDelayMs: number; + private _maxDelayMs: number; + private _nextDelayMs: number; + private _serverDirectedBaseMs?: number; + private readonly _normalInitialDelayMs: number; + private readonly _normalCeilingMs: number; + private readonly _extendedInitialDelayMs: number; + private readonly _extendedCeilingMs: number; + private readonly _operatingCadenceMs: number; + private readonly _resetPolicy: ResetPolicy; + private readonly _clock: () => number; + private readonly _random: () => number; + + constructor(config: RetryStateConfig) { + this._normalInitialDelayMs = config.normalInitialDelayMs; + this._normalCeilingMs = Math.max(config.normalCeilingMs, config.normalInitialDelayMs); + this._extendedInitialDelayMs = config.extendedInitialDelayMs; + this._extendedCeilingMs = Math.max(config.extendedCeilingMs, config.extendedInitialDelayMs); + this._operatingCadenceMs = config.operatingCadenceMs ?? 0; + this._resetPolicy = config.resetPolicy; + this._clock = config.clock ?? defaultClock(); + this._random = config.random ?? Math.random; + this._minDelayMs = this._normalInitialDelayMs; + this._maxDelayMs = this._normalCeilingMs; + // Before any outcome is recorded, the next wait is the ordinary interval. + this._nextDelayMs = this._operatingCadenceMs; + } + + /** + * The wait, in milliseconds, before the next operation, as the last + * recorded outcome decided it. Reading it has no side effects. + */ + get nextDelay(): number { + return this._nextDelayMs; + } + + /** + * Records a failed attempt and decides the wait before the next one. + * + * The state advances before the wait is computed, so {@link nextDelay} + * always reflects the failure just recorded. + */ + recordFailure(kind: FailureKind): void { + const nowMs = this._clock(); + // A reset that fell due during healthy operation is applied before the + // new failure is counted, so the failure computes from a fresh sequence. + this._resetIfDue(nowMs); + this._resetPolicy.noteFailure(); + + if (kind === 'unexpected' && !this._extended) { + // Moving to the extended regime raises both bounds and starts the delay + // sequence over. Only the move does this: a later unexpected failure + // keeps counting up rather than re-pinning the initial delay. + this._extended = true; + this._minDelayMs = this._extendedInitialDelayMs; + this._maxDelayMs = this._extendedCeilingMs; + this._attempts = 1; + } else { + this._attempts += 1; + } + + const base = this._serverDirectedBaseMs ?? this._minDelayMs; + const exponent = Math.min(Math.max(this._attempts - 1, 0), MAX_BACKOFF_EXPONENT); + const target = Math.min(base * 2 ** exponent, this._maxDelayMs); + const jitter = (this._random() * target) / 2; + this._nextDelayMs = Math.max(target - jitter, this._operatingCadenceMs); + } + + /** + * Records a successful operation and resets the retry state if that is now + * enough to satisfy the reset policy. + */ + recordSuccess(): void { + const nowMs = this._clock(); + this._resetPolicy.noteHealthy(nowMs); + this._resetIfDue(nowMs); + // A backoff wait applies to a retry, not to every operation, so after a + // success the next wait is the ordinary interval even while the retry + // state is raised. + this._nextDelayMs = this._operatingCadenceMs; + } + + /** + * Applies a server-directed retry time. + * + * The value replaces the base of the delay computation and restarts the + * doubling sequence. The regime and its ceiling are unaffected, and the + * value stays in effect until another one arrives, including across a + * reset. Non-finite or negative values are ignored; callers are expected to + * have validated and capped the value at its point of entry. + */ + applyServerDirectedRetry(delayMs: number): void { + if (typeof delayMs !== 'number' || !Number.isFinite(delayMs) || delayMs < 0) { + return; + } + this._serverDirectedBaseMs = delayMs; + this._attempts = 0; + } + + private _resetIfDue(nowMs: number): void { + if (!this._resetPolicy.isSatisfied(nowMs)) { + return; + } + this._attempts = 0; + this._extended = false; + this._minDelayMs = this._normalInitialDelayMs; + this._maxDelayMs = this._normalCeilingMs; + } +} + +/** + * Builds the retry state for a streaming data source: no wait during healthy + * operation, a normal regime running from the configured initial delay up to + * a 30 second ceiling, an extended regime of 5 minutes up to 1 hour, and a + * reset once a connection has been healthy for 60 seconds. + * + * The initial delay is validated here; the documented default of 1 second + * stands in for anything that is not a positive, finite number. The extended + * regime never starts below the configured delay. + */ +export function forStreaming(initialReconnectDelayMs: number, logger?: LDLogger): RetryState { + const validated = positiveFiniteOrDefault( + initialReconnectDelayMs, + DEFAULT_STREAMING_INITIAL_DELAY_MS, + 'initialReconnectDelayMs', + logger, + ); + return new RetryState({ + normalInitialDelayMs: validated, + normalCeilingMs: STREAMING_NORMAL_CEILING_MS, + extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), + extendedCeilingMs: EXTENDED_CEILING_MS, + resetPolicy: new AfterHealthyFor(STREAMING_RESET_INTERVAL_MS), + operatingCadenceMs: 0, + }); +} + +/** + * Builds the retry state for a polling data source: the poll interval is the + * operating cadence, so no wait is ever shorter than it; in the normal regime + * both bounds are the interval itself, so a normal failure simply polls again + * on schedule; the extended regime runs from 5 minutes (or the interval, if + * longer) up to 1 hour (or the interval, if longer); and the state resets + * after two successful polls in a row. + * + * The poll interval is validated here; the documented default of 30 seconds + * stands in for anything that is not a positive, finite number. + */ +export function forPolling(pollIntervalMs: number, logger?: LDLogger): RetryState { + const validated = positiveFiniteOrDefault( + pollIntervalMs, + DEFAULT_POLL_INTERVAL_MS, + 'pollIntervalMs', + logger, + ); + return new RetryState({ + normalInitialDelayMs: validated, + normalCeilingMs: validated, + extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), + extendedCeilingMs: Math.max(EXTENDED_CEILING_MS, validated), + resetPolicy: new AfterConsecutiveSuccesses(POLLING_RESET_SUCCESSES), + operatingCadenceMs: validated, + }); +} diff --git a/packages/shared/common/src/datasource/retry/classification.ts b/packages/shared/common/src/datasource/retry/classification.ts new file mode 100644 index 0000000000..2254132d86 --- /dev/null +++ b/packages/shared/common/src/datasource/retry/classification.ts @@ -0,0 +1,41 @@ +/** + * How a failure is classified, which decides how long the next wait is. + * + * A 'normal' failure is one the service is expected to recover from without + * intervention, so the wait stays short. An 'unexpected' failure suggests a + * problem that requires a person to fix, such as a rejected credential, so the + * wait becomes much longer. Neither classification ever means giving up; there + * is always a next attempt. + */ +export type FailureKind = 'normal' | 'unexpected'; + +// HTTP statuses in the 4xx range that are still normal failures. Every other +// 4xx is unexpected. +const NORMAL_4XX_STATUSES = [400, 408, 429]; + +/** + * Classifies an HTTP status. + * + * 400, 408, and 429 are normal, as is any 5xx. Every other 4xx, including 401 + * and 403, is unexpected. Anything outside the error ranges, including 0 for + * "no response", is normal. + */ +export function classifyHttpStatus(status: number): FailureKind { + if (status >= 400 && status < 500 && !NORMAL_4XX_STATUSES.includes(status)) { + return 'unexpected'; + } + return 'normal'; +} + +/** + * Classifies a transport-level failure (connection refused or dropped, DNS + * failure, timeout, TLS negotiation failure). + * + * Always normal: in an all-HTTPS system every transport error surfaces through + * the TLS layer, so a certificate misconfiguration cannot be reliably + * distinguished from a transient network fault, and treating transient faults + * as unexpected would hold ordinary recoveries to multi-minute waits. + */ +export function classifyTransportFailure(): FailureKind { + return 'normal'; +} diff --git a/packages/shared/common/src/datasource/retry/index.ts b/packages/shared/common/src/datasource/retry/index.ts new file mode 100644 index 0000000000..e2f2a86282 --- /dev/null +++ b/packages/shared/common/src/datasource/retry/index.ts @@ -0,0 +1,16 @@ +import { classifyHttpStatus, classifyTransportFailure, FailureKind } from './classification'; +import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; +import { forPolling, forStreaming, RetryState, RetryStateConfig } from './RetryState'; + +export { + AfterConsecutiveSuccesses, + AfterHealthyFor, + classifyHttpStatus, + classifyTransportFailure, + FailureKind, + forPolling, + forStreaming, + ResetPolicy, + RetryState, + RetryStateConfig, +}; From 81d471beb4eed1b8e2a34d9f35285adc87647cf5 Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Thu, 24 Sep 2026 16:02:15 -0400 Subject: [PATCH 02/10] docs: remove configuration-comparison sentence from the ResetPolicy doc --- packages/shared/common/src/datasource/retry/ResetPolicy.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/packages/shared/common/src/datasource/retry/ResetPolicy.ts b/packages/shared/common/src/datasource/retry/ResetPolicy.ts index ee1f2473d1..6326c17fb6 100644 --- a/packages/shared/common/src/datasource/retry/ResetPolicy.ts +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -1,7 +1,6 @@ /** * Decides when a component has operated well enough, for long enough, that its - * retry state should reset. This is the only behavioral difference between the - * streaming and polling retry configurations. + * retry state should reset. */ export interface ResetPolicy { /** From d4f41affeb9b9138b8835a0390ca5b6ecbd612dc Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Thu, 24 Sep 2026 16:43:45 -0400 Subject: [PATCH 03/10] refactor: make classifyHttpStatus the canonical partition beside the legacy helper --- .../classification.test.ts => errors.test.ts} | 19 +++++-- .../shared/common/src/datasource/index.ts | 6 --- .../common/src/datasource/retry/RetryState.ts | 2 +- .../src/datasource/retry/classification.ts | 41 --------------- .../common/src/datasource/retry/index.ts | 4 -- packages/shared/common/src/errors.ts | 51 +++++++++++++++++-- 6 files changed, 63 insertions(+), 60 deletions(-) rename packages/shared/common/__tests__/{datasource/retry/classification.test.ts => errors.test.ts} (62%) delete mode 100644 packages/shared/common/src/datasource/retry/classification.ts diff --git a/packages/shared/common/__tests__/datasource/retry/classification.test.ts b/packages/shared/common/__tests__/errors.test.ts similarity index 62% rename from packages/shared/common/__tests__/datasource/retry/classification.test.ts rename to packages/shared/common/__tests__/errors.test.ts index ec7b307493..6d07cfad5c 100644 --- a/packages/shared/common/__tests__/datasource/retry/classification.test.ts +++ b/packages/shared/common/__tests__/errors.test.ts @@ -1,7 +1,4 @@ -import { - classifyHttpStatus, - classifyTransportFailure, -} from '../../../src/datasource/retry/classification'; +import { classifyHttpStatus, classifyTransportFailure, isHttpRecoverable } from '../src/errors'; it.each([400, 408, 429])('classifies %i as a normal failure', (status) => { expect(classifyHttpStatus(status)).toEqual('normal'); @@ -25,3 +22,17 @@ it.each([0, 100, 200, 301, 304, 399, 600])( it('classifies transport failures as normal', () => { expect(classifyTransportFailure()).toEqual('normal'); }); + +it.each([400, 408, 429, 500, 503, 200, 0])( + 'reports %i as recoverable through the retained legacy helper', + (status) => { + expect(isHttpRecoverable(status)).toEqual(true); + }, +); + +it.each([401, 403, 404, 451])( + 'reports %i as unrecoverable through the retained legacy helper', + (status) => { + expect(isHttpRecoverable(status)).toEqual(false); + }, +); diff --git a/packages/shared/common/src/datasource/index.ts b/packages/shared/common/src/datasource/index.ts index db45c2778e..6d57d5a913 100644 --- a/packages/shared/common/src/datasource/index.ts +++ b/packages/shared/common/src/datasource/index.ts @@ -11,9 +11,6 @@ import { import { AfterConsecutiveSuccesses, AfterHealthyFor, - classifyHttpStatus, - classifyTransportFailure, - FailureKind, forPolling, forStreaming, ResetPolicy, @@ -25,12 +22,9 @@ export { AfterConsecutiveSuccesses, AfterHealthyFor, Backoff, - classifyHttpStatus, - classifyTransportFailure, CompositeDataSource, DefaultBackoff, DataSourceErrorKind, - FailureKind, forPolling, forStreaming, LDFileDataSourceError, diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts index d663e0ce6a..e90498334d 100644 --- a/packages/shared/common/src/datasource/retry/RetryState.ts +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -1,5 +1,5 @@ import { LDLogger } from '../../api/logging/LDLogger'; -import { FailureKind } from './classification'; +import { FailureKind } from '../../errors'; import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; // The delay bounds of the extended regime, in milliseconds. A component enters diff --git a/packages/shared/common/src/datasource/retry/classification.ts b/packages/shared/common/src/datasource/retry/classification.ts deleted file mode 100644 index 2254132d86..0000000000 --- a/packages/shared/common/src/datasource/retry/classification.ts +++ /dev/null @@ -1,41 +0,0 @@ -/** - * How a failure is classified, which decides how long the next wait is. - * - * A 'normal' failure is one the service is expected to recover from without - * intervention, so the wait stays short. An 'unexpected' failure suggests a - * problem that requires a person to fix, such as a rejected credential, so the - * wait becomes much longer. Neither classification ever means giving up; there - * is always a next attempt. - */ -export type FailureKind = 'normal' | 'unexpected'; - -// HTTP statuses in the 4xx range that are still normal failures. Every other -// 4xx is unexpected. -const NORMAL_4XX_STATUSES = [400, 408, 429]; - -/** - * Classifies an HTTP status. - * - * 400, 408, and 429 are normal, as is any 5xx. Every other 4xx, including 401 - * and 403, is unexpected. Anything outside the error ranges, including 0 for - * "no response", is normal. - */ -export function classifyHttpStatus(status: number): FailureKind { - if (status >= 400 && status < 500 && !NORMAL_4XX_STATUSES.includes(status)) { - return 'unexpected'; - } - return 'normal'; -} - -/** - * Classifies a transport-level failure (connection refused or dropped, DNS - * failure, timeout, TLS negotiation failure). - * - * Always normal: in an all-HTTPS system every transport error surfaces through - * the TLS layer, so a certificate misconfiguration cannot be reliably - * distinguished from a transient network fault, and treating transient faults - * as unexpected would hold ordinary recoveries to multi-minute waits. - */ -export function classifyTransportFailure(): FailureKind { - return 'normal'; -} diff --git a/packages/shared/common/src/datasource/retry/index.ts b/packages/shared/common/src/datasource/retry/index.ts index e2f2a86282..7e2dc819ab 100644 --- a/packages/shared/common/src/datasource/retry/index.ts +++ b/packages/shared/common/src/datasource/retry/index.ts @@ -1,13 +1,9 @@ -import { classifyHttpStatus, classifyTransportFailure, FailureKind } from './classification'; import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; import { forPolling, forStreaming, RetryState, RetryStateConfig } from './RetryState'; export { AfterConsecutiveSuccesses, AfterHealthyFor, - classifyHttpStatus, - classifyTransportFailure, - FailureKind, forPolling, forStreaming, ResetPolicy, diff --git a/packages/shared/common/src/errors.ts b/packages/shared/common/src/errors.ts index c97a9961ea..469371a62c 100644 --- a/packages/shared/common/src/errors.ts +++ b/packages/shared/common/src/errors.ts @@ -22,17 +22,60 @@ export class LDTimeoutError extends Error { } } +/** + * How a failure is classified, which decides how long the next wait is. + * + * A 'normal' failure is one the service is expected to recover from without + * intervention, so the wait stays short. An 'unexpected' failure suggests a + * problem that requires a person to fix, such as a rejected credential, so the + * wait becomes much longer. Neither classification ever means giving up; there + * is always a next attempt. + */ +export type FailureKind = 'normal' | 'unexpected'; + +// HTTP statuses in the 4xx range that are still normal failures. Every other +// 4xx is unexpected. +const NORMAL_4XX_STATUSES = [400, 408, 429]; + +/** + * Classifies an HTTP status. + * + * 400, 408, and 429 are normal, as is any 5xx. Every other 4xx, including 401 + * and 403, is unexpected. Anything outside the error ranges, including 0 for + * "no response", is normal. + */ +export function classifyHttpStatus(status: number): FailureKind { + if (status >= 400 && status < 500 && !NORMAL_4XX_STATUSES.includes(status)) { + return 'unexpected'; + } + return 'normal'; +} + +/** + * Classifies a transport-level failure (connection refused or dropped, DNS + * failure, timeout, TLS negotiation failure). + * + * Always normal: in an all-HTTPS system every transport error surfaces through + * the TLS layer, so a certificate misconfiguration cannot be reliably + * distinguished from a transient network fault, and treating transient faults + * as unexpected would hold ordinary recoveries to multi-minute waits. + */ +export function classifyTransportFailure(): FailureKind { + return 'normal'; +} + /** * Check if the HTTP error is recoverable. This will return false if a request * made with any payload could not recover. If the reason for the failure * is payload specific, for instance a payload that is too large, then * it could recover with a different payload. + * + * Superseded by {@link classifyHttpStatus}, which this delegates to. It is + * retained for the event-delivery pathway and is expected to be deprecated + * once that pathway migrates to failure classification. */ export function isHttpRecoverable(status: number) { - if (status >= 400 && status < 500) { - return status === 400 || status === 408 || status === 429; - } - return true; + return classifyHttpStatus(status) === 'normal'; } /** From d68d2cb84ae104ba1b24266b2962f569c5cbcfad Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Fri, 25 Sep 2026 11:11:09 -0400 Subject: [PATCH 04/10] refactor: make the reset policy interface clockless --- .../datasource/retry/RetryState.test.ts | 4 +- .../src/datasource/retry/ResetPolicy.ts | 46 +++++++++++++------ .../common/src/datasource/retry/RetryState.ts | 29 ++---------- 3 files changed, 37 insertions(+), 42 deletions(-) diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts index 93a32c4d9e..a552eeb3bc 100644 --- a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -27,9 +27,8 @@ function streamingState(overrides: Partial = {}): RetryState { normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 5 * MINUTE, extendedCeilingMs: HOUR, - resetPolicy: new AfterHealthyFor(MINUTE), + resetPolicy: new AfterHealthyFor(MINUTE, clock), operatingCadenceMs: 0, - clock, random: noJitter, ...overrides, }); @@ -43,7 +42,6 @@ function pollingState(intervalMs: number, overrides: Partial = extendedCeilingMs: Math.max(HOUR, intervalMs), resetPolicy: new AfterConsecutiveSuccesses(2), operatingCadenceMs: intervalMs, - clock, random: noJitter, ...overrides, }); diff --git a/packages/shared/common/src/datasource/retry/ResetPolicy.ts b/packages/shared/common/src/datasource/retry/ResetPolicy.ts index 6326c17fb6..3b5d4c583d 100644 --- a/packages/shared/common/src/datasource/retry/ResetPolicy.ts +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -1,3 +1,11 @@ +function defaultClock(): () => number { + const perf = (globalThis as any)?.performance; + if (perf && typeof perf.now === 'function') { + return () => perf.now(); + } + return Date.now; +} + /** * Decides when a component has operated well enough, for long enough, that its * retry state should reset. @@ -5,11 +13,8 @@ export interface ResetPolicy { /** * Records that the component is operating normally. - * - * @param nowMs The current time, from the clock of the RetryState that owns - * this policy. */ - noteHealthy(nowMs: number): void; + noteHealthy(): void; /** * Records a failure, which ends any healthy stretch in progress. @@ -18,11 +23,8 @@ export interface ResetPolicy { /** * Reports whether the reset condition is met. - * - * @param nowMs The current time, from the clock of the RetryState that owns - * this policy. */ - isSatisfied(nowMs: number): boolean; + isSatisfied(): boolean; } /** @@ -32,12 +34,24 @@ export interface ResetPolicy { */ export class AfterHealthyFor implements ResetPolicy { private _healthySinceMs?: number; + private readonly _clock: () => number; - constructor(private readonly _healthyForMs: number) {} + /** + * @param _healthyForMs How long the component must operate without failing, + * in milliseconds. + * @param clock The time source used to measure the healthy stretch; defaults + * to a monotonic clock. Primarily for testing. + */ + constructor( + private readonly _healthyForMs: number, + clock?: () => number, + ) { + this._clock = clock ?? defaultClock(); + } - noteHealthy(nowMs: number): void { + noteHealthy(): void { if (this._healthySinceMs === undefined) { - this._healthySinceMs = nowMs; + this._healthySinceMs = this._clock(); } } @@ -45,8 +59,10 @@ export class AfterHealthyFor implements ResetPolicy { this._healthySinceMs = undefined; } - isSatisfied(nowMs: number): boolean { - return this._healthySinceMs !== undefined && nowMs - this._healthySinceMs >= this._healthyForMs; + isSatisfied(): boolean { + return ( + this._healthySinceMs !== undefined && this._clock() - this._healthySinceMs >= this._healthyForMs + ); } } @@ -58,7 +74,7 @@ export class AfterConsecutiveSuccesses implements ResetPolicy { constructor(private readonly _count: number) {} - noteHealthy(_nowMs: number): void { + noteHealthy(): void { this._successes += 1; } @@ -66,7 +82,7 @@ export class AfterConsecutiveSuccesses implements ResetPolicy { this._successes = 0; } - isSatisfied(_nowMs: number): boolean { + isSatisfied(): boolean { return this._successes >= this._count; } } diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts index e90498334d..2433302eb6 100644 --- a/packages/shared/common/src/datasource/retry/RetryState.ts +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -21,14 +21,6 @@ const DEFAULT_POLL_INTERVAL_MS = 30 * 1000; // the delay computation. Any real ceiling is reached long before this. const MAX_BACKOFF_EXPONENT = 30; -function defaultClock(): () => number { - const perf = (globalThis as any)?.performance; - if (perf && typeof perf.now === 'function') { - return () => perf.now(); - } - return Date.now; -} - function positiveFiniteOrDefault( value: number, defaultValueMs: number, @@ -67,13 +59,6 @@ export interface RetryStateConfig { */ operatingCadenceMs?: number; - /** - * The time source used for all time-sensitive decisions. Defaults to a - * monotonic clock. All state in an instance lives on this clock's timeline; - * timestamps from other sources must not be mixed in. - */ - clock?: () => number; - /** The random source used for jitter. */ random?: () => number; } @@ -111,7 +96,6 @@ export class RetryState { private readonly _extendedCeilingMs: number; private readonly _operatingCadenceMs: number; private readonly _resetPolicy: ResetPolicy; - private readonly _clock: () => number; private readonly _random: () => number; constructor(config: RetryStateConfig) { @@ -121,7 +105,6 @@ export class RetryState { this._extendedCeilingMs = Math.max(config.extendedCeilingMs, config.extendedInitialDelayMs); this._operatingCadenceMs = config.operatingCadenceMs ?? 0; this._resetPolicy = config.resetPolicy; - this._clock = config.clock ?? defaultClock(); this._random = config.random ?? Math.random; this._minDelayMs = this._normalInitialDelayMs; this._maxDelayMs = this._normalCeilingMs; @@ -144,10 +127,9 @@ export class RetryState { * always reflects the failure just recorded. */ recordFailure(kind: FailureKind): void { - const nowMs = this._clock(); // A reset that fell due during healthy operation is applied before the // new failure is counted, so the failure computes from a fresh sequence. - this._resetIfDue(nowMs); + this._resetIfDue(); this._resetPolicy.noteFailure(); if (kind === 'unexpected' && !this._extended) { @@ -174,9 +156,8 @@ export class RetryState { * enough to satisfy the reset policy. */ recordSuccess(): void { - const nowMs = this._clock(); - this._resetPolicy.noteHealthy(nowMs); - this._resetIfDue(nowMs); + this._resetPolicy.noteHealthy(); + this._resetIfDue(); // A backoff wait applies to a retry, not to every operation, so after a // success the next wait is the ordinary interval even while the retry // state is raised. @@ -200,8 +181,8 @@ export class RetryState { this._attempts = 0; } - private _resetIfDue(nowMs: number): void { - if (!this._resetPolicy.isSatisfied(nowMs)) { + private _resetIfDue(): void { + if (!this._resetPolicy.isSatisfied()) { return; } this._attempts = 0; From c46a038f1f5b710fa22abe16b2d22a8005e624ae Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Fri, 25 Sep 2026 11:56:59 -0400 Subject: [PATCH 05/10] refactor: bound the backoff by the ceiling and harden the default clock fallback --- .../datasource/retry/RetryState.test.ts | 17 +++++++++++++++++ .../common/src/datasource/retry/ResetPolicy.ts | 5 +++-- .../common/src/datasource/retry/RetryState.ts | 14 ++++++++------ 3 files changed, 28 insertions(+), 8 deletions(-) diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts index a552eeb3bc..9a37a7a3e2 100644 --- a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -262,6 +262,23 @@ it('remains at the ceiling without overflowing after very many failures', () => expect(Number.isFinite(state.nextDelay)).toEqual(true); }); +it('stays finite when a server-directed retry time of zero is followed by very many failures', () => { + const state = streamingState(); + state.applyServerDirectedRetry(0); + for (let i = 0; i < 1100; i += 1) { + state.recordFailure('normal'); + } + expect(state.nextDelay).toEqual(0); + expect(Number.isFinite(state.nextDelay)).toEqual(true); +}); + +it('floors a zero server-directed retry time at the operating cadence', () => { + const state = pollingState(30 * 1000); + state.applyServerDirectedRetry(0); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(30 * 1000); +}); + it('binds streaming defaults through the factory', () => { const state = forStreaming(1000); expect(state.nextDelay).toEqual(0); diff --git a/packages/shared/common/src/datasource/retry/ResetPolicy.ts b/packages/shared/common/src/datasource/retry/ResetPolicy.ts index 3b5d4c583d..2a91319fac 100644 --- a/packages/shared/common/src/datasource/retry/ResetPolicy.ts +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -3,7 +3,7 @@ function defaultClock(): () => number { if (perf && typeof perf.now === 'function') { return () => perf.now(); } - return Date.now; + return () => Date.now(); } /** @@ -61,7 +61,8 @@ export class AfterHealthyFor implements ResetPolicy { isSatisfied(): boolean { return ( - this._healthySinceMs !== undefined && this._clock() - this._healthySinceMs >= this._healthyForMs + this._healthySinceMs !== undefined && + this._clock() - this._healthySinceMs >= this._healthyForMs ); } } diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts index 2433302eb6..cc88e21cb8 100644 --- a/packages/shared/common/src/datasource/retry/RetryState.ts +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -17,10 +17,6 @@ const DEFAULT_STREAMING_INITIAL_DELAY_MS = 1000; const POLLING_RESET_SUCCESSES = 2; const DEFAULT_POLL_INTERVAL_MS = 30 * 1000; -// An upper bound on the backoff exponent so that a long outage cannot overflow -// the delay computation. Any real ceiling is reached long before this. -const MAX_BACKOFF_EXPONENT = 30; - function positiveFiniteOrDefault( value: number, defaultValueMs: number, @@ -145,8 +141,14 @@ export class RetryState { } const base = this._serverDirectedBaseMs ?? this._minDelayMs; - const exponent = Math.min(Math.max(this._attempts - 1, 0), MAX_BACKOFF_EXPONENT); - const target = Math.min(base * 2 ** exponent, this._maxDelayMs); + // Compare against the ceiling scaled down rather than the base scaled up, so + // the computed value can never overflow the ceiling. A base of zero doubles + // to zero forever, so it short-circuits. + let target = 0; + if (base > 0) { + const exponent = this._attempts - 1; + target = base >= this._maxDelayMs / 2 ** exponent ? this._maxDelayMs : base * 2 ** exponent; + } const jitter = (this._random() * target) / 2; this._nextDelayMs = Math.max(target - jitter, this._operatingCadenceMs); } From 803367c69702ac9db1708a16c96f2e6d0629504a Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Fri, 25 Sep 2026 13:43:09 -0400 Subject: [PATCH 06/10] test: close controller test-parity gaps against the reference implementations --- .../datasource/retry/ResetPolicy.test.ts | 68 ++++++++++++++++ .../datasource/retry/RetryState.test.ts | 81 +++++++++++++++++++ 2 files changed, 149 insertions(+) create mode 100644 packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts diff --git a/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts new file mode 100644 index 0000000000..3462bd7b45 --- /dev/null +++ b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts @@ -0,0 +1,68 @@ +import { + AfterConsecutiveSuccesses, + AfterHealthyFor, +} from '../../../src/datasource/retry/ResetPolicy'; + +const MINUTE = 60 * 1000; + +let now = 0; +const clock = (): number => now; + +beforeEach(() => { + now = 0; +}); + +it('is not satisfied before any healthy report', () => { + const policy = new AfterHealthyFor(MINUTE, clock); + now = 10 * MINUTE; + expect(policy.isSatisfied()).toEqual(false); +}); + +it('is satisfied once the healthy stretch reaches the threshold, and not a moment sooner', () => { + const policy = new AfterHealthyFor(MINUTE, clock); + policy.noteHealthy(); + now = MINUTE - 1; + expect(policy.isSatisfied()).toEqual(false); + now = MINUTE; + expect(policy.isSatisfied()).toEqual(true); +}); + +it('does not move the start of the stretch on repeated healthy reports', () => { + const policy = new AfterHealthyFor(MINUTE, clock); + policy.noteHealthy(); + for (let i = 1; i < 60; i += 1) { + now = i * 1000; + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(false); + } + now = MINUTE; + expect(policy.isSatisfied()).toEqual(true); +}); + +it('clears the healthy stretch on a failure', () => { + const policy = new AfterHealthyFor(MINUTE, clock); + policy.noteHealthy(); + now = 30 * 1000; + policy.noteFailure(); + now = 15 * MINUTE; + expect(policy.isSatisfied()).toEqual(false); +}); + +it('is satisfied by consecutive successes and not by fewer', () => { + const policy = new AfterConsecutiveSuccesses(2); + expect(policy.isSatisfied()).toEqual(false); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(false); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(true); +}); + +it('starts the success count over on a failure', () => { + const policy = new AfterConsecutiveSuccesses(2); + policy.noteHealthy(); + policy.noteFailure(); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(false); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(true); +}); diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts index 9a37a7a3e2..c305b66978 100644 --- a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -198,6 +198,7 @@ it('never waits less than the poll interval in the normal regime', () => { }); it('jitters the wait into the upper half of the target delay', () => { + const delays: number[] = []; for (let i = 0; i < 100; i += 1) { const state = streamingState({ random: Math.random }); state.recordFailure('normal'); @@ -205,7 +206,44 @@ it('jitters the wait into the upper half of the target delay', () => { // Target is 2000; the wait must be in (1000, 2000]. expect(state.nextDelay).toBeGreaterThan(1000); expect(state.nextDelay).toBeLessThanOrEqual(2000); + delays.push(state.nextDelay); + } + // A disabled jitter would pass the range checks with a constant value. + expect(new Set(delays).size).toBeGreaterThan(1); +}); + +it('keeps the wait above half the target under the largest jitter draw', () => { + const state = streamingState({ random: () => 0.9999999 }); + state.recordFailure('normal'); + state.recordFailure('normal'); + // Target is 2000; the largest draw removes just under half of it. + expect(state.nextDelay).toBeGreaterThan(1000); + expect(state.nextDelay).toBeLessThan(1001); +}); + +it('replaces an earlier server-directed retry time with a later one', () => { + const state = streamingState(); + state.applyServerDirectedRetry(2500); + state.applyServerDirectedRetry(4000); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(4000); +}); + +it('does not enter the extended regime no matter how fast a connection flaps', () => { + const state = streamingState(); + const delays: number[] = []; + // Twenty cycles of connect, deliver briefly, and drop: never enough healthy + // time to reset, and never anything but normal failures. The delay may climb + // to the normal ceiling but must never reach the extended regime. + for (let i = 0; i < 20; i += 1) { + state.recordSuccess(); + now += 5 * 1000; + state.recordFailure('normal'); + delays.push(state.nextDelay); + now += 1000; } + delays.forEach((delay) => expect(delay).toBeLessThanOrEqual(30 * 1000)); + expect(delays[delays.length - 1]).toEqual(30 * 1000); }); it('uses a server-directed retry time as the new base and restarts the doubling', () => { @@ -311,6 +349,49 @@ it.each([0, -5, Number.NaN, Number.POSITIVE_INFINITY])( }, ); +it('raises the normal ceiling to a configured delay that exceeds it', () => { + // A 10 minute configured delay is above the 30 second normal ceiling; the + // configured value wins rather than being cut down to the ceiling. + const state = forStreaming(10 * MINUTE); + state.recordFailure('normal'); + expect(state.nextDelay).toBeGreaterThan(5 * MINUTE); + expect(state.nextDelay).toBeLessThanOrEqual(10 * MINUTE); + state.recordFailure('normal'); + expect(state.nextDelay).toBeGreaterThan(5 * MINUTE); + expect(state.nextDelay).toBeLessThanOrEqual(10 * MINUTE); +}); + +it('raises the extended bounds to a configured delay that exceeds them', () => { + // A 2 hour configured delay is above both the 5 minute extended initial and + // the 1 hour extended ceiling. + const state = forStreaming(2 * HOUR); + state.recordFailure('unexpected'); + expect(state.nextDelay).toBeGreaterThan(HOUR); + expect(state.nextDelay).toBeLessThanOrEqual(2 * HOUR); +}); + +it('leaves a configured delay below the normal ceiling untouched', () => { + const state = forStreaming(1000); + state.recordFailure('normal'); + expect(state.nextDelay).toBeGreaterThan(500); + expect(state.nextDelay).toBeLessThanOrEqual(1000); +}); + +it('collapses the extended regime when the poll interval exceeds its bounds', () => { + // With a 2 hour interval the cadence floor makes every wait exactly the + // interval: failures of either kind cannot escalate past it, and a success + // returns to it. + const state = forPolling(2 * HOUR); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(2 * HOUR); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(2 * HOUR); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(2 * HOUR); + state.recordSuccess(); + expect(state.nextDelay).toEqual(2 * HOUR); +}); + it('does not warn for valid factory inputs', () => { const warn = jest.fn(); const logger = { error: jest.fn(), warn, info: jest.fn(), debug: jest.fn() }; From ce199ca153b7101d62eeb06ca908f5cfaa22fc2e Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Fri, 25 Sep 2026 17:23:58 -0400 Subject: [PATCH 07/10] fix: address multi-agent review findings for retry controller Export the retry controller (RetryState, factories, and reset policies) from the package entry point so consumers can reach it. Clamp the normal regime at its ceiling when a configured initial delay exceeds it, matching the majority of the SDK fleet and the literal spec; the mandated extended floor stays. Document the config trust boundary, the server-directed hint precedence and ceiling bound, and the read-in-same-turn contract, and pin the hint-precedence and jitter-boundary behaviors and the reset-policy and default-cadence paths with tests. Part of SDK-2790. --- .../datasource/retry/ResetPolicy.test.ts | 23 ++++++ .../datasource/retry/RetryState.test.ts | 77 ++++++++++++++++--- .../src/datasource/retry/ResetPolicy.ts | 15 ++-- .../common/src/datasource/retry/RetryState.ts | 44 +++++++---- packages/shared/common/src/index.ts | 14 ++++ 5 files changed, 142 insertions(+), 31 deletions(-) diff --git a/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts index 3462bd7b45..c63c34e10f 100644 --- a/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts @@ -66,3 +66,26 @@ it('starts the success count over on a failure', () => { policy.noteHealthy(); expect(policy.isSatisfied()).toEqual(true); }); + +it('measures the healthy stretch with the default clock', () => { + const policy = new AfterHealthyFor(60 * MINUTE); + expect(policy.isSatisfied()).toEqual(false); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(false); +}); + +it('falls back to the wall clock when no monotonic source exists', () => { + const savedPerformance = (globalThis as any).performance; + const nowSpy = jest.spyOn(Date, 'now').mockReturnValue(0); + delete (globalThis as any).performance; + try { + const policy = new AfterHealthyFor(MINUTE); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(false); + nowSpy.mockReturnValue(MINUTE); + expect(policy.isSatisfied()).toEqual(true); + } finally { + nowSpy.mockRestore(); + (globalThis as any).performance = savedPerformance; + } +}); diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts index c305b66978..112fb17a50 100644 --- a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -212,13 +212,15 @@ it('jitters the wait into the upper half of the target delay', () => { expect(new Set(delays).size).toBeGreaterThan(1); }); -it('keeps the wait above half the target under the largest jitter draw', () => { - const state = streamingState({ random: () => 0.9999999 }); +it('keeps the wait at or above half the target under the largest jitter draw', () => { + // The largest double below 1 against a power-of-two target is the worst + // case: the subtraction ties and rounds to exactly half the target, so the + // wait interval is closed at T/2. + const state = streamingState({ normalInitialDelayMs: 1024, random: () => 1 - 2 ** -53 }); state.recordFailure('normal'); state.recordFailure('normal'); - // Target is 2000; the largest draw removes just under half of it. - expect(state.nextDelay).toBeGreaterThan(1000); - expect(state.nextDelay).toBeLessThan(1001); + // Target is 2048; the wait lands on exactly half of it, never below. + expect(state.nextDelay).toEqual(1024); }); it('replaces an earlier server-directed retry time with a later one', () => { @@ -278,6 +280,24 @@ it('keeps a server-directed retry time across a reset', () => { expect(state.nextDelay).toEqual(2500); }); +it('computes from a server-directed base rather than the extended initial delay', () => { + const state = streamingState(); + state.applyServerDirectedRetry(1000); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(1000); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(2000); +}); + +it('applies a server-directed retry time received while already extended', () => { + const state = streamingState(); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(5 * MINUTE); + state.applyServerDirectedRetry(1000); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(1000); +}); + it.each([Number.NaN, -1, Number.POSITIVE_INFINITY])( 'ignores the invalid server-directed retry time %p', (value) => { @@ -332,9 +352,12 @@ it.each([0, -5, Number.NaN, Number.POSITIVE_INFINITY])( (value) => { const warn = jest.fn(); const logger = { error: jest.fn(), warn, info: jest.fn(), debug: jest.fn() }; - forStreaming(value, logger); + const state = forStreaming(value, logger); expect(warn).toHaveBeenCalledTimes(1); expect(warn.mock.calls[0][0]).toMatch(/initialReconnectDelayMs/); + state.recordFailure('normal'); + expect(state.nextDelay).toBeGreaterThan(500); + expect(state.nextDelay).toBeLessThanOrEqual(1000); }, ); @@ -349,16 +372,16 @@ it.each([0, -5, Number.NaN, Number.POSITIVE_INFINITY])( }, ); -it('raises the normal ceiling to a configured delay that exceeds it', () => { +it('clamps waits at the normal ceiling when the configured delay exceeds it', () => { // A 10 minute configured delay is above the 30 second normal ceiling; the - // configured value wins rather than being cut down to the ceiling. + // ceiling wins rather than being raised to the configured value. const state = forStreaming(10 * MINUTE); state.recordFailure('normal'); - expect(state.nextDelay).toBeGreaterThan(5 * MINUTE); - expect(state.nextDelay).toBeLessThanOrEqual(10 * MINUTE); + expect(state.nextDelay).toBeGreaterThan(15 * 1000); + expect(state.nextDelay).toBeLessThanOrEqual(30 * 1000); state.recordFailure('normal'); - expect(state.nextDelay).toBeGreaterThan(5 * MINUTE); - expect(state.nextDelay).toBeLessThanOrEqual(10 * MINUTE); + expect(state.nextDelay).toBeGreaterThan(15 * 1000); + expect(state.nextDelay).toBeLessThanOrEqual(30 * 1000); }); it('raises the extended bounds to a configured delay that exceeds them', () => { @@ -377,6 +400,36 @@ it('leaves a configured delay below the normal ceiling untouched', () => { expect(state.nextDelay).toBeLessThanOrEqual(1000); }); +it('defaults the operating cadence to zero for direct construction', () => { + const state = new RetryState({ + normalInitialDelayMs: 1000, + normalCeilingMs: 30 * 1000, + extendedInitialDelayMs: 5 * MINUTE, + extendedCeilingMs: HOUR, + resetPolicy: new AfterHealthyFor(MINUTE, clock), + random: noJitter, + }); + expect(state.nextDelay).toEqual(0); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(1000); +}); + +it('floors the extended ceiling at the extended initial delay on the transition', () => { + const state = new RetryState({ + normalInitialDelayMs: 1000, + normalCeilingMs: 30 * 1000, + extendedInitialDelayMs: 10 * MINUTE, + extendedCeilingMs: 5 * MINUTE, + resetPolicy: new AfterHealthyFor(MINUTE, clock), + operatingCadenceMs: 0, + random: noJitter, + }); + state.recordFailure('unexpected'); + expect(state.nextDelay).toEqual(10 * MINUTE); + state.recordFailure('normal'); + expect(state.nextDelay).toEqual(10 * MINUTE); +}); + it('collapses the extended regime when the poll interval exceeds its bounds', () => { // With a 2 hour interval the cadence floor makes every wait exactly the // interval: failures of either kind cannot escalate past it, and a success diff --git a/packages/shared/common/src/datasource/retry/ResetPolicy.ts b/packages/shared/common/src/datasource/retry/ResetPolicy.ts index 2a91319fac..57e2313549 100644 --- a/packages/shared/common/src/datasource/retry/ResetPolicy.ts +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -34,18 +34,17 @@ export interface ResetPolicy { */ export class AfterHealthyFor implements ResetPolicy { private _healthySinceMs?: number; + private readonly _healthyForMs: number; private readonly _clock: () => number; /** - * @param _healthyForMs How long the component must operate without failing, - * in milliseconds. + * @param healthyForMs How long the component must operate without failing, + * in milliseconds; must be a positive, finite number. * @param clock The time source used to measure the healthy stretch; defaults * to a monotonic clock. Primarily for testing. */ - constructor( - private readonly _healthyForMs: number, - clock?: () => number, - ) { + constructor(healthyForMs: number, clock?: () => number) { + this._healthyForMs = healthyForMs; this._clock = clock ?? defaultClock(); } @@ -73,6 +72,10 @@ export class AfterHealthyFor implements ResetPolicy { export class AfterConsecutiveSuccesses implements ResetPolicy { private _successes = 0; + /** + * @param _count How many operations in a row must succeed; must be a + * positive integer. + */ constructor(private readonly _count: number) {} noteHealthy(): void { diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts index cc88e21cb8..c3547c1182 100644 --- a/packages/shared/common/src/datasource/retry/RetryState.ts +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -32,6 +32,15 @@ function positiveFiniteOrDefault( return defaultValueMs; } +/** + * Configuration for a {@link RetryState}. + * + * Values are trusted, not validated: each delay bound must be a positive, + * finite number of milliseconds, and the operating cadence a non-negative, + * finite one. The factories ({@link forStreaming}, {@link forPolling}) are the + * validated entry points for anything user-configurable; a caller constructing + * directly is expected to supply known-good values. + */ export interface RetryStateConfig { /** The delay before the first retry in the normal regime, in milliseconds. */ normalInitialDelayMs: number; @@ -42,7 +51,10 @@ export interface RetryStateConfig { /** The delay before the first retry in the extended regime, in milliseconds. */ extendedInitialDelayMs: number; - /** The longest extended-regime delay, in milliseconds. */ + /** + * The longest extended-regime delay, in milliseconds; never less than the + * extended initial delay. + */ extendedCeilingMs: number; /** Decides when the retry state resets. */ @@ -69,7 +81,9 @@ export interface RetryStateConfig { * ceiling, less a random jitter of up to half of it, and never less than the * operating cadence. An unexpected failure moves the state to the extended * regime, which raises both bounds; a normal failure that follows cannot lower - * them. The bounds stay raised until the reset policy is satisfied. + * them. The bounds stay raised until the reset policy is satisfied. A + * server-directed retry time, once applied, replaces the regime's initial + * delay as the base of the computation; the ceiling continues to apply. * * Recording a success sets the next wait back to the operating cadence, even * while the retry state is raised, because a backoff wait applies to a retry @@ -96,9 +110,9 @@ export class RetryState { constructor(config: RetryStateConfig) { this._normalInitialDelayMs = config.normalInitialDelayMs; - this._normalCeilingMs = Math.max(config.normalCeilingMs, config.normalInitialDelayMs); + this._normalCeilingMs = config.normalCeilingMs; this._extendedInitialDelayMs = config.extendedInitialDelayMs; - this._extendedCeilingMs = Math.max(config.extendedCeilingMs, config.extendedInitialDelayMs); + this._extendedCeilingMs = config.extendedCeilingMs; this._operatingCadenceMs = config.operatingCadenceMs ?? 0; this._resetPolicy = config.resetPolicy; this._random = config.random ?? Math.random; @@ -120,7 +134,8 @@ export class RetryState { * Records a failed attempt and decides the wait before the next one. * * The state advances before the wait is computed, so {@link nextDelay} - * always reflects the failure just recorded. + * always reflects the failure just recorded. Callers should read it in the + * same turn, before any other outcome is recorded. */ recordFailure(kind: FailureKind): void { // A reset that fell due during healthy operation is applied before the @@ -134,7 +149,7 @@ export class RetryState { // keeps counting up rather than re-pinning the initial delay. this._extended = true; this._minDelayMs = this._extendedInitialDelayMs; - this._maxDelayMs = this._extendedCeilingMs; + this._maxDelayMs = Math.max(this._extendedCeilingMs, this._extendedInitialDelayMs); this._attempts = 1; } else { this._attempts += 1; @@ -170,10 +185,12 @@ export class RetryState { * Applies a server-directed retry time. * * The value replaces the base of the delay computation and restarts the - * doubling sequence. The regime and its ceiling are unaffected, and the - * value stays in effect until another one arrives, including across a - * reset. Non-finite or negative values are ignored; callers are expected to - * have validated and capped the value at its point of entry. + * doubling sequence, taking precedence over the current regime's initial + * delay — including the extended regime's. The computed delay is still + * bounded by the regime's ceiling. The value stays in effect until another + * one arrives, including across a reset. Non-finite or negative values are + * ignored; callers are expected to have validated and capped the value at + * its point of entry. */ applyServerDirectedRetry(delayMs: number): void { if (typeof delayMs !== 'number' || !Number.isFinite(delayMs) || delayMs < 0) { @@ -201,8 +218,9 @@ export class RetryState { * reset once a connection has been healthy for 60 seconds. * * The initial delay is validated here; the documented default of 1 second - * stands in for anything that is not a positive, finite number. The extended - * regime never starts below the configured delay. + * stands in for anything that is not a positive, finite number. A configured + * delay above the normal ceiling is clamped to it, while the extended regime + * never starts below the configured delay. */ export function forStreaming(initialReconnectDelayMs: number, logger?: LDLogger): RetryState { const validated = positiveFiniteOrDefault( @@ -243,7 +261,7 @@ export function forPolling(pollIntervalMs: number, logger?: LDLogger): RetryStat normalInitialDelayMs: validated, normalCeilingMs: validated, extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), - extendedCeilingMs: Math.max(EXTENDED_CEILING_MS, validated), + extendedCeilingMs: EXTENDED_CEILING_MS, resetPolicy: new AfterConsecutiveSuccesses(POLLING_RESET_SUCCESSES), operatingCadenceMs: validated, }); diff --git a/packages/shared/common/src/index.ts b/packages/shared/common/src/index.ts index 9f88a96b07..e531cd5368 100644 --- a/packages/shared/common/src/index.ts +++ b/packages/shared/common/src/index.ts @@ -2,14 +2,21 @@ import AttributeReference from './AttributeReference'; import Context from './Context'; import ContextFilter from './ContextFilter'; import { + AfterConsecutiveSuccesses, + AfterHealthyFor, Backoff, CompositeDataSource, DataSourceErrorKind, DefaultBackoff, + forPolling, + forStreaming, LDFileDataSourceError, LDFlagDeliveryFallbackError, LDPollingError, LDStreamingError, + ResetPolicy, + RetryState, + RetryStateConfig, StreamingErrorHandler, } from './datasource'; @@ -35,4 +42,11 @@ export { StreamingErrorHandler, LDFileDataSourceError, LDFlagDeliveryFallbackError, + AfterConsecutiveSuccesses, + AfterHealthyFor, + forPolling, + forStreaming, + ResetPolicy, + RetryState, + RetryStateConfig, }; From b9ad84fcd3d6bf86b685a5733d169ecbce787e91 Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Fri, 25 Sep 2026 17:25:12 -0400 Subject: [PATCH 08/10] refactor: express RetryState as an interface with a factory Convert RetryState from an exported class to an interface plus a closure-backed createRetryState factory, per the prefer-interfaces guideline for publicly exposed types. State lives in closure locals, so the returned object is minifiable without the private-field convention and future mutators are additive. The forStreaming/forPolling factories return the interface; the reset-policy classes are unchanged, since the ResetPolicy interface already fronts them. No behavioral change. Part of SDK-2790. --- .../datasource/retry/RetryState.test.ts | 9 +- .../shared/common/src/datasource/index.ts | 2 + .../common/src/datasource/retry/RetryState.ts | 180 +++++++++--------- .../common/src/datasource/retry/index.ts | 9 +- packages/shared/common/src/index.ts | 2 + 5 files changed, 109 insertions(+), 93 deletions(-) diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts index 112fb17a50..593fa4ce34 100644 --- a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -3,6 +3,7 @@ import { AfterHealthyFor, } from '../../../src/datasource/retry/ResetPolicy'; import { + createRetryState, forPolling, forStreaming, RetryState, @@ -22,7 +23,7 @@ beforeEach(() => { }); function streamingState(overrides: Partial = {}): RetryState { - return new RetryState({ + return createRetryState({ normalInitialDelayMs: 1000, normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 5 * MINUTE, @@ -35,7 +36,7 @@ function streamingState(overrides: Partial = {}): RetryState { } function pollingState(intervalMs: number, overrides: Partial = {}): RetryState { - return new RetryState({ + return createRetryState({ normalInitialDelayMs: intervalMs, normalCeilingMs: intervalMs, extendedInitialDelayMs: Math.max(5 * MINUTE, intervalMs), @@ -401,7 +402,7 @@ it('leaves a configured delay below the normal ceiling untouched', () => { }); it('defaults the operating cadence to zero for direct construction', () => { - const state = new RetryState({ + const state = createRetryState({ normalInitialDelayMs: 1000, normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 5 * MINUTE, @@ -415,7 +416,7 @@ it('defaults the operating cadence to zero for direct construction', () => { }); it('floors the extended ceiling at the extended initial delay on the transition', () => { - const state = new RetryState({ + const state = createRetryState({ normalInitialDelayMs: 1000, normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 10 * MINUTE, diff --git a/packages/shared/common/src/datasource/index.ts b/packages/shared/common/src/datasource/index.ts index 6d57d5a913..947f5f7027 100644 --- a/packages/shared/common/src/datasource/index.ts +++ b/packages/shared/common/src/datasource/index.ts @@ -11,6 +11,7 @@ import { import { AfterConsecutiveSuccesses, AfterHealthyFor, + createRetryState, forPolling, forStreaming, ResetPolicy, @@ -23,6 +24,7 @@ export { AfterHealthyFor, Backoff, CompositeDataSource, + createRetryState, DefaultBackoff, DataSourceErrorKind, forPolling, diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts index c3547c1182..90b76d5217 100644 --- a/packages/shared/common/src/datasource/retry/RetryState.ts +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -89,97 +89,31 @@ export interface RetryStateConfig { * while the retry state is raised, because a backoff wait applies to a retry * and not to every operation. * - * This class never starts timers, performs I/O, or logs; callers schedule the - * waits it computes. Use one instance per component; instances are never + * A retry state never starts timers, performs I/O, or logs; callers schedule + * the waits it computes. Use one instance per component; instances are never * shared between components. */ -export class RetryState { - private _attempts = 0; - private _extended = false; - private _minDelayMs: number; - private _maxDelayMs: number; - private _nextDelayMs: number; - private _serverDirectedBaseMs?: number; - private readonly _normalInitialDelayMs: number; - private readonly _normalCeilingMs: number; - private readonly _extendedInitialDelayMs: number; - private readonly _extendedCeilingMs: number; - private readonly _operatingCadenceMs: number; - private readonly _resetPolicy: ResetPolicy; - private readonly _random: () => number; - - constructor(config: RetryStateConfig) { - this._normalInitialDelayMs = config.normalInitialDelayMs; - this._normalCeilingMs = config.normalCeilingMs; - this._extendedInitialDelayMs = config.extendedInitialDelayMs; - this._extendedCeilingMs = config.extendedCeilingMs; - this._operatingCadenceMs = config.operatingCadenceMs ?? 0; - this._resetPolicy = config.resetPolicy; - this._random = config.random ?? Math.random; - this._minDelayMs = this._normalInitialDelayMs; - this._maxDelayMs = this._normalCeilingMs; - // Before any outcome is recorded, the next wait is the ordinary interval. - this._nextDelayMs = this._operatingCadenceMs; - } - +export interface RetryState { /** * The wait, in milliseconds, before the next operation, as the last * recorded outcome decided it. Reading it has no side effects. */ - get nextDelay(): number { - return this._nextDelayMs; - } + readonly nextDelay: number; /** * Records a failed attempt and decides the wait before the next one. * - * The state advances before the wait is computed, so {@link nextDelay} + * The state advances before the wait is computed, so {@link RetryState.nextDelay} * always reflects the failure just recorded. Callers should read it in the * same turn, before any other outcome is recorded. */ - recordFailure(kind: FailureKind): void { - // A reset that fell due during healthy operation is applied before the - // new failure is counted, so the failure computes from a fresh sequence. - this._resetIfDue(); - this._resetPolicy.noteFailure(); - - if (kind === 'unexpected' && !this._extended) { - // Moving to the extended regime raises both bounds and starts the delay - // sequence over. Only the move does this: a later unexpected failure - // keeps counting up rather than re-pinning the initial delay. - this._extended = true; - this._minDelayMs = this._extendedInitialDelayMs; - this._maxDelayMs = Math.max(this._extendedCeilingMs, this._extendedInitialDelayMs); - this._attempts = 1; - } else { - this._attempts += 1; - } - - const base = this._serverDirectedBaseMs ?? this._minDelayMs; - // Compare against the ceiling scaled down rather than the base scaled up, so - // the computed value can never overflow the ceiling. A base of zero doubles - // to zero forever, so it short-circuits. - let target = 0; - if (base > 0) { - const exponent = this._attempts - 1; - target = base >= this._maxDelayMs / 2 ** exponent ? this._maxDelayMs : base * 2 ** exponent; - } - const jitter = (this._random() * target) / 2; - this._nextDelayMs = Math.max(target - jitter, this._operatingCadenceMs); - } + recordFailure(kind: FailureKind): void; /** * Records a successful operation and resets the retry state if that is now * enough to satisfy the reset policy. */ - recordSuccess(): void { - this._resetPolicy.noteHealthy(); - this._resetIfDue(); - // A backoff wait applies to a retry, not to every operation, so after a - // success the next wait is the ordinary interval even while the retry - // state is raised. - this._nextDelayMs = this._operatingCadenceMs; - } + recordSuccess(): void; /** * Applies a server-directed retry time. @@ -192,23 +126,93 @@ export class RetryState { * ignored; callers are expected to have validated and capped the value at * its point of entry. */ - applyServerDirectedRetry(delayMs: number): void { - if (typeof delayMs !== 'number' || !Number.isFinite(delayMs) || delayMs < 0) { - return; - } - this._serverDirectedBaseMs = delayMs; - this._attempts = 0; - } + applyServerDirectedRetry(delayMs: number): void; +} + +/** + * Builds a retry state from explicit bounds. See {@link RetryStateConfig} for + * the contract on the values. + */ +export function createRetryState(config: RetryStateConfig): RetryState { + const normalInitialDelayMs = config.normalInitialDelayMs; + const normalCeilingMs = config.normalCeilingMs; + const extendedInitialDelayMs = config.extendedInitialDelayMs; + const extendedCeilingMs = config.extendedCeilingMs; + const operatingCadenceMs = config.operatingCadenceMs ?? 0; + const { resetPolicy } = config; + const random = config.random ?? Math.random; + + let attempts = 0; + let extended = false; + let minDelayMs = normalInitialDelayMs; + let maxDelayMs = normalCeilingMs; + let serverDirectedBaseMs: number | undefined; + // Before any outcome is recorded, the next wait is the ordinary interval. + let nextDelayMs = operatingCadenceMs; - private _resetIfDue(): void { - if (!this._resetPolicy.isSatisfied()) { + function resetIfDue(): void { + if (!resetPolicy.isSatisfied()) { return; } - this._attempts = 0; - this._extended = false; - this._minDelayMs = this._normalInitialDelayMs; - this._maxDelayMs = this._normalCeilingMs; + attempts = 0; + extended = false; + minDelayMs = normalInitialDelayMs; + maxDelayMs = normalCeilingMs; } + + return { + get nextDelay(): number { + return nextDelayMs; + }, + + recordFailure(kind: FailureKind): void { + // A reset that fell due during healthy operation is applied before the + // new failure is counted, so the failure computes from a fresh sequence. + resetIfDue(); + resetPolicy.noteFailure(); + + if (kind === 'unexpected' && !extended) { + // Moving to the extended regime raises both bounds and starts the delay + // sequence over. Only the move does this: a later unexpected failure + // keeps counting up rather than re-pinning the initial delay. + extended = true; + minDelayMs = extendedInitialDelayMs; + maxDelayMs = Math.max(extendedCeilingMs, extendedInitialDelayMs); + attempts = 1; + } else { + attempts += 1; + } + + const base = serverDirectedBaseMs ?? minDelayMs; + // Compare against the ceiling scaled down rather than the base scaled up, so + // the computed value can never overflow the ceiling. A base of zero doubles + // to zero forever, so it short-circuits. + let target = 0; + if (base > 0) { + const exponent = attempts - 1; + target = base >= maxDelayMs / 2 ** exponent ? maxDelayMs : base * 2 ** exponent; + } + const jitter = (random() * target) / 2; + nextDelayMs = Math.max(target - jitter, operatingCadenceMs); + }, + + recordSuccess(): void { + resetPolicy.noteHealthy(); + resetIfDue(); + // A backoff wait applies to a retry, not to every operation, so after a + // success the next wait is the ordinary interval even while the retry + // state is raised. + nextDelayMs = operatingCadenceMs; + }, + + applyServerDirectedRetry(delayMs: number): void { + if (typeof delayMs !== 'number' || !Number.isFinite(delayMs) || delayMs < 0) { + return; + } + serverDirectedBaseMs = delayMs; + attempts = 0; + }, + }; } /** @@ -229,7 +233,7 @@ export function forStreaming(initialReconnectDelayMs: number, logger?: LDLogger) 'initialReconnectDelayMs', logger, ); - return new RetryState({ + return createRetryState({ normalInitialDelayMs: validated, normalCeilingMs: STREAMING_NORMAL_CEILING_MS, extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), @@ -257,7 +261,7 @@ export function forPolling(pollIntervalMs: number, logger?: LDLogger): RetryStat 'pollIntervalMs', logger, ); - return new RetryState({ + return createRetryState({ normalInitialDelayMs: validated, normalCeilingMs: validated, extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), diff --git a/packages/shared/common/src/datasource/retry/index.ts b/packages/shared/common/src/datasource/retry/index.ts index 7e2dc819ab..9d7da1cb08 100644 --- a/packages/shared/common/src/datasource/retry/index.ts +++ b/packages/shared/common/src/datasource/retry/index.ts @@ -1,9 +1,16 @@ import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; -import { forPolling, forStreaming, RetryState, RetryStateConfig } from './RetryState'; +import { + createRetryState, + forPolling, + forStreaming, + RetryState, + RetryStateConfig, +} from './RetryState'; export { AfterConsecutiveSuccesses, AfterHealthyFor, + createRetryState, forPolling, forStreaming, ResetPolicy, diff --git a/packages/shared/common/src/index.ts b/packages/shared/common/src/index.ts index e531cd5368..a6136ee3fd 100644 --- a/packages/shared/common/src/index.ts +++ b/packages/shared/common/src/index.ts @@ -6,6 +6,7 @@ import { AfterHealthyFor, Backoff, CompositeDataSource, + createRetryState, DataSourceErrorKind, DefaultBackoff, forPolling, @@ -44,6 +45,7 @@ export { LDFlagDeliveryFallbackError, AfterConsecutiveSuccesses, AfterHealthyFor, + createRetryState, forPolling, forStreaming, ResetPolicy, From 44034881490f51881c82d4b9fc958403e0cf8adc Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Mon, 28 Sep 2026 13:57:35 -0400 Subject: [PATCH 09/10] ci: raise js-sdk-common package size limit to 29500 --- .github/workflows/common.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/common.yml b/.github/workflows/common.yml index ae50fb8823..7734d26bca 100644 --- a/.github/workflows/common.yml +++ b/.github/workflows/common.yml @@ -35,4 +35,4 @@ jobs: target_file: 'packages/shared/common/dist/esm/index.mjs' package_name: '@launchdarkly/js-sdk-common' pr_number: ${{ github.event.number }} - size_limit: 29000 + size_limit: 29500 From 7f529cf560a43be79bb71eb8f4ed5fb16ab42267 Mon Sep 17 00:00:00 2001 From: Todd Anderson Date: Mon, 28 Sep 2026 15:40:16 -0400 Subject: [PATCH 10/10] refactor: express reset policies as factory functions --- .../datasource/retry/ResetPolicy.test.ts | 20 ++-- .../datasource/retry/RetryState.test.ts | 12 +-- .../shared/common/src/datasource/index.ts | 8 +- .../src/datasource/retry/ResetPolicy.ts | 94 +++++++++---------- .../common/src/datasource/retry/RetryState.ts | 6 +- .../common/src/datasource/retry/index.ts | 6 +- packages/shared/common/src/index.ts | 8 +- 7 files changed, 74 insertions(+), 80 deletions(-) diff --git a/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts index c63c34e10f..6372b30149 100644 --- a/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts @@ -1,6 +1,6 @@ import { - AfterConsecutiveSuccesses, - AfterHealthyFor, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, } from '../../../src/datasource/retry/ResetPolicy'; const MINUTE = 60 * 1000; @@ -13,13 +13,13 @@ beforeEach(() => { }); it('is not satisfied before any healthy report', () => { - const policy = new AfterHealthyFor(MINUTE, clock); + const policy = createAfterHealthyFor(MINUTE, clock); now = 10 * MINUTE; expect(policy.isSatisfied()).toEqual(false); }); it('is satisfied once the healthy stretch reaches the threshold, and not a moment sooner', () => { - const policy = new AfterHealthyFor(MINUTE, clock); + const policy = createAfterHealthyFor(MINUTE, clock); policy.noteHealthy(); now = MINUTE - 1; expect(policy.isSatisfied()).toEqual(false); @@ -28,7 +28,7 @@ it('is satisfied once the healthy stretch reaches the threshold, and not a momen }); it('does not move the start of the stretch on repeated healthy reports', () => { - const policy = new AfterHealthyFor(MINUTE, clock); + const policy = createAfterHealthyFor(MINUTE, clock); policy.noteHealthy(); for (let i = 1; i < 60; i += 1) { now = i * 1000; @@ -40,7 +40,7 @@ it('does not move the start of the stretch on repeated healthy reports', () => { }); it('clears the healthy stretch on a failure', () => { - const policy = new AfterHealthyFor(MINUTE, clock); + const policy = createAfterHealthyFor(MINUTE, clock); policy.noteHealthy(); now = 30 * 1000; policy.noteFailure(); @@ -49,7 +49,7 @@ it('clears the healthy stretch on a failure', () => { }); it('is satisfied by consecutive successes and not by fewer', () => { - const policy = new AfterConsecutiveSuccesses(2); + const policy = createAfterConsecutiveSuccesses(2); expect(policy.isSatisfied()).toEqual(false); policy.noteHealthy(); expect(policy.isSatisfied()).toEqual(false); @@ -58,7 +58,7 @@ it('is satisfied by consecutive successes and not by fewer', () => { }); it('starts the success count over on a failure', () => { - const policy = new AfterConsecutiveSuccesses(2); + const policy = createAfterConsecutiveSuccesses(2); policy.noteHealthy(); policy.noteFailure(); policy.noteHealthy(); @@ -68,7 +68,7 @@ it('starts the success count over on a failure', () => { }); it('measures the healthy stretch with the default clock', () => { - const policy = new AfterHealthyFor(60 * MINUTE); + const policy = createAfterHealthyFor(60 * MINUTE); expect(policy.isSatisfied()).toEqual(false); policy.noteHealthy(); expect(policy.isSatisfied()).toEqual(false); @@ -79,7 +79,7 @@ it('falls back to the wall clock when no monotonic source exists', () => { const nowSpy = jest.spyOn(Date, 'now').mockReturnValue(0); delete (globalThis as any).performance; try { - const policy = new AfterHealthyFor(MINUTE); + const policy = createAfterHealthyFor(MINUTE); policy.noteHealthy(); expect(policy.isSatisfied()).toEqual(false); nowSpy.mockReturnValue(MINUTE); diff --git a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts index 593fa4ce34..27824c7ee7 100644 --- a/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -1,6 +1,6 @@ import { - AfterConsecutiveSuccesses, - AfterHealthyFor, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, } from '../../../src/datasource/retry/ResetPolicy'; import { createRetryState, @@ -28,7 +28,7 @@ function streamingState(overrides: Partial = {}): RetryState { normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 5 * MINUTE, extendedCeilingMs: HOUR, - resetPolicy: new AfterHealthyFor(MINUTE, clock), + resetPolicy: createAfterHealthyFor(MINUTE, clock), operatingCadenceMs: 0, random: noJitter, ...overrides, @@ -41,7 +41,7 @@ function pollingState(intervalMs: number, overrides: Partial = normalCeilingMs: intervalMs, extendedInitialDelayMs: Math.max(5 * MINUTE, intervalMs), extendedCeilingMs: Math.max(HOUR, intervalMs), - resetPolicy: new AfterConsecutiveSuccesses(2), + resetPolicy: createAfterConsecutiveSuccesses(2), operatingCadenceMs: intervalMs, random: noJitter, ...overrides, @@ -407,7 +407,7 @@ it('defaults the operating cadence to zero for direct construction', () => { normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 5 * MINUTE, extendedCeilingMs: HOUR, - resetPolicy: new AfterHealthyFor(MINUTE, clock), + resetPolicy: createAfterHealthyFor(MINUTE, clock), random: noJitter, }); expect(state.nextDelay).toEqual(0); @@ -421,7 +421,7 @@ it('floors the extended ceiling at the extended initial delay on the transition' normalCeilingMs: 30 * 1000, extendedInitialDelayMs: 10 * MINUTE, extendedCeilingMs: 5 * MINUTE, - resetPolicy: new AfterHealthyFor(MINUTE, clock), + resetPolicy: createAfterHealthyFor(MINUTE, clock), operatingCadenceMs: 0, random: noJitter, }); diff --git a/packages/shared/common/src/datasource/index.ts b/packages/shared/common/src/datasource/index.ts index 947f5f7027..ba7337801e 100644 --- a/packages/shared/common/src/datasource/index.ts +++ b/packages/shared/common/src/datasource/index.ts @@ -9,8 +9,8 @@ import { StreamingErrorHandler, } from './errors'; import { - AfterConsecutiveSuccesses, - AfterHealthyFor, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, createRetryState, forPolling, forStreaming, @@ -20,10 +20,10 @@ import { } from './retry'; export { - AfterConsecutiveSuccesses, - AfterHealthyFor, Backoff, CompositeDataSource, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, createRetryState, DefaultBackoff, DataSourceErrorKind, diff --git a/packages/shared/common/src/datasource/retry/ResetPolicy.ts b/packages/shared/common/src/datasource/retry/ResetPolicy.ts index 57e2313549..28ba62c6d6 100644 --- a/packages/shared/common/src/datasource/retry/ResetPolicy.ts +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -28,65 +28,59 @@ export interface ResetPolicy { } /** - * Resets once the component has operated without failing for the given - * duration. Repeated healthy reports do not move the starting point; only a - * failure clears it. + * Creates a policy that resets once the component has operated without failing + * for the given duration. Repeated healthy reports do not move the starting + * point; only a failure clears it. + * + * @param healthyForMs How long the component must operate without failing, + * in milliseconds; must be a positive, finite number. + * @param clock The time source used to measure the healthy stretch; defaults + * to a monotonic clock. Primarily for testing. */ -export class AfterHealthyFor implements ResetPolicy { - private _healthySinceMs?: number; - private readonly _healthyForMs: number; - private readonly _clock: () => number; +export function createAfterHealthyFor( + healthyForMs: number, + clock: () => number = defaultClock(), +): ResetPolicy { + let healthySinceMs: number | undefined; - /** - * @param healthyForMs How long the component must operate without failing, - * in milliseconds; must be a positive, finite number. - * @param clock The time source used to measure the healthy stretch; defaults - * to a monotonic clock. Primarily for testing. - */ - constructor(healthyForMs: number, clock?: () => number) { - this._healthyForMs = healthyForMs; - this._clock = clock ?? defaultClock(); - } + return { + noteHealthy(): void { + if (healthySinceMs === undefined) { + healthySinceMs = clock(); + } + }, - noteHealthy(): void { - if (this._healthySinceMs === undefined) { - this._healthySinceMs = this._clock(); - } - } + noteFailure(): void { + healthySinceMs = undefined; + }, - noteFailure(): void { - this._healthySinceMs = undefined; - } - - isSatisfied(): boolean { - return ( - this._healthySinceMs !== undefined && - this._clock() - this._healthySinceMs >= this._healthyForMs - ); - } + isSatisfied(): boolean { + return healthySinceMs !== undefined && clock() - healthySinceMs >= healthyForMs; + }, + }; } /** - * Resets once the given number of operations in a row have succeeded. + * Creates a policy that resets once the given number of operations in a row + * have succeeded. + * + * @param count How many operations in a row must succeed; must be a positive + * integer. */ -export class AfterConsecutiveSuccesses implements ResetPolicy { - private _successes = 0; +export function createAfterConsecutiveSuccesses(count: number): ResetPolicy { + let successes = 0; - /** - * @param _count How many operations in a row must succeed; must be a - * positive integer. - */ - constructor(private readonly _count: number) {} - - noteHealthy(): void { - this._successes += 1; - } + return { + noteHealthy(): void { + successes += 1; + }, - noteFailure(): void { - this._successes = 0; - } + noteFailure(): void { + successes = 0; + }, - isSatisfied(): boolean { - return this._successes >= this._count; - } + isSatisfied(): boolean { + return successes >= count; + }, + }; } diff --git a/packages/shared/common/src/datasource/retry/RetryState.ts b/packages/shared/common/src/datasource/retry/RetryState.ts index 90b76d5217..cc022cb857 100644 --- a/packages/shared/common/src/datasource/retry/RetryState.ts +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -1,6 +1,6 @@ import { LDLogger } from '../../api/logging/LDLogger'; import { FailureKind } from '../../errors'; -import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; +import { createAfterConsecutiveSuccesses, createAfterHealthyFor, ResetPolicy } from './ResetPolicy'; // The delay bounds of the extended regime, in milliseconds. A component enters // the extended regime after an unexpected failure and leaves it when its reset @@ -238,7 +238,7 @@ export function forStreaming(initialReconnectDelayMs: number, logger?: LDLogger) normalCeilingMs: STREAMING_NORMAL_CEILING_MS, extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), extendedCeilingMs: EXTENDED_CEILING_MS, - resetPolicy: new AfterHealthyFor(STREAMING_RESET_INTERVAL_MS), + resetPolicy: createAfterHealthyFor(STREAMING_RESET_INTERVAL_MS), operatingCadenceMs: 0, }); } @@ -266,7 +266,7 @@ export function forPolling(pollIntervalMs: number, logger?: LDLogger): RetryStat normalCeilingMs: validated, extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), extendedCeilingMs: EXTENDED_CEILING_MS, - resetPolicy: new AfterConsecutiveSuccesses(POLLING_RESET_SUCCESSES), + resetPolicy: createAfterConsecutiveSuccesses(POLLING_RESET_SUCCESSES), operatingCadenceMs: validated, }); } diff --git a/packages/shared/common/src/datasource/retry/index.ts b/packages/shared/common/src/datasource/retry/index.ts index 9d7da1cb08..41e9e5c248 100644 --- a/packages/shared/common/src/datasource/retry/index.ts +++ b/packages/shared/common/src/datasource/retry/index.ts @@ -1,4 +1,4 @@ -import { AfterConsecutiveSuccesses, AfterHealthyFor, ResetPolicy } from './ResetPolicy'; +import { createAfterConsecutiveSuccesses, createAfterHealthyFor, ResetPolicy } from './ResetPolicy'; import { createRetryState, forPolling, @@ -8,8 +8,8 @@ import { } from './RetryState'; export { - AfterConsecutiveSuccesses, - AfterHealthyFor, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, createRetryState, forPolling, forStreaming, diff --git a/packages/shared/common/src/index.ts b/packages/shared/common/src/index.ts index a6136ee3fd..64c9a0cc86 100644 --- a/packages/shared/common/src/index.ts +++ b/packages/shared/common/src/index.ts @@ -2,10 +2,10 @@ import AttributeReference from './AttributeReference'; import Context from './Context'; import ContextFilter from './ContextFilter'; import { - AfterConsecutiveSuccesses, - AfterHealthyFor, Backoff, CompositeDataSource, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, createRetryState, DataSourceErrorKind, DefaultBackoff, @@ -43,8 +43,8 @@ export { StreamingErrorHandler, LDFileDataSourceError, LDFlagDeliveryFallbackError, - AfterConsecutiveSuccesses, - AfterHealthyFor, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, createRetryState, forPolling, forStreaming,