Skip to content

Commit a13b977

Browse files
committed
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.
1 parent 5fd29f0 commit a13b977

4 files changed

Lines changed: 193 additions & 16 deletions

File tree

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
import { authMockFns, permissionsMock, permissionsMockFns } from '@sim/testing'
2+
import { mcpPubsubMock, mcpPubsubMockFns } from '@sim/testing/mocks/mcp-pubsub.mock'
3+
import { NextRequest } from 'next/server'
4+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
5+
import { OPENED_COMMENT } from '@/lib/events/sse-endpoint'
6+
7+
vi.mock('@/lib/mcp/pubsub', () => mcpPubsubMock)
8+
vi.mock('@/lib/mcp/connection-manager', () => ({
9+
mcpConnectionManager: { subscribe: () => () => {} },
10+
}))
11+
vi.mock('@/lib/workspaces/permissions/utils', () => permissionsMock)
12+
13+
import { GET } from '@/app/api/mcp/events/route'
14+
15+
describe('MCP tool-change event stream', () => {
16+
beforeEach(() => {
17+
vi.useFakeTimers()
18+
authMockFns.mockGetSession.mockResolvedValue({ user: { id: 'user-1' } })
19+
permissionsMockFns.mockGetUserEntityPermissions.mockResolvedValue('read')
20+
})
21+
22+
afterEach(() => {
23+
vi.useRealTimers()
24+
})
25+
26+
it('opens once tool-change events reach this process', async () => {
27+
let live: () => void = () => {}
28+
const ready = new Promise<void>((resolve) => {
29+
live = resolve
30+
})
31+
mcpPubsubMockFns.mockReady.mockReturnValue(ready)
32+
const response = await GET(new NextRequest('http://localhost/api/mcp/events?workspaceId=ws-1'))
33+
let first: string | undefined
34+
if (!response.body) throw new Error('The event stream has no body')
35+
void response.body
36+
.getReader()
37+
.read()
38+
.then(({ value }) => {
39+
first = new TextDecoder().decode(value)
40+
})
41+
42+
await vi.advanceTimersByTimeAsync(1_000)
43+
expect(first).toBeUndefined()
44+
expect(mcpPubsubMockFns.mockReady).toHaveBeenCalledTimes(2)
45+
46+
live()
47+
await vi.advanceTimersByTimeAsync(0)
48+
expect(first).toBe(OPENED_COMMENT)
49+
})
50+
})

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

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -174,6 +174,34 @@ describe('Mothership owner-scoped event stream', () => {
174174
expect(chunks).toEqual([])
175175
})
176176

