Skip to content

Commit 46de8d4

Browse files
committed
fix: fallback variables need to be read on stream errors
1 parent 1400310 commit 46de8d4

2 files changed

Lines changed: 288 additions & 13 deletions

File tree

packages/shared/sdk-client/__tests__/datasource/fdv2/StreamingFDv2Base.test.ts

Lines changed: 223 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -595,7 +595,7 @@ it('emits terminal_error for a goodbye when a deferred open directive is pending
595595
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '45' },
596596
});
597597

598-
// Goodbye fires before any payload the pending directive triggers the fallback
598+
// Goodbye fires before any payload; the pending directive triggers the fallback
599599
simulateEvent(mockEventSource, 'goodbye', { reason: 'bye' });
600600

601601
const result = await base.takeResult();
@@ -614,13 +614,13 @@ it('clears a pending fallback directive when onopen fires without the fallback h
614614
const base = createBase(mockRequests, logger);
615615
base.start();
616616

617-
// First open carries fallback headers — arms the deferred directive
617+
// First open carries fallback headers, arming the deferred directive
618618
mockEventSource.onopen({
619619
type: 'open',
620620
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '60' },
621621
});
622622

623-
// Reconnect without fallback header must clear the stale directive
623+
// Reconnect without fallback header; must clear the stale directive
624624
mockEventSource.onopen({ type: 'open', headers: {} });
625625

626626
// A payload arriving after the clean reconnect must NOT carry fdv1Fallback=true
@@ -649,11 +649,230 @@ it('close resets the pending fallback directive', async () => {
649649
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '30' },
650650
});
651651

