Skip to content

Commit 5bb98aa

Browse files
committed
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.
1 parent 11f8190 commit 5bb98aa

14 files changed

Lines changed: 166 additions & 13 deletions

File tree

‎apps/sim/app/api/desktop/inbox/stream/route.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
3232
return createSSEStream(request, {
3333
label: 'desktop-inbox',
3434
revalidate: inbox.revalidate,
35-
subscriptions: [{ subscribe: inbox.subscribe }],
35+
subscriptions: [{ subscribe: inbox.subscribe, ready: inbox.ready }],
3636
})
3737
} catch (error) {
3838
if (error instanceof InternalUnauthenticatedError)

‎apps/sim/app/api/mcp/events/route.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ const mcpEventsHandler = createWorkspaceSSE({
3232
})
3333
})
3434
},
35+
ready: async () => mcpPubSub?.ready(),
3536
},
3637
{
3738
subscribe: (workspaceId, send) => {
@@ -45,6 +46,7 @@ const mcpEventsHandler = createWorkspaceSSE({
4546
})
4647
})
4748
},
49+
ready: async () => mcpPubSub?.ready(),
4850
},
4951
],
5052
})

‎apps/sim/app/api/mothership/events/route.test.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ import {
1616
import { NextRequest } from 'next/server'
1717
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
1818
import { OrchestrationError } from '@/lib/core/orchestration/types'
19-
import { HEARTBEAT_INTERVAL_MS } from '@/lib/events/sse-endpoint'
19+
import { HEARTBEAT_INTERVAL_MS, OPENED_COMMENT } from '@/lib/events/sse-endpoint'
2020
import type { ChatStatusEvent } from '@/lib/mothership/chat-status'
2121
import { PermissionGroupCapabilityError } from '@/lib/permission-groups/capability-error'
2222

@@ -41,13 +41,15 @@ function emit(event: ChatStatusEvent) {
4141
handler(event)
4242
}
4343

44+
/** Every chunk after the stream's opening comment. */
4445
async function collect(body: ReadableStream<Uint8Array>, chunks: string[]) {
4546
const reader = body.getReader()
4647
const decoder = new TextDecoder()
4748
while (true) {
4849
const { done, value } = await reader.read()
4950
if (done) return
50-
chunks.push(decoder.decode(value))
51+
const chunk = decoder.decode(value)
52+
if (chunk !== OPENED_COMMENT) chunks.push(chunk)
5153
}
5254
}
5355

‎apps/sim/app/api/mothership/events/route.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ const mothershipEventsHandler = createWorkspaceSSE({
4242
})
4343
})
4444
},
45+
ready: async () => chatPubSub?.ready(),
4546
},
4647
],
4748
})
@@ -80,6 +81,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
8081
timestamp: Date.now(),
8182
})
8283
}) ?? (() => {}),
84+
ready: async () => chatPubSub?.ready(),
8385
},
8486
],
8587
})

‎apps/sim/lib/desktop/application/executor.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
} from '@/lib/desktop/executor/constants'
1616
import {
1717
type DesktopInboxChangeReason,
18+
desktopInboxDoorbellReady,
1819
onDesktopInboxDoorbell,
1920
} from '@/lib/desktop/executor/doorbell'
2021
import {
@@ -201,6 +202,7 @@ export const openDesktopInboxStream = defineAuthorizedCredentialUserUseCase({
201202
onDesktopInboxDoorbell(deviceId, (reason: DesktopInboxChangeReason) =>
202203
send('inbox_changed', { reason })
203204
),
205+
ready: desktopInboxDoorbellReady,
204206
}
205207
},
206208
})

‎apps/sim/lib/desktop/executor/doorbell.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,11 @@ export function ringDesktopInbox(deviceId: string, reason: DesktopInboxChangeRea
4343
}
4444
}
4545