177+
it.each(['workspaceId=ws-1', 'organizationId=org-1'])(
178+
'opens the %s stream once chat status events reach this process',
179+
async (query) => {
180+
let live: () => void = () => {}
181+
mothershipChatStatusMockFns.mockReady.mockReturnValueOnce(
182+
new Promise<void>((resolve) => {
183+
live = resolve
184+
})
185+
)
186+
const response = await GET(request(query))
187+
let first: string | undefined
188+
if (!response.body) throw new Error('The event stream has no body')
189+
void response.body
190+
.getReader()
191+
.read()
192+
.then(({ value }) => {
193+
first = new TextDecoder().decode(value)
194+
})
195+
196+
await vi.advanceTimersByTimeAsync(1_000)
197+
expect(first).toBeUndefined()
198+
199+
live()
200+
await vi.advanceTimersByTimeAsync(0)
201+
expect(first).toBe(OPENED_COMMENT)
202+
}
203+
)
204+
177205
it('preserves workspace status events and excludes organization events', async () => {
178206
const abort = new AbortController()
179207
const response = await GET(request('workspaceId=ws-1', abort.signal))

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

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,7 @@
1+
import { setFlagsFromString } from 'node:v8'
2+
import { runInNewContext } from 'node:vm'
13
import { authMockFns, permissionsMock, permissionsMockFns } from '@sim/testing'
4+
import { sleep } from '@sim/utils/helpers'
25
import { NextRequest } from 'next/server'
36
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
47
import {
@@ -248,6 +251,62 @@ describe('createSSEStream', () => {
248251
expect(chunks).toEqual([OPENED_COMMENT])
249252
})
250253

254+
it('writes revalidated events in the order they arrive', async () => {
255+
let live: () => void = () => {}
256+
let authorized: () => void = () => {}
257+
let publish: (eventName: string, data: Record<string, unknown>) => void = () => {}
258+
const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), {
259+
label: 'test',
260+
revalidate: () =>
261+
new Promise<void>((resolve) => {
262+
authorized = resolve
263+
}),
264+
subscriptions: [
265+
{
266+
subscribe: (send) => {
267+
publish = send
268+
return () => {}
269+
},
270+
ready: () =>
271+
new Promise<void>((resolve) => {
272+
live = resolve
273+
}),
274+
},
275+
],
276+
})
277+
const chunks: string[] = []
278+
void collect(response.body as ReadableStream<Uint8Array>, chunks)
279+
280+
publish('changed', { n: 1 })
281+
live()
282+
await vi.advanceTimersByTimeAsync(0)
283+
publish('changed', { n: 2 })
284+
authorized()
285+
await vi.advanceTimersByTimeAsync(0)
286+
287+
expect(chunks).toEqual([
288+
OPENED_COMMENT,
289+
'event: changed\ndata: {"n":1}\n\n',
290+
'event: changed\ndata: {"n":2}\n\n',
291+
])
292+
})
293+
294+
it('clears its timers when it closes before it opens', async () => {
295+
const controller = new AbortController()
296+
createSSEStream(
297+
new NextRequest(new URL('https://sim.test/api/test/stream'), { signal: controller.signal }),
298+
{
299+
label: 'test',
300+
subscriptions: [{ subscribe: () => () => {}, ready: () => new Promise<void>(() => {}) }],
301+
}
302+
)
303+
expect(vi.getTimerCount()).toBe(2)
304+
305+
controller.abort()
306+
307+
expect(vi.getTimerCount()).toBe(0)
308+
})
309+
251310
it('delivers a revalidated event as soon as it is published', async () => {
252311
let publish: (eventName: string, data: Record<string, unknown>) => void = () => {}
253312
const response = createSSEStream(new NextRequest(new URL('https://sim.test/api/test/stream')), {
@@ -272,3 +331,29 @@ describe('createSSEStream', () => {
272331
expect(chunks).toEqual([OPENED_COMMENT, 'event: inbox_changed\ndata: {"reason":"call"}\n\n'])
273332
})
274333
})
334+
335+
describe('createSSEStream retention', () => {
336+
/** Opens a stream and closes it before it opens; only a weak reference to it escapes. */
337+
function closedBeforeOpening(ready: Promise<void>): WeakRef<object> {
338+
const controller = new AbortController()
339+
const response = createSSEStream(
340+
new NextRequest(new URL('https://sim.test/api/test/stream'), { signal: controller.signal }),
341+
{ label: 'test', subscriptions: [{ subscribe: () => () => {}, ready: () => ready }] }
342+
)
343+
controller.abort()
344+
return new WeakRef(response.body as object)
345+
}
346+
347+
it('does not stay reachable from a subscription that never becomes ready', async () => {
348+
setFlagsFromString('--expose_gc')
349+
const collectGarbage = runInNewContext('gc') as () => void
350+
const neverReady = new Promise<void>(() => {})
351+
352+
const stream = closedBeforeOpening(neverReady)
353+
await sleep(0)
354+
collectGarbage()
355+
356+
expect(stream.deref()).toBeUndefined()
357+
expect(neverReady).toBeInstanceOf(Promise)
358+
})
359+
})

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

Lines changed: 30 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
import type { SessionPrincipal } from '@sim/auth/principal'
99
import { createLogger } from '@sim/logger'
1010
import { getErrorMessage } from '@sim/utils/errors'
11+
import { noop } from '@sim/utils/helpers'
1112
import { randomFloat } from '@sim/utils/random'
1213
import type { NextRequest } from 'next/server'
1314
import { getSession } from '@/lib/auth'
@@ -164,14 +165,17 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
164165
})
165166
return authorization
166167
}
167-
/** Settles when the stream announces itself; events wait for it so none opens it early. */
168-
let opening: Promise<void> = Promise.resolve()
168+
/**
169+
* The stream's writes in order: the opening, then each event once it is authorized. An event
170+
* waits here while the stream has not opened or an earlier event is still being written.
171+
*/
172+
let writes: Promise<void> = Promise.resolve()
169173
let opened = false
170174
let pendingEvents = 0
171175
const send = (eventName: string, data: Record<string, unknown>) => {
172176
if (cleaned) return
173177
const payload = `event: ${eventName}\ndata: ${JSON.stringify(data)}\n\n`
174-
if (opened && !config.revalidate) {
178+
if (opened && pendingEvents === 0 && !config.revalidate) {
175179
enqueue(payload)
176180
return
177181
}
@@ -180,31 +184,41 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
180184
return
181185
}
182186
pendingEvents += 1
187+
// Authorized as soon as it arrives, so concurrent events share one check, but written only
188+
// after every event before it.
183189
const authorized = opened
184190
? revalidate()
185-
: opening.then(() => (cleaned ? undefined : revalidate()))
186-
void authorized.then(
187-
() => {
188-
pendingEvents -= 1
189-
enqueue(payload)
190-
},
191-
() => {
192-
pendingEvents -= 1
193-
close('authorization_lost')
194-
}
195-
)
191+
: writes.then(() => (cleaned ? undefined : revalidate()))
192+
authorized.catch(noop)
193+
writes = writes
194+
.then(() => authorized)
195+
.then(
196+
() => {
197+
pendingEvents -= 1
198+
enqueue(payload)
199+
},
200+
() => {
201+
pendingEvents -= 1
202+
close('authorization_lost')
203+
}
204+
)
196205
}
197206

198207
try {
199208
// The runtime sends the status and headers with the first body chunk, so the stream writes
200209
// one as soon as it opens. A client reads its state once the stream opens, so it opens once
201210
// every subscription receives events, or at the deadline when one is in an outage. A
202211
// heartbeat cannot open it first: the deadline is shorter than the heartbeat interval.
203-
opening = Promise.race([
212+
// Closing settles it too, so a stream closed before it opened is not kept alive by a
213+
// subscription that never becomes ready.
214+
writes = Promise.race([
204215
Promise.all(config.subscriptions.map((subscription) => subscription.ready?.())),
205216
new Promise<void>((resolve) => {
206217
const deadline = setTimeout(resolve, OPEN_DEADLINE_MS)
207-
teardowns.push(() => clearTimeout(deadline))
218+
teardowns.push(() => {
219+
clearTimeout(deadline)
220+
resolve()
221+
})
208222
}),
209223
]).then(
210224
() => {

0 commit comments

Comments
 (0)