652-
// Close before any payload arrives the shutdown result must NOT carry TTL
652+
// Close before any payload arrives; the shutdown result must NOT carry TTL
653653
base.close();
654654
const result = await base.takeResult();
655655
expect(result.type).toBe('status');
656656
if (result.type !== 'status') return;
657657
expect(result.state).toBe('shutdown');
658658
expect(result.fdv1FallbackTtlMs).toBeUndefined();
659659
});
660+
661+
it('surfaces a deferred fallback directive on the serverError path', async () => {
662+
const mockEventSource = createMockEventSource();
663+
const mockRequests = createMockRequests(mockEventSource);
664+
const base = createBase(mockRequests, logger);
665+
base.start();
666+
667+
// Arm the deferred directive at open.
668+
mockEventSource.onopen({
669+
type: 'open',
670+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
671+
});
672+
673+
// A server-side FDv2 'error' event arrives before any payload/goodbye.
674+
simulateEvent(mockEventSource, 'error', { reason: 'server error', payload_id: 'p1' });
675+
676+
const result = await base.takeResult();
677+
expect(result.type).toBe('status');
678+
if (result.type !== 'status') return;
679+
expect(result.state).toBe('interrupted');
680+
expect(result.fdv1Fallback).toBe(true);
681+
expect(result.fdv1FallbackTtlMs).toBe(75000);
682+
683+
base.close();
684+
});
685+
686+
it('surfaces a deferred fallback directive on the protocol error path', async () => {
687+
const mockEventSource = createMockEventSource();
688+
const mockRequests = createMockRequests(mockEventSource);
689+
const base = createBase(mockRequests, logger);
690+
base.start();
691+
692+
mockEventSource.onopen({
693+
type: 'open',
694+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
695+
});
696+
697+
// A server-intent with no payloads yields a MISSING_PAYLOAD protocol error
698+
// before any payload/goodbye.
699+
simulateEvent(mockEventSource, 'server-intent', { payloads: [] });
700+
701+
const result = await base.takeResult();
702+
expect(result.type).toBe('status');
703+
if (result.type !== 'status') return;
704+
expect(result.state).toBe('interrupted');
705+
expect(result.fdv1Fallback).toBe(true);
706+
expect(result.fdv1FallbackTtlMs).toBe(75000);
707+
708+
base.close();
709+
});
710+
711+
it('surfaces a deferred fallback directive on the malformed-JSON path', async () => {
712+
const mockEventSource = createMockEventSource();
713+
const mockRequests = createMockRequests(mockEventSource);
714+
const base = createBase(mockRequests, logger);
715+
base.start();
716+
717+
mockEventSource.onopen({
718+
type: 'open',
719+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
720+
});
721+
722+
// Invoke a registered listener directly with unparseable data.
723+
const { calls } = mockEventSource.addEventListener.mock;
724+
const listener = calls.find((c: any[]) => c[0] === 'server-intent')?.[1];
725+
listener({ data: 'not-valid-json{{{' });
726+
727+
const result = await base.takeResult();
728+
expect(result.type).toBe('status');
729+
if (result.type !== 'status') return;
730+
expect(result.state).toBe('interrupted');
731+
expect(result.errorInfo?.message).toContain('Malformed JSON');
732+
expect(result.fdv1Fallback).toBe(true);
733+
expect(result.fdv1FallbackTtlMs).toBe(75000);
734+
735+
base.close();
736+
});
737+
738+
it('surfaces a deferred fallback directive on the network-error path', async () => {
739+
const mockEventSource = createMockEventSource();
740+
const mockRequests = createMockRequests(mockEventSource);
741+
const base = createBase(mockRequests, logger);
742+
base.start();
743+
744+
mockEventSource.onopen({
745+
type: 'open',
746+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
747+
});
748+
749+
// Network error with no numeric status routes through es.onerror
750+
// (a numeric status would be handled by the error filter instead).
751+
mockEventSource.onerror({ message: 'IO Error' });
752+
753+
const result = await base.takeResult();
754+
expect(result.type).toBe('status');
755+
if (result.type !== 'status') return;
756+
expect(result.state).toBe('interrupted');
757+
expect(result.fdv1Fallback).toBe(true);
758+
expect(result.fdv1FallbackTtlMs).toBe(75000);
759+
760+
base.close();
761+
});
762+
763+
it('merges a deferred fallback directive into a successful ping-triggered poll result', async () => {
764+
const mockEventSource = createMockEventSource();
765+
const mockRequests = createMockRequests(mockEventSource);
766+
const pingHandler: PingHandler = {
767+
handlePing: jest.fn().mockResolvedValue({
768+
type: 'changeSet',
769+
payload: { events: [], selector: undefined },
770+
fdv1Fallback: false,
771+
}),
772+
};
773+
const base = createBase(mockRequests, logger, { pingHandler });
774+
base.start();
775+
776+
mockEventSource.onopen({
777+
type: 'open',
778+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
779+
});
780+
781+
const { calls } = mockEventSource.addEventListener.mock;
782+
const pingListener = calls.find((c: any[]) => c[0] === 'ping')?.[1];
783+
await pingListener();
784+
785+
const result = await base.takeResult();
786+
expect(result.fdv1Fallback).toBe(true);
787+
expect(result.fdv1FallbackTtlMs).toBe(75000);
788+
789+
base.close();
790+
});
791+
792+
it('surfaces a deferred fallback directive when a ping-triggered poll throws', async () => {
793+
const mockEventSource = createMockEventSource();
794+
const mockRequests = createMockRequests(mockEventSource);
795+
const pingHandler: PingHandler = {
796+
handlePing: jest.fn().mockRejectedValue(new Error('poll failed')),
797+
};
798+
const base = createBase(mockRequests, logger, { pingHandler });
799+
base.start();
800+
801+
mockEventSource.onopen({
802+
type: 'open',
803+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
804+
});
805+
806+
const { calls } = mockEventSource.addEventListener.mock;
807+
const pingListener = calls.find((c: any[]) => c[0] === 'ping')?.[1];
808+
await pingListener();
809+
810+
const result = await base.takeResult();
811+
expect(result.type).toBe('status');
812+
if (result.type !== 'status') return;
813+
expect(result.state).toBe('interrupted');
814+
expect(result.fdv1Fallback).toBe(true);
815+
expect(result.fdv1FallbackTtlMs).toBe(75000);
816+
817+
base.close();
818+
});
819+
820+
it('surfaces a deferred fallback directive on a non-retryable errorFilter error', async () => {
821+
const mockEventSource = createMockEventSource();
822+
const mockRequests = createMockRequests(mockEventSource);
823+
const base = createBase(mockRequests, logger);
824+
base.start();
825+
826+
// Arm a directive at onopen (a prior successful connection observed the header).
827+
mockEventSource.onopen({
828+
type: 'open',
829+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
830+
});
831+
832+
// A later reconnect attempt fails with a non-retryable status that does NOT
833+
// carry its own fallback header.
834+
const willRetry = simulateErrorFilter(mockRequests, {
835+
status: 401,
836+
message: 'Error 401',
837+
});
838+
expect(willRetry).toBe(false);
839+
840+
const result = await base.takeResult();
841+
expect(result.type).toBe('status');
842+
if (result.type !== 'status') return;
843+
expect(result.state).toBe('terminal_error');
844+
expect(result.fdv1Fallback).toBe(true);
845+
expect(result.fdv1FallbackTtlMs).toBe(75000);
846+
847+
base.close();
848+
});
849+
850+
it('surfaces a deferred fallback directive on a retryable errorFilter error', async () => {
851+
const mockEventSource = createMockEventSource();
852+
const mockRequests = createMockRequests(mockEventSource);
853+
const base = createBase(mockRequests, logger);
854+
base.start();
855+
856+
// Arm a directive at onopen (a prior successful connection observed the header).
857+
mockEventSource.onopen({
858+
type: 'open',
859+
headers: { 'x-ld-fd-fallback': 'true', 'x-ld-fd-fallback-ttl': '75' },
860+
});
861+
862+
// A later reconnect attempt fails with a retryable status that does NOT
863+
// carry its own fallback header.
864+
const willRetry = simulateErrorFilter(mockRequests, {
865+
status: 500,
866+
message: 'Error 500',
867+
});
868+
expect(willRetry).toBe(true);
869+
870+
const result = await base.takeResult();
871+
expect(result.type).toBe('status');
872+
if (result.type !== 'status') return;
873+
expect(result.state).toBe('interrupted');
874+
expect(result.fdv1Fallback).toBe(true);
875+
expect(result.fdv1FallbackTtlMs).toBe(75000);
876+
877+
base.close();
878+
});

