diff --git a/.github/scripts/desktop-live-changes.sh b/.github/scripts/desktop-live-changes.sh new file mode 100755 index 00000000000..dab0e69688a --- /dev/null +++ b/.github/scripts/desktop-live-changes.sh @@ -0,0 +1,24 @@ +#!/usr/bin/env bash +# Prints `changed=false` only when every file a pull request changes is clearly unrelated to the +# live desktop suite, and `changed=true` otherwise, including when the diff cannot be worked out. +# +# Usage: desktop-live-changes.sh +set -u + +unrelated='^(apps/docs/|apps/pii/|apps/sim/content/|packages/(python-sdk|ts-sdk)/)|\.mdx?$|(^|/)LICENSE$' + +base=${1:-} +run() { + echo "changed=true" + echo "Running the live desktop suite: $1" >&2 + exit 0 +} +[ -n "$base" ] || run 'no base commit' +git fetch --quiet --depth=1 origin "$base" || run "could not fetch $base" +names=$(git diff --name-only "$base" HEAD) || run "could not diff against $base" +[ -n "$names" ] || run 'no changed files listed' +if printf '%s\n' "$names" | grep -qvE "$unrelated"; then + run 'a change may affect it' +fi +echo "changed=false" +echo 'Skipping the live desktop suite: every change is unrelated to it' >&2 diff --git a/.github/workflows/test-build.yml b/.github/workflows/test-build.yml index e93dd0201b8..6309d81b51f 100644 --- a/.github/workflows/test-build.yml +++ b/.github/workflows/test-build.yml @@ -442,6 +442,148 @@ jobs: if-no-files-found: ignore retention-days: 7 + # Pull requests skip the live desktop suite only when every change is clearly unrelated to the + # app it drives (docs, other apps, published content). Anything else, and any failure to work + # out the diff, runs it: a pull request that skipped it wrongly would first fail on staging. + desktop-live-changes: + name: Detect desktop tool changes + runs-on: ${{ (vars.CI_PROVIDER == '' || vars.CI_PROVIDER == 'blacksmith') && 'blacksmith-2vcpu-ubuntu-2404' || 'ubuntu-latest' }} + timeout-minutes: 5 + outputs: + changed: ${{ github.event_name != 'pull_request' || steps.diff.outputs.changed != 'false' }} + steps: + - uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6 + if: github.event_name == 'pull_request' + with: + fetch-depth: 2 + - name: Diff against the pull request's base + id: diff + if: github.event_name == 'pull_request' + env: + BASE: ${{ github.event.pull_request.base.sha }} + run: bash .github/scripts/desktop-live-changes.sh "$BASE" >> "$GITHUB_OUTPUT" + + # Desktop tools in the real Electron app against a local app, on its own runner: the + # Electron app, the dev app and its realtime server together outgrow the http-e2e runner. + desktop-live-e2e: + name: Desktop tools against a local app + needs: desktop-live-changes + if: needs.desktop-live-changes.outputs.changed == 'true' + runs-on: ${{ (vars.CI_PROVIDER == '' || vars.CI_PROVIDER == 'blacksmith') && 'blacksmith-16vcpu-ubuntu-2404' || 'ubuntu-latest' }} + timeout-minutes: 30 + services: + postgres: + image: pgvector/pgvector:pg17 + env: + POSTGRES_USER: postgres + POSTGRES_PASSWORD: postgres + POSTGRES_DB: sim_test + ports: + - 5432:5432 + options: >- + --health-cmd "pg_isready -U postgres -d sim_test" + --health-interval 5s + --health-timeout 5s + --health-retries 10 + redis: + image: redis:7-alpine + ports: + - 6379:6379 + options: >- + --health-cmd "redis-cli ping" + --health-interval 5s + --health-timeout 5s + --health-retries 10 + env: + DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5432/sim_test + BETTER_AUTH_SECRET: desktop-live-e2e-ci-secret-at-least-32-characters + ENCRYPTION_KEY: '0000000000000000000000000000000000000000000000000000000000000000' + + steps: + - name: Checkout code + uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6 + + - name: Setup workspace + uses: ./.github/actions/setup-workspace + with: + provider: ${{ vars.CI_PROVIDER }} + + - name: Provision the database through migrations + working-directory: packages/db + run: bun run db:migrate + + # Chat switches, Stop, sign-out, approval and the flag-off foreground round trip. The spec + # runs the recording proxy (the app's public origin) and the stand-in worker. + - name: Verify desktop tools in the Electron app against a local app + env: + NEXT_PUBLIC_APP_URL: http://127.0.0.1:3020 + BETTER_AUTH_URL: http://127.0.0.1:3020 + REDIS_URL: redis://127.0.0.1:6379 + SIM_AGENT_API_URL: http://127.0.0.1:3022 + NEXT_PUBLIC_SOCKET_URL: http://127.0.0.1:3023 + SOCKET_SERVER_URL: http://127.0.0.1:3023 + NEXT_PUBLIC_FORCE_HOSTED: 'false' + COPILOT_API_KEY: desktop-tools-e2e-ci-local-copilot-key + COPILOT_TOOL_PERMISSIONS_ENABLED: 'true' + MOTHERSHIP_SIM_TRANSPORT: direct + INTERNAL_API_SECRET: desktop-tools-e2e-ci-local-secret-at-least-32-characters + DISABLE_TELEMETRY: 'true' + NEXT_TELEMETRY_DISABLED: '1' + READY_TIMEOUT_SECONDS: 300 + run: | + report_dir="$RUNNER_TEMP/e2e" + mkdir -p "$report_dir" + sudo apt-get update -q + sudo apt-get install -yq xvfb libgtk-3-0t64 libnss3 libasound2t64 libgbm1 libxss1 \ + libxtst6 libatk-bridge2.0-0t64 libxkbcommon0 > /dev/null + # Bundle only: `bun run build` also fetches the macOS node-pty prebuilds for packaging, + # which a Linux run does not use. + (cd apps/desktop && bun run scripts/build.ts) + # Each app runs in its own session under an E2E_APP tag, and stop-session.sh returns once + # every process it started has exited. + realtime_tag="desktop-realtime-$GITHUB_RUN_ID-$GITHUB_RUN_ATTEMPT-$$" + server_tag="desktop-tools-$GITHUB_RUN_ID-$GITHUB_RUN_ATTEMPT-$$" + (cd apps/realtime && PORT=3023 SIM_DB_ROLE=realtime ALLOWED_ORIGINS="$NEXT_PUBLIC_APP_URL" \ + E2E_APP="$realtime_tag" exec setsid bun src/index.ts > "$report_dir/desktop-tools-realtime.log" 2>&1) & + realtime_pid=$! + (cd apps/sim && E2E_APP="$server_tag" exec setsid node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 \ + --port 3021 > "$report_dir/desktop-tools-next.log" 2>&1) & + server_pid=$! + finish() { + status=$? + bash "$GITHUB_WORKSPACE/.github/scripts/stop-session.sh" "$server_pid" "$server_tag" || status=1 + bash "$GITHUB_WORKSPACE/.github/scripts/stop-session.sh" "$realtime_pid" "$realtime_tag" || status=1 + wait "$server_pid" "$realtime_pid" 2>/dev/null || true + exit "$status" + } + trap finish EXIT + started=$SECONDS + until curl --fail --silent --max-time 10 http://127.0.0.1:3021/api/health > /dev/null && + curl --fail --silent --max-time 10 http://127.0.0.1:3023/health > /dev/null; do + kill -0 "$server_pid" 2>/dev/null || { tail -n 200 "$report_dir/desktop-tools-next.log"; exit 1; } + kill -0 "$realtime_pid" 2>/dev/null || { tail -n 200 "$report_dir/desktop-tools-realtime.log"; exit 1; } + [ $((SECONDS - started)) -lt "$READY_TIMEOUT_SECONDS" ] || { echo '::error::Local app did not become ready'; exit 1; } + sleep 2 + done + cd apps/desktop + SIM_DESKTOP_E2E_SIM_URL=http://127.0.0.1:3021 \ + SIM_DESKTOP_E2E_PROXY_PORT=3020 \ + SIM_DESKTOP_E2E_AGENT_PORT=3022 \ + SIM_DESKTOP_E2E_DATABASE_URL="$DATABASE_URL" \ + SIM_DESKTOP_E2E_REDIS_URL="$REDIS_URL" \ + SIM_DESKTOP_E2E_AUTH_SECRET="$BETTER_AUTH_SECRET" \ + xvfb-run -a -s '-screen 0 1920x1200x24' bunx playwright test e2e/desktop-tools-live-sim.spec.ts \ + --output "$report_dir/desktop-tools-results" --retries=0 + + - name: Upload Electron E2E results and server logs + if: failure() + uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 + with: + name: desktop-live-e2e-results + path: ${{ runner.temp }}/e2e/ + if-no-files-found: ignore + retention-days: 7 + test-build: name: Lint and Test runs-on: ${{ (vars.CI_PROVIDER == '' || vars.CI_PROVIDER == 'blacksmith') && 'blacksmith-8vcpu-ubuntu-2404' || 'ubuntu-latest' }} diff --git a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts new file mode 100644 index 00000000000..e53ce45d7ae --- /dev/null +++ b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts @@ -0,0 +1,725 @@ +import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { dirname, join } from 'node:path' +import { fileURLToPath } from 'node:url' +import { + type ElectronApplication, + _electron as electron, + expect, + type Locator, + type Page, + test, +} from '@playwright/test' +import type { SimDesktopApi } from '@sim/desktop-bridge' +import { generateId } from '@sim/utils/id' +import { toRecord } from '@sim/utils/object' +import { + type LiveSimConfig, + liveSimConfig, + RedisMonitor, + ScriptedAgent, + type SeededUser, + SimDatabase, + SimProxy, +} from './fixtures/live-sim' + +/** + * Desktop tools in the real Electron app against a real local Sim: the renderer is Sim's own chat + * view, every claim and result crosses Sim's routes into PostgreSQL and Redis, and only the + * model's decisions are scripted (a stand-in worker at `SIM_AGENT_API_URL`). Each test checks + * what the user or the model would observe: the result the model is resumed with, what landed in + * the workspace, what Sim persisted, and which requests reached Sim. + */ + +const DESKTOP_DIR = fileURLToPath(new URL('..', import.meta.url)) +const config = liveSimConfig() +const PICKUP_GRACE_MS = 15_000 +/** How long a held request may take to arrive once the step that sends it ran. */ +const ARRIVAL_MS = 60_000 +/** First requests to a route compile it, which takes minutes on a cold dev app. */ +const COMPILE_MS = 300_000 + +type DesktopWindow = typeof globalThis & { simDesktop: SimDesktopApi } + +test.describe('desktop tools against a live Sim', () => { + test.skip(typeof config === 'string', typeof config === 'string' ? config : '') + test.describe.configure({ timeout: 360_000 }) + + let sim: LiveSimConfig + let proxy: SimProxy + let agent: ScriptedAgent + let db: SimDatabase + let app: ElectronApplication | undefined + let scratch: string + /** Errors the current window reported, shown when a click is blocked. */ + let pageErrors: string[] = [] + let testStartedAt = 0 + + test.beforeAll(async () => { + test.setTimeout(1_500_000) + if (typeof config === 'string') throw new Error(config) + sim = config + proxy = new SimProxy(sim) + agent = new ScriptedAgent(sim.agentPort) + db = new SimDatabase(sim) + await Promise.all([proxy.start(), agent.start()]) + await test.step('warm up the routes the tests use', warmUp) + }) + + test.afterAll(async () => { + await Promise.all([proxy?.stop(), agent?.stop(), db?.close()]) + }) + + test.beforeEach(() => { + testStartedAt = Date.now() + scratch = mkdtempSync(join(tmpdir(), 'sim-desktop-tools-live-')) + }) + + test.afterEach(async () => { + const testInfo = test.info() + if (testInfo.status !== testInfo.expectedStatus) { + const page = app?.windows()[0] + await page?.screenshot({ path: testInfo.outputPath('failure.png') }).catch(() => {}) + writeFileSync( + testInfo.outputPath('diagnostics.json'), + JSON.stringify( + { + url: page?.url(), + composer: await page + ?.getByRole('textbox') + .last() + .inputValue({ timeout: 2_000 }) + .catch(() => null), + pageErrors, + requests: proxy + .seen(testStartedAt) + .filter((entry) => !entry.path.startsWith('/_next/')), + }, + null, + 2 + ) + ) + } + proxy.clearHolds() + await app?.close().catch(() => {}) + app = undefined + proxy.rewriteChatBody(undefined) + rmSync(scratch, { recursive: true, force: true }) + }) + + /** + * A dev app compiles each route on its first request, one at a time, so a route first reached + * mid-test stalls every other request (Electron gives a tool authorization 8 s). The warm-up runs + * the tests' own flows once: a read and an import, a chat switch during a live turn and back, + * Stop, and the login page. It then waits until every request the app made has been answered, + * so no route the tests reach is left to compile. It ends with a marker request + * (`/api/health?e2e=warm-up-done`) so the dev server's log shows anything compiled after it. + */ + async function warmUp(): Promise { + scratch = mkdtempSync(join(tmpdir(), 'sim-desktop-tools-warm-')) + try { + const user = await db.seedUser(['Warm chat', 'Warm other chat']) + const headers = { + 'Content-Type': 'application/json', + Cookie: `better-auth.session_token=${user.cookie}`, + Origin: proxy.origin, + } + const compile = (path: string, method: 'GET' | 'POST' | 'PUT' = 'GET') => + fetch(new URL(path, sim.upstream), { + method, + headers, + body: method === 'GET' ? undefined : '{}', + redirect: 'manual', + signal: AbortSignal.timeout(COMPILE_MS), + }).then((response) => response.arrayBuffer()) + // Routes the tests reach that a warm-up turn alone would not: Stop's, registration, and the + // ones a running app loads in the background. + for (const path of [ + chatPath(user, 'Warm chat'), + '/login', + `/api/mothership/chats/${user.chats['Warm other chat']}`, + `/api/mothership/chat/stream?chatId=${user.chats['Warm chat']}`, + `/api/workspaces/${user.workspaceId}/files/folders`, + '/api/copilot/chats', + '/api/users/me/settings', + '/api/auth/oauth/connections', + ]) + await compile(path) + for (const path of [ + '/api/mothership/chat/stop', + '/api/mothership/chat/abort', + '/api/desktop/devices', + '/api/desktop/tool/authorize', + '/api/copilot/confirm', + '/api/files/uploads', + '/api/files/uploads/warm-up/parts', + '/api/files/uploads/warm-up/complete', + `/api/workspaces/${user.workspaceId}/files/folders`, + ]) + await compile(path, 'POST') + await compile('/api/v2/uploads/warm-up', 'PUT') + + const file = writeFile(join(scratch, 'warm.txt'), 'warm') + const folder = join(scratch, 'Warm folder') + writeFile(join(folder, 'warm.txt'), 'warm') + let proceed!: () => void + const proceeding = new Promise((resolve) => { + proceed = resolve + }) + agent.script( + '[warm-up]', + async (turn) => { + turn.text('Warming up.') + await proceeding + turn.toolCall({ toolName: 'read_local_file', args: { path: file } }) + turn.toolCall({ + toolName: 'import_local_files', + args: { path: folder, targetWorkspaceId: user.workspaceId }, + }) + turn.pause() + }, + (_resume, turn) => turn.complete('Warmed up.') + ) + agent.script('[warm-stop]', async (turn) => { + turn.text('Stopping soon.') + // The leg stays open until Stop ends it. + await turn.closed + }) + const page = await openApp(user, 'Warm chat', COMPILE_MS) + await send(page, '[warm-up] read and import', COMPILE_MS) + await expect(page.getByText('Warming up.')).toBeVisible({ timeout: COMPILE_MS }) + // Leaving and reopening a chat while its turn runs re-attaches to the turn's stream. + await openChat(page, user, 'Warm other chat', COMPILE_MS) + await openChat(page, user, 'Warm chat', COMPILE_MS) + proceed() + await expect(page.getByText('Warmed up.')).toBeVisible({ timeout: 2 * COMPILE_MS }) + await send(page, '[warm-stop] wait for Stop', COMPILE_MS) + await expect(page.getByText('Stopping soon.')).toBeVisible({ timeout: COMPILE_MS }) + await click(page, page.getByRole('button', { name: 'Stop generation' }), COMPILE_MS) + await expect(page.getByRole('button', { name: 'Stop generation' })).toBeHidden({ + timeout: COMPILE_MS, + }) + await proxy.settled(COMPILE_MS) + await compile('/api/health?e2e=warm-up-done') + } finally { + await app?.close().catch(() => {}) + app = undefined + rmSync(scratch, { recursive: true, force: true }) + } + } + + const chatPath = (user: SeededUser, title: string) => + `/workspace/${user.workspaceId}/chat/${user.chats[title]}` + + /** Launches the app signed in as `user`, showing the chat titled `title`. */ + async function openApp(user: SeededUser, title: string, timeout = 120_000): Promise { + app = await electron.launch({ + args: [process.env.SIM_DESKTOP_E2E_MAIN ?? '.'], + cwd: DESKTOP_DIR, + env: { + ...process.env, + SIM_DESKTOP_ORIGIN: proxy.origin, + SIM_DESKTOP_USER_DATA: join(scratch, 'profile'), + }, + }) + const page = await app.firstWindow({ timeout }) + pageErrors = [] + page.on('pageerror', (error) => pageErrors.push(error.message)) + page.on('console', (message) => { + if (message.type() === 'error') pageErrors.push(message.text()) + }) + const signIn = new URL('/__e2e/sign-in', proxy.origin) + signIn.searchParams.set('cookie', user.cookie) + signIn.searchParams.set('to', chatPath(user, title)) + await page.goto(signIn.toString(), { waitUntil: 'commit', timeout }) + await expect(composer(page)).toBeVisible({ timeout }) + return page + } + + const composer = (page: Page) => page.getByRole('textbox').last() + + /** The dev app's error overlay, if it is up, with the errors the window reported. */ + async function devOverlay(page: Page): Promise<{ text: string; consoleOnly: boolean } | null> { + const overlay = page.locator('nextjs-portal [data-nextjs-dialog]') + if ((await overlay.count()) === 0) return null + const text = await overlay + .first() + .innerText() + .catch(() => '') + return { + text: `${text}\n${pageErrors.join('\n')}`, + consoleOnly: /^\s*Console Error/.test(text), + } + } + + /** + * Clicks `target`. When the dev app's error overlay intercepts the click, a console-error notice + * (dev-only chrome over a logged error) is dismissed and the click retried; a runtime error fails + * at once with its message instead of waiting out a blocked click. + */ + async function click(page: Page, target: Locator, timeout = 15_000): Promise { + for (let attempt = 0; ; attempt++) { + try { + await target.click({ timeout: attempt === 0 ? Math.min(timeout, 5_000) : timeout }) + return + } catch (error) { + const overlay = await devOverlay(page) + if (overlay && !overlay.consoleOnly) + throw new Error(`Next.js error overlay: ${overlay.text}`) + if (overlay) await page.keyboard.press('Escape') + else if (attempt > 0) throw error + } + } + } + + /** + * Sends `message` and waits for its turn to reach Sim. Before the page hydrates, typing or the + * click can be lost: the Send button is missing or does nothing and the composer is not emptied, + * and only then is the message typed and sent again. Once the UI takes the message (the + * composer empties after the click), a turn that never reaches Sim is a lost send and fails. + */ + async function send(page: Page, message: string, timeout = 60_000): Promise { + const since = Date.now() + const deadline = since + timeout + const sent = () => + proxy + .seen(since) + .some((entry) => entry.method === 'POST' && entry.path === '/api/mothership/chat') + while (!sent()) { + if (Date.now() > deadline) + throw new Error( + `The UI never took the message: ${message} (errors: ${pageErrors.join(' | ')})` + ) + if ((await composer(page).inputValue()) !== message) await composer(page).fill(message) + const clicked = await click(page, page.getByRole('button', { name: 'Send message' }), 5_000) + .then(() => true) + .catch((error: unknown) => { + if (String(error).includes('Next.js error overlay')) throw error + return false + }) + if (!clicked) continue + const taken = await expect + .poll(() => composer(page).inputValue(), { timeout: 5_000 }) + .toBe('') + .then( + () => true, + () => false + ) + if (!taken) continue + await expect + .poll(sent, { + timeout: Math.max(deadline - Date.now(), 30_000), + message: `The UI took the message but its turn never reached Sim: ${message} (errors: ${pageErrors.join(' | ')})`, + }) + .toBe(true) + } + } + + /** Switches chats in-app, the way the sidebar does, without reloading the page. */ + async function openChat( + page: Page, + user: SeededUser, + title: string, + timeout = 30_000 + ): Promise { + await click(page, page.getByRole('link', { name: title }).first(), timeout) + await expect(page).toHaveURL(new RegExp(`${user.chats[title]}$`), { timeout }) + } + + /** The chat's first call as `status: error`, so a wrong terminal state says why. */ + async function callState(chatId: string): Promise { + const [call] = await db.toolCalls(chatId) + return call ? `${call.status}: ${call.error ?? ''}` : 'no call' + } + + function writeFile(path: string, contents: string): string { + mkdirSync(dirname(path), { recursive: true }) + writeFileSync(path, contents) + return path + } + + /** A folder to import: a file, then a subfolder holding another file. */ + function importSource(): string { + const root = join(scratch, 'Reports') + writeFile(join(root, 'a.txt'), 'first file') + writeFile(join(root, 'later', 'b.txt'), 'second file') + return root + } + + /** The first request of a file's upload session. */ + const isUploadStart = (method: string, path: string) => + method === 'POST' && path === '/api/files/uploads' + const isFolderCreate = (method: string, path: string) => + method === 'POST' && /^\/api\/workspaces\/[^/]+\/files\/folders$/.test(path) + const isDesktopClaim = (method: string, path: string) => + method === 'POST' && path === '/api/desktop/tool/authorize' + /** The chat turn's response stream, as the chat view reads it. */ + const isChatStream = (entry: { method: string; path: string }) => + entry.method === 'POST' && entry.path === '/api/mothership/chat' + /** A client tool's report of its own result. */ + const isToolReport = (method: string, path: string) => + method === 'POST' && path === '/api/copilot/confirm' + + test('a browser call issued after the user switched chats fails as not started after the pickup grace', async () => { + const user = await db.seedUser(['Browser chat', 'Other chat']) + let issue!: () => void + const issued = new Promise((resolve) => { + issue = resolve + }) + let callId = '' + let issuedAt = 0 + agent.script('[pickup-grace]', async (turn) => { + turn.text('Checking your browser tabs.') + await issued + callId = turn.toolCall({ toolName: 'browser_list_tabs', args: {} }) + issuedAt = Date.now() + turn.pause() + }) + const page = await openApp(user, 'Browser chat') + await send(page, '[pickup-grace] which tabs are open?') + await expect(page.getByText('Checking your browser tabs.')).toBeVisible({ timeout: 60_000 }) + await openChat(page, user, 'Other chat') + issue() + + await agent.waitForResume(() => Boolean(callId && agent.resultFor(callId)), 60_000) + const result = agent.resultFor(callId) + expect(result?.success).toBe(false) + expect(result?.data).toMatchObject({ notStarted: true }) + const waited = (result?.at ?? 0) - issuedAt + expect(waited).toBeGreaterThanOrEqual(PICKUP_GRACE_MS - 1_000) + expect(waited).toBeLessThan(PICKUP_GRACE_MS + 20_000) + const [call] = await db.toolCalls(user.chats['Browser chat']) + expect(call).toMatchObject({ toolName: 'browser_list_tabs', status: 'failed', claimedBy: null }) + }) + + test('leaving the chat does not end a local file read already under way', async () => { + const user = await db.seedUser(['Read chat', 'Other chat']) + const marker = generateId() + const file = writeFile(join(scratch, 'notes.txt'), `notes from disk ${marker}`) + let callId = '' + agent.script('[leave-read]', (turn) => { + callId = turn.toolCall({ toolName: 'read_local_file', args: { path: file } }) + turn.pause() + }) + const page = await openApp(user, 'Read chat') + // Visit the other chat once, so leaving for it later is a quick client-side switch. + await openChat(page, user, 'Other chat') + await openChat(page, user, 'Read chat') + const claim = proxy.hold(isDesktopClaim) + const since = Date.now() + await send(page, '[leave-read] read my notes') + await claim.arrival(ARRIVAL_MS, 'The read’s claim') + const turnStream = proxy.seen(since).find(isChatStream) + if (!turnStream) throw new Error('The turn’s stream never reached Sim') + await click(page, page.getByRole('link', { name: 'Other chat' }).first()) + // The view let go of the turn's stream, the moment a read tied to that view would end. + // Electron gives the held claim 8 s, so this waits only a few. + await expect.poll(() => turnStream.clientClosedAt, { timeout: 6_000 }).toBeDefined() + claim.release() + await expect(page).toHaveURL(new RegExp(`${user.chats['Other chat']}$`), { timeout: 30_000 }) + + await agent.waitForResume(() => Boolean(callId && agent.resultFor(callId)), 60_000) + const result = agent.resultFor(callId) + expect(result?.success).toBe(true) + expect(JSON.stringify(result?.data)).toContain(marker) + const [call] = await db.toolCalls(user.chats['Read chat']) + expect(call).toMatchObject({ toolName: 'read_local_file', status: 'completed' }) + }) + + test('an import keeps running across chat switches, and Stop from the reopened chat ends it', async () => { + const user = await db.seedUser(['Import chat', 'Other chat']) + const chatId = user.chats['Import chat'] + const source = importSource() + agent.script('[stop-import]', (turn) => { + turn.toolCall({ + toolName: 'import_local_files', + args: { path: source, targetWorkspaceId: user.workspaceId }, + }) + turn.pause() + }) + const page = await openApp(user, 'Import chat') + const firstUpload = proxy.hold(isUploadStart) + await send(page, '[stop-import] import my reports') + await firstUpload.arrival(ARRIVAL_MS, 'The first upload') + + await openChat(page, user, 'Other chat') + await openChat(page, user, 'Import chat') + const laterFolder = proxy.hold(isFolderCreate) + firstUpload.release() + // The import outlived both view changes: its first file landed. + await expect + .poll( + async () => + `${(await db.workspaceFileNames(user.workspaceId)).join(',')} | ${await callState(chatId)}`, + { timeout: 60_000 } + ) + .toMatch(/^a\.txt /) + await laterFolder.arrival(ARRIVAL_MS, 'The next folder') + + const stoppedAt = Date.now() + await click(page, page.getByRole('button', { name: 'Stop generation' })) + // Stop cancels the import's request in flight, which ends the import, and records the call as + // cancelled; the stopped import reports nothing that could contest that record. + await expect.poll(() => laterFolder.isAbandoned, { timeout: 15_000 }).toBe(true) + await expect.poll(() => callState(chatId), { timeout: 30_000 }).toMatch(/^cancelled/) + laterFolder.release() + expect(proxy.seen(stoppedAt).filter((entry) => isToolReport(entry.method, entry.path))).toEqual( + [] + ) + const [call] = await db.toolCalls(chatId) + expect(call).toMatchObject({ toolName: 'import_local_files', status: 'cancelled' }) + expect(await db.workspaceFolderNames(user.workspaceId)).not.toContain('later') + expect(await db.workspaceFileNames(user.workspaceId)).toEqual(['a.txt']) + }) + + test('signing out ends a desktop tool still running', async () => { + const user = await db.seedUser(['Import chat']) + const chatId = user.chats['Import chat'] + const source = importSource() + agent.script('[sign-out-import]', (turn) => { + turn.toolCall({ + toolName: 'import_local_files', + args: { path: source, targetWorkspaceId: user.workspaceId }, + }) + turn.pause() + }) + const page = await openApp(user, 'Import chat') + const firstUpload = proxy.hold(isUploadStart) + await send(page, '[sign-out-import] import my reports') + await firstUpload.arrival(ARRIVAL_MS, 'The first upload') + + // The desktop reloads into the login page once signing out completes, which would end any + // tool; holding the sign-out shows the tool ends at sign-out itself, not at that reload. + const signOut = proxy.hold((method, path) => method === 'POST' && path === '/api/auth/sign-out') + await click(page, page.getByRole('button', { name: user.name })) + await click(page, page.getByRole('menuitem', { name: 'Sign out' })) + await signOut.arrival(ARRIVAL_MS, 'The sign-out') + await expect.poll(() => firstUpload.isAbandoned, { timeout: 15_000 }).toBe(true) + const abandonedAt = Date.now() + signOut.release() + // The reload into the login page replaces the document, so nothing of the import can run after. + await expect(page).toHaveURL(/\/login/, { timeout: 30_000 }) + firstUpload.release() + // Between sign-out and that reload the cancelled import sent nothing more, not even a report. + const toolRequests = proxy + .seen(abandonedAt) + .filter( + (entry) => + isUploadStart(entry.method, entry.path) || + isFolderCreate(entry.method, entry.path) || + isToolReport(entry.method, entry.path) + ) + expect(toolRequests).toEqual([]) + expect(await callState(chatId)).toMatch(/^running/) + expect(await db.workspaceFileNames(user.workspaceId)).toEqual([]) + expect(await db.workspaceFolderNames(user.workspaceId)).not.toContain('later') + }) + + /** Sends `message` and returns the `toolName` call Sim persisted for it, held for approval. */ + async function gatedCall(page: Page, chatId: string, message: string, toolName: string) { + await send(page, message) + await expect + .poll(async () => (await db.toolCalls(chatId)).some((call) => call.toolName === toolName), { + timeout: 60_000, + }) + .toBe(true) + const call = (await db.toolCalls(chatId)).find((entry) => entry.toolName === toolName) + if (!call) throw new Error(`Missing ${toolName} call`) + expect(call.status).toBe('pending') + expect(call.permissionRequestedAt).not.toBeNull() + return call + } + + test("a local read awaiting the user's approval cannot be claimed", async () => { + const user = await db.seedUser(['Approval chat']) + const secret = generateId() + const file = writeFile(join(scratch, 'secret.txt'), `private ${secret}`) + agent.script('[unapproved-read]', (turn) => { + turn.toolCall({ + toolName: 'read_local_file', + args: { path: file }, + status: 'awaiting_approval', + }) + turn.pause() + }) + const page = await openApp(user, 'Approval chat') + const chatId = user.chats['Approval chat'] + const read = await gatedCall( + page, + chatId, + '[unapproved-read] read my secret', + 'read_local_file' + ) + + // A renderer acting on a call the user has not allowed (a replayed event) is refused. + const response = await page.evaluate( + (toolCallId) => + (globalThis as DesktopWindow).simDesktop.localFiles?.({ operation: 'read', toolCallId }), + read.toolCallId + ) + expect(response?.ok).toBe(false) + expect(JSON.stringify(response)).not.toContain(secret) + const [after] = await db.toolCalls(chatId) + expect(after).toMatchObject({ status: 'pending', claimedBy: null }) + }) + + test("a terminal command awaiting the user's approval cannot be claimed", async () => { + const user = await db.seedUser(['Approval chat']) + const command = `touch '${join(scratch, 'terminal-ran')}'` + agent.script('[unapproved-run]', (turn) => { + turn.toolCall({ toolName: 'terminal', args: { operation: 'run', command } }) + turn.pause() + }) + const page = await openApp(user, 'Approval chat') + const chatId = user.chats['Approval chat'] + const run = await gatedCall(page, chatId, '[unapproved-run] run a command', 'terminal') + + // The guard is the claim: Sim must not hand an unapproved command to the desktop, so the call + // stays pending and unclaimed. Whether a handed-over command then runs depends on a terminal + // being open for the chat, which this test does not set up, so it checks the claim only. + const response = await page.evaluate( + ({ toolCallId, command }) => + (globalThis as DesktopWindow).simDesktop.terminal + .executeTool(toolCallId, 'run', { command }, 'unapproved-e2e') + .catch((error: unknown) => ({ ok: false, error: String(error) })), + { toolCallId: run.toolCallId, command } + ) + expect(response).toMatchObject({ ok: false }) + const [after] = await db.toolCalls(chatId) + expect(after).toMatchObject({ status: 'pending', claimedBy: null }) + }) + + test('a claim that reaches Sim after the user pressed Stop is refused', async () => { + const user = await db.seedUser(['Stop chat']) + const file = writeFile(join(scratch, 'notes.txt'), 'notes') + agent.script('[stopped-claim]', (turn) => { + turn.toolCall({ toolName: 'read_local_file', args: { path: file } }) + turn.pause() + }) + const page = await openApp(user, 'Stop chat') + const chatId = user.chats['Stop chat'] + // Delivered to Sim on release even should Electron have given up on it meanwhile. + const lateClaim = proxy.hold(isDesktopClaim, { deliverIfAbandoned: true }) + await send(page, '[stopped-claim] read my secret') + const claim = await lateClaim.arrival(ARRIVAL_MS, 'The read’s claim') + await click(page, page.getByRole('button', { name: 'Stop generation' })) + // Stop settles the call nobody has claimed yet as never started. + await expect.poll(() => callState(chatId), { timeout: 30_000 }).toMatch(/^cancelled/) + lateClaim.release() + // The claim held across Stop reaches Sim after it and is refused, and so is a replay of it. + await expect.poll(() => claim.status, { timeout: 15_000 }).toBe(410) + const [call] = await db.toolCalls(chatId) + const replay = await page.evaluate( + (toolCallId) => + (globalThis as DesktopWindow).simDesktop.localFiles?.({ operation: 'read', toolCallId }), + call.toolCallId + ) + expect(replay?.ok).toBe(false) + const answered = () => + proxy.seen(claim.at, '/api/desktop/tool/authorize').filter((entry) => entry.status) + await expect.poll(() => answered().length, { timeout: 15_000 }).toBeGreaterThan(0) + for (const entry of answered()) expect(entry.status).toBe(410) + const [after] = await db.toolCalls(chatId) + expect(after).toMatchObject({ status: 'cancelled', claimedBy: null }) + }) + + test('with the background executor off, a foreground desktop round trip never binds a device or rings a doorbell', async () => { + const user = await db.seedUser(['Round trip']) + const marker = generateId() + const file = writeFile(join(scratch, 'plan.txt'), `plan ${marker}`) + const deviceId = generateId() + const monitor = new RedisMonitor(sim.redisUrl) + await monitor.start() + try { + // A desktop that speaks the executor protocol registers and offers itself for the turn. + const registration = await fetch(new URL('/api/desktop/devices', sim.upstream), { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + Cookie: `better-auth.session_token=${user.cookie}`, + Origin: proxy.origin, + 'User-Agent': 'Sim Desktop', + }, + body: JSON.stringify({ + deviceId, + name: 'E2E desktop', + appVersion: '0.9.0', + platform: `${process.platform}-${process.arch}`, + capabilities: { executor: 1, browser: true, terminal: true, localFiles: true }, + }), + signal: AbortSignal.timeout(COMPILE_MS), + }) + // Each layer of the dormant executor is checked on its own (soft), so a regression shows + // every layer it reaches: the answer to registration, the device record, the turn's + // binding, the routes the app calls, and the doorbell. + expect(registration.status).toBe(200) + expect + .soft(await registration.json(), 'registration answer') + .toMatchObject({ enabled: false }) + proxy.rewriteChatBody((body) => { + const desktop = toRecord(body.desktopCapabilities) + // The app offers its own install when it speaks the executor protocol; otherwise offer + // the device registered above, as such a desktop would. + body.desktopCapabilities = { deviceId, executor: 1, ...desktop } + }) + + let callId = '' + agent.script( + '[round-trip]', + (turn) => { + turn.text('Reading the plan.') + callId = turn.toolCall({ toolName: 'read_local_file', args: { path: file } }) + turn.pause() + }, + (resume, turn) => { + const result = resume.results.find((entry) => entry.callId === callId) + turn.complete( + JSON.stringify(result?.data).includes(marker) ? 'The plan says go.' : 'Could not read.' + ) + } + ) + const since = Date.now() + const page = await openApp(user, 'Round trip') + await send(page, '[round-trip] what does my plan say?') + await expect + .soft(page.getByText('The plan says go.'), 'foreground round trip') + .toBeVisible({ timeout: 60_000 }) + expect.soft(agent.resultFor(callId)?.success, 'read result').toBe(true) + expect(proxy.rewrittenChatBodies).toBeGreaterThan(0) + + // The app registers once signed out (refused) and again on sign-in. + const registeredSignedIn = () => + proxy + .seen(since, '/api/desktop/devices') + .some((entry) => entry.method === 'POST' && entry.status === 200) + await expect.poll(registeredSignedIn, { timeout: 30_000 }).toBe(true) + expect.soft(await db.desktopDeviceCount(user.userId), 'device records').toBe(0) + + const runs = await db.runs(user.chats['Round trip']) + expect(runs.length).toBeGreaterThan(0) + expect + .soft( + runs.map((run) => run.desktopDeviceId), + 'turn binding' + ) + .toEqual(runs.map(() => null)) + const [call] = await db.toolCalls(user.chats['Round trip']) + expect + .soft(call, 'foreground call') + .toMatchObject({ toolName: 'read_local_file', status: 'completed' }) + expect.soft(call?.persistSeq, 'persist order').not.toBeNull() + + // Only registration and the foreground claim: no inbox, doorbell stream, executor claim, + // lease or completion. + const desktopRoutes = new Set(proxy.seen(since, '/api/desktop/').map((entry) => entry.path)) + desktopRoutes.delete('/api/desktop/devices') + expect + .soft(desktopRoutes, 'desktop routes the app called') + .toEqual(new Set(['/api/desktop/tool/authorize'])) + expect(monitor.lines.length).toBeGreaterThan(0) + expect.soft(monitor.publishesTo('desktop:inbox'), 'doorbell').toEqual([]) + } finally { + monitor.stop() + } + }) +}) diff --git a/apps/desktop/e2e/fixtures/live-sim.ts b/apps/desktop/e2e/fixtures/live-sim.ts new file mode 100644 index 00000000000..6a23821f013 --- /dev/null +++ b/apps/desktop/e2e/fixtures/live-sim.ts @@ -0,0 +1,741 @@ +import { createHmac } from 'node:crypto' +import { + createServer, + request as httpRequest, + type IncomingHttpHeaders, + type IncomingMessage, + type Server, + type ServerResponse, +} from 'node:http' +import { connect, type Socket } from 'node:net' +import type { Duplex } from 'node:stream' +import { sleep } from '@sim/utils/helpers' +import { generateId, generateShortId } from '@sim/utils/id' +import { toArray } from '@sim/utils/object' +import postgres from 'postgres' + +/** + * A real local Sim app for the desktop E2E suite: the Electron app talks to it through a + * recording proxy, and Sim talks to a scripted agent standing in for the Mothership worker at + * `SIM_AGENT_API_URL`. Both sides of the desktop tools are the shipped code; only the model's + * decisions are scripted. + * + * Start Sim (for example `next dev`) with `NEXT_PUBLIC_APP_URL` and `BETTER_AUTH_URL` set to the + * proxy's origin, `SIM_AGENT_API_URL` set to the agent's, `REDIS_URL`, and + * `COPILOT_TOOL_PERMISSIONS_ENABLED=true`, then provide the environment `liveSimConfig` reads. + */ +export interface LiveSimConfig { + /** Where Sim itself listens. */ + upstream: string + /** The proxy's port: the origin Sim and Electron both use. */ + proxyPort: number + /** The scripted agent's port. */ + agentPort: number + databaseUrl: string + redisUrl: string + /** Sim's `BETTER_AUTH_SECRET`, to sign the session cookie of a seeded session. */ + authSecret: string +} + +/** The live Sim this run was given, or why the suite is skipped. */ +export function liveSimConfig(): LiveSimConfig | string { + const env = process.env + const missing = [ + 'SIM_DESKTOP_E2E_SIM_URL', + 'SIM_DESKTOP_E2E_PROXY_PORT', + 'SIM_DESKTOP_E2E_AGENT_PORT', + 'SIM_DESKTOP_E2E_DATABASE_URL', + 'SIM_DESKTOP_E2E_REDIS_URL', + 'SIM_DESKTOP_E2E_AUTH_SECRET', + ].filter((name) => !env[name]) + if (missing.length > 0) return `Needs a local Sim app: set ${missing.join(', ')}` + const databaseUrl = env.SIM_DESKTOP_E2E_DATABASE_URL ?? '' + if (!/(?:^|[_-])test$/i.test(new URL(databaseUrl).pathname.slice(1))) + return 'SIM_DESKTOP_E2E_DATABASE_URL must name a dedicated test database' + return { + upstream: env.SIM_DESKTOP_E2E_SIM_URL ?? '', + proxyPort: Number(env.SIM_DESKTOP_E2E_PROXY_PORT), + agentPort: Number(env.SIM_DESKTOP_E2E_AGENT_PORT), + databaseUrl, + redisUrl: env.SIM_DESKTOP_E2E_REDIS_URL ?? '', + authSecret: env.SIM_DESKTOP_E2E_AUTH_SECRET ?? '', + } +} + +async function readBody(request: IncomingMessage): Promise { + const chunks: Buffer[] = [] + for await (const chunk of request) + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)) + return Buffer.concat(chunks) +} + +function listen(server: Server, port: number): Promise { + return new Promise((resolve, reject) => { + server.once('error', reject) + server.listen(port, '127.0.0.1', () => { + server.off('error', reject) + resolve() + }) + }) +} + +function close(server: Server): Promise { + server.closeAllConnections() + return new Promise((resolve) => server.close(() => resolve())) +} + +/** One request the proxy saw, from the renderer or Electron's main process. */ +interface ProxiedRequest { + method: string + path: string + at: number + status?: number + /** When the client gave up on the response before it finished (a reader that let go). */ + clientClosedAt?: number +} + +/** + * A request the proxy holds before it reaches Sim, so a test can act while a tool's request is + * in flight. A request whose client gave up while held is dropped, never forwarded, as a + * cancelled request on a real network would be. + */ +class HeldRequest { + readonly reached: Promise + /** Settles once the client gives up on the held request (it cancelled it). */ + readonly abandoned: Promise + isAbandoned = false + private reach!: (request: ProxiedRequest) => void + private abandon!: () => void + private releaseHeld: (() => void) | undefined + private released = false + + constructor( + readonly matches: (method: string, path: string) => boolean, + /** Deliver the request to Sim on release even if its client gave up, as a late request would arrive. */ + readonly deliverIfAbandoned = false + ) { + this.reached = new Promise((resolve) => { + this.reach = resolve + }) + this.abandoned = new Promise((resolve) => { + this.abandon = resolve + }) + } + + /** Called by the proxy when the matching request arrives; resolves when the test releases it. */ + hold(request: ProxiedRequest, response: ServerResponse): Promise { + response.once('close', () => { + if (response.writableFinished || this.released) return + this.isAbandoned = true + this.abandon() + }) + this.reach(request) + return new Promise((resolve) => { + this.releaseHeld = () => resolve(!this.isAbandoned) + if (this.released) this.releaseHeld() + }) + } + + /** The held request, once it arrives; fails if it has not within `timeoutMs`. */ + async arrival(timeoutMs: number, label: string): Promise { + let timer: ReturnType | undefined + const timeout = new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error(`${label} did not arrive in ${timeoutMs} ms`)), + timeoutMs + ) + }) + try { + return await Promise.race([this.reached, timeout]) + } finally { + clearTimeout(timer) + } + } + + release(): void { + this.released = true + this.releaseHeld?.() + } +} + +/** + * The origin Electron and Sim share. It forwards everything to Sim unchanged, records each + * request, can hold one until the test releases it, and serves one test-only route that + * installs a seeded session's cookie the way Sim's own sign-in response would. + */ +export class SimProxy { + readonly requests: ProxiedRequest[] = [] + readonly origin: string + /** How many chat turns `rewriteChatBody` rewrote. */ + rewrittenChatBodies = 0 + private readonly server: Server + private holds: HeldRequest[] = [] + /** Requests a hold is keeping from Sim right now. */ + private readonly heldEntries = new Set() + private chatBodyRewrite: ((body: Record) => void) | undefined + private readonly sockets = new Set() + + constructor(private readonly config: LiveSimConfig) { + this.origin = `http://127.0.0.1:${config.proxyPort}` + this.server = createServer((request, response) => { + void this.handle(request, response).catch((error: unknown) => { + if (!response.headersSent) response.writeHead(502) + response.end(String(error)) + }) + }) + // Next's dev server pushes over a websocket; pass upgrades straight through. + this.server.on('upgrade', (request, socket, head) => { + const target = new URL(config.upstream) + const upstream = connect(Number(target.port), target.hostname, () => { + const lines = [`${request.method} ${request.url} HTTP/1.1`] + for (let i = 0; i < request.rawHeaders.length; i += 2) + lines.push(`${request.rawHeaders[i]}: ${request.rawHeaders[i + 1]}`) + upstream.write(`${lines.join('\r\n')}\r\n\r\n`) + upstream.write(head) + upstream.pipe(socket) + socket.pipe(upstream) + }) + const client = socket + this.sockets.add(upstream) + this.sockets.add(client) + const end = () => { + upstream.destroy() + client.destroy() + this.sockets.delete(upstream) + this.sockets.delete(client) + } + for (const side of [upstream, client]) { + side.on('error', end) + side.on('close', end) + } + }) + } + + start(): Promise { + return listen(this.server, this.config.proxyPort) + } + + async stop(): Promise { + for (const socket of this.sockets) socket.destroy() + await close(this.server) + } + + /** Holds the next request that matches until the returned handle releases it. */ + hold( + matches: (method: string, path: string) => boolean, + options: { deliverIfAbandoned?: boolean } = {} + ): HeldRequest { + const held = new HeldRequest(matches, options.deliverIfAbandoned) + this.holds.push(held) + return held + } + + /** Drops holds that never matched, so none outlives the test that set it. */ + clearHolds(): void { + for (const held of this.holds) held.release() + this.holds = [] + } + + /** Rewrites the JSON body of chat turns the renderer sends, as a newer client would send it. */ + rewriteChatBody(rewrite: ((body: Record) => void) | undefined): void { + this.chatBodyRewrite = rewrite + } + + /** + * Resolves once every request the app made has been answered (held ones aside), staying so for + * a moment, which on a dev app means none is waiting on a route to compile. + */ + async settled(timeoutMs: number): Promise { + const deadline = Date.now() + timeoutMs + let quietSince = Date.now() + for (;;) { + const waiting = this.requests.filter( + (entry) => + entry.status === undefined && + entry.clientClosedAt === undefined && + !this.heldEntries.has(entry) + ) + if (waiting.length > 0) quietSince = Date.now() + else if (Date.now() - quietSince >= 3_000) return + if (Date.now() > deadline) + throw new Error( + `Requests still unanswered: ${waiting.map((entry) => entry.path).join(', ')}` + ) + await sleep(250) + } + } + + /** The requests seen since `since`, optionally only those under a path prefix. */ + seen(since = 0, prefix = '/'): ProxiedRequest[] { + return this.requests.filter((entry) => entry.at >= since && entry.path.startsWith(prefix)) + } + + private async handle(request: IncomingMessage, response: ServerResponse): Promise { + const url = new URL(request.url ?? '/', this.origin) + const method = request.method ?? 'GET' + const entry: ProxiedRequest = { method, path: url.pathname, at: Date.now() } + this.requests.push(entry) + response.once('close', () => { + if (!response.writableFinished) entry.clientClosedAt = Date.now() + }) + if (url.pathname === '/__e2e/sign-in') { + const cookie = url.searchParams.get('cookie') ?? '' + response.writeHead(302, { + 'Set-Cookie': `better-auth.session_token=${cookie}; Path=/; HttpOnly; SameSite=Lax`, + Location: url.searchParams.get('to') ?? '/', + }) + response.end() + entry.status = 302 + return + } + let body = await readBody(request) + const held = this.holds.find((candidate) => candidate.matches(method, url.pathname)) + if (held) { + this.holds = this.holds.filter((candidate) => candidate !== held) + this.heldEntries.add(entry) + const deliver = await held.hold(entry, response) + this.heldEntries.delete(entry) + if (!deliver && !held.deliverIfAbandoned) return + } + if (this.chatBodyRewrite && method === 'POST' && url.pathname === '/api/mothership/chat') { + const parsed: Record = JSON.parse(body.toString('utf8')) + this.chatBodyRewrite(parsed) + this.rewrittenChatBodies += 1 + body = Buffer.from(JSON.stringify(parsed)) + } + const target = new URL(url.pathname + url.search, this.config.upstream) + // The body is forwarded whole, so it is sent with a length rather than chunked. + const { 'transfer-encoding': _chunked, ...forwarded } = request.headers + const headers: IncomingHttpHeaders = { ...forwarded, 'content-length': String(body.length) } + await new Promise((resolve, reject) => { + const clientGone = response.destroyed + const upstream = httpRequest(target, { method, headers }, (upstreamResponse) => { + entry.status = upstreamResponse.statusCode + upstreamResponse.on('end', resolve) + upstreamResponse.on('error', reject) + // A request delivered after its client gave up is answered to no one. + if (clientGone) { + upstreamResponse.resume() + return + } + response.writeHead(upstreamResponse.statusCode ?? 502, upstreamResponse.headers) + upstreamResponse.pipe(response) + }) + if (!clientGone) response.on('close', () => upstream.destroy()) + upstream.on('error', reject) + upstream.end(body) + }) + } +} + +/** One tool call a scripted turn issues. */ +interface ScriptedToolCall { + toolName: string + args: Record + /** `awaiting_approval` asks Sim to hold the call for the user's decision. */ + status?: 'awaiting_approval' +} + +/** One result Sim sends back when it resumes the agent after a checkpoint. */ +interface ResumedResult { + callId: string + name?: string + success?: boolean + data?: unknown +} + +/** A resume request Sim sent, with when it arrived. */ +interface Resume { + at: number + streamId: string + results: ResumedResult[] +} + +/** A chat turn Sim opened against the agent, written to as the scripted model acts. */ +class AgentTurn { + readonly toolCallIds: string[] = [] + /** Settles once this leg's stream has ended, from either side. */ + readonly closed: Promise + private markClosed!: () => void + private seq = 0 + private readonly keepAlive: ReturnType + private ended = false + + constructor( + readonly body: Record, + readonly streamId: string, + private readonly response: ServerResponse + ) { + this.closed = new Promise((resolve) => { + this.markClosed = resolve + }) + response.writeHead(200, { + 'Content-Type': 'text/event-stream', + 'Cache-Control': 'no-cache', + Connection: 'keep-alive', + }) + // The worker keeps an idle stream alive with SSE comments; so does this one. + this.keepAlive = setInterval(() => { + if (!this.ended) response.write(': keepalive\n\n') + }, 2_000) + response.on('close', () => this.finish()) + this.emit('session', { kind: 'start' }) + } + + get message(): string { + return typeof this.body.message === 'string' ? this.body.message : '' + } + + emit(type: string, payload: Record): void { + if (this.ended) return + this.seq += 1 + const chatId = typeof this.body.chatId === 'string' ? this.body.chatId : undefined + const envelope = { + v: 1, + type, + seq: this.seq, + ts: new Date().toISOString(), + stream: { streamId: this.streamId, ...(chatId ? { chatId } : {}), cursor: String(this.seq) }, + trace: { requestId: `agent-${this.streamId}` }, + payload, + } + this.response.write(`data: ${JSON.stringify(envelope)}\n\n`) + } + + text(text: string): void { + this.emit('text', { channel: 'assistant', text }) + } + + /** Issues a client-executed tool call and returns its id. */ + toolCall(call: ScriptedToolCall): string { + const toolCallId = `toolu_${generateShortId()}` + this.toolCallIds.push(toolCallId) + this.emit('tool', { + phase: 'call', + toolCallId, + toolName: call.toolName, + executor: 'client', + mode: 'async', + arguments: call.args, + ui: { clientExecutable: true }, + ...(call.status ? { status: call.status } : {}), + }) + return toolCallId + } + + /** Ends this leg waiting on the turn's tool calls, as the worker does at a checkpoint. */ + pause(): void { + this.emit('run', { + kind: 'checkpoint_pause', + checkpointId: generateId(), + executionId: generateId(), + runId: generateId(), + pendingToolCallIds: [...this.toolCallIds], + }) + this.finish() + } + + complete(text?: string): void { + if (text) this.text(text) + this.emit('complete', { status: 'complete' }) + this.finish() + } + + private finish(): void { + if (this.ended) return + this.ended = true + clearInterval(this.keepAlive) + this.response.end() + this.markClosed() + } +} + +type TurnScript = (turn: AgentTurn) => void | Promise +type ResumeScript = (resume: Resume, turn: AgentTurn) => void | Promise + +/** + * Stands in for the Mothership worker at Sim's `SIM_AGENT_API_URL`: each chat turn runs the + * script registered for its message, and each resume after a checkpoint records the tool results + * Sim delivers, which is what the real model would read. + */ +export class ScriptedAgent { + readonly turns: AgentTurn[] = [] + readonly resumes: Resume[] = [] + readonly unexpected: string[] = [] + private readonly server: Server + private readonly scripts = new Map() + private readonly resumeScripts = new Map() + private readonly waiters: (() => void)[] = [] + + constructor(private readonly port: number) { + this.server = createServer((request, response) => { + void this.handle(request, response).catch((error: unknown) => { + this.unexpected.push(String(error)) + if (!response.headersSent) response.writeHead(500) + response.end() + }) + }) + } + + get url(): string { + return `http://127.0.0.1:${this.port}` + } + + start(): Promise { + return listen(this.server, this.port) + } + + stop(): Promise { + return close(this.server) + } + + /** + * Runs `script` for the turn whose message contains `marker`; `onResume` answers each resume of + * that turn (by default the turn completes with a short reply). + */ + script(marker: string, script: TurnScript, onResume?: ResumeScript): void { + this.scripts.set(marker, script) + if (onResume) this.resumeScripts.set(marker, onResume) + } + + /** Resolves with the first resume matching `predicate`, failing after `timeoutMs`. */ + async waitForResume(predicate: (resume: Resume) => boolean, timeoutMs: number): Promise { + const deadline = Date.now() + timeoutMs + for (;;) { + const found = this.resumes.find(predicate) + if (found) return found + const remaining = deadline - Date.now() + if (remaining <= 0) throw new Error(`No matching resume within ${timeoutMs} ms`) + await new Promise((resolve) => { + const timer = setTimeout(resolve, remaining) + this.waiters.push(() => { + clearTimeout(timer) + resolve() + }) + }) + } + } + + /** The result Sim delivered for `toolCallId`, if any resume carried it. */ + resultFor(toolCallId: string): (ResumedResult & { at: number }) | undefined { + for (const resume of this.resumes) { + const result = resume.results.find((entry) => entry.callId === toolCallId) + if (result) return { ...result, at: resume.at } + } + return undefined + } + + private async handle(request: IncomingMessage, response: ServerResponse): Promise { + const path = new URL(request.url ?? '/', this.url).pathname + const raw = (await readBody(request)).toString('utf8') + const body: Record = raw ? JSON.parse(raw) : {} + if (path === '/api/mothership' || path === '/api/copilot') { + const message = typeof body.message === 'string' ? body.message : '' + const streamId = typeof body.messageId === 'string' ? body.messageId : generateId() + const turn = new AgentTurn(body, streamId, response) + this.turns.push(turn) + const entry = [...this.scripts].find(([marker]) => message.includes(marker)) + if (!entry) { + turn.complete('No script for this message.') + return + } + await entry[1](turn) + return + } + if (path === '/api/tools/resume') { + const streamId = typeof body.streamId === 'string' ? body.streamId : '' + const results = toArray(body.results) + const resume: Resume = { at: Date.now(), streamId, results } + this.resumes.push(resume) + for (const wake of this.waiters.splice(0)) wake() + const turn = this.turns.find((candidate) => candidate.streamId === streamId) + const resumed = new AgentTurn(turn?.body ?? body, streamId, response) + const marker = turn + ? [...this.resumeScripts.keys()].find((key) => turn.message.includes(key)) + : undefined + const onResume = marker ? this.resumeScripts.get(marker) : undefined + if (onResume) await onResume(resume, resumed) + else resumed.complete('Done.') + return + } + if (path === '/api/generate-chat-title') { + response.writeHead(200, { 'Content-Type': 'application/json' }) + response.end(JSON.stringify({ title: 'Desktop tools E2E' })) + return + } + // Stop, cleanup and every other worker callback: acknowledge. + response.writeHead(200, { 'Content-Type': 'application/json' }) + response.end(JSON.stringify({ settled: true })) + } +} + +/** A seeded user signed in to Sim, with one workspace and named chats. */ +export interface SeededUser { + userId: string + name: string + workspaceId: string + sessionId: string + /** The `better-auth.session_token` cookie value for the seeded session. */ + cookie: string + chats: Record +} + +/** Direct database access for seeding and for asserting what Sim persisted. */ +export class SimDatabase { + readonly sql: postgres.Sql + private readonly userIds: string[] = [] + + constructor(private readonly config: LiveSimConfig) { + this.sql = postgres(config.databaseUrl, { max: 4, onnotice: () => {} }) + } + + async close(): Promise { + if (this.userIds.length > 0) + await this.sql`delete from "user" where id in ${this.sql(this.userIds)}`.catch(() => {}) + await this.sql.end({ timeout: 5 }) + } + + /** A user with a workspace, a session as the desktop sign-in creates, and one chat per title. */ + async seedUser(chatTitles: string[]): Promise { + const userId = generateId() + const workspaceId = generateId() + const sessionId = generateId() + const token = generateShortId() + const name = `Desktop E2E ${userId.slice(0, 6)}` + const email = `${userId}@desktop-tools-e2e.test` + const chats: Record = {} + await this.sql.begin(async (tx) => { + await tx`insert into "user" (id, name, email, normalized_email, email_verified, created_at, updated_at) + values (${userId}, ${name}, ${email}, ${email}, true, now(), now())` + await tx`insert into user_stats (id, user_id) values (${generateId()}, ${userId})` + await tx`insert into workspace (id, name, owner_id, billed_account_user_id) + values (${workspaceId}, 'Desktop tools E2E', ${userId}, ${userId})` + await tx`insert into permissions (id, user_id, entity_type, entity_id, permission_type) + values (${generateId()}, ${userId}, 'workspace', ${workspaceId}, 'admin')` + await tx`insert into session (id, token, user_id, user_agent, expires_at, created_at, updated_at) + values (${sessionId}, ${token}, ${userId}, 'Sim Desktop', now() + interval '1 day', now(), now())` + for (const title of chatTitles) { + const chatId = generateId() + chats[title] = chatId + await tx`insert into copilot_chats (id, user_id, workspace_id, type, title) + values (${chatId}, ${userId}, ${workspaceId}, 'mothership', ${title})` + } + }) + this.userIds.push(userId) + const signature = createHmac('sha256', this.config.authSecret).update(token).digest('base64') + return { + userId, + name, + workspaceId, + sessionId, + cookie: encodeURIComponent(`${token}.${signature}`), + chats, + } + } + + /** The persisted desktop calls of a chat's runs, oldest first. */ + toolCalls(chatId: string) { + return this.sql< + { + toolCallId: string + toolName: string + status: string + claimedBy: string | null + error: string | null + result: unknown + permissionRequestedAt: Date | null + persistSeq: number | null + runId: string + }[] + >`select c.tool_call_id as "toolCallId", c.tool_name as "toolName", c.status, + c.claimed_by as "claimedBy", c.error, c.result, + c.permission_requested_at as "permissionRequestedAt", c.persist_seq as "persistSeq", + c.run_id as "runId" + from copilot_async_tool_calls c join copilot_runs r on r.id = c.run_id + where r.chat_id = ${chatId} order by c.created_at` + } + + runs(chatId: string) { + return this.sql<{ id: string; status: string; desktopDeviceId: string | null }[]>` + select id, status, desktop_device_id as "desktopDeviceId" from copilot_runs + where chat_id = ${chatId} order by created_at` + } + + async desktopDeviceCount(userId?: string): Promise { + const [row] = userId + ? await this.sql<{ count: number }[]>` + select count(*)::int as count from desktop_devices where user_id = ${userId}` + : await this.sql<{ count: number }[]>`select count(*)::int as count from desktop_devices` + return row?.count ?? 0 + } + + /** Names of the files a workspace holds. */ + async workspaceFileNames(workspaceId: string): Promise { + const rows = await this.sql<{ name: string }[]>` + select original_name as name from workspace_files + where workspace_id = ${workspaceId} and context = 'workspace' and deleted_at is null + order by original_name` + return rows.map((row) => row.name) + } + + /** Names of the file folders a workspace holds. */ + async workspaceFolderNames(workspaceId: string): Promise { + const rows = await this.sql<{ name: string }[]>` + select name from folder where workspace_id = ${workspaceId} and resource_type = 'file' + order by name` + return rows.map((row) => row.name) + } +} + +/** + * Records every command Sim sends Redis (`MONITOR`), so a test can assert that a turn rang no + * desktop doorbell. Speaks just enough RESP for that. + */ +export class RedisMonitor { + readonly lines: string[] = [] + private socket: Socket | undefined + + constructor(private readonly url: string) {} + + async start(): Promise { + const target = new URL(this.url) + const socket = connect(Number(target.port || 6379), target.hostname) + this.socket = socket + await new Promise((resolve, reject) => { + socket.once('connect', resolve) + socket.once('error', reject) + }) + // Each command sent before MONITOR answers `+OK`; MONITOR's own `+OK` means it is recording. + let pendingAcks = target.password ? 2 : 1 + let buffer = '' + const acknowledged = new Promise((resolve, reject) => { + socket.on('data', (chunk) => { + buffer += chunk.toString('utf8') + const lines = buffer.split('\r\n') + buffer = lines.pop() ?? '' + for (const line of lines) { + if (pendingAcks > 0) { + if (line.startsWith('-')) reject(new Error(`Redis refused MONITOR: ${line}`)) + else if (line === '+OK' && --pendingAcks === 0) resolve() + continue + } + if (line.startsWith('+')) this.lines.push(line) + } + }) + }) + if (target.password) socket.write(`AUTH ${decodeURIComponent(target.password)}\r\n`) + socket.write('MONITOR\r\n') + await acknowledged + } + + /** Commands seen that publish to a channel whose name contains `channel`. */ + publishesTo(channel: string): string[] { + return this.lines.filter( + (line) => /"publish"/i.test(line) && line.toLowerCase().includes(channel.toLowerCase()) + ) + } + + stop(): void { + this.socket?.destroy() + } +} diff --git a/apps/desktop/package.json b/apps/desktop/package.json index a4484720362..3a6dbd01c5e 100644 --- a/apps/desktop/package.json +++ b/apps/desktop/package.json @@ -60,6 +60,7 @@ "jsdom": "^26.0.0", "postcss": "^8", "postcss-load-config": "6.0.1", + "postgres": "^3.4.5", "react": "19.2.4", "react-dom": "19.2.4", "typescript": "^7.0.2", diff --git a/bun.lock b/bun.lock index 4f462187545..70df25efd9d 100644 --- a/bun.lock +++ b/bun.lock @@ -76,6 +76,7 @@ "jsdom": "^26.0.0", "postcss": "^8", "postcss-load-config": "6.0.1", + "postgres": "^3.4.5", "react": "19.2.4", "react-dom": "19.2.4", "typescript": "^7.0.2",