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 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..6372b30149 --- /dev/null +++ b/packages/shared/common/__tests__/datasource/retry/ResetPolicy.test.ts @@ -0,0 +1,91 @@ +import { + createAfterConsecutiveSuccesses, + createAfterHealthyFor, +} 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 = 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 = createAfterHealthyFor(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 = createAfterHealthyFor(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 = createAfterHealthyFor(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 = createAfterConsecutiveSuccesses(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 = createAfterConsecutiveSuccesses(2); + policy.noteHealthy(); + policy.noteFailure(); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(false); + policy.noteHealthy(); + expect(policy.isSatisfied()).toEqual(true); +}); + +it('measures the healthy stretch with the default clock', () => { + const policy = createAfterHealthyFor(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 = createAfterHealthyFor(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 new file mode 100644 index 0000000000..27824c7ee7 --- /dev/null +++ b/packages/shared/common/__tests__/datasource/retry/RetryState.test.ts @@ -0,0 +1,455 @@ +import { + createAfterConsecutiveSuccesses, + createAfterHealthyFor, +} from '../../../src/datasource/retry/ResetPolicy'; +import { + createRetryState, + 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 createRetryState({ + normalInitialDelayMs: 1000, + normalCeilingMs: 30 * 1000, + extendedInitialDelayMs: 5 * MINUTE, + extendedCeilingMs: HOUR, + resetPolicy: createAfterHealthyFor(MINUTE, clock), + operatingCadenceMs: 0, + random: noJitter, + ...overrides, + }); +} + +function pollingState(intervalMs: number, overrides: Partial = {}): RetryState { + return createRetryState({ + normalInitialDelayMs: intervalMs, + normalCeilingMs: intervalMs, + extendedInitialDelayMs: Math.max(5 * MINUTE, intervalMs), + extendedCeilingMs: Math.max(HOUR, intervalMs), + resetPolicy: createAfterConsecutiveSuccesses(2), + operatingCadenceMs: intervalMs, + 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', () => { + const delays: number[] = []; + 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); + 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 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 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', () => { + 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', () => { + 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('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) => { + 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('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); +}); + +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() }; + 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); + }, +); + +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('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 + // ceiling wins rather than being raised to the configured value. + const state = forStreaming(10 * MINUTE); + state.recordFailure('normal'); + expect(state.nextDelay).toBeGreaterThan(15 * 1000); + expect(state.nextDelay).toBeLessThanOrEqual(30 * 1000); + state.recordFailure('normal'); + expect(state.nextDelay).toBeGreaterThan(15 * 1000); + expect(state.nextDelay).toBeLessThanOrEqual(30 * 1000); +}); + +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('defaults the operating cadence to zero for direct construction', () => { + const state = createRetryState({ + normalInitialDelayMs: 1000, + normalCeilingMs: 30 * 1000, + extendedInitialDelayMs: 5 * MINUTE, + extendedCeilingMs: HOUR, + resetPolicy: createAfterHealthyFor(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 = createRetryState({ + normalInitialDelayMs: 1000, + normalCeilingMs: 30 * 1000, + extendedInitialDelayMs: 10 * MINUTE, + extendedCeilingMs: 5 * MINUTE, + resetPolicy: createAfterHealthyFor(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 + // 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() }; + forStreaming(1000, logger); + forPolling(30 * 1000, logger); + expect(warn).not.toHaveBeenCalled(); +}); diff --git a/packages/shared/common/__tests__/errors.test.ts b/packages/shared/common/__tests__/errors.test.ts new file mode 100644 index 0000000000..6d07cfad5c --- /dev/null +++ b/packages/shared/common/__tests__/errors.test.ts @@ -0,0 +1,38 @@ +import { classifyHttpStatus, classifyTransportFailure, isHttpRecoverable } from '../src/errors'; + +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'); +}); + +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 497a7588fe..ba7337801e 100644 --- a/packages/shared/common/src/datasource/index.ts +++ b/packages/shared/common/src/datasource/index.ts @@ -8,15 +8,33 @@ import { LDStreamingError, StreamingErrorHandler, } from './errors'; +import { + createAfterConsecutiveSuccesses, + createAfterHealthyFor, + createRetryState, + forPolling, + forStreaming, + ResetPolicy, + RetryState, + RetryStateConfig, +} from './retry'; export { Backoff, CompositeDataSource, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, + createRetryState, DefaultBackoff, DataSourceErrorKind, + 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..28ba62c6d6 --- /dev/null +++ b/packages/shared/common/src/datasource/retry/ResetPolicy.ts @@ -0,0 +1,86 @@ +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. + */ +export interface ResetPolicy { + /** + * Records that the component is operating normally. + */ + noteHealthy(): void; + + /** + * Records a failure, which ends any healthy stretch in progress. + */ + noteFailure(): void; + + /** + * Reports whether the reset condition is met. + */ + isSatisfied(): boolean; +} + +/** + * 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 function createAfterHealthyFor( + healthyForMs: number, + clock: () => number = defaultClock(), +): ResetPolicy { + let healthySinceMs: number | undefined; + + return { + noteHealthy(): void { + if (healthySinceMs === undefined) { + healthySinceMs = clock(); + } + }, + + noteFailure(): void { + healthySinceMs = undefined; + }, + + isSatisfied(): boolean { + return healthySinceMs !== undefined && clock() - healthySinceMs >= healthyForMs; + }, + }; +} + +/** + * 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 function createAfterConsecutiveSuccesses(count: number): ResetPolicy { + let successes = 0; + + return { + noteHealthy(): void { + successes += 1; + }, + + noteFailure(): void { + successes = 0; + }, + + 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 new file mode 100644 index 0000000000..cc022cb857 --- /dev/null +++ b/packages/shared/common/src/datasource/retry/RetryState.ts @@ -0,0 +1,272 @@ +import { LDLogger } from '../../api/logging/LDLogger'; +import { FailureKind } from '../../errors'; +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 +// 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; + +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; +} + +/** + * 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; + + /** 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; never less than the + * extended initial delay. + */ + 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 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. 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 + * and not to every operation. + * + * 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 interface RetryState { + /** + * The wait, in milliseconds, before the next operation, as the last + * recorded outcome decided it. Reading it has no side effects. + */ + 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 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; + + /** + * Records a successful operation and resets the retry state if that is now + * enough to satisfy the reset policy. + */ + recordSuccess(): void; + + /** + * Applies a server-directed retry time. + * + * The value replaces the base of the delay computation and restarts the + * 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; +} + +/** + * 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; + + function resetIfDue(): void { + if (!resetPolicy.isSatisfied()) { + return; + } + 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; + }, + }; +} + +/** + * 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. 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( + initialReconnectDelayMs, + DEFAULT_STREAMING_INITIAL_DELAY_MS, + 'initialReconnectDelayMs', + logger, + ); + return createRetryState({ + normalInitialDelayMs: validated, + normalCeilingMs: STREAMING_NORMAL_CEILING_MS, + extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), + extendedCeilingMs: EXTENDED_CEILING_MS, + resetPolicy: createAfterHealthyFor(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 createRetryState({ + normalInitialDelayMs: validated, + normalCeilingMs: validated, + extendedInitialDelayMs: Math.max(EXTENDED_INITIAL_DELAY_MS, validated), + extendedCeilingMs: EXTENDED_CEILING_MS, + 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 new file mode 100644 index 0000000000..41e9e5c248 --- /dev/null +++ b/packages/shared/common/src/datasource/retry/index.ts @@ -0,0 +1,19 @@ +import { createAfterConsecutiveSuccesses, createAfterHealthyFor, ResetPolicy } from './ResetPolicy'; +import { + createRetryState, + forPolling, + forStreaming, + RetryState, + RetryStateConfig, +} from './RetryState'; + +export { + createAfterConsecutiveSuccesses, + createAfterHealthyFor, + createRetryState, + forPolling, + forStreaming, + ResetPolicy, + RetryState, + RetryStateConfig, +}; 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'; } /** diff --git a/packages/shared/common/src/index.ts b/packages/shared/common/src/index.ts index 9f88a96b07..64c9a0cc86 100644 --- a/packages/shared/common/src/index.ts +++ b/packages/shared/common/src/index.ts @@ -4,12 +4,20 @@ import ContextFilter from './ContextFilter'; import { Backoff, CompositeDataSource, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, + createRetryState, DataSourceErrorKind, DefaultBackoff, + forPolling, + forStreaming, LDFileDataSourceError, LDFlagDeliveryFallbackError, LDPollingError, LDStreamingError, + ResetPolicy, + RetryState, + RetryStateConfig, StreamingErrorHandler, } from './datasource'; @@ -35,4 +43,12 @@ export { StreamingErrorHandler, LDFileDataSourceError, LDFlagDeliveryFallbackError, + createAfterConsecutiveSuccesses, + createAfterHealthyFor, + createRetryState, + forPolling, + forStreaming, + ResetPolicy, + RetryState, + RetryStateConfig, };