Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/sim/app/api/desktop/inbox/stream/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
50 changes: 50 additions & 0 deletions apps/sim/app/api/mcp/events/route.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>((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)
})
})
2 changes: 2 additions & 0 deletions apps/sim/app/api/mcp/events/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ const mcpEventsHandler = createWorkspaceSSE({
})
})
},
ready: async () => mcpPubSub?.ready(),
},
{
subscribe: (workspaceId, send) => {
Expand All @@ -45,6 +46,7 @@ const mcpEventsHandler = createWorkspaceSSE({
})
})
},
ready: async () => mcpPubSub?.ready(),
},
],
})
Expand Down
36 changes: 34 additions & 2 deletions apps/sim/app/api/mothership/events/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand All @@ -41,13 +41,15 @@ function emit(event: ChatStatusEvent) {
handler(event)
}

/** Every chunk after the stream's opening comment. */
async function collect(body: ReadableStream<Uint8Array>, 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)
}
}

Expand Down Expand Up @@ -154,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<void>((resolve) => {
Expand All @@ -171,11 +174,40 @@ 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<void>((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))
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' })
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/app/api/mothership/events/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ const mothershipEventsHandler = createWorkspaceSSE({
})
})
},
ready: async () => chatPubSub?.ready(),
},
],
})
Expand Down Expand Up @@ -80,6 +81,7 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
timestamp: Date.now(),
})
}) ?? (() => {}),
ready: async () => chatPubSub?.ready(),
},
],
})
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/lib/desktop/application/executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {
} from '@/lib/desktop/executor/constants'
import {
type DesktopInboxChangeReason,
desktopInboxDoorbellReady,
onDesktopInboxDoorbell,
} from '@/lib/desktop/executor/doorbell'
import {
Expand Down Expand Up @@ -201,6 +202,7 @@ export const openDesktopInboxStream = defineAuthorizedCredentialUserUseCase({
onDesktopInboxDoorbell(deviceId, (reason: DesktopInboxChangeReason) =>
send('inbox_changed', { reason })
),
ready: desktopInboxDoorbellReady,
Comment thread
waleedlatif1 marked this conversation as resolved.
}
},
})
Expand Down
5 changes: 5 additions & 0 deletions apps/sim/lib/desktop/executor/doorbell.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,11 @@ export function ringDesktopInbox(deviceId: string, reason: DesktopInboxChangeRea
}
}

/** Settles once this process hears every ring. */
export function desktopInboxDoorbellReady(): Promise<void> {
return channel().ready()
}

/** Subscribes to one device's doorbell; returns the unsubscribe. */
export function onDesktopInboxDoorbell(
deviceId: string,
Expand Down
97 changes: 97 additions & 0 deletions apps/sim/lib/events/pubsub.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
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<EventEmitter & { subscribed: Array<(err: Error | null) => 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<void>): Promise<boolean> {
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('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)
})
})
52 changes: 47 additions & 5 deletions apps/sim/lib/events/pubsub.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,12 @@ const logger = createLogger('PubSub')
export interface PubSubChannel<T> {
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. Stays pending while the channel cannot subscribe, so a caller
* that must not wait indefinitely bounds the wait.
*/
ready(): Promise<void>
dispose(): void
}

Expand All @@ -29,6 +35,12 @@ class RedisPubSubChannel<T> implements PubSubChannel<T> {
private sub: Redis
private handlers = new Set<(event: T) => void>()
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<void> = Promise.resolve()
private markSubscribed: () => void = noop

constructor(
redisUrl: string,
Expand Down Expand Up @@ -56,12 +68,28 @@ class RedisPubSubChannel<T> implements PubSubChannel<T> {
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 {
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. 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)
return
Comment thread
waleedlatif1 marked this conversation as resolved.
}
this.listening = true
logger.info(`Subscribed to ${config.label} channel`)
}
this.markSubscribed()
})
})
this.sub.on('close', () => {
this.closedConnections += 1
if (!this.listening) return
Comment thread
waleedlatif1 marked this conversation as resolved.
this.listening = false
this.awaitSubscription()
})

this.sub.on('message', (channel: string, message: string) => {
Expand Down Expand Up @@ -95,6 +123,16 @@ class RedisPubSubChannel<T> implements PubSubChannel<T> {
}
}

ready(): Promise<void> {
return this.subscribed
}
Comment thread
waleedlatif1 marked this conversation as resolved.

private awaitSubscription(): void {
this.subscribed = new Promise((resolve) => {
this.markSubscribed = resolve
})
}

dispose(): void {
this.disposed = true
this.handlers.clear()
Expand Down Expand Up @@ -130,6 +168,10 @@ class LocalPubSubChannel<T> implements PubSubChannel<T> {
}
}

ready(): Promise<void> {
return Promise.resolve()
}

dispose(): void {
this.emitter.removeAllListeners()
logger.info(`${this.config.label} local pub/sub disposed`)
Expand Down
Loading
Loading