Skip to content

Commit ab588d4

Browse files
authored
fix(events): authorize the events queued before a stream opens once (#8718)
* fix(events): authorize the events queued before a stream opens once Events that arrived before the stream opened each queued their own authorization behind the write chain, so ten of them made ten serial checks. They now share one that starts when the stream opens, while each is still written in order. - Comments explain the two defensive guards no test can reach: the fast path's empty-queue check and the rejection handler on each authorization. - The mothership and MCP event route tests reset their ready mocks between tests. * test(events): prove the shared pre-open authorization through the stream's writes
1 parent dfe8ecf commit ab588d4

4 files changed

Lines changed: 54 additions & 4 deletions

File tree

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ describe('MCP tool-change event stream', () => {
1717
vi.useFakeTimers()
1818
authMockFns.mockGetSession.mockResolvedValue({ user: { id: 'user-1' } })
1919
permissionsMockFns.mockGetUserEntityPermissions.mockResolvedValue('read')
20+
mcpPubsubMockFns.mockReady.mockReset().mockResolvedValue(undefined)
2021
})
2122

2223
afterEach(() => {

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ describe('Mothership owner-scoped event stream', () => {
6464
permissionsMockFns.mockGetUserEntityPermissions.mockResolvedValue('read')
6565
authorize.mockResolvedValue({ organizationId: 'org-1', userId: 'user-1', role: 'member' })
6666
subscribe.mockReturnValue(unsubscribe)
67+
mothershipChatStatusMockFns.mockReady.mockReset().mockResolvedValue(undefined)
6768
})
6869
afterEach(() => {
6970
vi.useRealTimers()

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

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -291,6 +291,44 @@ describe('createSSEStream', () => {
291291
])
292292
})
293293

294+
it('writes every event queued before it opened once one authorization passes', async () => {
295+
let live: () => void = () => {}
296+
let publish: (eventName: string, data: Record<string, unknown>) => void = () => {}
297+
const authorizations: Array<() => void> = []
298+
const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), {
299+
label: 'test',
300+
revalidate: () =>
301+
new Promise<void>((resolve) => {
302+
authorizations.push(resolve)
303+
}),
304+
subscriptions: [
305+
{
306+
subscribe: (send) => {
307+
publish = send
308+
return () => {}
309+
},
310+
ready: () =>
311+
new Promise<void>((resolve) => {
312+
live = resolve
313+
}),
314+
},
315+
],
316+
})
317+
const chunks: string[] = []
318+
void collect(response.body as ReadableStream<Uint8Array>, chunks)
319+
for (let n = 1; n <= 10; n += 1) publish('changed', { n })
320+
live()
321+
await vi.advanceTimersByTimeAsync(0)
322+
323+
authorizations[0]()
324+
await vi.advanceTimersByTimeAsync(0)
325+
326+
expect(chunks).toEqual([
327+
OPENED_COMMENT,
328+
...Array.from({ length: 10 }, (_, index) => `event: changed\ndata: {"n":${index + 1}}\n\n`),
329+
])
330+
})
331+
294332
it('clears its timers when it closes before it opens', async () => {
295333
const controller = new AbortController()
296334
createSSEStream(

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

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -170,11 +170,17 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
170170
* waits here while the stream has not opened or an earlier event is still being written.
171171
*/
172172
let writes: Promise<void> = Promise.resolve()
173+
/** Settles when the stream opens, or when it closes first. */
174+
let opening: Promise<void> = Promise.resolve()
175+
/** The one authorization every event that arrived before the stream opened waits for. */
176+
let preOpenAuthorization: Promise<void> | undefined
173177
let opened = false
174178
let pendingEvents = 0
175179
const send = (eventName: string, data: Record<string, unknown>) => {
176180
if (cleaned) return
177181
const payload = `event: ${eventName}\ndata: ${JSON.stringify(data)}\n\n`
182+
// Defensive: the opening and the writes queued before it settle in the same run of
183+
// microtasks, so an event sent from a microtask in between must still queue behind them.
178184
if (opened && pendingEvents === 0 && !config.revalidate) {
179185
enqueue(payload)
180186
return
@@ -184,11 +190,14 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
184190
return
185191
}
186192
pendingEvents += 1
187-
// Authorized as soon as it arrives, so concurrent events share one check, but written only
188-
// after every event before it.
193+
// Authorized as soon as the stream is open, so events in flight together share one check:
194+
// those that arrived before it opened share one that starts when it opens, and later ones
195+
// share whichever check is in flight. Each is written only after every event before it.
189196
const authorized = opened
190197
? revalidate()
191-
: writes.then(() => (cleaned ? undefined : revalidate()))
198+
: (preOpenAuthorization ??= opening.then(() => (cleaned ? undefined : revalidate())))
199+
// Defensive: a rejection is handled once the write chain reaches this event, which can be
200+
// after the authorization settles; this keeps it from surfacing as unhandled meanwhile.
192201
authorized.catch(noop)
193202
writes = writes
194203
.then(() => authorized)
@@ -211,7 +220,7 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
211220
// heartbeat cannot open it first: the deadline is shorter than the heartbeat interval.
212221
// Closing settles it too, so a stream closed before it opened is not kept alive by a
213222
// subscription that never becomes ready.
214-
writes = Promise.race([
223+
opening = Promise.race([
215224
Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())),
216225
new Promise<void>((resolve) => {
217226
const deadline = setTimeout(resolve, OPEN_DEADLINE_MS)
@@ -227,6 +236,7 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
227236
},
228237
() => close('subscription_failed')
229238
)
239+
writes = opening
230240
for (const subscription of config.subscriptions) {
231241
teardowns.push(subscription.subscribe(send))
232242
}

0 commit comments

Comments
 (0)