packages/shared/sdk-client/src/datasource/fdv2/StreamingFDv2Base.ts

Lines changed: 65 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,28 @@ export function createStreamingBase(config: {
133133
connectionAttemptStartTime = undefined;
134134
}
135135

136+
/**
137+
* Promotes a directive deferred from `onopen` (`pendingFallback`) into the
138+
* committed `fdv1Fallback`/`fdv1FallbackTtlMs` state, clears the pending
139+
* pair, and returns the current committed fallback state.
140+
*
141+
* Every path that can produce a result, a stream error, a network failure,
142+
* or a ping-triggered poll (including one that succeeds without its own
143+
* fallback signal), calls this instead of reading the closure variables
144+
* directly, so a directive deferred at onopen surfaces no matter which
145+
* path fires next. Safe to call when nothing is pending; it just returns
146+
* the current state unchanged.
147+
*/
148+
function resolveFallback(): { fdv1Fallback: boolean; fdv1FallbackTtlMs: number | undefined } {
149+
if (pendingFallback) {
150+
fdv1Fallback = true;
151+
fdv1FallbackTtlMs = pendingFallbackTtlMs;
152+
pendingFallback = false;
153+
pendingFallbackTtlMs = undefined;
154+
}
155+
return { fdv1Fallback, fdv1FallbackTtlMs };
156+
}
157+
136158
function handleAction(action: internal.ProtocolAction, rawData?: unknown): void {
137159
switch (action.type) {
138160
case 'payload':
@@ -169,15 +191,22 @@ export function createStreamingBase(config: {
169191
break;
170192
}
171193

172-
case 'serverError':
173-
resultQueue.put(interrupted(errorInfoFromUnknown(action.reason), fdv1Fallback, fdv1FallbackTtlMs));
194+
case 'serverError': {
195+
const fallback = resolveFallback();
196+
resultQueue.put(
197+
interrupted(errorInfoFromUnknown(action.reason), fallback.fdv1Fallback, fallback.fdv1FallbackTtlMs),
198+
);
174199
break;
200+
}
175201

176202
case 'error':
177203
// Only actionable errors are queued; informational ones (UNKNOWN_EVENT)
178204
// are logged by the protocol handler.
179205
if (action.kind === 'MISSING_PAYLOAD' || action.kind === 'PROTOCOL_ERROR') {
180-
resultQueue.put(interrupted(errorInfoFromInvalidData(action.message), fdv1Fallback, fdv1FallbackTtlMs));
206+
const fallback = resolveFallback();
207+
resultQueue.put(
208+
interrupted(errorInfoFromInvalidData(action.message), fallback.fdv1Fallback, fallback.fdv1FallbackTtlMs),
209+
);
181210
}
182211
break;
183212

@@ -203,17 +232,23 @@ export function createStreamingBase(config: {
203232
return false;
204233
}
205234

235+
const fallback = resolveFallback();
236+
206237
if (!shouldRetry(err)) {
207238
config.logger?.error(httpErrorMessage(err, 'streaming request'));
208239
logConnectionResult(false);
209-
resultQueue.put(terminalError(errorInfoFromHttpError(err.status ?? 0), fdv1Fallback, fdv1FallbackTtlMs));
240+
resultQueue.put(
241+
terminalError(errorInfoFromHttpError(err.status ?? 0), fallback.fdv1Fallback, fallback.fdv1FallbackTtlMs),
242+
);
210243
return false;
211244
}
212245

213246
config.logger?.warn(httpErrorMessage(err, 'streaming request', 'will retry'));
214247
logConnectionResult(false);
215248
logConnectionAttempt();
216-
resultQueue.put(interrupted(errorInfoFromHttpError(err.status ?? 0), fdv1Fallback, fdv1FallbackTtlMs));
249+
resultQueue.put(
250+
interrupted(errorInfoFromHttpError(err.status ?? 0), fallback.fdv1Fallback, fallback.fdv1FallbackTtlMs),
251+
);
217252
return true;
218253
}
219254

@@ -243,8 +278,13 @@ export function createStreamingBase(config: {
243278
`Stream received data that was unable to be parsed in "${eventName}" message`,
244279
);
245280
config.logger?.debug(`Data follows: ${event.data}`);
281+
const fallback = resolveFallback();
246282
resultQueue.put(
247-
interrupted(errorInfoFromInvalidData('Malformed JSON in EventStream'), fdv1Fallback, fdv1FallbackTtlMs),
283+
interrupted(
284+
errorInfoFromInvalidData('Malformed JSON in EventStream'),
285+
fallback.fdv1Fallback,
286+
fallback.fdv1FallbackTtlMs,
287+
),
248288
);
249289
return;
250290
}
@@ -276,18 +316,29 @@ export function createStreamingBase(config: {
276316
return;
277317
}
278318

319+
// Trust the poll's own fallback signal (e.g. from its own HTTP
320+
// response headers) over a directive deferred at onopen. Only
321+
// backfill from the deferred directive when the poll result doesn't
322+
// already indicate fallback, so an onopen directive still surfaces
323+
// even while ping-triggered polls keep succeeding.
324+
const fallback = resolveFallback();
325+
if (!result.fdv1Fallback && fallback.fdv1Fallback) {
326+
result.fdv1Fallback = true;
327+
result.fdv1FallbackTtlMs = fallback.fdv1FallbackTtlMs;
328+
}
279329
resultQueue.put(result);
280330
} catch (err: any) {
281331
if (stopped) {
282332
return;
283333
}
284334

285335
config.logger?.error(`Error handling ping: ${err?.message ?? err}`);
336+
const fallback = resolveFallback();
286337
resultQueue.put(
287338
interrupted(
288339
errorInfoFromNetworkError(err?.message ?? 'Error during ping poll'),
289-
fdv1Fallback,
290-
fdv1FallbackTtlMs,
340+
fallback.fdv1Fallback,
341+
fallback.fdv1FallbackTtlMs,
291342
),
292343
);
293344
}
@@ -329,8 +380,13 @@ export function createStreamingBase(config: {
329380
// This condition will be handled by the error filter.
330381
return;
331382
}
383+
const fallback = resolveFallback();
332384
resultQueue.put(
333-
interrupted(errorInfoFromNetworkError(err?.message ?? 'IO Error'), fdv1Fallback, fdv1FallbackTtlMs),
385+
interrupted(
386+
errorInfoFromNetworkError(err?.message ?? 'IO Error'),
387+
fallback.fdv1Fallback,
388+
fallback.fdv1FallbackTtlMs,
389+
),
334390
);
335391
};
336392

0 commit comments

Comments
 (0)