diff --git a/.github/scripts/http-e2e.sh b/.github/scripts/http-e2e.sh index 202a9a2f2a2..3523bbe8a1a 100755 --- a/.github/scripts/http-e2e.sh +++ b/.github/scripts/http-e2e.sh @@ -137,7 +137,6 @@ case "$group" in desktop-inbox) export REDIS_URL=redis://127.0.0.1:6379 export NEXT_PUBLIC_FORCE_HOSTED=false - export MSHIP_DESKTOP_BACKGROUND_EXECUTOR=true export COPILOT_TOOL_PERMISSIONS_ENABLED=true export INTERNAL_API_SECRET=desktop-inbox-http-ci-local-secret-at-least-32-characters export DB_TX_TRIPWIRE=throw diff --git a/.github/workflows/checks.yml b/.github/workflows/checks.yml index 817b76b7d2a..f680a832415 100644 --- a/.github/workflows/checks.yml +++ b/.github/workflows/checks.yml @@ -307,8 +307,9 @@ jobs: key: ${{ steps.next-cache.outputs.key }} path: ./apps/sim/.next/dev - # 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. + # Chat switches, Stop, sign-out, approval, the dormant-executor round trip, and background + # runs across a chat switch and a network cut. 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 diff --git a/apps/desktop/e2e/background-executor.spec.ts b/apps/desktop/e2e/background-executor.spec.ts index 72ec81b2018..5eef9b4e582 100644 --- a/apps/desktop/e2e/background-executor.spec.ts +++ b/apps/desktop/e2e/background-executor.spec.ts @@ -5,6 +5,7 @@ import { type ElectronApplication, expect, test } from '@playwright/test' import type { SimDesktopApi } from '@sim/desktop-bridge' import { getErrorMessage } from '@sim/utils/errors' import { sleep } from '@sim/utils/helpers' +import { recordCheck } from './check-report' import { type FixtureCall, FixtureSim, @@ -24,37 +25,32 @@ import { * device protocol (register, inbox, doorbell, claim, lease, complete, import) the way Sim's * routes do. * The window navigates, reloads and leaves the chats while their calls run: nothing in this - * suite depends on a chat view, which is the point. Each scenario's checks land in a JSON - * report at BACKGROUND_EXECUTOR_REPORT_PATH. + * suite depends on a chat view, which is the point. Each check lands in a JSON report at + * BACKGROUND_EXECUTOR_REPORT_PATH as it finishes. */ const CHAT_A = 'chat-browser-a' const CHAT_B = 'chat-terminal-b' const CHAT_C = 'chat-idle-c' -interface ReportCheck { - name: string - status: 'passed' | 'failed' - durationMs: number - error?: string -} - -const report: ReportCheck[] = [] - const sim = new FixtureSim() +/** Runs one check and adds its outcome to the report as soon as it is known. */ async function check(name: string, body: () => Promise): Promise { const startedAt = Date.now() - try { - await body() - report.push({ name, status: 'passed', durationMs: Date.now() - startedAt }) - } catch (error) { - report.push({ + const record = (status: 'passed' | 'failed', error?: unknown) => + recordCheck(process.env.BACKGROUND_EXECUTOR_REPORT_PATH, 'background-executor', { name, - status: 'failed', + status, durationMs: Date.now() - startedAt, - error: getErrorMessage(error), + retry: test.info().retry, + ...(error === undefined ? {} : { error: getErrorMessage(error) }), }) + try { + await body() + record('passed') + } catch (error) { + record('failed', error) throw error } } @@ -74,13 +70,6 @@ test.describe('background executor', () => { test.afterAll(async () => { await sim.stop() - const reportPath = process.env.BACKGROUND_EXECUTOR_REPORT_PATH - if (reportPath) { - writeFileSync( - reportPath, - JSON.stringify({ suite: 'background-executor', checks: report }, null, 2) - ) - } }) test('A: two chats run browser and terminal work while the user is elsewhere and reloads', async () => { @@ -168,7 +157,7 @@ test.describe('background executor', () => { }) }) - test('B: a result produced while offline is delivered once after reconnecting', async () => { + test('B: a result produced while the network is cut is delivered once after reconnecting', async () => { const userData = mkdtempSync(join(tmpdir(), 'sim-executor-b-')) app = (await launch(sim, userData)).app const deviceId = await registeredDevice(sim) @@ -178,21 +167,33 @@ test.describe('background executor', () => { args: { command: 'sleep 2; echo offline-done', waitSeconds: 30 }, }) await expect.poll(() => sim.requireCall(run).claims).toBe(1) - sim.offline = true - await sleep(6_000) + sim.disconnect() - await check('B: nothing reached Sim while offline', async () => { - expect(sim.requireCall(run).completions).toHaveLength(0) - expect(sim.droppedWhileOffline).toBeGreaterThan(0) - }) - sim.offline = false + await check( + 'B: the cut drops the doorbell and nothing reaches Sim while it lasts', + async () => { + await expect.poll(() => sim.streams.get(deviceId)?.size ?? 0).toBe(0) + await expect.poll(() => sim.droppedWhileOffline, { timeout: 15_000 }).toBeGreaterThan(0) + await sleep(4_000) + expect(sim.requireCall(run).completions).toHaveLength(0) + } + ) + sim.reconnect() await check('B: the result arrives once after reconnecting', async () => { const completion = await settled(sim, run, 60_000) - expect(completion.status).toBe('success') + expect(completion.status, completion.message).toBe('success') expect(JSON.stringify(completion.data)).toContain('offline-done') await sleep(3_000) expect(sim.requireCall(run).completions).toHaveLength(1) + expect(sim.requireCall(run).claims).toBe(1) + }) + + await check('B: the device reopens its doorbell and picks up new work', async () => { + await expect.poll(() => sim.streams.get(deviceId)?.size ?? 0, { timeout: 45_000 }).toBe(1) + const after = sim.issue(deviceId, CHAT_B, 'browser_list_tabs', {}) + const completion = await settled(sim, after) + expect(completion.status, completion.message).toBe('success') }) }) @@ -376,6 +377,34 @@ test.describe('background executor', () => { }) }) + test('F: a call the user declines in a background chat never runs', async () => { + const userData = mkdtempSync(join(tmpdir(), 'sim-executor-f-declined-')) + const marker = join(userData, 'declined-marker.txt') + app = (await launch(sim, userData)).app + const deviceId = await registeredDevice(sim) + const pulls = () => + sim.requests.filter((request) => request.startsWith('GET /api/desktop/inbox')).length + const pullsBefore = pulls() + + const gated = sim.issue( + deviceId, + CHAT_B, + 'terminal', + { operation: 'run', args: { command: `echo ran >> '${marker}'`, waitSeconds: 30 } }, + 'awaiting_approval' + ) + // The device has pulled its inbox and seen the call waiting for approval. + await expect.poll(pulls).toBeGreaterThan(pullsBefore) + sim.decline(gated) + + await check('F: the declined call is never claimed, and its command never runs', async () => { + await sleep(RECONCILE_MS * 2) + expect(sim.requireCall(gated).claims).toBe(0) + expect(sim.requireCall(gated).completions).toEqual([]) + expect(readFileSafe(marker)).toBe('') + }) + }) + test('G: only the device a turn is bound to claims its calls', async () => { const first = await launch(sim, mkdtempSync(join(tmpdir(), 'sim-executor-g1-'))) const firstDevice = await registeredDevice(sim) diff --git a/apps/desktop/e2e/check-report.spec.ts b/apps/desktop/e2e/check-report.spec.ts new file mode 100644 index 00000000000..dc2a6378255 --- /dev/null +++ b/apps/desktop/e2e/check-report.spec.ts @@ -0,0 +1,49 @@ +import { spawnSync } from 'node:child_process' +import { mkdtempSync, readFileSync, writeFileSync } from 'node:fs' +import { createRequire } from 'node:module' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { fileURLToPath } from 'node:url' +import { expect, test } from '@playwright/test' +import { omit } from '@sim/utils/object' +import type { ReportCheck } from './check-report' + +/** + * The suites' JSON reports against a real worker restart: a nested Playwright run of a fixture + * whose first test fails, so the second runs in a fresh worker with fresh module state. The report + * must keep the failure, and must not carry checks from an earlier run. + */ +const CLI = createRequire(import.meta.url).resolve('@playwright/test/cli') +const CONFIG = fileURLToPath(new URL('./fixtures/worker-restart.config.ts', import.meta.url)) + +test('a check that failed before the worker restarted stays in the report', () => { + const reportPath = join(mkdtempSync(join(tmpdir(), 'sim-check-report-')), 'report.json') + writeFileSync( + reportPath, + JSON.stringify({ + suite: 'worker-restart', + run: -1, + checks: [{ name: 'left by an earlier run', status: 'passed', durationMs: 0 }], + }) + ) + // The outer run's worker variables would make the nested runner think it is a worker. + const env = omit( + process.env, + Object.keys(process.env).filter((key) => /^(TEST_|PW_)/.test(key)) + ) + + const run = spawnSync(process.execPath, [CLI, 'test', '--config', CONFIG], { + env: { ...env, WORKER_RESTART_REPORT_PATH: reportPath }, + encoding: 'utf8', + timeout: 60_000, + }) + + expect(run.status, run.stdout + run.stderr).toBe(1) + const checks = (JSON.parse(readFileSync(reportPath, 'utf8')) as { checks: ReportCheck[] }).checks + expect(checks.map(({ name, status }) => ({ name, status }))).toEqual([ + { name: 'fails first', status: 'failed' }, + { name: 'passes after the worker restarts', status: 'passed' }, + ]) + // Proof the second check ran in a replacement worker, not the one that failed. + expect(checks[0]?.error).not.toBe(checks[1]?.error) +}) diff --git a/apps/desktop/e2e/check-report.ts b/apps/desktop/e2e/check-report.ts new file mode 100644 index 00000000000..5fd3692f70c --- /dev/null +++ b/apps/desktop/e2e/check-report.ts @@ -0,0 +1,44 @@ +import { readFileSync, writeFileSync } from 'node:fs' + +/** One check's outcome in a suite's JSON report. */ +export interface ReportCheck { + name: string + status: string + durationMs: number + error?: string + retry?: number +} + +interface Report { + suite: string + /** The Playwright run that wrote it: its workers' parent process. */ + run: number + checks: ReportCheck[] +} + +/** + * Adds a check to the suite's report at `reportPath` as soon as its outcome is known, keeping the + * checks already there from the same run. + * + * Playwright replaces a worker after a failure, and the new one starts with empty module state, so + * a report kept in memory and written at the end would drop everything before the failure, the + * failure included, and could read green. The file holds it instead. Every worker of one run is a + * child of the same runner, so a report left by an earlier run is replaced rather than added to. + */ +export function recordCheck( + reportPath: string | undefined, + suite: string, + check: ReportCheck +): void { + if (!reportPath) return + const run = process.ppid + let checks: ReportCheck[] = [] + try { + const existing = JSON.parse(readFileSync(reportPath, 'utf8')) as Report + if (existing.suite === suite && existing.run === run) checks = existing.checks + } catch { + // No report yet from this run. + } + const report: Report = { suite, run, checks: [...checks, check] } + writeFileSync(reportPath, JSON.stringify(report, null, 2)) +} diff --git a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts index 10aa0c1a030..be9dd5041d7 100644 --- a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts +++ b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts @@ -13,7 +13,6 @@ import { import type { SimDesktopApi } from '@sim/desktop-bridge' import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' -import { toRecord } from '@sim/utils/object' import { type LiveSimConfig, liveSimConfig, @@ -30,6 +29,11 @@ import { * 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. + * + * Sim runs the background executor wherever it has Redis, as it does here. Most tests cover the + * chat view running desktop tools itself, as on an install without Redis: the proxy answers the + * app's registration as such an install does, so the app stays dormant and no turn binds to it. + * The tests marked as running in the background let Sim's own answer through. */ const DESKTOP_DIR = fileURLToPath(new URL('..', import.meta.url)) @@ -43,6 +47,11 @@ const LEASE_MS = 60_000 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 +const REGISTRATION_PATH = '/api/desktop/devices' +/** Sim's answer to registration on an install that cannot run the background executor. */ +const executorUnavailable = (answer: Record) => { + answer.enabled = false +} type DesktopWindow = typeof globalThis & { simDesktop: SimDesktopApi } @@ -68,6 +77,7 @@ test.describe('desktop tools against a live Sim', () => { agent = new ScriptedAgent(sim.agentPort) db = new SimDatabase(sim) await Promise.all([proxy.start(), agent.start()]) + proxy.rewriteAnswer(REGISTRATION_PATH, executorUnavailable) await test.step('warm up the routes the tests use', warmUp) }) @@ -106,9 +116,10 @@ test.describe('desktop tools against a live Sim', () => { ) } proxy.clearHolds() + proxy.restoreNetwork() await app?.close().catch(() => {}) app = undefined - proxy.rewriteChatBody(undefined) + proxy.rewriteAnswer(REGISTRATION_PATH, executorUnavailable) rmSync(scratch, { recursive: true, force: true }) }) @@ -148,6 +159,7 @@ test.describe('desktop tools against a live Sim', () => { '/api/copilot/chats', '/api/users/me/settings', '/api/auth/oauth/connections', + '/api/desktop/inbox', ]) await compile(path) for (const path of [ @@ -155,6 +167,8 @@ test.describe('desktop tools against a live Sim', () => { '/api/mothership/chat/abort', '/api/desktop/devices', '/api/desktop/tool/authorize', + '/api/desktop/tool/claim', + '/api/desktop/tool/complete', '/api/copilot/confirm', '/api/files/uploads', '/api/files/uploads/warm-up/parts', @@ -361,6 +375,9 @@ test.describe('desktop tools against a live Sim', () => { /** 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' + /** The background executor's report of a call's result. */ + const isDesktopCompletion = (method: string, path: string) => + method === 'POST' && path === '/api/desktop/tool/complete' /** A client tool's report of its own result. */ const isToolReport = (method: string, path: string) => method === 'POST' && path === '/api/copilot/confirm' @@ -691,46 +708,13 @@ test.describe('desktop tools against a live Sim', () => { 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 () => { + test('where Sim cannot run the background executor, the app stays dormant and a desktop round trip runs in the chat view', 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]', @@ -749,19 +733,19 @@ test.describe('desktop tools against a live Sim', () => { const since = Date.now() const page = await openApp(user, 'Round trip') await send(page, '[round-trip] what does my plan say?') + // Each layer is checked on its own (soft), so a regression shows every layer it reaches: + // the round trip, the turn's binding, the routes the app calls, and the doorbell. 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') + .seen(since, REGISTRATION_PATH) .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) @@ -778,9 +762,11 @@ test.describe('desktop tools against a live Sim', () => { expect.soft(call?.persistSeq, 'persist order').not.toBeNull() // Only registration and the foreground claim: no inbox, doorbell stream, executor claim, - // lease or completion. + // lease or completion. The workspace's activity poll is Sim's page, not the app: this Sim + // has Redis, so its page asks. const desktopRoutes = new Set(proxy.seen(since, '/api/desktop/').map((entry) => entry.path)) - desktopRoutes.delete('/api/desktop/devices') + desktopRoutes.delete(REGISTRATION_PATH) + desktopRoutes.delete('/api/desktop/activity') expect .soft(desktopRoutes, 'desktop routes the app called') .toEqual(new Set(['/api/desktop/tool/authorize'])) @@ -790,4 +776,90 @@ test.describe('desktop tools against a live Sim', () => { monitor.stop() } }) + + test('in the background, a call issued after the user switched chats runs on the desktop', async () => { + proxy.rewriteAnswer(REGISTRATION_PATH, undefined) + const user = await db.seedUser(['Background chat', 'Other chat']) + const marker = generateId() + const file = writeFile(join(scratch, 'notes.txt'), `notes from disk ${marker}`) + let issue!: () => void + const issued = new Promise((resolve) => { + issue = resolve + }) + let callId = '' + let issuedAt = 0 + agent.script('[background-read]', async (turn) => { + turn.text('Reading your notes.') + await issued + callId = turn.toolCall({ toolName: 'read_local_file', args: { path: file } }) + issuedAt = Date.now() + turn.pause() + }) + const page = await openApp(user, 'Background chat') + await send(page, '[background-read] read my notes') + await expect(page.getByText('Reading your notes.')).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(true) + expect(JSON.stringify(result?.data)).toContain(marker) + // No pickup grace: the desktop, not the chat view the user left, ran it. + expect((result?.at ?? 0) - issuedAt).toBeLessThan(PICKUP_GRACE_MS) + const chatId = user.chats['Background chat'] + const runs = await db.runs(chatId) + expect(runs.some((run) => run.desktopDeviceId !== null)).toBe(true) + const [call] = await db.toolCalls(chatId) + expect(call).toMatchObject({ toolName: 'read_local_file', status: 'completed' }) + expect(proxy.seen(issuedAt, '/api/desktop/tool/authorize')).toEqual([]) + }) + + test('in the background, a result reported across a network cut reaches the agent exactly once', async () => { + proxy.rewriteAnswer(REGISTRATION_PATH, undefined) + const user = await db.seedUser(['Cut chat']) + const chatId = user.chats['Cut chat'] + const marker = generateId() + const file = writeFile(join(scratch, 'notes.txt'), `notes from disk ${marker}`) + let callId = '' + agent.script('[network-cut]', (turn) => { + callId = turn.toolCall({ toolName: 'read_local_file', args: { path: file } }) + turn.pause() + }) + const page = await openApp(user, 'Cut chat') + // Reaches Sim on release although the cut made the app give up on it, as a report already on + // the wire would: the app cannot know it landed, so it reports again once back online. + const completion = proxy.hold(isDesktopCompletion, { deliverIfAbandoned: true }) + await send(page, '[network-cut] read my notes') + await completion.arrival(ARRIVAL_MS, 'The result report') + + proxy.cutNetwork() + await expect.poll(() => completion.isAbandoned, { timeout: 15_000 }).toBe(true) + completion.release() + await expect.poll(() => callState(chatId), { timeout: 30_000 }).toMatch(/^completed/) + await sleep(5_000) + const restoredAt = Date.now() + proxy.restoreNetwork() + + // Back online, the app reopens its doorbell and reports the result again. + await expect + .poll(() => proxy.seen(restoredAt, '/api/desktop/inbox/stream').length, { timeout: 60_000 }) + .toBeGreaterThan(0) + const retried = () => + proxy + .seen(restoredAt) + .filter((entry) => isDesktopCompletion(entry.method, entry.path) && entry.status) + await expect.poll(() => retried().length, { timeout: 60_000 }).toBeGreaterThan(0) + for (const entry of retried()) expect(entry.status).toBeLessThan(300) + + await agent.waitForResume(() => Boolean(agent.resultFor(callId)), 60_000) + await proxy.settled(30_000) + const delivered = agent.resumes.filter((resume) => + resume.results.some((entry) => entry.callId === callId) + ) + expect(delivered).toHaveLength(1) + expect(JSON.stringify(agent.resultFor(callId)?.data)).toContain(marker) + const [call] = await db.toolCalls(chatId) + expect(call).toMatchObject({ toolName: 'read_local_file', status: 'completed' }) + }) }) diff --git a/apps/desktop/e2e/executor-sim.ts b/apps/desktop/e2e/executor-sim.ts index af7a6d3da12..0fd7f3d3153 100644 --- a/apps/desktop/e2e/executor-sim.ts +++ b/apps/desktop/e2e/executor-sim.ts @@ -1,7 +1,10 @@ import { execFileSync } from 'node:child_process' import { createHash } from 'node:crypto' -import { readFileSync } from 'node:fs' +import { mkdtempSync, readFileSync } from 'node:fs' import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'node:http' +import type { Socket } from 'node:net' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import { fileURLToPath } from 'node:url' import { type ElectronApplication, @@ -68,8 +71,10 @@ export class FixtureSim { readonly streams = new Map>() readonly imported: ImportedEntry[] = [] enabled = true - offline = false + /** Connections reset while the network was cut. */ droppedWhileOffline = 0 + private offline = false + private readonly sockets = new Set() private server: Server | null = null origin = '' private nextToken = 0 @@ -78,6 +83,15 @@ export class FixtureSim { this.server = createServer((request, response) => { void this.handle(request, response) }) + this.server.on('connection', (socket) => { + if (this.offline) { + this.droppedWhileOffline += 1 + socket.destroy() + return + } + this.sockets.add(socket) + socket.on('close', () => this.sockets.delete(socket)) + }) await new Promise((resolve) => this.server?.listen(0, '127.0.0.1', resolve)) const address = this.server.address() if (!address || typeof address === 'string') throw new Error('Missing fixture address') @@ -90,6 +104,19 @@ export class FixtureSim { await new Promise((resolve) => this.server?.close(() => resolve())) } + /** + * Cuts the network between the app and Sim: every open connection drops, the doorbell stream and + * kept-alive request sockets included, and every new one is reset until {@link reconnect}. + */ + disconnect(): void { + this.offline = true + for (const socket of this.sockets) socket.destroy() + } + + reconnect(): void { + this.offline = false + } + reset(): void { this.calls.clear() this.devices.clear() @@ -140,6 +167,13 @@ export class FixtureSim { this.ring(call.deviceId, 'approval') } + /** The user declines a held call: Sim settles it without the device and rings again. */ + decline(toolCallId: string): void { + const call = this.requireCall(toolCallId) + call.status = 'cancelled' + this.ring(call.deviceId, 'approval') + } + requireCall(toolCallId: string): FixtureCall { const call = this.calls.get(toolCallId) if (!call) throw new Error(`No fixture call ${toolCallId}`) @@ -219,11 +253,6 @@ export class FixtureSim { const url = new URL(request.url ?? '/', this.origin) const path = url.pathname this.requests.push(`${request.method} ${path}`) - if (path.startsWith('/api/desktop/') && this.offline) { - this.droppedWhileOffline += 1 - request.socket.destroy() - return - } if (path === '/api/auth/get-session') { this.json( response, @@ -404,6 +433,12 @@ export class FixtureSim { } } +/** + * Launches the app against the fixture Sim. Its shells are zsh with an empty config directory, so + * what they do is the app's and not this machine's: a developer's `.zshrc` can stop at a question + * (oh-my-zsh asks before updating) and hold the first prompt for as long as nobody answers it. + * The terminal's handling of slow and stalled startup files has its own tests. + */ export async function launch( sim: FixtureSim, userData: string, @@ -414,6 +449,8 @@ export async function launch( cwd: DESKTOP_DIR, env: { ...process.env, + SHELL: '/bin/zsh', + ZDOTDIR: mkdtempSync(join(tmpdir(), 'sim-e2e-zdotdir-')), SIM_DESKTOP_ORIGIN: sim.origin, SIM_DESKTOP_USER_DATA: userData, ...env, diff --git a/apps/desktop/e2e/fixtures/live-sim.ts b/apps/desktop/e2e/fixtures/live-sim.ts index 6a23821f013..a27873c9ac2 100644 --- a/apps/desktop/e2e/fixtures/live-sim.ts +++ b/apps/desktop/e2e/fixtures/live-sim.ts @@ -159,21 +159,23 @@ class HeldRequest { } /** - * 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. + * The origin Electron and Sim share. It forwards everything to Sim, records each request, can + * hold one until the test releases it, rewrite an answer or cut the network, 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 answerRewrites = new Map) => void>() private readonly sockets = new Set() + /** Every client connection open now, HTTP and upgraded alike. */ + private readonly connections = new Set() + private networkCut = false constructor(private readonly config: LiveSimConfig) { this.origin = `http://127.0.0.1:${config.proxyPort}` @@ -183,6 +185,15 @@ export class SimProxy { response.end(String(error)) }) }) + // While the network is cut, a new connection is reset as an unreachable host's would be. + this.server.on('connection', (socket: Socket) => { + if (this.networkCut) { + socket.destroy() + return + } + this.connections.add(socket) + socket.once('close', () => this.connections.delete(socket)) + }) // 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) @@ -236,9 +247,29 @@ export class SimProxy { 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 + /** + * Rewrites Sim's JSON answer to `path`, as a Sim configured otherwise would answer it; `undefined` + * stops rewriting it. + */ + rewriteAnswer( + path: string, + rewrite: ((body: Record) => void) | undefined + ): void { + if (rewrite) this.answerRewrites.set(path, rewrite) + else this.answerRewrites.delete(path) + } + + /** + * Cuts the network between the app and Sim: every open connection drops mid-flight, the + * doorbell stream included, and new ones are reset until `restoreNetwork`. + */ + cutNetwork(): void { + this.networkCut = true + for (const socket of this.connections) socket.destroy() + } + + restoreNetwork(): void { + this.networkCut = false } /** @@ -288,7 +319,7 @@ export class SimProxy { entry.status = 302 return } - let body = await readBody(request) + const 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) @@ -297,16 +328,20 @@ export class SimProxy { 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) } + const rewriteAnswer = this.answerRewrites.get(url.pathname) + // The body is forwarded whole, so it is sent with a length rather than chunked. An answer to + // rewrite is asked for uncompressed. + const { + 'transfer-encoding': _chunked, + 'accept-encoding': acceptEncoding, + ...forwarded + } = request.headers + const headers: IncomingHttpHeaders = { + ...forwarded, + ...(rewriteAnswer ? {} : { 'accept-encoding': acceptEncoding }), + 'content-length': String(body.length), + } await new Promise((resolve, reject) => { const clientGone = response.destroyed const upstream = httpRequest(target, { method, headers }, (upstreamResponse) => { @@ -318,6 +353,20 @@ export class SimProxy { upstreamResponse.resume() return } + if (rewriteAnswer && upstreamResponse.statusCode === 200) { + void readBody(upstreamResponse).then((raw) => { + const answer: Record = JSON.parse(raw.toString('utf8')) + rewriteAnswer(answer) + const rewritten = Buffer.from(JSON.stringify(answer)) + const { 'transfer-encoding': _chunked, ...answerHeaders } = upstreamResponse.headers + response.writeHead(200, { + ...answerHeaders, + 'content-length': String(rewritten.length), + }) + response.end(rewritten) + }, reject) + return + } response.writeHead(upstreamResponse.statusCode ?? 502, upstreamResponse.headers) upstreamResponse.pipe(response) }) @@ -661,14 +710,6 @@ export class SimDatabase { 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 }[]>` diff --git a/apps/desktop/e2e/fixtures/worker-restart.config.ts b/apps/desktop/e2e/fixtures/worker-restart.config.ts new file mode 100644 index 00000000000..d0f0e3a8b2d --- /dev/null +++ b/apps/desktop/e2e/fixtures/worker-restart.config.ts @@ -0,0 +1,10 @@ +import { defineConfig } from '@playwright/test' + +/** Runs only the worker-restart fixture, for `check-report.spec.ts`. */ +export default defineConfig({ + testDir: '.', + testMatch: 'worker-restart.ts', + workers: 1, + retries: 0, + reporter: [['line']], +}) diff --git a/apps/desktop/e2e/fixtures/worker-restart.ts b/apps/desktop/e2e/fixtures/worker-restart.ts new file mode 100644 index 00000000000..8a8c3921787 --- /dev/null +++ b/apps/desktop/e2e/fixtures/worker-restart.ts @@ -0,0 +1,26 @@ +import { expect, test } from '@playwright/test' +import { recordCheck } from '../check-report' + +/** + * Two checks in one file, the first failing, run by `check-report.spec.ts` in a Playwright run of + * its own. The failure makes Playwright start a new worker for the second, as it does mid-suite. + */ +const reportPath = process.env.WORKER_RESTART_REPORT_PATH + +function record(name: string, status: string): void { + recordCheck(reportPath, 'worker-restart', { + name, + status, + durationMs: 0, + error: `worker ${test.info().workerIndex}`, + }) +} + +test('fails first', () => { + record('fails first', 'failed') + expect('failed').toBe('passed') +}) + +test('passes after the worker restarts', () => { + record('passes after the worker restarts', 'passed') +}) diff --git a/apps/desktop/e2e/terminal-cancel.spec.ts b/apps/desktop/e2e/terminal-cancel.spec.ts index 9b1e08a9d04..1d639b7dacc 100644 --- a/apps/desktop/e2e/terminal-cancel.spec.ts +++ b/apps/desktop/e2e/terminal-cancel.spec.ts @@ -1,11 +1,12 @@ import { execFileSync } from 'node:child_process' -import { chmodSync, mkdtempSync, readdirSync, readFileSync, writeFileSync } from 'node:fs' +import { chmodSync, mkdtempSync, readdirSync, writeFileSync } from 'node:fs' import { tmpdir } from 'node:os' import { join } from 'node:path' import { type ElectronApplication, expect, type Page, test } from '@playwright/test' import type { SimDesktopApi } from '@sim/desktop-bridge' import { sleep } from '@sim/utils/helpers' import { randomInt } from '@sim/utils/random' +import { recordCheck } from './check-report' import { FixtureSim, launch, processRunning, registeredDevice, settled } from './executor-sim' /** @@ -31,32 +32,6 @@ const STUBBORN = "trap '' HUP INT TERM" const sim = new FixtureSim() /** The app's own log output for the current test, attached when it fails. */ const appOutput: string[] = [] -interface ReportCheck { - name: string - status: string - durationMs: number - retry: number -} - -/** - * Adds a scenario's outcome to the report at `TERMINAL_CANCEL_REPORT_PATH` as soon as it is known. - * Playwright replaces the worker after a failure, so the file, not the worker's memory, holds what - * came before: a failure and its retry both stay in it. - */ -function reportCheck(check: ReportCheck): void { - const reportPath = process.env.TERMINAL_CANCEL_REPORT_PATH - if (!reportPath) return - let checks: ReportCheck[] = [] - try { - checks = (JSON.parse(readFileSync(reportPath, 'utf8')) as { checks: ReportCheck[] }).checks - } catch { - // The first check of the run. - } - writeFileSync( - reportPath, - JSON.stringify({ suite: 'terminal-cancel', checks: [...checks, check] }, null, 2) - ) -} /** The shells, sleeps and tmux processes running, for a failure about who started what. */ function processes(): string { @@ -290,7 +265,7 @@ test.describe('terminal cancel', () => { test.afterEach(async () => { const testInfo = test.info() - reportCheck({ + recordCheck(process.env.TERMINAL_CANCEL_REPORT_PATH, 'terminal-cancel', { name: testInfo.title, status: testInfo.status ?? 'unknown', durationMs: testInfo.duration, diff --git a/apps/desktop/src/main/desktop-executor/service.ts b/apps/desktop/src/main/desktop-executor/service.ts index 6e27d8e1459..95b856be5e3 100644 --- a/apps/desktop/src/main/desktop-executor/service.ts +++ b/apps/desktop/src/main/desktop-executor/service.ts @@ -345,14 +345,14 @@ export function createDesktopExecutorService( ? { deviceId: id, protocolVersion: DESKTOP_EXECUTOR_PROTOCOL_VERSION } : null logger.info('Desktop executor registered', { enabled: nextTiming.enabled }) - // Off for this user, and nothing here to finish: stay dormant. No inbox, no doorbell; only + // Off on this Sim, and nothing here to finish: stay dormant. No inbox, no doorbell; only // a slow recheck, so switching the executor on in Sim reaches this device without a relaunch. if (!device && !executor && (await journal.load()).length === 0) { if (registrationGeneration !== generation) return scheduleRegistration(DORMANT_RECHECK_MS) return } - // Off for this user mid-session, or results from a previous run to deliver: a turn already + // Off on this Sim mid-session, or results from a previous run to deliver: a turn already // bound here still finishes, so the executor serves the inbox; no new turn binds. await startExecutor(id, nextTiming, registrationGeneration) } catch (error) { diff --git a/apps/desktop/src/main/terminal/index.ts b/apps/desktop/src/main/terminal/index.ts index 6dc00680ffe..20ee810f9ee 100644 --- a/apps/desktop/src/main/terminal/index.ts +++ b/apps/desktop/src/main/terminal/index.ts @@ -39,7 +39,7 @@ import { } from '@/main/resource-shortcuts' import { readForegroundProcessGroup, signalProcessGroup } from '@/main/terminal/process-group' import type { RunLedger } from '@/main/terminal/run-ledger' -import { elide, TerminalSession } from '@/main/terminal/session' +import { elide, type ShellStartupBounds, TerminalSession } from '@/main/terminal/session' import { activePane, awaitRun, @@ -65,11 +65,15 @@ import { const logger = createLogger('DesktopTerminal') /** - * How long to let a just-spawned shell finish its startup files before - * concluding it has no integration. Generous because a heavy `.zshrc` - * (nvm, pyenv, starship) can take a while on a cold start. + * How long a just-spawned shell gets to reach its first prompt. It begins our startup files + * before anything of the user's runs, so one that has not within seconds was never instrumented. + * The user's own files get far longer, since a heavy `.zshrc` (oh-my-zsh, nvm, pyenv) is slow on + * a cold start and slower on a busy machine; past that, it is waiting on something. */ -const SHELL_INTEGRATION_TIMEOUT_MS = 8_000 +const SHELL_STARTUP_BOUNDS: ShellStartupBounds = { unstartedMs: 8_000, startingMs: 30_000 } + +/** Enough of a stalled startup's screen to show what it is waiting on. */ +const STARTUP_SCREEN_LINES = 20 /** Grace for a program to react to input before its screen is worth reading. */ const INPUT_ECHO_MS = 250 @@ -1388,13 +1392,7 @@ export class TerminalService { throw new TerminalError('INVALID_REQUEST', 'run needs a `command`.') } if (!session.hasShellIntegration) { - await session.waitForShellIntegration(SHELL_INTEGRATION_TIMEOUT_MS) - } - if (!session.hasShellIntegration) { - throw new TerminalError( - 'NO_SHELL_INTEGRATION', - 'This shell did not load Sim shell integration, so command boundaries and exit codes cannot be determined. Ask the user to run the command themselves, or use a bash/zsh session.' - ) + await this.awaitShellStartup(session, latch) } if (session.isBusy) { throw new TerminalError( @@ -1408,6 +1406,31 @@ export class TerminalService { return session.runCommand(command, toolCallId, resolveRunWaitMs(args.waitSeconds)) } + /** Waits for a shell's first prompt, refusing the run with what it is doing if none comes. */ + private async awaitShellStartup(session: TerminalSession, latch: StopLatch): Promise { + const readiness = await session.waitForShellIntegration(SHELL_STARTUP_BOUNDS, latch.signal) + switch (readiness) { + case 'ready': + return + case 'stopped': + throw stoppedBeforeStart() + case 'exited': + throw new TerminalError('SESSION_CLOSED', 'The shell exited before it reached a prompt.') + case 'starting': { + const screen = (await session.readScrollback(STARTUP_SCREEN_LINES)).output.trim() + throw new TerminalError( + 'NO_SHELL_INTEGRATION', + `The shell has been running its startup files for over ${SHELL_STARTUP_BOUNDS.startingMs / 1000} s without reaching a prompt, so nothing was run. Its screen:\n${screen || '(empty)'}\nIf it is waiting for an answer, ask the user to answer it in that terminal (terminalId ${session.terminalId}), then run the command again.` + ) + } + case 'not-instrumented': + throw new TerminalError( + 'NO_SHELL_INTEGRATION', + 'This shell did not load Sim shell integration, so command boundaries and exit codes cannot be determined. Ask the user to run the command themselves, or use a bash/zsh session.' + ) + } + } + private spawn( cwd: string, cols: number, diff --git a/apps/desktop/src/main/terminal/session.ts b/apps/desktop/src/main/terminal/session.ts index 00f816f6ff6..6c64e85d3af 100644 --- a/apps/desktop/src/main/terminal/session.ts +++ b/apps/desktop/src/main/terminal/session.ts @@ -261,6 +261,21 @@ interface PendingCommand { resolve(result: TerminalRunResult): void } +/** How long a new shell gets to reach its first prompt. */ +export interface ShellStartupBounds { + /** From spawn, to begin the startup files we generated; a shell that has not by then never will. */ + unstartedMs: number + /** From the moment they began, for those startup files to reach a prompt. */ + startingMs: number +} + +/** + * How a wait for the shell's first prompt ended: `ready` with integration live; `not-instrumented` + * when our startup files never ran; `starting` when they are still running at the bound, often + * because something in them is waiting for input; `exited`; or `stopped` by the caller. + */ +export type ShellReadiness = 'ready' | 'not-instrumented' | 'starting' | 'exited' | 'stopped' + export interface TerminalSessionCallbacks { onData(terminalId: string, data: string): void onState(): void @@ -313,6 +328,10 @@ export class TerminalSession { private columns: number private lines: number private shellIntegration = false + /** The shell has begun the startup files we generated, so its integration is on the way. */ + /** When the shell began the startup files we generated, or null while it has not. */ + private startupBegunAt: number | null = null + private readonly spawnedAt = Date.now() private altScreen = false private foregroundCommand: string | null = null private foregroundToolCallId: string | null = null @@ -325,7 +344,8 @@ export class TerminalSession { private pendingCommand: PendingCommand | null = null /** Command line reported by the shell but not yet bracketed by output-start. */ private announcedCommand: string | null = null - private integrationWaiters: Array<() => void> = [] + /** Pending shell-readiness waits, rechecked whenever the shell's startup moves on. */ + private readonly readinessWaiters = new Set<() => void>() private constructor( options: TerminalSessionOptions, @@ -333,7 +353,8 @@ export class TerminalSession { integrationDir: string, nonce: string, shellName: string, - shellEnv: NodeJS.ProcessEnv + shellEnv: NodeJS.ProcessEnv, + private readonly instrumented: boolean ) { this.callbacks = options.callbacks this.terminalId = options.terminalId @@ -378,7 +399,15 @@ export class TerminalSession { }) logger.info('Started terminal session', { shell: shellPath, instrumented: shell !== null }) - return new TerminalSession(options, pty, integrationDir, nonce, shell ?? shellPath, shellEnv) + return new TerminalSession( + options, + pty, + integrationDir, + nonce, + shell ?? shellPath, + shellEnv, + shell !== null + ) } /** @@ -486,27 +515,57 @@ export class TerminalSession { } /** - * Resolves once the shell has emitted its first integration marker, or when - * `timeoutMs` elapses. A shell takes a few hundred milliseconds to run its - * startup files, so a command issued immediately after spawn would otherwise - * be refused for having no integration when it is merely early. + * Resolves once the shell reaches its first prompt with integration live, or once it is clear + * it will not: it never began our startup files within `bounds.unstartedMs` of spawning, it is + * still running them `bounds.startingMs` after they began, it exited, or `signal` stopped the + * wait. A shell's startup files take a while, and longer on a busy machine, so a command issued + * soon after spawn would otherwise be refused when the shell is merely early. */ - waitForShellIntegration(timeoutMs: number): Promise { - if (this.shellIntegration) return Promise.resolve(true) - if (this.disposed) return Promise.resolve(false) + waitForShellIntegration( + bounds: ShellStartupBounds, + signal?: AbortSignal + ): Promise { return new Promise((resolve) => { - const notify = () => { - clearTimeout(timer) - resolve(this.shellIntegration) + let timer: NodeJS.Timeout | null = null + const check = () => { + if (timer) clearTimeout(timer) + timer = null + const readiness = signal?.aborted ? 'stopped' : this.readiness(bounds) + if (readiness === null) { + const deadline = + this.startupBegunAt === null + ? this.spawnedAt + bounds.unstartedMs + : this.startupBegunAt + bounds.startingMs + timer = setTimeout(check, Math.max(0, deadline - Date.now())) + return + } + this.readinessWaiters.delete(check) + signal?.removeEventListener('abort', check) + resolve(readiness) } - const timer = setTimeout(() => { - this.integrationWaiters = this.integrationWaiters.filter((entry) => entry !== notify) - resolve(this.shellIntegration) - }, timeoutMs) - this.integrationWaiters.push(notify) + this.readinessWaiters.add(check) + signal?.addEventListener('abort', check) + check() }) } + /** Where the shell's startup stands against `bounds`, or null while it may still get there. */ + private readiness(bounds: ShellStartupBounds): ShellReadiness | null { + if (this.shellIntegration) return 'ready' + if (this.disposed) return 'exited' + const now = Date.now() + if (this.startupBegunAt === null) { + return !this.instrumented || now - this.spawnedAt >= bounds.unstartedMs + ? 'not-instrumented' + : null + } + return now - this.startupBegunAt >= bounds.startingMs ? 'starting' : null + } + + private notifyReadinessWaiters(): void { + for (const check of [...this.readinessWaiters]) check() + } + write(data: string): void { if (this.disposed) return this.pty.write(data) @@ -739,6 +798,7 @@ export class TerminalSession { if (this.flushTimer) clearTimeout(this.flushTimer) this.flushTimer = null this.finishCommand(null) + this.notifyReadinessWaiters() try { this.pty.kill() } catch { @@ -794,13 +854,17 @@ export class TerminalSession { marker: ReturnType['markers'][number] ): void { switch (marker.kind) { + case 'startup': + if (this.startupBegunAt === null) { + this.startupBegunAt = Date.now() + this.notifyReadinessWaiters() + } + break case 'prompt-start': if (!this.shellIntegration) { this.shellIntegration = true this.emitState() - const waiters = this.integrationWaiters - this.integrationWaiters = [] - for (const notify of waiters) notify() + this.notifyReadinessWaiters() } break case 'command-line': @@ -1030,9 +1094,7 @@ export class TerminalSession { private handleExit(): void { if (this.disposed) return this.disposed = true - const waiters = this.integrationWaiters - this.integrationWaiters = [] - for (const notify of waiters) notify() + this.notifyReadinessWaiters() if (this.flushTimer) clearTimeout(this.flushTimer) this.flushTimer = null this.flush() diff --git a/apps/desktop/src/main/terminal/shell-integration.test.ts b/apps/desktop/src/main/terminal/shell-integration.test.ts index ab47bf5a8ab..573f601484a 100644 --- a/apps/desktop/src/main/terminal/shell-integration.test.ts +++ b/apps/desktop/src/main/terminal/shell-integration.test.ts @@ -1,5 +1,13 @@ +import { spawnSync } from 'node:child_process' +import { existsSync, mkdtempSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import { describe, expect, it } from 'vitest' -import { detectShell, ShellIntegrationParser } from '@/main/terminal/shell-integration' +import { + buildShellLaunch, + detectShell, + ShellIntegrationParser, +} from '@/main/terminal/shell-integration' const NONCE = 'testnonce' @@ -65,6 +73,14 @@ describe('ShellIntegrationParser', () => { expect(markers).toEqual([{ kind: 'cwd', cwd: '/tmp/some dir' }]) }) + it("recognises the startup marker sent before the user's startup files", () => { + const parser = new ShellIntegrationParser(NONCE) + const { text, markers } = parser.parse(`${osc(`SimStartup;${NONCE}`)}loading`) + + expect(text).toBe('loading') + expect(markers).toEqual([{ kind: 'startup' }]) + }) + it('accepts ST as well as BEL as a terminator', () => { const parser = new ShellIntegrationParser(NONCE) const { text, markers } = parser.parse(`a\u001b]633;A;${NONCE}\u001b\\b`) @@ -93,3 +109,48 @@ describe('detectShell', () => { expect(detectShell('/bin/sh')).toBeNull() }) }) + +/** + * The generated startup files run in a real shell, with the user's own files standing in as a + * temp home whose startup prints a line. The startup marker has to reach the terminal before that + * line, and only from the interactive shell the terminal runs. + */ +describe('startup marker', () => { + const STARTUP = `\u001b]633;SimStartup;${NONCE}\u0007` + + function userHome(rcFile: string): string { + const home = mkdtempSync(join(tmpdir(), 'sim-shell-home-')) + writeFileSync(join(home, rcFile), 'echo user-startup\n') + return home + } + + it.skipIf(!existsSync('/bin/zsh'))( + "zsh sends it before the user's files, interactively only", + () => { + const home = userHome('.zshrc') + const env = { PATH: process.env.PATH ?? '/usr/bin:/bin', HOME: home } + const launch = buildShellLaunch('zsh', mkdtempSync(join(tmpdir(), 'sim-zsh-')), NONCE, env) + const run = (args: string[]) => + spawnSync('/bin/zsh', args, { env: { ...env, ...launch.env }, encoding: 'utf8' }).stdout + + const interactive = run([...launch.args, '-i', '-c', 'true']) + expect(interactive.indexOf(STARTUP)).toBeGreaterThanOrEqual(0) + expect(interactive.indexOf(STARTUP)).toBeLessThan(interactive.indexOf('user-startup')) + expect(run(['-c', 'echo script'])).toBe('script\n') + } + ) + + it.skipIf(!existsSync('/bin/bash'))("bash sends it before the user's files", () => { + const home = userHome('.bashrc') + const env = { PATH: process.env.PATH ?? '/usr/bin:/bin', HOME: home } + const launch = buildShellLaunch('bash', mkdtempSync(join(tmpdir(), 'sim-bash-')), NONCE, env) + const output = spawnSync('/bin/bash', launch.args, { + env: { ...env, ...launch.env }, + input: 'exit\n', + encoding: 'utf8', + }).stdout + + expect(output.indexOf(STARTUP)).toBeGreaterThanOrEqual(0) + expect(output.indexOf(STARTUP)).toBeLessThan(output.indexOf('user-startup')) + }) +}) diff --git a/apps/desktop/src/main/terminal/shell-integration.ts b/apps/desktop/src/main/terminal/shell-integration.ts index f69e86e1d2d..b0f770f430e 100644 --- a/apps/desktop/src/main/terminal/shell-integration.ts +++ b/apps/desktop/src/main/terminal/shell-integration.ts @@ -36,9 +36,13 @@ export function createNonce(): string { /** * Marker kinds we act on. `A` (prompt start) doubles as the "integration is * live" signal; `C`/`D` bracket a command's output; `E` reports the exact - * command line; `P` tracks the working directory across `cd`. + * command line; `P` tracks the working directory across `cd`. `SimStartup` is + * ours, not VS Code's: the shell sends it before running the user's startup + * files, so a shell that is still starting can be told from one we never + * instrumented. */ export type ShellMarker = + | { kind: 'startup' } | { kind: 'prompt-start' } | { kind: 'command-line'; command: string } | { kind: 'output-start' } @@ -123,6 +127,8 @@ export class ShellIntegrationParser { if (nonce !== this.nonce) return null switch (kind) { + case 'SimStartup': + return { kind: 'startup' } case 'A': return { kind: 'prompt-start' } case 'C': @@ -163,9 +169,14 @@ function writeZshFiles(dir: string, nonce: string, originalZdotdir: string): voi const sourceOriginal = (file: string) => `[ -f "$SIM_ZDOTDIR_ORIG/${file}" ] && builtin source "$SIM_ZDOTDIR_ORIG/${file}"` + // `.zshenv` is the first file zsh reads, so the startup marker goes out before any of the user's + // files run. Only from an interactive shell: a script's output must not carry it. writeFileSync( join(dir, '.zshenv'), - `SIM_ZDOTDIR_ORIG="\${SIM_ZDOTDIR_ORIG:-${originalZdotdir}}"\n${sourceOriginal('.zshenv')}\n` + `[[ -o interactive ]] && builtin printf '\\e]633;SimStartup;%s\\a' '${nonce}' +SIM_ZDOTDIR_ORIG="\${SIM_ZDOTDIR_ORIG:-${originalZdotdir}}" +${sourceOriginal('.zshenv')} +` ) writeFileSync(join(dir, '.zprofile'), `${sourceOriginal('.zprofile')}\n`) writeFileSync(join(dir, '.zlogin'), `${sourceOriginal('.zlogin')}\n`) @@ -225,7 +236,8 @@ function writeBashFile(dir: string, nonce: string): string { const rcPath = join(dir, 'sim-bash-rc.sh') writeFileSync( rcPath, - `[ -f "$HOME/.bashrc" ] && builtin source "$HOME/.bashrc" + `builtin printf '\\e]633;SimStartup;%s\\a' '${nonce}' +[ -f "$HOME/.bashrc" ] && builtin source "$HOME/.bashrc" __sim_nonce='${nonce}' __sim_in_cmd='' diff --git a/apps/desktop/src/main/terminal/shell-startup.test.ts b/apps/desktop/src/main/terminal/shell-startup.test.ts new file mode 100644 index 00000000000..4df9be581ab --- /dev/null +++ b/apps/desktop/src/main/terminal/shell-startup.test.ts @@ -0,0 +1,155 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { TerminalService } from '@/main/terminal' + +/** + * A new shell's first prompt, as the service waits for it before running the agent's command. + * node-pty is stubbed so each test plays the shell's side: the startup marker our generated files + * send before the user's run, whatever the user's files print, and the first prompt. The service + * and the session are real. + */ +const pty = vi.hoisted(() => ({ + emit: null as ((data: string) => void) | null, + writes: [] as string[], +})) + +vi.mock('@lydell/node-pty', () => ({ + spawn: () => ({ + pid: 4321, + onData: (handler: (data: string) => void) => { + pty.emit = handler + }, + onExit: () => {}, + write: (data: string) => pty.writes.push(data), + resize: vi.fn(), + kill: vi.fn(), + pause: vi.fn(), + resume: vi.fn(), + }), +})) + +vi.mock('@/main/terminal/shell-integration', async (importOriginal) => { + const actual = await importOriginal() + return { ...actual, createNonce: () => 'nonce' } +}) + +vi.mock('@/main/terminal/tmux', async (importOriginal) => { + const actual = await importOriginal() + return { ...actual, isTmuxUnavailable: () => true } +}) + +const marker = (body: string) => `\u001b]633;${body};nonce\u0007` +const STARTUP = marker('SimStartup') +const PROMPT = marker('A') + +function shell(data: string): void { + pty.emit?.(data) +} + +/** Plays a command the shell was asked to run through to its prompt, once it was typed. */ +async function answerCommand(output: string): Promise { + await vi.waitFor(() => expect(pty.writes.some((write) => write.endsWith('\r'))).toBe(true)) + shell(`${marker('C')}${output}\r\n${marker('D;0')}${PROMPT}`) +} + +/** Lets timers run until `promise` settles: the screen is read through an emulator that parses on them. */ +async function settle(promise: Promise): Promise { + let done = false + void promise.finally(() => { + done = true + }) + for (let step = 0; step < 1_000 && !done; step++) await vi.advanceTimersByTimeAsync(10) + return promise +} + +beforeEach(() => { + vi.stubEnv('SHELL', '/bin/zsh') + vi.useFakeTimers({ + toFake: ['setTimeout', 'clearTimeout', 'setInterval', 'clearInterval', 'Date'], + }) +}) + +afterEach(() => { + vi.useRealTimers() + pty.emit = null + pty.writes.length = 0 + vi.unstubAllEnvs() +}) + +describe('a shell that is still starting', () => { + it('runs the command once slow startup files reach a prompt', async () => { + // A heavy .zshrc on a busy machine: our files began at once, the user's took 20 s. + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const running = terminal.executeTool('call-slow', 'run', { + command: 'echo hi', + waitSeconds: 30, + }) + shell(STARTUP) + await vi.advanceTimersByTimeAsync(20_000) + expect(pty.writes).toEqual([]) + + shell(PROMPT) + await answerCommand('hi') + const response = await running + + expect(response).toMatchObject({ ok: true, result: { exitCode: 0 } }) + terminal.dispose() + }) + + it('gives startup files their full bound from when they began, however late that was', async () => { + // A machine so busy our files began 7 s after spawn, and the user's then took 25 s more. + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const running = terminal.executeTool('call-late', 'run', { command: 'echo hi' }) + await vi.advanceTimersByTimeAsync(7_000) + shell(STARTUP) + await vi.advanceTimersByTimeAsync(25_000) + expect(pty.writes).toEqual([]) + + shell(PROMPT) + await answerCommand('hi') + const response = await running + + expect(response).toMatchObject({ ok: true, result: { exitCode: 0 } }) + terminal.dispose() + }) + + it('refuses with the screen when startup files stall on a question', async () => { + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const running = terminal.executeTool('call-stalled', 'run', { command: 'echo hi' }) + shell(`${STARTUP}[oh-my-zsh] Would you like to update? [Y/n] `) + await vi.advanceTimersByTimeAsync(30_000) + const response = await settle(running) + + expect(response).toMatchObject({ ok: false, code: 'NO_SHELL_INTEGRATION' }) + expect(response.error).toContain('[oh-my-zsh] Would you like to update? [Y/n]') + expect(response.error).toContain('ask the user to answer it') + expect(pty.writes).toEqual([]) + terminal.dispose() + }) + + it('refuses a shell that never began our startup files at the short bound', async () => { + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const running = terminal.executeTool('call-plain', 'run', { command: 'echo hi' }) + await vi.advanceTimersByTimeAsync(8_000) + const response = await running + + expect(response).toMatchObject({ ok: false, code: 'NO_SHELL_INTEGRATION' }) + expect(response.error).toContain('did not load Sim shell integration') + terminal.dispose() + }) + + it('ends the wait as soon as the call is stopped, without running anything', async () => { + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const running = terminal.executeTool('call-stop', 'run', { command: 'echo hi' }) + shell(STARTUP) + await vi.advanceTimersByTimeAsync(1_000) + + expect(await terminal.cancelTool('call-stop')).toBe(true) + const response = await running + + expect(response).toMatchObject({ ok: false, code: 'CANCELLED' }) + shell(PROMPT) + await vi.advanceTimersByTimeAsync(1_000) + expect(pty.writes).toEqual([]) + terminal.dispose() + }) +}) diff --git a/apps/sim/app/o/[organizationId]/home/components/composer/composer.test.tsx b/apps/sim/app/o/[organizationId]/home/components/composer/composer.test.tsx index e78a1c35fbe..aa1466b9499 100644 --- a/apps/sim/app/o/[organizationId]/home/components/composer/composer.test.tsx +++ b/apps/sim/app/o/[organizationId]/home/components/composer/composer.test.tsx @@ -235,7 +235,6 @@ async function render( 'table-row-ttl': false, 'mothership-model-selector': mocks.advanced, 'mothership-plan-mode': mocks.plan, - 'mothership-desktop-background-executor': false, }} > @@ -314,7 +313,6 @@ it.each([ 'table-row-ttl': false, 'mothership-model-selector': false, 'mothership-plan-mode': planEnabled, - 'mothership-desktop-background-executor': false, }} > @@ -410,7 +408,6 @@ it('keeps restored queued skills scoped when replacing a draft', async () => { 'table-row-ttl': false, 'mothership-model-selector': mocks.advanced, 'mothership-plan-mode': mocks.plan, - 'mothership-desktop-background-executor': false, }} > diff --git a/apps/sim/app/o/[organizationId]/layout.tsx b/apps/sim/app/o/[organizationId]/layout.tsx index 14ea8b747ca..3ad21c07b64 100644 --- a/apps/sim/app/o/[organizationId]/layout.tsx +++ b/apps/sim/app/o/[organizationId]/layout.tsx @@ -79,7 +79,6 @@ export default async function OrganizationLayout({ 'table-row-ttl': tableRowTtlEnabled, 'mothership-model-selector': modelSelectorEnabled, 'mothership-plan-mode': planModeEnabled, - 'mothership-desktop-background-executor': false, }} > diff --git a/apps/sim/app/workspace/[workspaceId]/layout.test.tsx b/apps/sim/app/workspace/[workspaceId]/layout.test.tsx index 9a884ac45b9..6cdd0319416 100644 --- a/apps/sim/app/workspace/[workspaceId]/layout.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/layout.test.tsx @@ -9,6 +9,8 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' const { mockBrandingProvider, + mockIsDesktopPresenceAvailable, + mockWorkspaceChrome, mockGetOrgWhitelabelSettings, mockPrefetchWorkspaceHostContext, mockPrefetchWorkspaceSidebar, @@ -16,6 +18,10 @@ const { mockPrefetchWorkspaceForkAvailability, } = vi.hoisted(() => ({ mockBrandingProvider: vi.fn(({ children }: { children: ReactNode }) => children), + mockIsDesktopPresenceAvailable: vi.fn(() => false), + mockWorkspaceChrome: vi.fn( + ({ children }: { children: ReactNode; sidebar: ReactNode }) => children + ), mockGetOrgWhitelabelSettings: vi.fn(), mockPrefetchWorkspaceHostContext: vi.fn(), mockPrefetchWorkspaceSidebar: vi.fn(), @@ -64,7 +70,11 @@ vi.mock('@/app/workspace/[workspaceId]/components/session-expired', () => ({ })) vi.mock('@/app/workspace/[workspaceId]/components/workspace-chrome', () => ({ - WorkspaceChrome: ({ children }: { children: ReactNode }) => children, + WorkspaceChrome: mockWorkspaceChrome, +})) + +vi.mock('@/lib/desktop/executor/presence', () => ({ + isDesktopPresenceAvailable: mockIsDesktopPresenceAvailable, })) vi.mock('@/app/workspace/[workspaceId]/w/components/sidebar/sidebar', () => ({ @@ -202,4 +212,21 @@ describe('WorkspaceLayout host context', () => { expect(mockPrefetchWorkspaceAccess).not.toHaveBeenCalled() expect(mockGetOrgWhitelabelSettings).not.toHaveBeenCalled() }) + + it.each([true, false])( + 'tells the sidebar the desktop executor runs only where presence is tracked (%s)', + async (presenceAvailable) => { + mockIsDesktopPresenceAvailable.mockReturnValue(presenceAvailable) + mockWorkspaceChrome.mockClear() + + const element = await WorkspaceLayout({ + children:
Workspace child
, + params: Promise.resolve({ workspaceId: 'workspace-b' }), + }) + renderToStaticMarkup(element) + + const { sidebar } = mockWorkspaceChrome.mock.calls[0][0] + expect(sidebar).toMatchObject({ props: { desktopExecutorAvailable: presenceAvailable } }) + } + ) }) diff --git a/apps/sim/app/workspace/[workspaceId]/layout.tsx b/apps/sim/app/workspace/[workspaceId]/layout.tsx index 006daa63cfd..fe983f71dc1 100644 --- a/apps/sim/app/workspace/[workspaceId]/layout.tsx +++ b/apps/sim/app/workspace/[workspaceId]/layout.tsx @@ -5,7 +5,7 @@ import { SettingsNavigationProvider } from '@/components/settings/settings-navig import { getSession } from '@/lib/auth' import { getActiveOrganizationId } from '@/lib/auth/session-response' import { isDashboardsEnabled } from '@/lib/dashboards/feature-flag' -import { isDesktopBackgroundExecutorEnabled } from '@/lib/desktop/executor/flag' +import { isDesktopBackgroundExecutorAvailable } from '@/lib/desktop/executor/availability' import { isMothershipModelSelectorEnabled, isPlanModeEnabled } from '@/lib/mothership/feature-flags' import { resolveOrganizationEntryPath } from '@/lib/navigation/resolve-app-entry' import { isTableRowTtlEnabled } from '@/lib/table/ttl-availability' @@ -68,7 +68,6 @@ export default async function WorkspaceLayout({ planModeEnabled, organizationHref, dashboardsEnabled, - desktopBackgroundExecutorEnabled, ] = await Promise.all([ cookies(), hostContext.hostOrganizationId @@ -86,7 +85,6 @@ export default async function WorkspaceLayout({ isPlanModeEnabled(), resolveOrganizationEntryPath(session), isDashboardsEnabled(hostContext.hostOrganizationId), - isDesktopBackgroundExecutorEnabled(session.user.id), prefetchWorkspaceAccess(queryClient, workspaceId, principal), prefetchWorkspaceForkAvailability(queryClient, workspaceId, principal, hostContext), ]) @@ -100,7 +98,6 @@ export default async function WorkspaceLayout({ 'table-row-ttl': tableRowTtlEnabled, 'mothership-model-selector': modelSelectorEnabled, 'mothership-plan-mode': planModeEnabled, - 'mothership-desktop-background-executor': desktopBackgroundExecutorEnabled, }} > @@ -122,7 +119,12 @@ export default async function WorkspaceLayout({ } + sidebar={ + + } initialSidebarCollapsed={initialSidebarCollapsed} > {children} diff --git a/apps/sim/app/workspace/[workspaceId]/providers/feature-flags-provider.tsx b/apps/sim/app/workspace/[workspaceId]/providers/feature-flags-provider.tsx index ac6eb96fcae..d27ba042643 100644 --- a/apps/sim/app/workspace/[workspaceId]/providers/feature-flags-provider.tsx +++ b/apps/sim/app/workspace/[workspaceId]/providers/feature-flags-provider.tsx @@ -7,7 +7,6 @@ export interface WorkspaceFeatureFlags { 'table-row-ttl': boolean 'mothership-model-selector': boolean 'mothership-plan-mode': boolean - 'mothership-desktop-background-executor': boolean } const FeatureFlagsContext = createContext(null) diff --git a/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx b/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx index 2f8353dfb03..8010fb27410 100644 --- a/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx +++ b/apps/sim/app/workspace/[workspaceId]/w/components/sidebar/sidebar.tsx @@ -372,6 +372,8 @@ const DRAG_EXEMPT_CLASS = '[-webkit-app-region:no-drag]' interface SidebarProps { organizationHref: string | null + /** Whether this install runs the desktop background executor, so chats show desktop activity. */ + desktopExecutorAvailable: boolean } /** @@ -390,7 +392,10 @@ interface SidebarProps { * * @returns Sidebar with workflows panel */ -export const Sidebar = memo(function Sidebar({ organizationHref }: SidebarProps) { +export const Sidebar = memo(function Sidebar({ + organizationHref, + desktopExecutorAvailable, +}: SidebarProps) { const { isCollapsed: isCollapsedProp, isPeeking } = useSidebarChrome() const isCollapsed = isCollapsedProp && !isPeeking const params = useParams() @@ -885,15 +890,14 @@ export const Sidebar = memo(function Sidebar({ organizationHref }: SidebarProps) { enabled: chatEnabled && !permissionConfig.hideCopilot } ) - const desktopExecutorEnabled = useFeatureFlag('mothership-desktop-background-executor') useMothershipChatEvents( workspaceId, chatEnabled && !permissionConfig.hideCopilot, - desktopExecutorEnabled + desktopExecutorAvailable ) const { data: desktopActivity } = useDesktopActivity( workspaceId, - desktopExecutorEnabled && chatEnabled && !permissionConfig.hideCopilot + desktopExecutorAvailable && chatEnabled && !permissionConfig.hideCopilot ) const desktopActivityByChat = useMemo( () => new Map((desktopActivity ?? []).map((activity) => [activity.chatId, activity])), diff --git a/apps/sim/hooks/queries/desktop-activity.ts b/apps/sim/hooks/queries/desktop-activity.ts index a583838ab88..494f2de62ae 100644 --- a/apps/sim/hooks/queries/desktop-activity.ts +++ b/apps/sim/hooks/queries/desktop-activity.ts @@ -11,8 +11,8 @@ import { desktopActivityKeys } from '@/hooks/queries/utils/desktop-activity-keys const DESKTOP_ACTIVITY_STALE_TIME = 10 * 1000 /** * Presence, approvals and new background turns change without a chat event this query hears, so - * it is re-read on a timer: often while a desktop runs a chat, rarely otherwise. Only users in the - * executor's rollout ever run it. + * it is re-read on a timer: often while a desktop runs a chat, rarely otherwise. Only installs that + * run the executor ever read it. */ const DESKTOP_ACTIVITY_ACTIVE_REFETCH_MS = 15 * 1000 const DESKTOP_ACTIVITY_IDLE_REFETCH_MS = 30 * 1000 diff --git a/apps/sim/hooks/use-mothership-chat-events.ts b/apps/sim/hooks/use-mothership-chat-events.ts index 145da3b4eae..02166feff8b 100644 --- a/apps/sim/hooks/use-mothership-chat-events.ts +++ b/apps/sim/hooks/use-mothership-chat-events.ts @@ -258,7 +258,7 @@ export function reflectBackgroundChatStatus( export function useMothershipChatEvents( owner: MothershipChatOwner | undefined, chatEnabled: boolean, - /** Announce chats that finish in the background; only with the background executor on. */ + /** Announce chats that finish in the background; only where the background executor runs. */ announceBackgroundCompletions = false ) { const queryClient = useQueryClient() diff --git a/apps/sim/lib/api/contracts/desktop-executor.ts b/apps/sim/lib/api/contracts/desktop-executor.ts index fbc10ea159f..deba49ea953 100644 --- a/apps/sim/lib/api/contracts/desktop-executor.ts +++ b/apps/sim/lib/api/contracts/desktop-executor.ts @@ -30,8 +30,8 @@ const registerDesktopDeviceBodySchema = z.object({ export type RegisterDesktopDeviceBody = z.input /** - * `enabled: false` is the kill switch: the device keeps finishing calls on runs already bound to - * it, but no new turn binds to it. + * `enabled: false` means no new turn binds to the device, as on an install that cannot track + * desktop presence; the device still finishes calls on runs already bound to it. */ export const registerDesktopDeviceResponseSchema = z.object({ enabled: z.boolean(), diff --git a/apps/sim/lib/core/config/env.ts b/apps/sim/lib/core/config/env.ts index 32c747e15e1..0d1d4df54d8 100644 --- a/apps/sim/lib/core/config/env.ts +++ b/apps/sim/lib/core/config/env.ts @@ -601,7 +601,6 @@ export const env = createEnv({ DASHBOARDS: z.boolean().optional(), PROJECT_API_ENABLED: z.boolean().optional(), // Fallback for the `projects` feature flag off AppConfig MSHIP_MODEL_SELECTOR: z.boolean().optional(), - MSHIP_DESKTOP_BACKGROUND_EXECUTOR: z.boolean().optional(), // Fallback for the `mothership-desktop-background-executor` feature flag off AppConfig INBOX_ENABLED: z.boolean().optional(), // Enable inbox (Sim Mailer) on self-hosted (bypasses hosted requirements) SANDBOXES_ENABLED: z.boolean().optional(), // Enable custom sandboxes on self-hosted (bypasses hosted requirements) diff --git a/apps/sim/lib/core/config/feature-flags.test.ts b/apps/sim/lib/core/config/feature-flags.test.ts index a5e5ca1e1fd..ad1f84dde1d 100644 --- a/apps/sim/lib/core/config/feature-flags.test.ts +++ b/apps/sim/lib/core/config/feature-flags.test.ts @@ -44,7 +44,6 @@ setEnv({ TABLE_ROW_TTL: undefined, MSHIP_MODEL_SELECTOR: undefined, MSHIP_PLAN_MODE: undefined, - MSHIP_DESKTOP_BACKGROUND_EXECUTOR: undefined, AGENT_MEMORY_HISTORY: undefined, CREDENTIAL_GROUPS: undefined, KNOWLEDGE_MEMBER_ACCESS: undefined, diff --git a/apps/sim/lib/core/config/feature-flags.ts b/apps/sim/lib/core/config/feature-flags.ts index 4a823d1c6eb..566d3481832 100644 --- a/apps/sim/lib/core/config/feature-flags.ts +++ b/apps/sim/lib/core/config/feature-flags.ts @@ -63,14 +63,6 @@ const FEATURE_FLAGS = { 'and workspace surfaces.', fallback: 'MSHIP_PLAN_MODE', }, - 'mothership-desktop-background-executor': { - description: - "Let the Sim desktop app run a chat's desktop tools in the background, so work continues " + - 'after the user leaves the chat. Gates device registration and binding new turns to a ' + - 'device; runs already bound keep finishing when it is turned off. Supports userId and ' + - 'admin targeting; off-AppConfig falls back to MSHIP_DESKTOP_BACKGROUND_EXECUTOR.', - fallback: 'MSHIP_DESKTOP_BACKGROUND_EXECUTOR', - }, 'agent-memory-history': { description: 'Capture durable Workflow Agent tool history and continue existing retries. Supports workspace rollout targeting; version-aware memory storage remains active when capture is disabled.', diff --git a/apps/sim/lib/desktop/application/activity.integration.ts b/apps/sim/lib/desktop/application/activity.integration.ts index b587caf094c..855013e53da 100644 --- a/apps/sim/lib/desktop/application/activity.integration.ts +++ b/apps/sim/lib/desktop/application/activity.integration.ts @@ -2,7 +2,7 @@ * Background desktop activity against real PostgreSQL and Redis: which of a user's chats a desktop * runs, and whether each is running, waiting on the user's approval, or blocked by an offline desktop. */ -import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' +import { afterAll, describe, expect, it, vi } from 'vitest' const { redisUrl } = await vi.hoisted(async () => { const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') @@ -13,8 +13,6 @@ const { redisUrl } = await vi.hoisted(async () => { return { redisUrl: url } }) -vi.mock('@/lib/core/config/feature-flags', () => featureFlagsMock) - import type { SessionPrincipal } from '@sim/auth/principal' import { db } from '@sim/db' import { @@ -25,7 +23,6 @@ import { user, workspace, } from '@sim/db/schema' -import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock' import { generateId } from '@sim/utils/id' import { inArray } from 'drizzle-orm' import { closeRedisConnection } from '@/lib/core/config/redis' @@ -43,10 +40,6 @@ describe.runIf(Boolean(redisUrl))('background desktop activity', () => { const workspaceIds: string[] = [] const deviceIds: string[] = [] - beforeEach(() => { - featureFlagsMockFns.mockIsFeatureEnabled.mockResolvedValue(true) - }) - afterAll(async () => { if (workspaceIds.length) await db.delete(workspace).where(inArray(workspace.id, workspaceIds)) if (deviceIds.length) diff --git a/apps/sim/lib/desktop/application/executor.integration.ts b/apps/sim/lib/desktop/application/executor.integration.ts index dcd057447b7..7ef9a70eccf 100644 --- a/apps/sim/lib/desktop/application/executor.integration.ts +++ b/apps/sim/lib/desktop/application/executor.integration.ts @@ -6,15 +6,19 @@ */ import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' -const { redisUrl } = await vi.hoisted(async () => { +const { redisUrl, presence } = await vi.hoisted(async () => { const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') const url = readTestRedisUrl() /** The real Redis module reads this at import. */ if (url) process.env.REDIS_URL = url - return { redisUrl: url } + return { redisUrl: url, presence: { available: true } } }) -vi.mock('@/lib/core/config/feature-flags', () => featureFlagsMock) +/** Redis is always configured here; `presence.available` stands in for an install without it. */ +vi.mock('@/lib/desktop/executor/presence', async (importOriginal) => ({ + ...(await importOriginal()), + isDesktopPresenceAvailable: () => presence.available, +})) import type { SessionPrincipal } from '@sim/auth/principal' import { db } from '@sim/db' @@ -29,7 +33,6 @@ import { workspace, } from '@sim/db/schema' import { createDeferred } from '@sim/testing/helpers/deferred' -import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock' import { generateId } from '@sim/utils/id' import { compareStrings } from '@sim/utils/string' import { eq, inArray, sql } from 'drizzle-orm' @@ -69,8 +72,8 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () => * A signed-in desktop: its user (created unless `sameUserAs` names one), its Better Auth * session, and its install id, registered the way the app registers after sign-in. */ - async function signedInDesktop(sameUserAs?: string, newUserId?: string) { - const userId = sameUserAs ?? newUserId ?? generateId() + async function signedInDesktop(sameUserAs?: string) { + const userId = sameUserAs ?? generateId() const now = new Date() if (!sameUserAs) { userIds.push(userId) @@ -229,18 +232,8 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () => return settled } - /** The flag targets users; `excluded` is outside its rollout. */ - const excluded = new Set() beforeEach(() => { - featureFlagsMockFns.mockIsFeatureEnabled.mockImplementation( - async (flag, context) => - flag === 'mothership-desktop-background-executor' && - typeof context === 'object' && - context !== null && - 'userId' in context && - typeof context.userId === 'string' && - !excluded.has(context.userId) - ) + presence.available = true }) afterAll(async () => { @@ -255,10 +248,15 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () => }) describe('registration', () => { - it('writes nothing and reports the executor off for a user outside the rollout', async () => { - const userId = generateId() - excluded.add(userId) - const desktop = await signedInDesktop(undefined, userId) + it('reports the executor on wherever Sim can track presence', async () => { + const desktop = await signedInDesktop() + expect(desktop.enabled).toBe(true) + await expect(inbox(desktop)).resolves.toEqual({ items: [] }) + }) + + it('writes nothing and reports the executor off where Sim cannot track presence', async () => { + presence.available = false + const desktop = await signedInDesktop() expect(desktop.enabled).toBe(false) const [stored] = await db .select() diff --git a/apps/sim/lib/desktop/application/executor.ts b/apps/sim/lib/desktop/application/executor.ts index f70c844f05e..cc0d2e19bae 100644 --- a/apps/sim/lib/desktop/application/executor.ts +++ b/apps/sim/lib/desktop/application/executor.ts @@ -9,6 +9,7 @@ import { type CredentialUserAuditEntry, defineAuthorizedCredentialUserUseCase, } from '@/lib/credentials/application/authorized-user-use-case' +import { isDesktopBackgroundExecutorAvailable } from '@/lib/desktop/executor/availability' import { DESKTOP_EXECUTOR_PROTOCOL_VERSION, DESKTOP_INBOX_RECONCILE_MS, @@ -22,7 +23,6 @@ import { DesktopCallRevokedError, DesktopDeviceUnrecognizedError, } from '@/lib/desktop/executor/errors' -import { isDesktopBackgroundExecutorEnabled } from '@/lib/desktop/executor/flag' import { classifyDesktopInbox, type DesktopInboxEntry } from '@/lib/desktop/executor/inbox' import { markDesktopPresent } from '@/lib/desktop/executor/presence' import { @@ -72,15 +72,15 @@ async function requireBoundDevice(principal: SessionPrincipal, deviceId: string) } /** - * The device a new turn binds to: the composer's own, but only while the executor is on for this - * user and the device is registered to this very session as an executor. Anything else leaves + * The device a new turn binds to: the composer's own, but only while this install can run the + * executor and the device is registered to this very session as an executor. Anything else leaves * the turn to the chat view, as before the executor existed. */ export async function resolveTurnDesktopDevice( principal: SessionPrincipal, deviceId: string ): Promise { - if (!(await isDesktopBackgroundExecutorEnabled(principal.userId))) return null + if (!isDesktopBackgroundExecutorAvailable()) return null const device = await getBoundDesktopDevice( { deviceId, userId: principal.userId, sessionId: principal.sessionId }, { executor: true } @@ -103,8 +103,8 @@ interface RegisterDesktopDeviceInput extends DeviceInput { } /** - * Binds the install to this user and session. With the executor turned off nothing is written, - * and the device stays dormant. + * Binds the install to this user and session. Where the executor is unavailable nothing is + * written, and the device stays dormant. */ export const registerDesktopDevice = defineAuthorizedCredentialUserUseCase({ // permission-group-exempt: registering grants no access; each call's run was admitted under the Chat capability. @@ -120,7 +120,7 @@ export const registerDesktopDevice = defineAuthorizedCredentialUserUseCase({ principal: SessionPrincipal input: RegisterDesktopDeviceInput }) { - const enabled = await isDesktopBackgroundExecutorEnabled(principal.userId) + const enabled = isDesktopBackgroundExecutorAvailable() if (enabled) { const registered = await upsertDesktopDevice({ id: input.deviceId, diff --git a/apps/sim/lib/desktop/executor/availability.ts b/apps/sim/lib/desktop/executor/availability.ts new file mode 100644 index 00000000000..59ef02e81f3 --- /dev/null +++ b/apps/sim/lib/desktop/executor/availability.ts @@ -0,0 +1,10 @@ +import { isDesktopPresenceAvailable } from '@/lib/desktop/executor/presence' + +/** + * Whether new turns may bind to a desktop, and whether the workspace shows background desktop + * activity. The executor needs presence tracking, so an install without Redis keeps desktop tools + * in the chat window. Never consulted for a run already bound: it keeps running on its desktop. + */ +export function isDesktopBackgroundExecutorAvailable(): boolean { + return isDesktopPresenceAvailable() +} diff --git a/apps/sim/lib/desktop/executor/bound-turn.integration.ts b/apps/sim/lib/desktop/executor/bound-turn.integration.ts index f0830b84911..6a13a2be313 100644 --- a/apps/sim/lib/desktop/executor/bound-turn.integration.ts +++ b/apps/sim/lib/desktop/executor/bound-turn.integration.ts @@ -8,7 +8,7 @@ import { authMock, authMockFns } from '@sim/testing/mocks/auth.mock' import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' -const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => { +const { redisUrl, inheritedEnv, worker, presence } = await vi.hoisted(async () => { const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') const { createServer } = await import('node:http') /** Stands in for the agent worker, which Stop also tells to end the stream. */ @@ -30,11 +30,15 @@ const { redisUrl, inheritedEnv, worker } = await vi.hoisted(async () => { if (url) process.env.REDIS_URL = url process.env.COPILOT_TOOL_PERMISSIONS_ENABLED = 'true' process.env.SIM_AGENT_API_URL = `http://127.0.0.1:${port}` - return { redisUrl: url, inheritedEnv, worker: { server } } + return { redisUrl: url, inheritedEnv, worker: { server }, presence: { available: true } } }) vi.mock('@/lib/auth', () => authMock) -vi.mock('@/lib/core/config/feature-flags', () => featureFlagsMock) +/** Redis is always configured here; `presence.available` stands in for an install without it. */ +vi.mock('@/lib/desktop/executor/presence', async (importOriginal) => ({ + ...(await importOriginal()), + isDesktopPresenceAvailable: () => presence.available, +})) import type { SessionPrincipal } from '@sim/auth/principal' import { db } from '@sim/db' @@ -49,7 +53,6 @@ import { user, workspace, } from '@sim/db/schema' -import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock' import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' import { eq, inArray, sql } from 'drizzle-orm' @@ -357,9 +360,7 @@ describe.runIf(Boolean(redisUrl))("a turn bound to a desktop's background execut }) beforeEach(() => { - featureFlagsMockFns.mockIsFeatureEnabled.mockImplementation( - async (flag) => flag === 'mothership-desktop-background-executor' - ) + presence.available = true }) afterAll(async () => { @@ -372,7 +373,7 @@ describe.runIf(Boolean(redisUrl))("a turn bound to a desktop's background execut await db.delete(user).where(eq(user.id, userId)) }) - it('binds a turn only to an executor registered to this very session, with the flag on', async () => { + it('binds a turn only to an executor registered to this very session, where Sim tracks presence', async () => { const desktop = await signedInDesktop() const viewer = await signedInDesktop({ executor: 0 }) const otherSession: SessionPrincipal = { ...desktop.principal, sessionId: generateId() } @@ -382,7 +383,7 @@ describe.runIf(Boolean(redisUrl))("a turn bound to a desktop's background execut ) expect(await resolveTurnDesktopDevice(otherSession, desktop.deviceId)).toBeNull() expect(await resolveTurnDesktopDevice(viewer.principal, viewer.deviceId)).toBeNull() - featureFlagsMockFns.mockIsFeatureEnabled.mockResolvedValue(false) + presence.available = false expect(await resolveTurnDesktopDevice(desktop.principal, desktop.deviceId)).toBeNull() }) diff --git a/apps/sim/lib/desktop/executor/flag.ts b/apps/sim/lib/desktop/executor/flag.ts deleted file mode 100644 index 9b0b5fc1348..00000000000 --- a/apps/sim/lib/desktop/executor/flag.ts +++ /dev/null @@ -1,11 +0,0 @@ -import { isFeatureEnabled } from '@/lib/core/config/feature-flags' -import { isDesktopPresenceAvailable } from '@/lib/desktop/executor/presence' - -/** - * Whether this user's new turns may bind to a desktop, and whether the UI shows background desktop - * activity. Never consulted for a run already bound: it keeps running on its desktop until it ends. - */ -export async function isDesktopBackgroundExecutorEnabled(userId: string): Promise { - if (!isDesktopPresenceAvailable()) return false - return isFeatureEnabled('mothership-desktop-background-executor', { userId }) -} diff --git a/apps/sim/lib/desktop/executor/presence.ts b/apps/sim/lib/desktop/executor/presence.ts index 059ff43dd22..407b98f8b1c 100644 --- a/apps/sim/lib/desktop/executor/presence.ts +++ b/apps/sim/lib/desktop/executor/presence.ts @@ -14,7 +14,7 @@ function presenceKey(deviceId: string): string { return `desktop:presence:${deviceId}` } -/** Whether presence can be tracked at all; without Redis the executor stays off. */ +/** Whether presence can be tracked at all, which needs Redis. */ export function isDesktopPresenceAvailable(): boolean { return getRedisClient() !== null } diff --git a/apps/sim/scripts/test-desktop-inbox-e2e.ts b/apps/sim/scripts/test-desktop-inbox-e2e.ts index 8ce6c3da3c9..0f107054653 100644 --- a/apps/sim/scripts/test-desktop-inbox-e2e.ts +++ b/apps/sim/scripts/test-desktop-inbox-e2e.ts @@ -34,9 +34,8 @@ import { * real HTTP boundary: registration, the SSE doorbell, the inbox pull, claim, lease renewal and * completion. No chat view is open at any point. * - * Start the app with the executor flag on and Redis configured, for example: - * MSHIP_DESKTOP_BACKGROUND_EXECUTOR=true COPILOT_TOOL_PERMISSIONS_ENABLED=true \ - * REDIS_URL=redis://127.0.0.1:6379 bun run dev + * Start the app with Redis configured, which the executor needs, for example: + * COPILOT_TOOL_PERMISSIONS_ENABLED=true REDIS_URL=redis://127.0.0.1:6379 bun run dev * then run: * DESKTOP_INBOX_E2E_BASE_URL=http://127.0.0.1:3000 \ * DESKTOP_INBOX_E2E_DATABASE_URL=postgresql://postgres@127.0.0.1:5432/sim_test \ @@ -416,11 +415,7 @@ async function run() { await check('registers the desktop and receives the executor timing contract', async () => { const registration = await register(desktop) - assert.equal( - registration.enabled, - true, - 'Start the app with MSHIP_DESKTOP_BACKGROUND_EXECUTOR=true' - ) + assert.equal(registration.enabled, true, 'Start the app with REDIS_URL set') assert(registration.leaseRenewMs < registration.leaseMs) assert(registration.reconcileMs < PICKUP_GRACE_SECONDS * 1000) }) diff --git a/knip.jsonc b/knip.jsonc index 68c63368a0e..e3cbb2641cd 100644 --- a/knip.jsonc +++ b/knip.jsonc @@ -75,14 +75,16 @@ "apps/realtime": { "includeEntryExports": false, "entry": ["src/bootstrap.ts"] }, "apps/desktop": { "includeEntryExports": false, - // scripts/build.ts and e2e/updater.spec.ts supply these to esbuild by path. + // scripts/build.ts and e2e/updater.spec.ts supply these to esbuild by path; the nested + // Playwright run in e2e/check-report.spec.ts loads worker-restart.ts through its config. "entry": [ "src/main/index.ts", "src/preload/*.ts", "src/preload/browser/index.ts", "src/renderer/*/index.tsx", "scripts/*.ts", - "e2e/fixtures/updater.ts" + "e2e/fixtures/updater.ts", + "e2e/fixtures/worker-restart.ts" ], // ensure-pty-prebuilds.ts assembles these package names for universal builds. "ignoreDependencies": ["@lydell/node-pty-darwin-arm64", "@lydell/node-pty-darwin-x64"]