|
| 1 | +/** |
| 2 | + * One Chat open in the desktop app and in a plain web tab at once. Both tail the same run, so a |
| 3 | + * desktop-only tool call (the agent browser, the terminal) reaches both. Runs against real |
| 4 | + * PostgreSQL and Redis: the web tab is the production stream tool-event path and client executor, |
| 5 | + * the desktop speaks its production protocol (claim through the authorize route, then report |
| 6 | + * through the confirm route), and the agent's answer is what the server-side waiter resolves. |
| 7 | + */ |
| 8 | +import { isCurrentBrowserToolName } from '@sim/browser-protocol' |
| 9 | +import { authMock, authMockFns } from '@sim/testing/mocks/auth.mock' |
| 10 | +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' |
| 11 | + |
| 12 | +const { redisUrl, inheritedEnv } = await vi.hoisted(async () => { |
| 13 | + const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') |
| 14 | + const url = readTestRedisUrl() |
| 15 | + const inheritedEnv = { REDIS_URL: process.env.REDIS_URL } |
| 16 | + /** The real Redis module and the confirmation channel read this at import. */ |
| 17 | + process.env.REDIS_URL = url |
| 18 | + return { redisUrl: url, inheritedEnv } |
| 19 | +}) |
| 20 | + |
| 21 | +vi.mock('@/lib/auth', () => authMock) |
| 22 | + |
| 23 | +import { db } from '@sim/db' |
| 24 | +import { copilotAsyncToolCalls, copilotChats, copilotRuns, user, workspace } from '@sim/db/schema' |
| 25 | +import { sleep } from '@sim/utils/helpers' |
| 26 | +import { generateId } from '@sim/utils/id' |
| 27 | +import { eq } from 'drizzle-orm' |
| 28 | +import { NextRequest } from 'next/server' |
| 29 | +import { closeRedisConnection } from '@/lib/core/config/redis' |
| 30 | +import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' |
| 31 | +import type { PersistedStreamEventEnvelope } from '@/lib/mothership/request/session/contract' |
| 32 | +import { waitForClientToolCompletion } from '@/lib/mothership/request/tools/client' |
| 33 | +import { executeBrowserToolOnClient } from '@/lib/mothership/tools/client/browser-tool-execution' |
| 34 | +import { executeTerminalToolOnClient } from '@/lib/mothership/tools/client/terminal-tool-execution' |
| 35 | +import { POST as confirmPOST } from '@/app/api/copilot/confirm/route' |
| 36 | +import { POST as authorizePOST } from '@/app/api/desktop/tool/authorize/route' |
| 37 | +import { dispatchStreamEvent } from '@/app/workspace/[workspaceId]/home/hooks/stream/dispatch-stream-event' |
| 38 | +import { createStreamLoopContext } from '@/app/workspace/[workspaceId]/home/hooks/stream/stream-context' |
| 39 | +import { makeStreamLoopDeps } from '@/app/workspace/[workspaceId]/home/hooks/stream/stream-test-helpers' |
| 40 | + |
| 41 | +const APP_ORIGIN = 'http://localhost:3000' |
| 42 | +const ROUTES: Record<string, (request: NextRequest) => Promise<Response>> = { |
| 43 | + '/api/copilot/confirm': async (request) => confirmPOST(request, {}), |
| 44 | + '/api/desktop/tool/authorize': async (request) => authorizePOST(request, {}), |
| 45 | +} |
| 46 | + |
| 47 | +/** Every route request a client has started, so a test can wait for fire-and-forget reports. */ |
| 48 | +const inFlight: Promise<Response>[] = [] |
| 49 | + |
| 50 | +function callRoute(path: string, { method, headers, body }: RequestInit): Promise<Response> { |
| 51 | + const handler = ROUTES[path] |
| 52 | + if (!handler) throw new Error(`No route fixture for ${path}`) |
| 53 | + const response = handler(new NextRequest(new URL(path, APP_ORIGIN), { method, headers, body })) |
| 54 | + inFlight.push(response) |
| 55 | + return response |
| 56 | +} |
| 57 | + |
| 58 | +/** The desktop app's server protocol: claim the pending call, run it, report its result. */ |
| 59 | +async function desktopExecutes(toolCallId: string, result: Record<string, unknown>) { |
| 60 | + const authorization = await callRoute('/api/desktop/tool/authorize', { |
| 61 | + method: 'POST', |
| 62 | + headers: { 'Content-Type': 'application/json' }, |
| 63 | + body: JSON.stringify({ toolCallId }), |
| 64 | + }) |
| 65 | + if (!authorization.ok) return { authorized: authorization.status, confirmed: null } |
| 66 | + const confirmation = await callRoute('/api/copilot/confirm', { |
| 67 | + method: 'POST', |
| 68 | + headers: { 'Content-Type': 'application/json' }, |
| 69 | + body: JSON.stringify({ toolCallId, status: 'success', message: 'Done', data: result }), |
| 70 | + }) |
| 71 | + return { authorized: authorization.status, confirmed: confirmation.status } |
| 72 | +} |
| 73 | + |
| 74 | +/** |
| 75 | + * A plain web tab receiving the call frame on its live tail: the production stream handler with |
| 76 | + * the executors `useChat` wires into it. |
| 77 | + */ |
| 78 | +async function webTabReceives( |
| 79 | + chatId: string, |
| 80 | + toolCallId: string, |
| 81 | + toolName: string, |
| 82 | + args: Record<string, unknown> |
| 83 | +) { |
| 84 | + const ctx = createStreamLoopContext( |
| 85 | + makeStreamLoopDeps({ |
| 86 | + startClientBrowserTool: (id, name, toolArgs, eventTs) => { |
| 87 | + if (isCurrentBrowserToolName(name)) |
| 88 | + executeBrowserToolOnClient(id, name, toolArgs, chatId, eventTs) |
| 89 | + }, |
| 90 | + startClientTerminalTool: (id, _name, toolArgs, eventTs) => |
| 91 | + executeTerminalToolOnClient(id, toolArgs, chatId, eventTs), |
| 92 | + }) |
| 93 | + ) |
| 94 | + const envelope: PersistedStreamEventEnvelope = { |
| 95 | + type: 'tool', |
| 96 | + v: 1, |
| 97 | + seq: 1, |
| 98 | + ts: new Date().toISOString(), |
| 99 | + stream: { streamId: generateId(), cursor: '1' }, |
| 100 | + payload: { |
| 101 | + phase: 'call', |
| 102 | + executor: 'client', |
| 103 | + mode: 'async', |
| 104 | + toolCallId, |
| 105 | + toolName, |
| 106 | + arguments: args, |
| 107 | + status: 'executing', |
| 108 | + }, |
| 109 | + } |
| 110 | + dispatchStreamEvent(ctx, envelope) |
| 111 | + // The client executors are fire-and-forget: let any report they start reach the server. |
| 112 | + await sleep(250) |
| 113 | + await Promise.allSettled(inFlight) |
| 114 | +} |
| 115 | + |
| 116 | +afterAll(async () => { |
| 117 | + const channels = globalThis as typeof globalThis & { |
| 118 | + _toolConfirmationChannel?: { dispose(): void } |
| 119 | + } |
| 120 | + channels._toolConfirmationChannel?.dispose() |
| 121 | + channels._toolConfirmationChannel = undefined |
| 122 | + await closeRedisConnection() |
| 123 | + for (const [key, value] of Object.entries(inheritedEnv)) { |
| 124 | + if (value === undefined) delete process.env[key] |
| 125 | + else process.env[key] = value |
| 126 | + } |
| 127 | +}) |
| 128 | + |
| 129 | +describe.runIf(Boolean(redisUrl))('a desktop tool call watched by a web tab', () => { |
| 130 | + const userId = generateId() |
| 131 | + const workspaceId = generateId() |
| 132 | + const chatId = generateId() |
| 133 | + const runId = generateId() |
| 134 | + |
| 135 | + beforeAll(async () => { |
| 136 | + const now = new Date() |
| 137 | + await db.insert(user).values({ |
| 138 | + id: userId, |
| 139 | + name: 'Desktop executor fixture', |
| 140 | + email: `${userId}@desktop-executor.test`, |
| 141 | + emailVerified: true, |
| 142 | + createdAt: now, |
| 143 | + updatedAt: now, |
| 144 | + }) |
| 145 | + await db.insert(workspace).values({ |
| 146 | + id: workspaceId, |
| 147 | + name: 'Desktop executor fixture', |
| 148 | + ownerId: userId, |
| 149 | + billedAccountUserId: userId, |
| 150 | + }) |
| 151 | + await db.insert(copilotChats).values({ |
| 152 | + id: chatId, |
| 153 | + userId, |
| 154 | + workspaceId, |
| 155 | + type: 'mothership', |
| 156 | + conversationId: generateId(), |
| 157 | + }) |
| 158 | + await db.insert(copilotRuns).values({ |
| 159 | + id: runId, |
| 160 | + executionId: generateId(), |
| 161 | + chatId, |
| 162 | + userId, |
| 163 | + workspaceId, |
| 164 | + streamId: generateId(), |
| 165 | + toolExecutionVersion: SIM_TOOL_EXECUTION_VERSION, |
| 166 | + status: 'paused_waiting_for_tool', |
| 167 | + requestContext: { source: 'headless_lifecycle' }, |
| 168 | + }) |
| 169 | + authMockFns.mockGetSession.mockResolvedValue({ |
| 170 | + user: { id: userId, email: `${userId}@desktop-executor.test`, name: 'Desktop executor' }, |
| 171 | + session: { id: generateId(), userId }, |
| 172 | + }) |
| 173 | + vi.stubGlobal('fetch', (input: RequestInfo | URL, init?: RequestInit) => |
| 174 | + callRoute(new URL(String(input), APP_ORIGIN).pathname, init ?? {}) |
| 175 | + ) |
| 176 | + }) |
| 177 | + |
| 178 | + afterAll(async () => { |
| 179 | + await db.delete(copilotChats).where(eq(copilotChats.id, chatId)) |
| 180 | + await db.delete(workspace).where(eq(workspace.id, workspaceId)) |
| 181 | + await db.delete(user).where(eq(user.id, userId)) |
| 182 | + }) |
| 183 | + |
| 184 | + it.each([ |
| 185 | + ['browser_find', { query: 'Sign in' }], |
| 186 | + ['terminal', { operation: 'run', args: { command: 'ls' } }], |
| 187 | + ])( |
| 188 | + 'delivers the desktop result of %s when the web tab receives the call first', |
| 189 | + async (toolName, args) => { |
| 190 | + const toolCallId = generateId() |
| 191 | + await db |
| 192 | + .insert(copilotAsyncToolCalls) |
| 193 | + .values({ runId, toolCallId, toolName, args, status: 'pending' }) |
| 194 | + const agentAnswer = waitForClientToolCompletion({ |
| 195 | + toolCallId, |
| 196 | + runId, |
| 197 | + userId, |
| 198 | + timeoutMs: 10_000, |
| 199 | + }) |
| 200 | + |
| 201 | + await webTabReceives(chatId, toolCallId, toolName, args) |
| 202 | + const desktop = await desktopExecutes(toolCallId, { matches: 1 }) |
| 203 | + expect(desktop).toEqual({ authorized: 200, confirmed: 200 }) |
| 204 | + expect(await agentAnswer).toMatchObject({ status: 'success' }) |
| 205 | + const [row] = await db |
| 206 | + .select({ status: copilotAsyncToolCalls.status }) |
| 207 | + .from(copilotAsyncToolCalls) |
| 208 | + .where(eq(copilotAsyncToolCalls.toolCallId, toolCallId)) |
| 209 | + expect(row.status).toBe('completed') |
| 210 | + } |
| 211 | + ) |
| 212 | +}) |
0 commit comments