From 5bb98aaee79d7dd95d57a4092567dfd4805e380d Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 15:58:12 -0700 Subject: [PATCH 1/4] fix(events): start every SSE response once its subscriptions are live Next sends a route's status and headers with the first body chunk, and createSSEStream wrote nothing until an event or the 30 s heartbeat. A client therefore saw no response for up to 30 s after connecting: the desktop inbox doorbell E2E waited 32-34 s for its stream to open in every run, and Sim desktop's doorbell gives up on a response that has not started within 15 s. - createSSEStream writes a `: connected` comment as soon as every subscription is live, which starts the response at once. - PubSubChannel exposes ready(), settled when the process's Redis SUBSCRIBE completes, so a client that reads its state on open cannot miss an event published before the channel was listening. The desktop doorbell, chat status and MCP streams wait on their channels. - The desktop inbox E2E asserts the doorbell opens within the desktop's 15 s handshake and that a ring arrives within 2 s of the open. --- .../sim/app/api/desktop/inbox/stream/route.ts | 2 +- apps/sim/app/api/mcp/events/route.ts | 2 + .../app/api/mothership/events/route.test.ts | 6 +- apps/sim/app/api/mothership/events/route.ts | 2 + apps/sim/lib/desktop/application/executor.ts | 2 + apps/sim/lib/desktop/executor/doorbell.ts | 5 ++ apps/sim/lib/events/pubsub.ts | 31 ++++++-- apps/sim/lib/events/sse-endpoint.test.ts | 71 +++++++++++++++++++ apps/sim/lib/events/sse-endpoint.ts | 15 ++++ apps/sim/lib/mcp/pubsub.ts | 5 ++ apps/sim/lib/mothership/chat-status.ts | 1 + apps/sim/scripts/test-desktop-inbox-e2e.ts | 29 +++++++- packages/testing/src/mocks/mcp-pubsub.mock.ts | 4 +- .../src/mocks/mothership-chat-status.mock.ts | 4 +- 14 files changed, 166 insertions(+), 13 deletions(-) diff --git a/apps/sim/app/api/desktop/inbox/stream/route.ts b/apps/sim/app/api/desktop/inbox/stream/route.ts index 4f6ae8911de..e2ddf11b2b9 100644 --- a/apps/sim/app/api/desktop/inbox/stream/route.ts +++ b/apps/sim/app/api/desktop/inbox/stream/route.ts @@ -32,7 +32,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => { return createSSEStream(request, { label: 'desktop-inbox', revalidate: inbox.revalidate, - subscriptions: [{ subscribe: inbox.subscribe }], + subscriptions: [{ subscribe: inbox.subscribe, ready: inbox.ready }], }) } catch (error) { if (error instanceof InternalUnauthenticatedError) diff --git a/apps/sim/app/api/mcp/events/route.ts b/apps/sim/app/api/mcp/events/route.ts index 73a1aaa55c5..b294e878957 100644 --- a/apps/sim/app/api/mcp/events/route.ts +++ b/apps/sim/app/api/mcp/events/route.ts @@ -32,6 +32,7 @@ const mcpEventsHandler = createWorkspaceSSE({ }) }) }, + ready: async () => mcpPubSub?.ready(), }, { subscribe: (workspaceId, send) => { @@ -45,6 +46,7 @@ const mcpEventsHandler = createWorkspaceSSE({ }) }) }, + ready: async () => mcpPubSub?.ready(), }, ], }) diff --git a/apps/sim/app/api/mothership/events/route.test.ts b/apps/sim/app/api/mothership/events/route.test.ts index cd4eaaeda5c..39ed74df41f 100644 --- a/apps/sim/app/api/mothership/events/route.test.ts +++ b/apps/sim/app/api/mothership/events/route.test.ts @@ -16,7 +16,7 @@ import { import { NextRequest } from 'next/server' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { OrchestrationError } from '@/lib/core/orchestration/types' -import { HEARTBEAT_INTERVAL_MS } from '@/lib/events/sse-endpoint' +import { HEARTBEAT_INTERVAL_MS, OPENED_COMMENT } from '@/lib/events/sse-endpoint' import type { ChatStatusEvent } from '@/lib/mothership/chat-status' import { PermissionGroupCapabilityError } from '@/lib/permission-groups/capability-error' @@ -41,13 +41,15 @@ function emit(event: ChatStatusEvent) { handler(event) } +/** Every chunk after the stream's opening comment. */ async function collect(body: ReadableStream, chunks: string[]) { const reader = body.getReader() const decoder = new TextDecoder() while (true) { const { done, value } = await reader.read() if (done) return - chunks.push(decoder.decode(value)) + const chunk = decoder.decode(value) + if (chunk !== OPENED_COMMENT) chunks.push(chunk) } } diff --git a/apps/sim/app/api/mothership/events/route.ts b/apps/sim/app/api/mothership/events/route.ts index 80e8850d7fe..4b0e5695c3e 100644 --- a/apps/sim/app/api/mothership/events/route.ts +++ b/apps/sim/app/api/mothership/events/route.ts @@ -42,6 +42,7 @@ const mothershipEventsHandler = createWorkspaceSSE({ }) }) }, + ready: async () => chatPubSub?.ready(), }, ], }) @@ -80,6 +81,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => { timestamp: Date.now(), }) }) ?? (() => {}), + ready: async () => chatPubSub?.ready(), }, ], }) diff --git a/apps/sim/lib/desktop/application/executor.ts b/apps/sim/lib/desktop/application/executor.ts index 11b76c6cba0..0dd8fc89b39 100644 --- a/apps/sim/lib/desktop/application/executor.ts +++ b/apps/sim/lib/desktop/application/executor.ts @@ -15,6 +15,7 @@ import { } from '@/lib/desktop/executor/constants' import { type DesktopInboxChangeReason, + desktopInboxDoorbellReady, onDesktopInboxDoorbell, } from '@/lib/desktop/executor/doorbell' import { @@ -201,6 +202,7 @@ export const openDesktopInboxStream = defineAuthorizedCredentialUserUseCase({ onDesktopInboxDoorbell(deviceId, (reason: DesktopInboxChangeReason) => send('inbox_changed', { reason }) ), + ready: desktopInboxDoorbellReady, } }, }) diff --git a/apps/sim/lib/desktop/executor/doorbell.ts b/apps/sim/lib/desktop/executor/doorbell.ts index 4fd07e241ca..8f8caf060c6 100644 --- a/apps/sim/lib/desktop/executor/doorbell.ts +++ b/apps/sim/lib/desktop/executor/doorbell.ts @@ -43,6 +43,11 @@ export function ringDesktopInbox(deviceId: string, reason: DesktopInboxChangeRea } } +/** Settles once this process hears every ring. */ +export function desktopInboxDoorbellReady(): Promise { + return channel().ready() +} + /** Subscribes to one device's doorbell; returns the unsubscribe. */ export function onDesktopInboxDoorbell( deviceId: string, diff --git a/apps/sim/lib/events/pubsub.ts b/apps/sim/lib/events/pubsub.ts index dd372ab488c..853e68bc0d1 100644 --- a/apps/sim/lib/events/pubsub.ts +++ b/apps/sim/lib/events/pubsub.ts @@ -16,6 +16,11 @@ const logger = createLogger('PubSub') export interface PubSubChannel { publish(event: T): void subscribe(handler: (event: T) => void): () => void + /** + * Settles once this process receives the channel's publications; anything published before + * then reaches no subscriber here. + */ + ready(): Promise dispose(): void } @@ -29,6 +34,7 @@ class RedisPubSubChannel implements PubSubChannel { private sub: Redis private handlers = new Set<(event: T) => void>() private disposed = false + private readonly subscribed: Promise constructor( redisUrl: string, @@ -56,12 +62,17 @@ class RedisPubSubChannel implements PubSubChannel { this.pub.on('connect', () => logger.info(`${config.label} publish client connected`)) this.sub.on('connect', () => logger.info(`${config.label} subscribe client connected`)) - this.sub.subscribe(config.channel, (err) => { - if (err) { - logger.error(`Failed to subscribe to ${config.label} channel:`, err) - } else { - logger.info(`Subscribed to ${config.label} channel`) - } + // Settles on failure too: nothing retries a failed subscribe, so waiting on it would only + // hold back every stream on this channel. + this.subscribed = new Promise((resolve) => { + this.sub.subscribe(config.channel, (err) => { + if (err) { + logger.error(`Failed to subscribe to ${config.label} channel:`, err) + } else { + logger.info(`Subscribed to ${config.label} channel`) + } + resolve() + }) }) this.sub.on('message', (channel: string, message: string) => { @@ -95,6 +106,10 @@ class RedisPubSubChannel implements PubSubChannel { } } + ready(): Promise { + return this.subscribed + } + dispose(): void { this.disposed = true this.handlers.clear() @@ -130,6 +145,10 @@ class LocalPubSubChannel implements PubSubChannel { } } + ready(): Promise { + return Promise.resolve() + } + dispose(): void { this.emitter.removeAllListeners() logger.info(`${this.config.label} local pub/sub disposed`) diff --git a/apps/sim/lib/events/sse-endpoint.test.ts b/apps/sim/lib/events/sse-endpoint.test.ts index b4da2936f93..a9b878c5f1a 100644 --- a/apps/sim/lib/events/sse-endpoint.test.ts +++ b/apps/sim/lib/events/sse-endpoint.test.ts @@ -2,10 +2,12 @@ import { authMockFns, permissionsMock, permissionsMockFns } from '@sim/testing' import { NextRequest } from 'next/server' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { + createSSEStream, createWorkspaceSSE, HEARTBEAT_INTERVAL_MS, MAX_CONNECTION_MS, MAX_UNDRAINED_CHUNKS, + OPENED_COMMENT, ROTATION_GRACE_MS, } from '@/lib/events/sse-endpoint' @@ -65,6 +67,16 @@ describe('createWorkspaceSSE', () => { vi.useRealTimers() }) + it('starts the response before the first heartbeat', async () => { + const { body } = await openConnection() + const chunks: string[] = [] + void collect(body, chunks) + + await vi.advanceTimersByTimeAsync(HEARTBEAT_INTERVAL_MS - 1) + + expect(chunks).toEqual([OPENED_COMMENT]) + }) + it('announces rotation before releasing the old connection', async () => { const { body, unsubscribe } = await openConnection() const chunks: string[] = [] @@ -155,3 +167,62 @@ describe('createWorkspaceSSE', () => { expect(unsubscribe).toHaveBeenCalledTimes(1) }) }) + +describe('createSSEStream', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('opens only once its subscriptions receive events', async () => { + let live: () => void = () => {} + const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { + label: 'test', + subscriptions: [ + { + subscribe: () => () => {}, + ready: () => + new Promise((resolve) => { + live = resolve + }), + }, + ], + }) + const chunks: string[] = [] + void collect(response.body as ReadableStream, chunks) + await vi.advanceTimersByTimeAsync(1_000) + expect(chunks).toEqual([]) + + live() + await vi.advanceTimersByTimeAsync(0) + + expect(chunks).toEqual([OPENED_COMMENT]) + }) + + it('delivers a revalidated event as soon as it is published', async () => { + let publish: (eventName: string, data: Record) => void = () => {} + const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { + label: 'test', + revalidate: async () => {}, + subscriptions: [ + { + subscribe: (send) => { + publish = send + return () => {} + }, + }, + ], + }) + const chunks: string[] = [] + void collect(response.body as ReadableStream, chunks) + await vi.advanceTimersByTimeAsync(0) + + publish('inbox_changed', { reason: 'call' }) + await vi.advanceTimersByTimeAsync(0) + + expect(chunks).toEqual([OPENED_COMMENT, 'event: inbox_changed\ndata: {"reason":"call"}\n\n']) + }) +}) diff --git a/apps/sim/lib/events/sse-endpoint.ts b/apps/sim/lib/events/sse-endpoint.ts index 9f110e58f08..5b52447d87f 100644 --- a/apps/sim/lib/events/sse-endpoint.ts +++ b/apps/sim/lib/events/sse-endpoint.ts @@ -19,6 +19,8 @@ interface SSESubscription { workspaceId: string, send: (eventName: string, data: Record) => void ): () => void + /** Settles once the subscription receives events; the stream is announced only after it. */ + ready?: () => Promise } interface WorkspaceSSEConfig { @@ -30,6 +32,9 @@ const encoder = new TextEncoder() export const HEARTBEAT_INTERVAL_MS = 30_000 +/** Written once a stream's subscriptions are live; clients ignore comments. */ +export const OPENED_COMMENT = ': connected\n\n' + /** * Starts a make-before-break rotation for one connection. Healthy clients open * a replacement before this stream closes; orphaned streams are released after @@ -83,6 +88,7 @@ export function createWorkspaceSSE(config: WorkspaceSSEConfig) { label: `${config.label}:workspace:${workspaceId}`, subscriptions: config.subscriptions.map((subscription) => ({ subscribe: (send) => subscription.subscribe(workspaceId, send), + ready: subscription.ready, })), }) } @@ -92,6 +98,8 @@ interface SSEStreamConfig { label: string subscriptions: Array<{ subscribe(send: (eventName: string, data: Record) => void): () => void + /** Settles once the subscription receives events; the stream is announced only after it. */ + ready?: () => Promise }> /** Rechecks a long-lived authorization before each publication and on heartbeats. */ revalidate?: () => Promise @@ -224,6 +232,13 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): }) teardowns.push(() => listenerScope.abort()) + // The runtime sends the status and headers with the first body chunk, so without this the + // response would not start until the first event or heartbeat. A client reads its state + // once the stream opens, so it opens only when no later event can be missed. + void Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())).then( + () => enqueue(OPENED_COMMENT), + () => close('subscription_failed') + ) logger.info(`SSE connection opened for ${config.label}`) } catch (error) { cleanup('setup_failed') diff --git a/apps/sim/lib/mcp/pubsub.ts b/apps/sim/lib/mcp/pubsub.ts index 43bae170979..26a8064dcc0 100644 --- a/apps/sim/lib/mcp/pubsub.ts +++ b/apps/sim/lib/mcp/pubsub.ts @@ -19,6 +19,8 @@ interface McpPubSubAdapter { publishWorkflowToolsChanged(event: WorkflowToolsChangedEvent): void onToolsChanged(handler: (event: ToolsChangedEvent) => void): () => void onWorkflowToolsChanged(handler: (event: WorkflowToolsChangedEvent) => void): () => void + /** Settles once this process receives both channels' events. */ + ready(): Promise dispose(): void } @@ -60,6 +62,9 @@ export const mcpPubSub: McpPubSubAdapter | null = publishWorkflowToolsChanged: (event) => workflowToolsChannel.publish(event), onToolsChanged: (handler) => toolsChannel.subscribe(handler), onWorkflowToolsChanged: (handler) => workflowToolsChannel.subscribe(handler), + ready: async () => { + await Promise.all([toolsChannel.ready(), workflowToolsChannel.ready()]) + }, dispose: () => { toolsChannel.dispose() workflowToolsChannel.dispose() diff --git a/apps/sim/lib/mothership/chat-status.ts b/apps/sim/lib/mothership/chat-status.ts index 3fb5edf163c..ac0b0a6fe03 100644 --- a/apps/sim/lib/mothership/chat-status.ts +++ b/apps/sim/lib/mothership/chat-status.ts @@ -40,6 +40,7 @@ export const chatPubSub = channel ? { publishStatusChanged: (event: ChatStatusEvent) => channel.publish(event), onStatusChanged: (handler: (event: ChatStatusEvent) => void) => channel.subscribe(handler), + ready: () => channel.ready(), dispose: () => channel.dispose(), } : null diff --git a/apps/sim/scripts/test-desktop-inbox-e2e.ts b/apps/sim/scripts/test-desktop-inbox-e2e.ts index 5ef02e09210..fb0ace60925 100644 --- a/apps/sim/scripts/test-desktop-inbox-e2e.ts +++ b/apps/sim/scripts/test-desktop-inbox-e2e.ts @@ -48,6 +48,10 @@ import { const logger = createLogger('DesktopInboxE2E') const PICKUP_GRACE_SECONDS = 15 const RECONCILE_MS = 2_000 +/** Sim desktop gives up on a doorbell whose response has not started by then, and reconnects. */ +const DESKTOP_DOORBELL_HANDSHAKE_MS = 15_000 +/** A ring is a Redis publish and one device check away from the open stream. */ +const RING_DELIVERY_MS = 2_000 function requiredEnvironment(name: string): string { const value = process.env[name] @@ -228,16 +232,18 @@ function openDoorbell(desktop: Desktop) { const controller = new AbortController() const events: string[] = [] const listeners = new Set<() => void>() + const startedAt = Date.now() // boundary-raw-fetch: reads the SSE doorbell as a stream. const opened = fetch(new URL(`/api/desktop/inbox/stream?deviceId=${desktop.deviceId}`, baseUrl), { headers: { Cookie: `better-auth.session_token=${desktop.cookie}`, Accept: 'text/event-stream' }, signal: controller.signal, }).then(async (response) => { + const startedInMs = Date.now() - startedAt http.push({ method: 'GET', path: '/api/desktop/inbox/stream', status: response.status, - durationMs: 0, + durationMs: startedInMs, }) assert.equal(response.status, 200, `inbox stream: ${response.status}`) assert(response.body, 'inbox stream has no body') @@ -265,8 +271,10 @@ function openDoorbell(desktop: Desktop) { } } catch {} })() + return startedInMs }) return { + /** Resolves with how long the response took to start. */ opened, events, onEvent(listener: () => void) { @@ -428,7 +436,11 @@ async function run() { await redis.del(`desktop:presence:${desktop.deviceId}`) const doorbell = openDoorbell(desktop) await check('counts the device online once it opens its doorbell stream', async () => { - await doorbell.opened + const startedInMs = await doorbell.opened + assert( + startedInMs < DESKTOP_DOORBELL_HANDSHAKE_MS, + `the doorbell response started after ${startedInMs} ms` + ) await waitFor( async () => (await redis.exists(`desktop:presence:${desktop.deviceId}`)) === 1, 10_000, @@ -436,6 +448,19 @@ async function run() { ) }) + await check('rings the doorbell as soon as the inbox changes', async () => { + const heardBefore = doorbell.events.length + await redis.publish( + 'desktop:inbox', + JSON.stringify({ deviceId: desktop.deviceId, reason: 'call' }) + ) + await waitFor( + () => doorbell.events.slice(heardBefore).includes('inbox_changed:call'), + RING_DELIVERY_MS, + 'the doorbell' + ) + }) + const executor = startExecutor(desktop, doorbell) try { await check( diff --git a/packages/testing/src/mocks/mcp-pubsub.mock.ts b/packages/testing/src/mocks/mcp-pubsub.mock.ts index bcd090651d2..8db7097e6e2 100644 --- a/packages/testing/src/mocks/mcp-pubsub.mock.ts +++ b/packages/testing/src/mocks/mcp-pubsub.mock.ts @@ -3,7 +3,7 @@ import { vi } from 'vitest' /** * Controllable mock functions for `@/lib/mcp/pubsub`. * - * `publish*` and `mockDispose` are bare no-ops; `mockOnToolsChanged` and + * `publish*` and `mockDispose` are bare no-ops, `mockReady` settles at once; `mockOnToolsChanged` and * `mockOnWorkflowToolsChanged` return a no-op unsubscribe (the real `channel.subscribe` contract). * * @example @@ -21,6 +21,7 @@ export const mcpPubsubMockFns = { mockPublishWorkflowToolsChanged: vi.fn(), mockOnToolsChanged: vi.fn((_handler: (event: unknown) => void): (() => void) => () => {}), mockOnWorkflowToolsChanged: vi.fn((_handler: (event: unknown) => void): (() => void) => () => {}), + mockReady: vi.fn(async (): Promise => {}), mockDispose: vi.fn(), } @@ -40,6 +41,7 @@ export const mcpPubsubMock = { publishWorkflowToolsChanged: mcpPubsubMockFns.mockPublishWorkflowToolsChanged, onToolsChanged: mcpPubsubMockFns.mockOnToolsChanged, onWorkflowToolsChanged: mcpPubsubMockFns.mockOnWorkflowToolsChanged, + ready: mcpPubsubMockFns.mockReady, dispose: mcpPubsubMockFns.mockDispose, }, } diff --git a/packages/testing/src/mocks/mothership-chat-status.mock.ts b/packages/testing/src/mocks/mothership-chat-status.mock.ts index 25c67c021a9..9785be123ed 100644 --- a/packages/testing/src/mocks/mothership-chat-status.mock.ts +++ b/packages/testing/src/mocks/mothership-chat-status.mock.ts @@ -4,7 +4,7 @@ import { vi } from 'vitest' * Controllable mock functions for `@/lib/mothership/chat-status`. * * Every function is a bare `vi.fn()` except `mockOnStatusChanged`, which returns a no-op - * unsubscribe (the real `channel.subscribe` contract). + * unsubscribe (the real `channel.subscribe` contract), and `mockReady`, which settles at once. * * @example * ```ts @@ -20,6 +20,7 @@ export const mothershipChatStatusMockFns = { mockPublishChatStatusChanged: vi.fn(), mockPublishStatusChanged: vi.fn(), mockOnStatusChanged: vi.fn((_handler: (event: unknown) => void): (() => void) => () => {}), + mockReady: vi.fn(async (): Promise => {}), mockDispose: vi.fn(), } @@ -36,6 +37,7 @@ export const mothershipChatStatusMock = { chatPubSub: { publishStatusChanged: mothershipChatStatusMockFns.mockPublishStatusChanged, onStatusChanged: mothershipChatStatusMockFns.mockOnStatusChanged, + ready: mothershipChatStatusMockFns.mockReady, dispose: mothershipChatStatusMockFns.mockDispose, }, publishChatStatusChanged: mothershipChatStatusMockFns.mockPublishChatStatusChanged, From a6460c4c55b0d982b1aeed134569c60b48790515 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 16:46:28 -0700 Subject: [PATCH 2/4] fix(events): hold events until the stream opens and bound the wait - Events wait for the opening gate too, so one published while another subscription is still connecting cannot start the response early. - The gate opens at OPEN_DEADLINE_MS (5 s) when a subscription is not live, so a Redis outage degrades to an open stream instead of a silent one. - PubSubChannel subscribes on every ready connection and is not ready again after a dropped one until it has resubscribed. - The desktop inbox E2E compiles the doorbell route with a refused request before timing its open. --- .../app/api/mothership/events/route.test.ts | 2 + apps/sim/lib/events/pubsub.test.ts | 80 +++++++++++++++++++ apps/sim/lib/events/pubsub.ts | 28 +++++-- apps/sim/lib/events/sse-endpoint.test.ts | 46 +++++++++++ apps/sim/lib/events/sse-endpoint.ts | 43 +++++++--- apps/sim/scripts/test-desktop-inbox-e2e.ts | 8 ++ 6 files changed, 192 insertions(+), 15 deletions(-) create mode 100644 apps/sim/lib/events/pubsub.test.ts diff --git a/apps/sim/app/api/mothership/events/route.test.ts b/apps/sim/app/api/mothership/events/route.test.ts index 39ed74df41f..59a718a7dcf 100644 --- a/apps/sim/app/api/mothership/events/route.test.ts +++ b/apps/sim/app/api/mothership/events/route.test.ts @@ -156,6 +156,7 @@ describe('Mothership owner-scoped event stream', () => { const response = await GET(request('organizationId=org-1')) const chunks: string[] = [] const collected = collect(response.body!, chunks) + await vi.advanceTimersByTimeAsync(0) let authorizeDone: (() => void) | undefined authorize.mockReturnValueOnce( new Promise((resolve) => { @@ -178,6 +179,7 @@ describe('Mothership owner-scoped event stream', () => { const response = await GET(request('workspaceId=ws-1', abort.signal)) const chunks: string[] = [] const collected = collect(response.body!, chunks) + await vi.advanceTimersByTimeAsync(0) emit({ organizationId: 'org-1', userId: 'user-1', chatId: 'org-chat', type: 'created' }) emit({ workspaceId: 'ws-2', chatId: 'other-workspace-chat', type: 'created' }) emit({ workspaceId: 'ws-1', chatId: 'workspace-chat', type: 'renamed' }) diff --git a/apps/sim/lib/events/pubsub.test.ts b/apps/sim/lib/events/pubsub.test.ts new file mode 100644 index 00000000000..bef5fec0435 --- /dev/null +++ b/apps/sim/lib/events/pubsub.test.ts @@ -0,0 +1,80 @@ +import { EventEmitter } from 'node:events' +import { redisConfigMockFns } from '@sim/testing/mocks/redis-config.mock' +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const { clients } = vi.hoisted(() => ({ + clients: [] as Array void> }>, +})) + +vi.mock('ioredis', () => ({ + default: class extends EventEmitter { + subscribed: Array<(err: Error | null) => void> = [] + + constructor() { + super() + clients.push(this) + } + + subscribe(_channel: string, done: (err: Error | null) => void) { + this.subscribed.push(done) + } + + publish = vi.fn(async () => 1) + unsubscribe = vi.fn(async () => undefined) + quit = vi.fn(async () => 'OK') + }, +})) + +import { createPubSubChannel } from '@/lib/events/pubsub' + +/** Whether the promise has settled by the time already-queued work runs. */ +async function settled(promise: Promise): Promise { + return Promise.race([promise.then(() => true), Promise.resolve().then(() => false)]) +} + +/** The subscriber connection: the channel opens its publisher first. */ +function subscriber() { + const client = clients.at(-1) + if (!client) throw new Error('no Redis client') + return client +} + +describe('createPubSubChannel over Redis', () => { + beforeEach(() => { + clients.length = 0 + redisConfigMockFns.mockGetConfiguredRedisUrl.mockReturnValue('redis://localhost:6379') + }) + + it('is ready only once its connection has subscribed', async () => { + const channel = createPubSubChannel({ channel: 'test', label: 'Test' }) + expect(await settled(channel.ready())).toBe(false) + + subscriber().emit('ready') + expect(await settled(channel.ready())).toBe(false) + + subscriber().subscribed[0](null) + expect(await settled(channel.ready())).toBe(true) + }) + + it('is not ready again until a dropped connection has resubscribed', async () => { + const channel = createPubSubChannel({ channel: 'test', label: 'Test' }) + subscriber().emit('ready') + subscriber().subscribed[0](null) + + subscriber().emit('close') + expect(await settled(channel.ready())).toBe(false) + + subscriber().emit('ready') + subscriber().subscribed[1](null) + expect(await settled(channel.ready())).toBe(true) + }) + + it('settles when the subscribe fails, so streams are not held back', async () => { + const channel = createPubSubChannel({ channel: 'test', label: 'Test' }) + subscriber().emit('ready') + + subscriber().subscribed[0](new Error('NOPERM')) + + expect(await settled(channel.ready())).toBe(true) + }) +}) diff --git a/apps/sim/lib/events/pubsub.ts b/apps/sim/lib/events/pubsub.ts index 853e68bc0d1..04b4f79ce21 100644 --- a/apps/sim/lib/events/pubsub.ts +++ b/apps/sim/lib/events/pubsub.ts @@ -34,7 +34,10 @@ class RedisPubSubChannel implements PubSubChannel { private sub: Redis private handlers = new Set<(event: T) => void>() private disposed = false - private readonly subscribed: Promise + /** Whether the current connection has subscribed; a dropped connection has to again. */ + private listening = false + private subscribed: Promise = Promise.resolve() + private markSubscribed: () => void = noop constructor( redisUrl: string, @@ -62,18 +65,27 @@ class RedisPubSubChannel implements PubSubChannel { this.pub.on('connect', () => logger.info(`${config.label} publish client connected`)) this.sub.on('connect', () => logger.info(`${config.label} subscribe client connected`)) - // Settles on failure too: nothing retries a failed subscribe, so waiting on it would only - // hold back every stream on this channel. - this.subscribed = new Promise((resolve) => { + this.awaitSubscription() + // Subscribes on every ready connection: ioredis resubscribes after a reconnect on its own but + // does not report when that lands, and SUBSCRIBE is idempotent. Readiness settles on failure + // too: nothing retries a failed subscribe, so waiting on it would only hold back every stream + // on this channel. + this.sub.on('ready', () => { this.sub.subscribe(config.channel, (err) => { if (err) { logger.error(`Failed to subscribe to ${config.label} channel:`, err) } else { + this.listening = true logger.info(`Subscribed to ${config.label} channel`) } - resolve() + this.markSubscribed() }) }) + this.sub.on('close', () => { + if (!this.listening) return + this.listening = false + this.awaitSubscription() + }) this.sub.on('message', (channel: string, message: string) => { if (channel !== config.channel) return @@ -110,6 +122,12 @@ class RedisPubSubChannel implements PubSubChannel { return this.subscribed } + private awaitSubscription(): void { + this.subscribed = new Promise((resolve) => { + this.markSubscribed = resolve + }) + } + dispose(): void { this.disposed = true this.handlers.clear() diff --git a/apps/sim/lib/events/sse-endpoint.test.ts b/apps/sim/lib/events/sse-endpoint.test.ts index a9b878c5f1a..86efe7fa020 100644 --- a/apps/sim/lib/events/sse-endpoint.test.ts +++ b/apps/sim/lib/events/sse-endpoint.test.ts @@ -7,6 +7,7 @@ import { HEARTBEAT_INTERVAL_MS, MAX_CONNECTION_MS, MAX_UNDRAINED_CHUNKS, + OPEN_DEADLINE_MS, OPENED_COMMENT, ROTATION_GRACE_MS, } from '@/lib/events/sse-endpoint' @@ -202,6 +203,51 @@ describe('createSSEStream', () => { expect(chunks).toEqual([OPENED_COMMENT]) }) + it('holds events until the stream opens', async () => { + let live: () => void = () => {} + let publish: (eventName: string, data: Record) => void = () => {} + const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { + label: 'test', + subscriptions: [ + { + subscribe: (send) => { + publish = send + return () => {} + }, + ready: () => + new Promise((resolve) => { + live = resolve + }), + }, + ], + }) + const chunks: string[] = [] + void collect(response.body as ReadableStream, chunks) + + publish('changed', { id: 1 }) + await vi.advanceTimersByTimeAsync(0) + expect(chunks).toEqual([]) + + live() + await vi.advanceTimersByTimeAsync(0) + expect(chunks).toEqual([OPENED_COMMENT, 'event: changed\ndata: {"id":1}\n\n']) + }) + + it('opens at the deadline when a subscription never becomes ready', async () => { + const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { + label: 'test', + subscriptions: [{ subscribe: () => () => {}, ready: () => new Promise(() => {}) }], + }) + const chunks: string[] = [] + void collect(response.body as ReadableStream, chunks) + + await vi.advanceTimersByTimeAsync(OPEN_DEADLINE_MS - 1) + expect(chunks).toEqual([]) + + await vi.advanceTimersByTimeAsync(1) + expect(chunks).toEqual([OPENED_COMMENT]) + }) + it('delivers a revalidated event as soon as it is published', async () => { let publish: (eventName: string, data: Record) => void = () => {} const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { diff --git a/apps/sim/lib/events/sse-endpoint.ts b/apps/sim/lib/events/sse-endpoint.ts index 5b52447d87f..0c6d5f78974 100644 --- a/apps/sim/lib/events/sse-endpoint.ts +++ b/apps/sim/lib/events/sse-endpoint.ts @@ -35,6 +35,13 @@ export const HEARTBEAT_INTERVAL_MS = 30_000 /** Written once a stream's subscriptions are live; clients ignore comments. */ export const OPENED_COMMENT = ': connected\n\n' +/** + * How long a stream waits for its subscriptions before opening anyway. A subscription still not + * live is in an outage, which open streams ride out the same way, and a client that hears nothing + * gives up on the connection: Sim desktop after 15 s. + */ +export const OPEN_DEADLINE_MS = 5_000 + /** * Starts a make-before-break rotation for one connection. Healthy clients open * a replacement before this stream closes; orphaned streams are released after @@ -157,20 +164,26 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): }) return authorization } + /** Settles when the stream announces itself; events wait for it so none opens it early. */ + let opening: Promise = Promise.resolve() + let opened = false let pendingEvents = 0 const send = (eventName: string, data: Record) => { if (cleaned) return const payload = `event: ${eventName}\ndata: ${JSON.stringify(data)}\n\n` - if (!config.revalidate) { + if (opened && !config.revalidate) { enqueue(payload) return } if (pendingEvents >= MAX_UNDRAINED_CHUNKS) { - close('authorization_backpressure') + close('pending_backpressure') return } pendingEvents += 1 - void revalidate().then( + const authorized = opened + ? revalidate() + : opening.then(() => (cleaned ? undefined : revalidate())) + void authorized.then( () => { pendingEvents -= 1 enqueue(payload) @@ -183,6 +196,23 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): } try { + // The runtime sends the status and headers with the first body chunk, so the stream writes + // one as soon as it opens. A client reads its state once the stream opens, so it opens once + // every subscription receives events, or at the deadline when one is in an outage. A + // heartbeat cannot open it first: the deadline is shorter than the heartbeat interval. + opening = Promise.race([ + Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())), + new Promise((resolve) => { + const deadline = setTimeout(resolve, OPEN_DEADLINE_MS) + teardowns.push(() => clearTimeout(deadline)) + }), + ]).then( + () => { + opened = true + enqueue(OPENED_COMMENT) + }, + () => close('subscription_failed') + ) for (const subscription of config.subscriptions) { teardowns.push(subscription.subscribe(send)) } @@ -232,13 +262,6 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): }) teardowns.push(() => listenerScope.abort()) - // The runtime sends the status and headers with the first body chunk, so without this the - // response would not start until the first event or heartbeat. A client reads its state - // once the stream opens, so it opens only when no later event can be missed. - void Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())).then( - () => enqueue(OPENED_COMMENT), - () => close('subscription_failed') - ) logger.info(`SSE connection opened for ${config.label}`) } catch (error) { cleanup('setup_failed') diff --git a/apps/sim/scripts/test-desktop-inbox-e2e.ts b/apps/sim/scripts/test-desktop-inbox-e2e.ts index fb0ace60925..c841c3707a3 100644 --- a/apps/sim/scripts/test-desktop-inbox-e2e.ts +++ b/apps/sim/scripts/test-desktop-inbox-e2e.ts @@ -432,6 +432,14 @@ async function run() { } }) + /** + * Compiles the doorbell route before its open is timed. Refused before the use case runs, so the + * open below is still the first to subscribe this process to the doorbell channel. + */ + await check('refuses a doorbell without a device', async () => { + await request(desktop, 'GET', '/api/desktop/inbox/stream', { expected: 400 }) + }) + /** Starts absent, so only the stream open below can mark the device present. */ await redis.del(`desktop:presence:${desktop.deviceId}`) const doorbell = openDoorbell(desktop) From 5fd29f0e4f127d80f82988996c849677bbbb6c6b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 17:29:40 -0700 Subject: [PATCH 3/4] fix(events): keep a channel not ready while its subscribe fails or answers a closed connection A failed SUBSCRIBE no longer settles ready(): the stream's opening deadline bounds the wait, and the next connection subscribes again. A subscribe answered after its connection closed is ignored, so it cannot mark a reconnected channel ready before that connection has subscribed. --- apps/sim/lib/events/pubsub.test.ts | 19 ++++++++++++++++++- apps/sim/lib/events/pubsub.ts | 19 ++++++++++++------- 2 files changed, 30 insertions(+), 8 deletions(-) diff --git a/apps/sim/lib/events/pubsub.test.ts b/apps/sim/lib/events/pubsub.test.ts index bef5fec0435..06dd0922dcc 100644 --- a/apps/sim/lib/events/pubsub.test.ts +++ b/apps/sim/lib/events/pubsub.test.ts @@ -69,12 +69,29 @@ describe('createPubSubChannel over Redis', () => { expect(await settled(channel.ready())).toBe(true) }) - it('settles when the subscribe fails, so streams are not held back', async () => { + it('is not ready while its subscribe fails, and retries on the next connection', async () => { const channel = createPubSubChannel({ channel: 'test', label: 'Test' }) subscriber().emit('ready') subscriber().subscribed[0](new Error('NOPERM')) + expect(await settled(channel.ready())).toBe(false) + + subscriber().emit('close') + subscriber().emit('ready') + subscriber().subscribed[1](null) + expect(await settled(channel.ready())).toBe(true) + }) + + it('ignores a subscribe answered after its connection closed', async () => { + const channel = createPubSubChannel({ channel: 'test', label: 'Test' }) + subscriber().emit('ready') + + subscriber().emit('close') + subscriber().subscribed[0](null) + expect(await settled(channel.ready())).toBe(false) + subscriber().emit('ready') + subscriber().subscribed[1](null) expect(await settled(channel.ready())).toBe(true) }) }) diff --git a/apps/sim/lib/events/pubsub.ts b/apps/sim/lib/events/pubsub.ts index 04b4f79ce21..9550ae4db8d 100644 --- a/apps/sim/lib/events/pubsub.ts +++ b/apps/sim/lib/events/pubsub.ts @@ -18,7 +18,8 @@ export interface PubSubChannel { subscribe(handler: (event: T) => void): () => void /** * Settles once this process receives the channel's publications; anything published before - * then reaches no subscriber here. + * then reaches no subscriber here. Stays pending while the channel cannot subscribe, so a caller + * that must not wait indefinitely bounds the wait. */ ready(): Promise dispose(): void @@ -36,6 +37,8 @@ class RedisPubSubChannel implements PubSubChannel { private disposed = false /** Whether the current connection has subscribed; a dropped connection has to again. */ private listening = false + /** Counts closed connections, so a subscribe answered after its connection closed is ignored. */ + private closedConnections = 0 private subscribed: Promise = Promise.resolve() private markSubscribed: () => void = noop @@ -67,21 +70,23 @@ class RedisPubSubChannel implements PubSubChannel { this.awaitSubscription() // Subscribes on every ready connection: ioredis resubscribes after a reconnect on its own but - // does not report when that lands, and SUBSCRIBE is idempotent. Readiness settles on failure - // too: nothing retries a failed subscribe, so waiting on it would only hold back every stream - // on this channel. + // does not report when that lands, and SUBSCRIBE is idempotent. A failed subscribe leaves the + // channel not ready; the next connection tries again. this.sub.on('ready', () => { + const connection = this.closedConnections this.sub.subscribe(config.channel, (err) => { + if (connection !== this.closedConnections) return if (err) { logger.error(`Failed to subscribe to ${config.label} channel:`, err) - } else { - this.listening = true - logger.info(`Subscribed to ${config.label} channel`) + return } + this.listening = true + logger.info(`Subscribed to ${config.label} channel`) this.markSubscribed() }) }) this.sub.on('close', () => { + this.closedConnections += 1 if (!this.listening) return this.listening = false this.awaitSubscription() From a13b9773dcf40b376404ccc56dfccfa8a34b77ea Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 18:37:37 -0700 Subject: [PATCH 4/4] fix(events): write each stream's events in order and release a stream closed before it opens - Every event goes through one ordered write chain: it is authorized as soon as it arrives, so concurrent events still share one check, but written only after every event before it. An event queued before the stream opened can no longer be overtaken by one that arrived during its authorization. - Closing settles the opening gate as well as clearing its deadline, so a stream closed before it opens is not kept alive by a subscription that never becomes ready. - Tests cover event order, retention after an early close, timer cleanup, and the ready wiring of the mothership and MCP event routes. --- apps/sim/app/api/mcp/events/route.test.ts | 50 +++++++++++ .../app/api/mothership/events/route.test.ts | 28 ++++++ apps/sim/lib/events/sse-endpoint.test.ts | 85 +++++++++++++++++++ apps/sim/lib/events/sse-endpoint.ts | 46 ++++++---- 4 files changed, 193 insertions(+), 16 deletions(-) create mode 100644 apps/sim/app/api/mcp/events/route.test.ts diff --git a/apps/sim/app/api/mcp/events/route.test.ts b/apps/sim/app/api/mcp/events/route.test.ts new file mode 100644 index 00000000000..4d1848cf084 --- /dev/null +++ b/apps/sim/app/api/mcp/events/route.test.ts @@ -0,0 +1,50 @@ +import { authMockFns, permissionsMock, permissionsMockFns } from '@sim/testing' +import { mcpPubsubMock, mcpPubsubMockFns } from '@sim/testing/mocks/mcp-pubsub.mock' +import { NextRequest } from 'next/server' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { OPENED_COMMENT } from '@/lib/events/sse-endpoint' + +vi.mock('@/lib/mcp/pubsub', () => mcpPubsubMock) +vi.mock('@/lib/mcp/connection-manager', () => ({ + mcpConnectionManager: { subscribe: () => () => {} }, +})) +vi.mock('@/lib/workspaces/permissions/utils', () => permissionsMock) + +import { GET } from '@/app/api/mcp/events/route' + +describe('MCP tool-change event stream', () => { + beforeEach(() => { + vi.useFakeTimers() + authMockFns.mockGetSession.mockResolvedValue({ user: { id: 'user-1' } }) + permissionsMockFns.mockGetUserEntityPermissions.mockResolvedValue('read') + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('opens once tool-change events reach this process', async () => { + let live: () => void = () => {} + const ready = new Promise((resolve) => { + live = resolve + }) + mcpPubsubMockFns.mockReady.mockReturnValue(ready) + const response = await GET(new NextRequest('http://localhost/api/mcp/events?workspaceId=ws-1')) + let first: string | undefined + if (!response.body) throw new Error('The event stream has no body') + void response.body + .getReader() + .read() + .then(({ value }) => { + first = new TextDecoder().decode(value) + }) + + await vi.advanceTimersByTimeAsync(1_000) + expect(first).toBeUndefined() + expect(mcpPubsubMockFns.mockReady).toHaveBeenCalledTimes(2) + + live() + await vi.advanceTimersByTimeAsync(0) + expect(first).toBe(OPENED_COMMENT) + }) +}) diff --git a/apps/sim/app/api/mothership/events/route.test.ts b/apps/sim/app/api/mothership/events/route.test.ts index 59a718a7dcf..2b379d6257c 100644 --- a/apps/sim/app/api/mothership/events/route.test.ts +++ b/apps/sim/app/api/mothership/events/route.test.ts @@ -174,6 +174,34 @@ describe('Mothership owner-scoped event stream', () => { expect(chunks).toEqual([]) }) + it.each(['workspaceId=ws-1', 'organizationId=org-1'])( + 'opens the %s stream once chat status events reach this process', + async (query) => { + let live: () => void = () => {} + mothershipChatStatusMockFns.mockReady.mockReturnValueOnce( + new Promise((resolve) => { + live = resolve + }) + ) + const response = await GET(request(query)) + let first: string | undefined + if (!response.body) throw new Error('The event stream has no body') + void response.body + .getReader() + .read() + .then(({ value }) => { + first = new TextDecoder().decode(value) + }) + + await vi.advanceTimersByTimeAsync(1_000) + expect(first).toBeUndefined() + + live() + await vi.advanceTimersByTimeAsync(0) + expect(first).toBe(OPENED_COMMENT) + } + ) + it('preserves workspace status events and excludes organization events', async () => { const abort = new AbortController() const response = await GET(request('workspaceId=ws-1', abort.signal)) diff --git a/apps/sim/lib/events/sse-endpoint.test.ts b/apps/sim/lib/events/sse-endpoint.test.ts index 86efe7fa020..603d0c95e88 100644 --- a/apps/sim/lib/events/sse-endpoint.test.ts +++ b/apps/sim/lib/events/sse-endpoint.test.ts @@ -1,4 +1,7 @@ +import { setFlagsFromString } from 'node:v8' +import { runInNewContext } from 'node:vm' import { authMockFns, permissionsMock, permissionsMockFns } from '@sim/testing' +import { sleep } from '@sim/utils/helpers' import { NextRequest } from 'next/server' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { @@ -248,6 +251,62 @@ describe('createSSEStream', () => { expect(chunks).toEqual([OPENED_COMMENT]) }) + it('writes revalidated events in the order they arrive', async () => { + let live: () => void = () => {} + let authorized: () => void = () => {} + let publish: (eventName: string, data: Record) => void = () => {} + const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { + label: 'test', + revalidate: () => + new Promise((resolve) => { + authorized = resolve + }), + subscriptions: [ + { + subscribe: (send) => { + publish = send + return () => {} + }, + ready: () => + new Promise((resolve) => { + live = resolve + }), + }, + ], + }) + const chunks: string[] = [] + void collect(response.body as ReadableStream, chunks) + + publish('changed', { n: 1 }) + live() + await vi.advanceTimersByTimeAsync(0) + publish('changed', { n: 2 }) + authorized() + await vi.advanceTimersByTimeAsync(0) + + expect(chunks).toEqual([ + OPENED_COMMENT, + 'event: changed\ndata: {"n":1}\n\n', + 'event: changed\ndata: {"n":2}\n\n', + ]) + }) + + it('clears its timers when it closes before it opens', async () => { + const controller = new AbortController() + createSSEStream( + new NextRequest(new URL('https://sim.test/api/test/stream'), { signal: controller.signal }), + { + label: 'test', + subscriptions: [{ subscribe: () => () => {}, ready: () => new Promise(() => {}) }], + } + ) + expect(vi.getTimerCount()).toBe(2) + + controller.abort() + + expect(vi.getTimerCount()).toBe(0) + }) + it('delivers a revalidated event as soon as it is published', async () => { let publish: (eventName: string, data: Record) => void = () => {} const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), { @@ -272,3 +331,29 @@ describe('createSSEStream', () => { expect(chunks).toEqual([OPENED_COMMENT, 'event: inbox_changed\ndata: {"reason":"call"}\n\n']) }) }) + +describe('createSSEStream retention', () => { + /** Opens a stream and closes it before it opens; only a weak reference to it escapes. */ + function closedBeforeOpening(ready: Promise): WeakRef { + const controller = new AbortController() + const response = createSSEStream( + new NextRequest(new URL('https://sim.test/api/test/stream'), { signal: controller.signal }), + { label: 'test', subscriptions: [{ subscribe: () => () => {}, ready: () => ready }] } + ) + controller.abort() + return new WeakRef(response.body as object) + } + + it('does not stay reachable from a subscription that never becomes ready', async () => { + setFlagsFromString('--expose_gc') + const collectGarbage = runInNewContext('gc') as () => void + const neverReady = new Promise(() => {}) + + const stream = closedBeforeOpening(neverReady) + await sleep(0) + collectGarbage() + + expect(stream.deref()).toBeUndefined() + expect(neverReady).toBeInstanceOf(Promise) + }) +}) diff --git a/apps/sim/lib/events/sse-endpoint.ts b/apps/sim/lib/events/sse-endpoint.ts index 0c6d5f78974..4e44e3ce673 100644 --- a/apps/sim/lib/events/sse-endpoint.ts +++ b/apps/sim/lib/events/sse-endpoint.ts @@ -8,6 +8,7 @@ import type { SessionPrincipal } from '@sim/auth/principal' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' +import { noop } from '@sim/utils/helpers' import { randomFloat } from '@sim/utils/random' import type { NextRequest } from 'next/server' import { getSession } from '@/lib/auth' @@ -164,14 +165,17 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): }) return authorization } - /** Settles when the stream announces itself; events wait for it so none opens it early. */ - let opening: Promise = Promise.resolve() + /** + * The stream's writes in order: the opening, then each event once it is authorized. An event + * waits here while the stream has not opened or an earlier event is still being written. + */ + let writes: Promise = Promise.resolve() let opened = false let pendingEvents = 0 const send = (eventName: string, data: Record) => { if (cleaned) return const payload = `event: ${eventName}\ndata: ${JSON.stringify(data)}\n\n` - if (opened && !config.revalidate) { + if (opened && pendingEvents === 0 && !config.revalidate) { enqueue(payload) return } @@ -180,19 +184,24 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): return } pendingEvents += 1 + // Authorized as soon as it arrives, so concurrent events share one check, but written only + // after every event before it. const authorized = opened ? revalidate() - : opening.then(() => (cleaned ? undefined : revalidate())) - void authorized.then( - () => { - pendingEvents -= 1 - enqueue(payload) - }, - () => { - pendingEvents -= 1 - close('authorization_lost') - } - ) + : writes.then(() => (cleaned ? undefined : revalidate())) + authorized.catch(noop) + writes = writes + .then(() => authorized) + .then( + () => { + pendingEvents -= 1 + enqueue(payload) + }, + () => { + pendingEvents -= 1 + close('authorization_lost') + } + ) } try { @@ -200,11 +209,16 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig): // one as soon as it opens. A client reads its state once the stream opens, so it opens once // every subscription receives events, or at the deadline when one is in an outage. A // heartbeat cannot open it first: the deadline is shorter than the heartbeat interval. - opening = Promise.race([ + // Closing settles it too, so a stream closed before it opened is not kept alive by a + // subscription that never becomes ready. + writes = Promise.race([ Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())), new Promise((resolve) => { const deadline = setTimeout(resolve, OPEN_DEADLINE_MS) - teardowns.push(() => clearTimeout(deadline)) + teardowns.push(() => { + clearTimeout(deadline) + resolve() + }) }), ]).then( () => {