From de033b7dbd96e37466250c3a7bb4646b09583953 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Sat, 1 Aug 2026 07:22:05 -0400 Subject: [PATCH] =?UTF-8?q?=F0=9F=A6=97=20fix:=20Ignore=20Sequenced=20Redi?= =?UTF-8?q?s=20Events=20Without=20SSE=20Subscribers=20(#14557)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../__tests__/RedisEventTransport.spec.ts | 104 ++++++++++++++++++ .../implementations/RedisEventTransport.ts | 7 ++ 2 files changed, 111 insertions(+) diff --git a/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts b/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts index 6b9ae43f7f..59d7baaaaf 100644 --- a/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts +++ b/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts @@ -157,6 +157,110 @@ describe('RedisEventTransport', () => { transport.destroy(); }); + it('ignores sequenced events while only internal abort listeners are attached', async () => { + jest.useFakeTimers(); + const mockPublisher = createMockPublisher(); + const mockSubscriber = createMockSubscriber(); + const transport = new RedisEventTransport( + mockPublisher as unknown as Redis, + mockSubscriber as unknown as Redis, + ); + const streamId = 'abort-only-sequenced-events'; + const onAbort = jest.fn(); + const onPreempt = jest.fn(); + const warn = jest.spyOn(logger, 'warn').mockImplementation(() => logger); + + try { + await transport.onAbort(streamId, onAbort); + await transport.onPreempt(streamId, onPreempt); + const messageHandler = getMessageHandler(mockSubscriber); + + deliverSequencedMessage(messageHandler, streamId, { + type: 'chunk', + seq: 16_119, + data: { index: 0 }, + }); + deliverSequencedMessage(messageHandler, streamId, { + type: 'done', + seq: 16_120, + data: { final: true }, + }); + deliverSequencedMessage(messageHandler, streamId, { + type: 'error', + seq: 16_121, + error: 'remote error', + }); + await jest.advanceTimersByTimeAsync(1_000); + + messageHandler( + `stream:{${streamId}}:events`, + JSON.stringify({ type: 'abort', generationId: 9 }), + ); + const preempt = { op: 'arm', createdAt: 10, steerIds: ['steer-1'] }; + messageHandler(`stream:{${streamId}}:events`, JSON.stringify({ type: 'preempt', preempt })); + + expect(onAbort).toHaveBeenCalledWith(9); + expect(onPreempt).toHaveBeenCalledWith(preempt); + expect(warn).not.toHaveBeenCalledWith(expect.stringContaining(`Stream ${streamId}:`)); + } finally { + warn.mockRestore(); + transport.destroy(); + jest.useRealTimers(); + } + }); + + it('synchronizes a resumed subscriber after ignoring detached stream traffic', async () => { + const mockPublisher = createMockPublisher(); + const mockSubscriber = createMockSubscriber(); + const transport = new RedisEventTransport( + mockPublisher as unknown as Redis, + mockSubscriber as unknown as Redis, + ); + const streamId = 'abort-only-then-sse'; + const messageHandler = getMessageHandler(mockSubscriber); + + await transport.onAbort(streamId, () => undefined); + const initial: object[] = []; + const initialSubscription = transport.subscribe(streamId, { + onChunk: (event) => initial.push(event as object), + }); + await initialSubscription.ready; + deliverSequencedMessage(messageHandler, streamId, { + type: 'chunk', + seq: 0, + data: { index: 0 }, + }); + initialSubscription.unsubscribe(); + + deliverSequencedMessage(messageHandler, streamId, { + type: 'chunk', + seq: 41, + data: { ignored: true }, + }); + + mockPublisher.get.mockResolvedValueOnce('42'); + const received: object[] = []; + const subscription = transport.subscribe( + streamId, + { onChunk: (event) => received.push(event as object) }, + { deferSequenceDelivery: true }, + ); + await subscription.ready; + await transport.syncReorderBuffer(streamId); + + deliverSequencedMessage(messageHandler, streamId, { + type: 'chunk', + seq: 42, + data: { index: 42 }, + }); + + expect(initial).toEqual([{ index: 0 }]); + expect(received).toEqual([{ index: 42 }]); + + subscription.unsubscribe(); + transport.destroy(); + }); + it('releases each generation abort subscription after successful completion', async () => { const mockPublisher = createMockPublisher(); const mockSubscriber = createMockSubscriber(); diff --git a/packages/api/src/stream/implementations/RedisEventTransport.ts b/packages/api/src/stream/implementations/RedisEventTransport.ts index 8b6bc005e3..2e17a1eeb1 100644 --- a/packages/api/src/stream/implementations/RedisEventTransport.ts +++ b/packages/api/src/stream/implementations/RedisEventTransport.ts @@ -505,6 +505,13 @@ export class RedisEventTransport implements IEventTransport { try { const parsed = JSON.parse(message) as PubSubMessage; + if ( + streamState.count === 0 && + parsed.type !== EventTypes.ABORT && + parsed.type !== EventTypes.PREEMPT + ) { + return; + } if (parsed.type === EventTypes.CHUNK && parsed.seq != null) { this.handleOrderedChunk(streamId, streamState, parsed); } else if (