diff --git a/.github/workflows/server-node-nightly.yaml b/.github/workflows/server-node-nightly.yaml new file mode 100644 index 0000000000..2908325885 --- /dev/null +++ b/.github/workflows/server-node-nightly.yaml @@ -0,0 +1,41 @@ +name: sdk/server-node nightly + +on: + schedule: + - cron: '0 3 * * *' + workflow_dispatch: + +jobs: + contract-tests: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + - uses: ./actions/setup-yarn + with: + node-version: 22 + registry-url: 'https://registry.npmjs.org' + - id: shared + name: Shared CI Steps + uses: ./actions/ci + with: + workspace_name: '@launchdarkly/node-server-sdk' + workspace_path: packages/sdk/server-node + - name: Install contract test service dependencies + env: + ELECTRON_SKIP_BINARY_DOWNLOAD: '1' + run: yarn workspace @launchdarkly/node-server-sdk-contract-tests install --no-immutable + - name: Build shared contract test utils (server) + run: yarn workspace @launchdarkly/js-contract-test-utils build:server + - name: Build the test service + run: yarn workspace @launchdarkly/node-server-sdk-contract-tests build + - name: Launch the test service in the background + run: yarn workspace @launchdarkly/node-server-sdk-contract-tests start 2>&1 & + # The retry-conformance suite is long-running (real backoff waits), so it + # runs here nightly rather than on every PR. + - name: Run contract tests (long-running) + uses: launchdarkly/gh-actions/actions/contract-tests@contract-tests-v1 + with: + test_service_port: 8000 + token: ${{ secrets.GITHUB_TOKEN }} + stop_service: 'false' + extra_params: '--enable-long-running-tests --skip-from=./packages/sdk/server-node/contract-tests/testharness-suppressions.txt' diff --git a/.yarnrc.yml b/.yarnrc.yml index 39bb43d906..4bb2d4d9ae 100644 --- a/.yarnrc.yml +++ b/.yarnrc.yml @@ -7,6 +7,12 @@ nodeLinker: node-modules npmMinimalAgeGate: 4320 +# Exempt this first-party package from the minimum-age gate so freshly published +# versions install without waiting it out. Exact name (not a glob) so it cannot +# match an attacker-squatted launchdarkly-* package. +npmPreapprovedPackages: + - 'launchdarkly-eventsource' + npmPublishAccess: public npmScopes: diff --git a/packages/sdk/browser/src/platform/DefaultBrowserEventSource.ts b/packages/sdk/browser/src/platform/DefaultBrowserEventSource.ts index 4a2ea36998..c4a6c099c1 100644 --- a/packages/sdk/browser/src/platform/DefaultBrowserEventSource.ts +++ b/packages/sdk/browser/src/platform/DefaultBrowserEventSource.ts @@ -35,8 +35,8 @@ export default class DefaultBrowserEventSource implements LDEventSource { options: EventSourceInitDict, ) { this._backoff = new DefaultBackoff( - options.initialRetryDelayMillis, - options.retryResetIntervalMillis, + options.initialRetryDelayMillis ?? 1000, + options.retryResetIntervalMillis ?? 60 * 1000, ); this._errorFilter = options.errorFilter; this._urlBuilder = options.urlBuilder; diff --git a/packages/sdk/server-node/contract-tests/src/index.ts b/packages/sdk/server-node/contract-tests/src/index.ts index 7a7f8c5b89..729e908b44 100644 --- a/packages/sdk/server-node/contract-tests/src/index.ts +++ b/packages/sdk/server-node/contract-tests/src/index.ts @@ -24,6 +24,8 @@ app.get('/', (req: Request, res: Response) => { capabilities: [ 'server-side-polling', 'server-side', + 'retry-conformance-fdv1-streaming', + 'retry-conformance-fdv1-polling', 'all-flags-client-side-only', 'all-flags-details-only-for-tracked-flags', 'all-flags-with-reasons', diff --git a/packages/sdk/server-node/package.json b/packages/sdk/server-node/package.json index a0cd2caa89..69b59f2620 100644 --- a/packages/sdk/server-node/package.json +++ b/packages/sdk/server-node/package.json @@ -49,7 +49,7 @@ "dependencies": { "@launchdarkly/js-server-sdk-common": "2.21.6", "https-proxy-agent": "^7.0.6", - "launchdarkly-eventsource": "2.2.0" + "launchdarkly-eventsource": "2.3.0" }, "devDependencies": { "@types/jest": "^29.4.0", diff --git a/packages/shared/common/src/api/platform/EventSource.ts b/packages/shared/common/src/api/platform/EventSource.ts index 559581a19e..1103b8cd18 100644 --- a/packages/shared/common/src/api/platform/EventSource.ts +++ b/packages/shared/common/src/api/platform/EventSource.ts @@ -17,14 +17,44 @@ export interface EventSource { close(): void; } +/** + * The strategy an {@link EventSource} uses to decide how long to wait before + * each reconnection attempt. It is injected so a data source can drive + * reconnection timing from its own retry state rather than the transport's + * built-in backoff. The `nowMs` arguments are the transport's wall clock; an + * implementation that keeps its own clock may ignore them. + */ +export interface EventSourceRetryDelayStrategy { + /** Returns the delay, in milliseconds, before the next reconnection. */ + nextRetryDelay(nowMs: number): number; + /** Records that the connection is currently healthy. */ + setGoodSince(nowMs: number): void; + /** Applies a server-directed base delay, in milliseconds. */ + setBaseDelay(baseDelayMs: number): void; +} + export interface EventSourceInitDict { method?: string; headers: { [key: string]: string | string[] }; body?: string; errorFilter: (err: HttpErrorResponse) => boolean; - initialRetryDelayMillis: number; readTimeoutMillis: number; - retryResetIntervalMillis: number; + + /** + * The initial delay and reset window for the built-in (default) retry-delay + * strategy. These configure the default behavior; they are ignored when a + * custom {@link retryDelayStrategy} is provided, since that strategy then + * owns all reconnection timing. + */ + initialRetryDelayMillis?: number; + retryResetIntervalMillis?: number; + + /** + * A custom strategy that replaces the built-in (default) one. When provided, + * the EventSource defers all reconnection timing to it and the built-in + * options (`initialRetryDelayMillis`, `retryResetIntervalMillis`) are ignored. + */ + retryDelayStrategy?: EventSourceRetryDelayStrategy; /** * Optional callback that returns a fresh URL on each reconnection attempt. * When provided, the EventSource implementation should call this instead of diff --git a/packages/shared/sdk-server/__tests__/data_sources/PollingProcessor.test.ts b/packages/shared/sdk-server/__tests__/data_sources/PollingProcessor.test.ts index 85cfa588a4..79c31d4749 100644 --- a/packages/shared/sdk-server/__tests__/data_sources/PollingProcessor.test.ts +++ b/packages/shared/sdk-server/__tests__/data_sources/PollingProcessor.test.ts @@ -164,7 +164,7 @@ describe('given a polling processor with a short poll duration', () => { processor.start(); expect(initSuccessHandler).not.toBeCalled(); - expect(errorHandler.mock.lastCall[0].message).toMatch(/malformed json/i); + expect(errorHandler).not.toBeCalled(); setTimeout(() => { expect(requestor.requestAllData.mock.calls.length).toBeGreaterThanOrEqual(2); @@ -173,26 +173,77 @@ describe('given a polling processor with a short poll duration', () => { }, 300); }); - it.each([401, 403])( - 'does not continue after non-recoverable error', - (status, done) => { - requestor.requestAllData = jest.fn((cb) => - cb( - { - status, - }, - undefined, - ), - ); + it('continues polling when deserialization throws on structurally invalid data', (done) => { + // Parses as JSON but throws during the revive step (here, a null flag + // entry). It must be handled like unparseable data, not escape and kill + // the poll loop. + requestor.requestAllData = jest.fn((cb) => cb(undefined, '{"flags":{"x":null},"segments":{}}')); + + processor.start(); + + expect(initSuccessHandler).not.toBeCalled(); + expect(errorHandler).not.toBeCalled(); + + setTimeout(() => { + expect(requestor.requestAllData.mock.calls.length).toBeGreaterThanOrEqual(2); + expect(testLogger.getCount(LogLevel.Error)).toBeGreaterThan(2); + (done as jest.DoneCallback)(); + }, 300); + }); + + it('cancels the scheduled poll when stopped before it fires', () => { + jest.useFakeTimers(); + try { + requestor.requestAllData = jest.fn((cb) => cb({ status: 500 }, undefined)); processor.start(); - expect(initSuccessHandler).not.toBeCalled(); - expect(errorHandler.mock.lastCall[0].message).toMatch(new RegExp(`${status}.*permanently`)); + expect(requestor.requestAllData).toHaveBeenCalledTimes(1); - setTimeout(() => { - expect(requestor.requestAllData.mock.calls.length).toBe(1); + // The failed poll armed the next-poll timer; stopping must cancel it. + processor.stop(); + jest.advanceTimersByTime(5 * 60 * 1000); + + expect(requestor.requestAllData).toHaveBeenCalledTimes(1); + } finally { + jest.useRealTimers(); + } + }); + + it('treats a status-less transport error as a normal, retryable failure', (done) => { + requestor.requestAllData = jest.fn((cb) => cb({ message: 'socket hang up' } as any, undefined)); + + processor.start(); + + expect(errorHandler).not.toBeCalled(); + setTimeout(() => { + expect(requestor.requestAllData.mock.calls.length).toBeGreaterThanOrEqual(2); + expect(testLogger.getCount(LogLevel.Error)).toBe(0); + expect(testLogger.getCount(LogLevel.Warn)).toBeGreaterThan(2); + (done as jest.DoneCallback)(); + }, 300); + }); + + it.each([401, 403])( + 'retries with extended backoff rather than stopping after error %p', + (status) => { + jest.useFakeTimers(); + try { + requestor.requestAllData = jest.fn((cb) => cb({ status }, undefined)); + processor.start(); + + expect(initSuccessHandler).not.toBeCalled(); + expect(requestor.requestAllData).toHaveBeenCalledTimes(1); + // Previously terminal; now an 'error'-level *log* that keeps retrying — + // like any recoverable failure, it is not surfaced as an 'error' event. + expect(errorHandler).not.toBeCalled(); expect(testLogger.getCount(LogLevel.Error)).toBe(1); - (done as jest.DoneCallback)(); - }, 300); + + // Not a permanent stop: an unexpected failure schedules the next poll in + // the extended regime (~5 minutes), so polling resumes after that wait. + jest.advanceTimersByTime(5 * 60 * 1000); + expect(requestor.requestAllData.mock.calls.length).toBeGreaterThanOrEqual(2); + } finally { + jest.useRealTimers(); + } }, ); }); diff --git a/packages/shared/sdk-server/__tests__/data_sources/StreamingProcessor.test.ts b/packages/shared/sdk-server/__tests__/data_sources/StreamingProcessor.test.ts index a91a90ffd2..d0864e9c97 100644 --- a/packages/shared/sdk-server/__tests__/data_sources/StreamingProcessor.test.ts +++ b/packages/shared/sdk-server/__tests__/data_sources/StreamingProcessor.test.ts @@ -1,11 +1,9 @@ import { - DataSourceErrorKind, defaultHeaders, EventName, Info, internal, LDLogger, - LDStreamingError, ProcessStreamResponse, subsystem, } from '@launchdarkly/js-sdk-common'; @@ -143,14 +141,17 @@ describe('given a stream processor with mock event source', () => { { errorFilter: expect.any(Function), headers: defaultHeaders(sdkKey, info, undefined), - initialRetryDelayMillis: 1000, readTimeoutMillis: 300000, - retryResetIntervalMillis: 60000, + retryDelayStrategy: { + nextRetryDelay: expect.any(Function), + setGoodSince: expect.any(Function), + setBaseDelay: expect.any(Function), + }, }, ); }); - it('sets streamInitialReconnectDelay correctly', () => { + it('applies streamInitialReconnectDelay to the retry backoff', () => { streamingProcessor = new StreamingProcessor( { basicConfiguration: getBasicConfiguration(logger), @@ -170,16 +171,12 @@ describe('given a stream processor with mock event source', () => { ); streamingProcessor.start(); - expect(basicPlatform.requests.createEventSource).toHaveBeenLastCalledWith( - `${serviceEndpoints.streaming}/all`, - { - errorFilter: expect.any(Function), - headers: defaultHeaders(sdkKey, info, undefined), - initialRetryDelayMillis: 22000, - readTimeoutMillis: 300000, - retryResetIntervalMillis: 60000, - }, - ); + // The configured 22s initial delay is no longer an init-dict field; it + // lives in the injected strategy. A first normal failure computes from it. + mockEventSource.options.errorFilter({ status: 500, message: 'err' }); + const delay = mockEventSource.options.retryDelayStrategy.nextRetryDelay(0); + expect(delay).toBeGreaterThan(11000); + expect(delay).toBeLessThanOrEqual(22000); }); it('adds listeners', () => { @@ -216,21 +213,109 @@ describe('given a stream processor with mock event source', () => { expect(mockListener.processJson).toHaveBeenNthCalledWith(1, expect.any(Object), headers); }); - it('passes error to callback if json data is malformed', async () => { + it('records a failure and reconnects when json data is malformed', () => { (mockListener.deserializeData as jest.Mock).mockReturnValue(false); + const createSpy = basicPlatform.requests.createEventSource as jest.Mock; + const callsBefore = createSpy.mock.calls.length; simulatePutEvent(); expect(logger.error).toHaveBeenCalledWith(expect.stringMatching(/invalid data in "put"/)); expect(logger.debug).toHaveBeenCalledWith(expect.stringMatching(/invalid json/i)); - expect(mockErrorHandler.mock.lastCall[0].message).toMatch(/malformed json/i); + // Recoverable now: surfaced as a log, not an 'error' event. + expect(mockErrorHandler).not.toHaveBeenCalled(); + // The connection is fine but the data is not, so the stream is torn down + // and re-established after the backoff wait. + expect(mockEventSource.close).toHaveBeenCalled(); + jest.advanceTimersByTime(1000); + expect(createSpy.mock.calls.length).toEqual(callsBefore + 1); + }); + + it('does not double-reconnect when two invalid payloads arrive in one parse pass', () => { + (mockListener.deserializeData as jest.Mock).mockReturnValue(false); + const createSpy = basicPlatform.requests.createEventSource as jest.Mock; + const callsBefore = createSpy.mock.calls.length; + + // Two malformed events delivered synchronously, before any reconnect timer + // fires. The second must be a no-op — the first already tore the stream + // down — so only one source is closed and only one reconnect is scheduled. + simulatePutEvent(); + simulatePutEvent(); + + expect(mockEventSource.close).toHaveBeenCalledTimes(1); + jest.advanceTimersByTime(1000); + expect(createSpy.mock.calls.length).toEqual(callsBefore + 1); + }); + + it('treats a deserialization that throws as invalid data and reconnects', () => { + (mockListener.deserializeData as jest.Mock).mockImplementation(() => { + throw new Error('structurally invalid'); + }); + const createSpy = basicPlatform.requests.createEventSource as jest.Mock; + const callsBefore = createSpy.mock.calls.length; + + // A throw during deserialization must not escape the listener; it is + // handled like any other invalid payload. + expect(() => simulatePutEvent()).not.toThrow(); + + expect(logger.error).toHaveBeenCalledWith(expect.stringMatching(/invalid data in "put"/)); + expect(mockErrorHandler).not.toHaveBeenCalled(); + expect(mockEventSource.close).toHaveBeenCalled(); + jest.advanceTimersByTime(1000); + expect(createSpy.mock.calls.length).toEqual(callsBefore + 1); + }); + + it('cancels a scheduled reconnect when stopped before it fires', () => { + (mockListener.deserializeData as jest.Mock).mockReturnValue(false); + const createSpy = basicPlatform.requests.createEventSource as jest.Mock; + const callsBefore = createSpy.mock.calls.length; + + // Malformed data arms a reconnect timer; stopping must cancel it. + simulatePutEvent(); + expect(mockEventSource.close).toHaveBeenCalled(); + + streamingProcessor.stop(); + jest.advanceTimersByTime(5 * 60 * 1000); + + expect(createSpy.mock.calls.length).toEqual(callsBefore); + }); + + it('classifies a status-less transport error as a normal, retryable failure', () => { + // No HTTP status → classifyTransportFailure() → normal: warn, keep retrying, + // and do not surface it through the error handler. + const willRetry = mockEventSource.options.errorFilter({ message: 'socket hang up' }); + + expect(willRetry).toBeTruthy(); + expect(mockErrorHandler).not.toHaveBeenCalled(); + expect(logger.warn).toHaveBeenCalledWith(expect.stringMatching(/will retry/)); + expect(logger.error).not.toHaveBeenCalled(); + }); + + it('counts a self-initiated restart after malformed data as a single failure', () => { + const recordFailure = jest.spyOn((streamingProcessor as any)['_retryState'], 'recordFailure'); + (mockListener.deserializeData as jest.Mock).mockReturnValue(false); + + simulatePutEvent(); + + // The malformed payload is one normal failure; the SDK's own close() + // during the restart must not be counted as a second. + expect(recordFailure).toHaveBeenCalledTimes(1); + expect(recordFailure).toHaveBeenCalledWith('normal'); }); - it('calls error handler if event.data prop is missing', async () => { + it('logs and restarts if event.data prop is missing', () => { simulatePutEvent({ flags: {} }); expect(mockListener.deserializeData).not.toHaveBeenCalled(); expect(mockListener.processJson).not.toHaveBeenCalled(); - expect(mockErrorHandler.mock.lastCall[0].message).toMatch(/unexpected payload/i); + expect(mockErrorHandler).not.toHaveBeenCalled(); + expect(logger.error).toHaveBeenCalledWith(expect.stringMatching(/unexpected payload/i)); + }); + + it('logs the reconnect delay via onretrying, unchanged from prior behavior', () => { + mockEventSource.onretrying({ delayMillis: 5000 }); + expect(logger.info).toHaveBeenCalledWith( + expect.stringMatching(/Will retry stream connection in 5000 milliseconds/), + ); }); it('closes and stops', async () => { @@ -275,18 +360,18 @@ describe('given a stream processor with mock event source', () => { }); }); - describe.each([401, 403])('given irrecoverable http errors', (status) => { - it(`stops retrying after error: ${status}`, () => { + describe.each([401, 403])('given unexpected http errors', (status) => { + it(`retries at error level rather than stopping after error: ${status}`, () => { const startTime = Date.now(); - const testError = { status, message: 'stopping. irrecoverable.' }; + const testError = { status, message: 'unexpected but still retried.' }; const willRetry = simulateError(testError); - expect(willRetry).toBeFalsy(); - expect(mockErrorHandler).toHaveBeenCalledWith( - new LDStreamingError(DataSourceErrorKind.Unknown, testError.message, testError.status), - ); + // RETRY conformance abolishes permanent stops: 401/403 now retry, and + // like any recoverable failure they surface as a log, not an 'error' event. + expect(willRetry).toBeTruthy(); + expect(mockErrorHandler).not.toHaveBeenCalled(); expect(logger.error).toHaveBeenCalledWith( - expect.stringMatching(new RegExp(`${status}.*permanently`)), + expect.stringMatching(new RegExp(`${status}.*will retry`)), ); const diagnosticEvent = diagnosticsManager.createStatsEventAndReset(0, 0, 0); @@ -296,5 +381,30 @@ describe('given a stream processor with mock event source', () => { expect(si.failed).toBeTruthy(); expect(si.durationMillis).toBeGreaterThanOrEqual(0); }); + + it(`enters the extended regime after error: ${status}`, () => { + simulateError({ status, message: 'unexpected' }); + const delay = mockEventSource.options.retryDelayStrategy.nextRetryDelay(0); + // Extended initial delay is 5 minutes, less up to half for jitter. + expect(delay).toBeGreaterThan(2.5 * 60 * 1000); + expect(delay).toBeLessThanOrEqual(5 * 60 * 1000); + }); + }); + + it('resets the reconnect delay to the operating cadence after a healthy event', () => { + const strategy = mockEventSource.options.retryDelayStrategy; + simulateError({ status: 500, message: 'transient' }); + expect(strategy.nextRetryDelay(0)).toBeGreaterThan(0); + strategy.setGoodSince(0); + expect(strategy.nextRetryDelay(0)).toEqual(0); + }); + + it('applies a server-directed retry time as the backoff base', () => { + const strategy = mockEventSource.options.retryDelayStrategy; + strategy.setBaseDelay(2500); + simulateError({ status: 500, message: 'transient' }); + const delay = strategy.nextRetryDelay(0); + expect(delay).toBeGreaterThan(1250); + expect(delay).toBeLessThanOrEqual(2500); }); }); diff --git a/packages/shared/sdk-server/src/data_sources/PollingProcessor.ts b/packages/shared/sdk-server/src/data_sources/PollingProcessor.ts index 2d0ff5230b..1613bd0d1e 100644 --- a/packages/shared/sdk-server/src/data_sources/PollingProcessor.ts +++ b/packages/shared/sdk-server/src/data_sources/PollingProcessor.ts @@ -1,10 +1,13 @@ import { - DataSourceErrorKind, + classifyHttpStatus, + classifyTransportFailure, + FailureKind, + forPolling, httpErrorMessage, internal, - isHttpRecoverable, LDLogger, LDPollingError, + RetryState, subsystem, VoidFunction, } from '@launchdarkly/js-sdk-common'; @@ -23,6 +26,7 @@ const { initMetadataFromHeaders } = internal; */ export default class PollingProcessor implements subsystem.LDStreamProcessor { private _stopped = false; + private readonly _retryState: RetryState; private _timeoutHandle: any; @@ -32,78 +36,90 @@ export default class PollingProcessor implements subsystem.LDStreamProcessor { private readonly _featureStore: LDDataSourceUpdates, private readonly _logger?: LDLogger, private readonly _initSuccessHandler: VoidFunction = () => {}, + // Reserved for a future terminal-failure channel; intentionally unused + // today private readonly _errorHandler?: PollingErrorHandler, - ) {} + ) { + this._retryState = forPolling(1000 * this._pollInterval); + } + + private _scheduleNextPoll() { + const delay = this._retryState.nextDelay; + this._logger?.debug('Scheduling next poll in %d ms', delay); + this._timeoutHandle = setTimeout(() => { + this._poll(); + }, delay); + } private _poll() { if (this._stopped) { return; } - const reportJsonError = (data: string) => { - this._logger?.error('Polling received invalid data'); - this._logger?.debug(`Invalid JSON follows: ${data}`); - this._errorHandler?.( - new LDPollingError( - DataSourceErrorKind.InvalidData, - 'Malformed JSON data in polling response', - ), - ); - }; - - const startTime = Date.now(); this._logger?.debug('Polling LaunchDarkly for feature flag updates'); this._requestor.requestAllData((err, body, headers) => { - const elapsed = Date.now() - startTime; - const sleepFor = Math.max(this._pollInterval * 1000 - elapsed, 0); + if (this._stopped) { + return; + } - this._logger?.debug('Elapsed: %d ms, sleeping for %d ms', elapsed, sleepFor); if (err) { const { status } = err; - if (status && !isHttpRecoverable(status)) { - const message = httpErrorMessage(err, 'polling request'); + const kind: FailureKind = + status !== undefined ? classifyHttpStatus(status) : classifyTransportFailure(); + this._retryState.recordFailure(kind); + const message = httpErrorMessage(err, 'polling request', 'will retry'); + // No failure is terminal now, so this is surfaced as a log only — the + // same outward treatment the SDK has always given a recoverable poll + // failure. The error-event channel stays reserved for a terminal case. + if (kind === 'unexpected') { this._logger?.error(message); - this._errorHandler?.( - new LDPollingError(DataSourceErrorKind.ErrorResponse, message, status), - ); - // It is not recoverable, return and do not trigger another - // poll. - return; + } else { + this._logger?.warn(message); + } + this._scheduleNextPoll(); + return; + } + + if (body) { + let parsed; + try { + parsed = deserializePoll(body); + } catch { + // Structurally invalid data can throw during deserialization; treat + // it the same as the unparseable payload handled below. + parsed = undefined; } - this._logger?.warn(httpErrorMessage(err, 'polling request', 'will retry')); - } else if (body) { - const parsed = deserializePoll(body); if (!parsed) { - // We could not parse this JSON. Report the problem and fallthrough to - // start another poll. - reportJsonError(body); - } else { - const initData = { - [VersionedDataKinds.Features.namespace]: parsed.flags, - [VersionedDataKinds.Segments.namespace]: parsed.segments, - }; - this._featureStore.init( - initData, - () => { - this._initSuccessHandler(); - // Triggering the next poll after the init has completed. - this._timeoutHandle = setTimeout(() => { - this._poll(); - }, sleepFor); - }, - initMetadataFromHeaders(headers), - ); - // The poll will be triggered by the feature store initialization - // completing. + // Unusable data is a normal failure: record it, report it, and poll + // again after the resulting wait. + this._retryState.recordFailure('normal'); + this._logger?.error('Polling received invalid data'); + this._logger?.debug(`Invalid JSON follows: ${body}`); + this._scheduleNextPoll(); return; } + + const initData = { + [VersionedDataKinds.Features.namespace]: parsed.flags, + [VersionedDataKinds.Segments.namespace]: parsed.segments, + }; + this._featureStore.init( + initData, + () => { + this._retryState.recordSuccess(); + this._initSuccessHandler(); + // Triggering the next poll after the init has completed. + this._scheduleNextPoll(); + }, + initMetadataFromHeaders(headers), + ); + return; } - // Falling through, there was some type of error and we need to trigger - // a new poll. - this._timeoutHandle = setTimeout(() => { - this._poll(); - }, sleepFor); + // No error and no body (for example, a not-modified response): a + // successful poll that delivered nothing new. + this._retryState.recordSuccess(); + this._scheduleNextPoll(); }); } diff --git a/packages/shared/sdk-server/src/data_sources/StreamingProcessor.ts b/packages/shared/sdk-server/src/data_sources/StreamingProcessor.ts index 41cd409174..82a3f44c5d 100644 --- a/packages/shared/sdk-server/src/data_sources/StreamingProcessor.ts +++ b/packages/shared/sdk-server/src/data_sources/StreamingProcessor.ts @@ -1,44 +1,37 @@ import { ClientContext, - DataSourceErrorKind, + classifyHttpStatus, + classifyTransportFailure, EventName, EventSource, + EventSourceRetryDelayStrategy, + FailureKind, + forStreaming, getStreamingUri, httpErrorMessage, HttpErrorResponse, internal, LDHeaders, LDLogger, - LDStreamingError, ProcessStreamResponse, Requests, - shouldRetry, + RetryState, StreamingErrorHandler, subsystem, } from '@launchdarkly/js-sdk-common'; -const reportJsonError = ( - type: string, - data: string, - logger?: LDLogger, - errorHandler?: StreamingErrorHandler, -) => { - logger?.error(`Stream received invalid data in "${type}" message`); - logger?.debug(`Invalid JSON follows: ${data}`); - errorHandler?.( - new LDStreamingError(DataSourceErrorKind.InvalidData, 'Malformed JSON data in event stream'), - ); -}; - export default class StreamingProcessor implements subsystem.LDStreamProcessor { private readonly _headers: { [key: string]: string | string[] }; private readonly _streamUri: string; private readonly _logger?: LDLogger; + private readonly _retryState: RetryState; private _eventSource?: EventSource; private _requests: Requests; private _connectionAttemptStartTime?: number; private _initHeaders?: { [key: string]: string }; + private _reconnectTimeout?: ReturnType; + private _stopped = false; constructor( clientContext: ClientContext, @@ -47,6 +40,8 @@ export default class StreamingProcessor implements subsystem.LDStreamProcessor { private readonly _listeners: Map, baseHeaders: LDHeaders, private readonly _diagnosticsManager?: internal.DiagnosticsManager, + // Reserved for a future terminal-failure channel; intentionally unused + // today private readonly _errorHandler?: StreamingErrorHandler, private readonly _streamInitialReconnectDelay = 1, ) { @@ -62,6 +57,7 @@ export default class StreamingProcessor implements subsystem.LDStreamProcessor { streamUriPath, parameters, ); + this._retryState = forStreaming(1000 * this._streamInitialReconnectDelay); } private _logConnectionStarted() { @@ -81,40 +77,73 @@ export default class StreamingProcessor implements subsystem.LDStreamProcessor { } /** - * This is a wrapper around the passed errorHandler which adds additional - * diagnostics and logging logic. + * Records a connection failure and logs it. The retry state decides the + * wait; a server-directed `retry:` value, if any, has already been applied + * to it through the injected strategy. * - * @param err The error to be logged and handled. - * @return boolean whether to retry the connection. + * @param err The error to be recorded and logged. * * @private */ - private _retryAndHandleError(err: HttpErrorResponse) { - if (!shouldRetry(err)) { - this._logConnectionResult(false); - this._errorHandler?.( - new LDStreamingError(DataSourceErrorKind.ErrorResponse, err.message, err.status), - ); - this._logger?.error(httpErrorMessage(err, 'streaming request')); - return false; + private _retryAndHandleError(err: HttpErrorResponse): void { + const kind: FailureKind = + err.status !== undefined ? classifyHttpStatus(err.status) : classifyTransportFailure(); + this._retryState.recordFailure(kind); + this._logConnectionResult(false); + + const message = httpErrorMessage(err, 'streaming request', 'will retry'); + // No failure is terminal now, so this is surfaced as a log only — the same + // outward treatment the SDK has always given a recoverable/interrupted + // failure. The error-event channel stays reserved for a terminal condition. + if (kind === 'unexpected') { + this._logger?.error(message); + } else { + this._logger?.warn(message); } - this._logger?.warn(httpErrorMessage(err, 'streaming request', 'will retry')); - this._logConnectionResult(false); this._logConnectionStarted(); - return true; + } + + private _restartForInvalidData() { + // A second invalid payload can arrive in the same parse pass — the event + // source emits every complete event in a chunk synchronously, even after + // close(). The first call tears the stream down, so any re-entry finds no + // live source and must return without scheduling again; otherwise it would + // orphan the first timer, leak the first EventSource, and duplicate the + // stream once both reconnects fired. + if (this._stopped || !this._eventSource) { + return; + } + + this._retryState.recordFailure('normal'); + + this._eventSource.close(); + this._eventSource = undefined; + this._reconnectTimeout = setTimeout(() => { + if (!this._stopped) { + this.start(); + } + }, this._retryState.nextDelay); } start() { this._logConnectionStarted(); + const retryDelayStrategy: EventSourceRetryDelayStrategy = { + nextRetryDelay: () => this._retryState.nextDelay, + setGoodSince: () => this._retryState.recordSuccess(), + setBaseDelay: (baseDelayMs: number) => this._retryState.applyServerDirectedRetry(baseDelayMs), + }; + // TLS is handled by the platform implementation. const eventSource = this._requests.createEventSource(this._streamUri, { headers: this._headers, - errorFilter: (error: HttpErrorResponse) => this._retryAndHandleError(error), - initialRetryDelayMillis: 1000 * this._streamInitialReconnectDelay, + errorFilter: (error: HttpErrorResponse) => { + this._retryAndHandleError(error); + return true; + }, readTimeoutMillis: 5 * 60 * 1000, - retryResetIntervalMillis: 60 * 1000, + retryDelayStrategy, }); this._eventSource = eventSource; @@ -142,26 +171,36 @@ export default class StreamingProcessor implements subsystem.LDStreamProcessor { if (event?.data) { this._logConnectionResult(true); const { data } = event; - const dataJson = deserializeData(data); + let dataJson; + try { + dataJson = deserializeData(data); + } catch { + // Structurally invalid data can throw during deserialization; treat + // it the same as the unparseable payload handled below. + dataJson = undefined; + } if (!dataJson) { - reportJsonError(eventName, data, this._logger, this._errorHandler); + this._logger?.error(`Stream received invalid data in "${eventName}" message`); + this._logger?.debug(`Invalid JSON follows: ${data}`); + this._restartForInvalidData(); return; } processJson(dataJson, this._initHeaders); } else { - this._errorHandler?.( - new LDStreamingError( - DataSourceErrorKind.Unknown, - 'Unexpected payload from event stream', - ), - ); + this._logger?.error('Unexpected payload from event stream'); + this._restartForInvalidData(); } }); }); } stop() { + if (this._reconnectTimeout) { + clearTimeout(this._reconnectTimeout); + this._reconnectTimeout = undefined; + } + this._stopped = true; this._eventSource?.close(); this._eventSource = undefined; }