46+
/** Settles once this process hears every ring. */
47+
export function desktopInboxDoorbellReady(): Promise<void> {
48+
return channel().ready()
49+
}
50+
4651
/** Subscribes to one device's doorbell; returns the unsubscribe. */
4752
export function onDesktopInboxDoorbell(
4853
deviceId: string,

‎apps/sim/lib/events/pubsub.ts‎

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,11 @@ const logger = createLogger('PubSub')
1616
export interface PubSubChannel<T> {
1717
publish(event: T): void
1818
subscribe(handler: (event: T) => void): () => void
19+
/**
20+
* Settles once this process receives the channel's publications; anything published before
21+
* then reaches no subscriber here.
22+
*/
23+
ready(): Promise<void>
1924
dispose(): void
2025
}
2126

@@ -29,6 +34,7 @@ class RedisPubSubChannel<T> implements PubSubChannel<T> {
2934
private sub: Redis
3035
private handlers = new Set<(event: T) => void>()
3136
private disposed = false
37+
private readonly subscribed: Promise<void>
3238

3339
constructor(
3440
redisUrl: string,
@@ -56,12 +62,17 @@ class RedisPubSubChannel<T> implements PubSubChannel<T> {
5662
this.pub.on('connect', () => logger.info(`${config.label} publish client connected`))
5763
this.sub.on('connect', () => logger.info(`${config.label} subscribe client connected`))
5864

59-
this.sub.subscribe(config.channel, (err) => {
60-
if (err) {
61-
logger.error(`Failed to subscribe to ${config.label} channel:`, err)
62-
} else {
63-
logger.info(`Subscribed to ${config.label} channel`)
64-
}
65+
// Settles on failure too: nothing retries a failed subscribe, so waiting on it would only
66+
// hold back every stream on this channel.
67+
this.subscribed = new Promise((resolve) => {
68+
this.sub.subscribe(config.channel, (err) => {
69+
if (err) {
70+
logger.error(`Failed to subscribe to ${config.label} channel:`, err)
71+
} else {
72+
logger.info(`Subscribed to ${config.label} channel`)
73+
}
74+
resolve()
75+
})
6576
})
6677

6778
this.sub.on('message', (channel: string, message: string) => {
@@ -95,6 +106,10 @@ class RedisPubSubChannel<T> implements PubSubChannel<T> {
95106
}
96107
}
97108

109+
ready(): Promise<void> {
110+
return this.subscribed
111+
}
112+
98113
dispose(): void {
99114
this.disposed = true
100115
this.handlers.clear()
@@ -130,6 +145,10 @@ class LocalPubSubChannel<T> implements PubSubChannel<T> {
130145
}
131146
}
132147

148+
ready(): Promise<void> {
149+
return Promise.resolve()
150+
}
151+
133152
dispose(): void {
134153
this.emitter.removeAllListeners()
135154
logger.info(`${this.config.label} local pub/sub disposed`)

‎apps/sim/lib/events/sse-endpoint.test.ts‎

Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,10 +2,12 @@ import { authMockFns, permissionsMock, permissionsMockFns } from '@sim/testing'
22
import { NextRequest } from 'next/server'
33
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
44
import {
5+
createSSEStream,
56
createWorkspaceSSE,
67
HEARTBEAT_INTERVAL_MS,
78
MAX_CONNECTION_MS,
89
MAX_UNDRAINED_CHUNKS,
10+
OPENED_COMMENT,
911
ROTATION_GRACE_MS,
1012
} from '@/lib/events/sse-endpoint'
1113

@@ -65,6 +67,16 @@ describe('createWorkspaceSSE', () => {
6567
vi.useRealTimers()
6668
})
6769

70+
it('starts the response before the first heartbeat', async () => {
71+
const { body } = await openConnection()
72+
const chunks: string[] = []
73+
void collect(body, chunks)
74+
75+
await vi.advanceTimersByTimeAsync(HEARTBEAT_INTERVAL_MS - 1)
76+
77+
expect(chunks).toEqual([OPENED_COMMENT])
78+
})
79+
6880
it('announces rotation before releasing the old connection', async () => {
6981
const { body, unsubscribe } = await openConnection()
7082
const chunks: string[] = []
@@ -155,3 +167,62 @@ describe('createWorkspaceSSE', () => {
155167
expect(unsubscribe).toHaveBeenCalledTimes(1)
156168
})
157169
})
170+
171+
describe('createSSEStream', () => {
172+
beforeEach(() => {
173+
vi.useFakeTimers()
174+
})
175+
176+
afterEach(() => {
177+
vi.useRealTimers()
178+
})
179+
180+
it('opens only once its subscriptions receive events', async () => {
181+
let live: () => void = () => {}
182+
const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), {
183+
label: 'test',
184+
subscriptions: [
185+
{
186+
subscribe: () => () => {},
187+
ready: () =>
188+
new Promise<void>((resolve) => {
189+
live = resolve
190+
}),
191+
},
192+
],
193+
})
194+
const chunks: string[] = []
195+
void collect(response.body as ReadableStream<Uint8Array>, chunks)
196+
await vi.advanceTimersByTimeAsync(1_000)
197+
expect(chunks).toEqual([])
198+
199+
live()
200+
await vi.advanceTimersByTimeAsync(0)
201+
202+
expect(chunks).toEqual([OPENED_COMMENT])
203+
})
204+
205+
it('delivers a revalidated event as soon as it is published', async () => {
206+
let publish: (eventName: string, data: Record<string, unknown>) => void = () => {}
207+
const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), {
208+
label: 'test',
209+
revalidate: async () => {},
210+
subscriptions: [
211+
{
212+
subscribe: (send) => {
213+
publish = send
214+
return () => {}
215+
},
216+
},
217+
],
218+
})
219+
const chunks: string[] = []
220+
void collect(response.body as ReadableStream<Uint8Array>, chunks)
221+
await vi.advanceTimersByTimeAsync(0)
222+
223+
publish('inbox_changed', { reason: 'call' })
224+
await vi.advanceTimersByTimeAsync(0)
225+
226+
expect(chunks).toEqual([OPENED_COMMENT, 'event: inbox_changed\ndata: {"reason":"call"}\n\n'])
227+
})
228+
})

‎apps/sim/lib/events/sse-endpoint.ts‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ interface SSESubscription {
1919
workspaceId: string,
2020
send: (eventName: string, data: Record<string, unknown>) => void
2121
): () => void
22+
/** Settles once the subscription receives events; the stream is announced only after it. */
23+
ready?: () => Promise<void>
2224
}
2325

2426
interface WorkspaceSSEConfig {
@@ -30,6 +32,9 @@ const encoder = new TextEncoder()
3032

3133
export const HEARTBEAT_INTERVAL_MS = 30_000
3234

35+
/** Written once a stream's subscriptions are live; clients ignore comments. */
36+
export const OPENED_COMMENT = ': connected\n\n'
37+
3338
/**
3439
* Starts a make-before-break rotation for one connection. Healthy clients open
3540
* a replacement before this stream closes; orphaned streams are released after
@@ -83,6 +88,7 @@ export function createWorkspaceSSE(config: WorkspaceSSEConfig) {
8388
label: `${config.label}:workspace:${workspaceId}`,
8489
subscriptions: config.subscriptions.map((subscription) => ({
8590
subscribe: (send) => subscription.subscribe(workspaceId, send),
91+
ready: subscription.ready,
8692
})),
8793
})
8894
}
@@ -92,6 +98,8 @@ interface SSEStreamConfig {
9298
label: string
9399
subscriptions: Array<{
94100
subscribe(send: (eventName: string, data: Record<string, unknown>) => void): () => void
101+
/** Settles once the subscription receives events; the stream is announced only after it. */
102+
ready?: () => Promise<void>
95103
}>
96104
/** Rechecks a long-lived authorization before each publication and on heartbeats. */
97105
revalidate?: () => Promise<void>
@@ -224,6 +232,13 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
224232
})
225233
teardowns.push(() => listenerScope.abort())
226234

235+
// The runtime sends the status and headers with the first body chunk, so without this the
236+
// response would not start until the first event or heartbeat. A client reads its state
237+
// once the stream opens, so it opens only when no later event can be missed.
238+
void Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())).then(
239+
() => enqueue(OPENED_COMMENT),
240+
() => close('subscription_failed')
241+
)
227242
logger.info(`SSE connection opened for ${config.label}`)
228243
} catch (error) {
229244
cleanup('setup_failed')

‎apps/sim/lib/mcp/pubsub.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@ interface McpPubSubAdapter {
1919
publishWorkflowToolsChanged(event: WorkflowToolsChangedEvent): void
2020
onToolsChanged(handler: (event: ToolsChangedEvent) => void): () => void
2121
onWorkflowToolsChanged(handler: (event: WorkflowToolsChangedEvent) => void): () => void
22+
/** Settles once this process receives both channels' events. */
23+
ready(): Promise<void>
2224
dispose(): void
2325
}
2426

@@ -60,6 +62,9 @@ export const mcpPubSub: McpPubSubAdapter | null =
6062
publishWorkflowToolsChanged: (event) => workflowToolsChannel.publish(event),
6163
onToolsChanged: (handler) => toolsChannel.subscribe(handler),
6264
onWorkflowToolsChanged: (handler) => workflowToolsChannel.subscribe(handler),
65+
ready: async () => {
66+
await Promise.all([toolsChannel.ready(), workflowToolsChannel.ready()])
67+
},
6368
dispose: () => {
6469
toolsChannel.dispose()
6570
workflowToolsChannel.dispose()

0 commit comments

Comments
 (0)