diff --git a/apps/desktop/src/main/desktop-executor/executor.test.ts b/apps/desktop/src/main/desktop-executor/executor.test.ts index f1ed7801d55..795d4df54de 100644 --- a/apps/desktop/src/main/desktop-executor/executor.test.ts +++ b/apps/desktop/src/main/desktop-executor/executor.test.ts @@ -168,6 +168,8 @@ function setup( const runner = new FakeRunner() const onUnregistered = vi.fn() const busy: boolean[] = [] + /** Calls whose result the executor reported as reaching the model, in order. */ + const delivered: string[] = [] const executor = new DesktopExecutor({ client: sim.client, journal, @@ -177,12 +179,13 @@ function setup( onUnregistered, onBusyChange: (value) => busy.push(value), onApprovals: (items) => approvals.push(items), + onResultDelivered: (toolCallId) => delivered.push(toolCallId), ...(options.maxHeldCalls ? { maxHeldCalls: options.maxHeldCalls } : {}), ...(options.deliveryAwakeLimitMs !== undefined ? { deliveryAwakeLimitMs: options.deliveryAwakeLimitMs } : {}), }) - return { sim, journal, runner, executor, onUnregistered, busy, approvals } + return { sim, journal, runner, executor, onUnregistered, busy, approvals, delivered } } describe('claiming', () => { @@ -203,6 +206,24 @@ describe('claiming', () => { expect(executor.heldCallCount()).toBe(0) }) + it('reports a result as reaching the model only once Sim takes it as the call own', async () => { + const { sim, journal, runner, executor, delivered } = setup() + runner.immediate = DONE + // Sim settled this call first: the result never reached the model. + sim.completionOutcome = 'superseded' + sim.inbox = [callItem('call-superseded', 'chat-a')] + await executor.reconcile() + await vi.waitFor(() => expect(sim.completions).toHaveLength(1)) + expect(delivered).toEqual([]) + + // Taken by Sim, even with nothing written locally (no OS encryption, say). + sim.completionOutcome = 'recorded' + journal.failOn = 'result' + sim.inbox = [callItem('call-recorded', 'chat-a')] + await executor.reconcile() + await vi.waitFor(() => expect(delivered).toEqual(['call-recorded'])) + }) + it('claims a whole backlog at once, before any of it runs', async () => { const { sim, runner, executor } = setup() sim.inbox = [callItem('a-1', 'chat-a'), callItem('a-2', 'chat-a'), callItem('a-3', 'chat-a')] diff --git a/apps/desktop/src/main/desktop-executor/executor.ts b/apps/desktop/src/main/desktop-executor/executor.ts index 0d56584c027..6ca794182cd 100644 --- a/apps/desktop/src/main/desktop-executor/executor.ts +++ b/apps/desktop/src/main/desktop-executor/executor.ts @@ -67,6 +67,11 @@ export interface DesktopExecutorOptions { onApprovals?: (items: DesktopApprovalItem[]) => void /** Called whenever the number of held calls changes between zero and more. */ onBusyChange?: (busy: boolean) => void + /** + * Called once Sim has taken a call's result as the call's own (recorded, or a duplicate of one + * it recorded): the model has it, so anything it hands back (a pane still running) is in use. + */ + onResultDelivered?: (toolCallId: string, completion: DesktopToolCompletion) => void maxHeldCalls?: number /** First delivery retry delay; tests shorten it. */ retryBaseMs?: number @@ -401,6 +406,8 @@ export class DesktopExecutor { completion: pending, }) logger.info('Desktop call result acknowledged', { toolCallId, outcome }) + // Superseded: Sim settled the call first, so this result never reached the model. + if (outcome !== 'superseded') this.options.onResultDelivered?.(toolCallId, pending) break } catch (error) { // Encoding failed on this machine, so nothing was sent; the same data would fail again. diff --git a/apps/desktop/src/main/desktop-executor/service.test.ts b/apps/desktop/src/main/desktop-executor/service.test.ts index 20ff6168e68..924ad2b551b 100644 --- a/apps/desktop/src/main/desktop-executor/service.test.ts +++ b/apps/desktop/src/main/desktop-executor/service.test.ts @@ -7,6 +7,7 @@ import { describe, expect, it, vi } from 'vitest' vi.mock('electron', () => import('@/test/electron-mock')) import { net } from 'electron' +import { createExecutorJournal } from '@/main/desktop-executor/journal' import { createDesktopExecutorService, deviceName } from '@/main/desktop-executor/service' /** Sim's device routes, with registration answers held until the test releases them. */ @@ -95,6 +96,36 @@ async function service(protocolVersion = 1, userDataPath?: string) { return { sim, desktopExecutor, busy } } +describe('results recovery will hand to the model', () => { + it('names the calls whose real result the journal holds for recovery to send', async () => { + const userData = await mkdtemp(join(tmpdir(), 'sim-executor-service-')) + const journal = createExecutorJournal(join(userData, 'desktop-executor-journal.json')) + await journal.put({ toolCallId: 'claimed', state: 'claimed', executionToken: 't1' }) + await journal.put({ toolCallId: 'started', state: 'started', executionToken: 't2' }) + await journal.put({ + toolCallId: 'unknown', + state: 'result', + executionToken: 't3', + completion: { status: 'error', message: 'x', data: { outcomeUnknown: true } }, + }) + await journal.put({ + toolCallId: 'not-started', + state: 'result', + executionToken: 't4', + completion: { status: 'error', message: 'x', data: { notStarted: true } }, + }) + await journal.put({ + toolCallId: 'handed-back', + state: 'result', + executionToken: 't5', + completion: { status: 'success', message: 'running', data: { status: 'running' } }, + }) + const { desktopExecutor } = await service(1, userData) + + expect([...(await desktopExecutor.pendingResults())]).toEqual(['handed-back']) + }) +}) + describe('desktop executor registration', () => { it('offers the device for binding once Sim enables it', async () => { const { sim, desktopExecutor } = await service() diff --git a/apps/desktop/src/main/desktop-executor/service.ts b/apps/desktop/src/main/desktop-executor/service.ts index 5812e06ae9f..7611dfd5c2f 100644 --- a/apps/desktop/src/main/desktop-executor/service.ts +++ b/apps/desktop/src/main/desktop-executor/service.ts @@ -7,6 +7,7 @@ import { hostname } from 'node:os' import { join } from 'node:path' import type { DesktopExecutorDevice } from '@sim/desktop-bridge' +import type { DesktopToolCompletion } from '@sim/desktop-bridge/tool-results' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' @@ -61,6 +62,8 @@ export interface DesktopExecutorServiceDeps { onApprovals?: (items: DesktopApprovalItem[]) => void /** Whether any chat has desktop work claimed on this machine changed. */ onBusyChange?: (busy: boolean) => void + /** Sim has taken a call's result as the call's own, so the model has it. */ + onResultDelivered?: (toolCallId: string, completion: DesktopToolCompletion) => void } export interface DesktopExecutorService { @@ -68,6 +71,12 @@ export interface DesktopExecutorService { /** Re-registers after a sign-in, a session change, or a change to what this device can run. */ refreshRegistration(): void getDevice(): DesktopExecutorDevice | null + /** + * The calls whose real result (not one reported as not started or outcome unknown) is in the + * journal and not yet acknowledged, so recovery will hand it to the model. Read before recovery + * changes the journal; empty when it cannot be read. + */ + pendingResults(): Promise> /** Stores one entry of a claimed import, as this device's registered session. */ importEntry( request: DesktopImportEntryRequest, @@ -239,6 +248,7 @@ export function createDesktopExecutorService( onUnregistered: handleUnrecognized, ...(deps.onApprovals ? { onApprovals: deps.onApprovals } : {}), ...(deps.onBusyChange ? { onBusyChange: deps.onBusyChange } : {}), + ...(deps.onResultDelivered ? { onResultDelivered: deps.onResultDelivered } : {}), }) await executor.recover() // Signed out while recovering: sign-out already disposed this executor. @@ -443,6 +453,22 @@ export function createDesktopExecutorService( getDevice() { return device }, + async pendingResults() { + const pending = new Set() + try { + for (const entry of await journal.load()) { + if (entry.state !== 'result') continue + const data = entry.completion.data + if (data?.outcomeUnknown === true || data?.notStarted === true) continue + pending.add(entry.toolCallId) + } + } catch (error) { + logger.warn('Could not read the executor journal for pending results', { + error: getErrorMessage(error), + }) + } + return pending + }, importEntry(request, signal) { if (!client) throw new Error('The Sim desktop app is not signed in to Sim.') return client.importEntry(request, signal) diff --git a/apps/desktop/src/main/index.ts b/apps/desktop/src/main/index.ts index 40d44cf2ce5..d777615929f 100644 --- a/apps/desktop/src/main/index.ts +++ b/apps/desktop/src/main/index.ts @@ -101,6 +101,7 @@ import { import { setShellTheme } from '@/main/shell-theme' import { attachTelemetryPolicy } from '@/main/telemetry-policy' import { TerminalRegistry } from '@/main/terminal/registry' +import { createRunLedger } from '@/main/terminal/run-ledger' import { installTray, type TrayHandle } from '@/main/tray' import { checkForUpdatesInteractive, initUpdater, type UpdaterHandle } from '@/main/updater' import { installBrowserUserAgent } from '@/main/user-agent' @@ -167,15 +168,20 @@ function main(): void { ), }) const scopeEvents = new ScopedEventRouter() - const terminal = new TerminalRegistry({ - load: (scopeId) => desktopChatSessions.getTerminal(processOrigin, scopeId) ?? undefined, - save: (scopeId, snapshot) => desktopChatSessions.setTerminal(processOrigin, scopeId, snapshot), - migrate: (fromScopeId, toScopeId) => - desktopChatSessions.migrateTerminal(processOrigin, fromScopeId, toScopeId), - disposeScope: (scopeId) => { - desktopChatSessions.deleteScope(processOrigin, scopeId) + const terminal = new TerminalRegistry( + { + load: (scopeId) => desktopChatSessions.getTerminal(processOrigin, scopeId) ?? undefined, + save: (scopeId, snapshot) => + desktopChatSessions.setTerminal(processOrigin, scopeId, snapshot), + migrate: (fromScopeId, toScopeId) => + desktopChatSessions.migrateTerminal(processOrigin, fromScopeId, toScopeId), + disposeScope: (scopeId) => { + desktopChatSessions.deleteScope(processOrigin, scopeId) + }, }, - }) + undefined, + createRunLedger(join(userDataPath, 'terminal-runs')) + ) const preloadPath = join(__dirname, 'preload.cjs') const windows = new Set() @@ -601,6 +607,13 @@ function main(): void { accountDataAvailable, onApprovals: (items) => approvalNotifier.update(items), onBusyChange: (busy) => sleepBlocker.setBusy(busy), + // A result the model has (not one reported as not started or outcome unknown) makes a tmux + // run it handed back as still going collectable across a restart. + onResultDelivered: (toolCallId, completion) => { + if (completion.data?.outcomeUnknown !== true && completion.data?.notStarted !== true) { + terminal.markRunDelivered(toolCallId) + } + }, runner: createDesktopToolRunner({ preferences: () => desktopSettings.getPreferences(), accountDataAvailable, @@ -804,6 +817,14 @@ function main(): void { } } + // The same user's tmux runs from a previous process: a run whose call never handed back its + // result (or whose result the journal will report as unknown) has nothing left to collect what + // it does, so it is stopped, while its pane still carries its tag. A run already handed back as + // still going, with its pane, is left to the model, which may come back to it. Read before the + // executor starts, since its recovery rewrites the journal. + const pendingResults = desktopExecutor.pendingResults() + void pendingResults.then((pending) => terminal.stopUncollectableRuns(pending)) + if (!accountDataAvailable()) { logger.warn( 'Account-bearing browser, terminal, and local filesystem APIs are unavailable until local recovery succeeds' @@ -953,6 +974,7 @@ function main(): void { ensureAppSession().cookies.on('changed', (_event, cookie, _cause, removed) => { if (!removed && isSessionCookieName(cookie.name)) desktopExecutor.refreshRegistration() }) + await pendingResults desktopExecutor.start() } await ensureMainWindow() diff --git a/apps/desktop/src/main/terminal/index.ts b/apps/desktop/src/main/terminal/index.ts index c8de703d177..0150434ef0a 100644 --- a/apps/desktop/src/main/terminal/index.ts +++ b/apps/desktop/src/main/terminal/index.ts @@ -38,6 +38,7 @@ import { resourceTabTargetIndex, } 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 { activePane, @@ -49,6 +50,7 @@ import { killPane, listPanes, pollRun, + type RecordedRun, resolveAttachment, runPaneState, sendKey, @@ -156,6 +158,8 @@ export interface TerminalServiceOptions { canSpawn?(): boolean /** Reads and signals a terminal's foreground process group; the OS's by default. */ processGroups?: TerminalProcessGroups + /** Where tagged tmux runs are recorded, so a later process can still stop them. */ + runLedger?: RunLedger } interface TerminalProcessGroups { @@ -490,7 +494,10 @@ export class TerminalService { private async reapFinishedRuns(terminalId: string, env: NodeJS.ProcessEnv): Promise { // A closed tab's run whose pane has since gone (its command ended) needs no stopping. for (const [handle, orphanEnv] of this.orphanedRuns) { - if ((await runPaneState(handle, orphanEnv)) === 'gone') this.orphanedRuns.delete(handle) + if ((await runPaneState(handle, orphanEnv)) === 'gone') { + this.orphanedRuns.delete(handle) + this.forgetRun(handle) + } } for (const handle of this.pendingRuns.get(terminalId) ?? []) { if (this.awaitedRuns.has(handle)) continue @@ -499,11 +506,17 @@ export class TerminalService { // A pane kept open after its command ended (`remain-on-exit`) closes with its run. if (complete) await closeRunPane(handle, env) this.untrackRun(terminalId, handle) + this.forgetRun(handle) handle.dispose() } } } + /** Drops a run's record once nothing of it is left for any process to stop. */ + private forgetRun(handle: TmuxRunHandle): void { + if (handle.runId) this.options.runLedger?.forget(handle.runId) + } + /** Removes a run's files now, or once the call still reading them is done with them. */ private releaseRun(handle: TmuxRunHandle): void { if (this.awaitedRuns.has(handle)) this.releasedAwaitedRuns.add(handle) @@ -524,7 +537,18 @@ export class TerminalService { const pending = this.pendingRuns.get(terminalId) if (!pending) return for (const handle of pending) { - // An untracked run is never stopped, so there is nothing to keep it for. + // A finished run's pane may still be open (`remain-on-exit`): it is closed, while still the + // run's, before the record goes, and its files go only after that check, which an untracked + // run needs them for. Without the shell's environment the record stays, and the next sweep + // closes it. An untracked run is never stopped, so it is not kept either. + if (isRunComplete(handle) && env) { + void closeRunPane(handle, env) + .then(async () => { + if ((await runPaneState(handle, env)) === 'gone') this.forgetRun(handle) + }) + .finally(() => this.releaseRun(handle)) + continue + } if (env && handle.runId !== null && !isRunComplete(handle)) this.orphanedRuns.set(handle, env) this.releaseRun(handle) } @@ -1035,7 +1059,7 @@ export class TerminalService { } case 'run': return tmux - ? this.runInTmux(session, tmux.session, args, latch) + ? this.runInTmux(toolCallId, session, tmux.session, args, latch) : this.run(toolCallId, session, args, latch) case 'read': { const requested = Number(args.lines) @@ -1269,6 +1293,7 @@ export class TerminalService { * see through tmux. */ private async runInTmux( + toolCallId: string, terminal: TerminalSession, session: string, args: TerminalToolArgs, @@ -1280,7 +1305,16 @@ export class TerminalService { const started = Date.now() await this.reapFinishedRuns(terminal.terminalId, terminal.env) - const handle = await startRun(session, command, terminal.currentCwd, terminal.env) + const ledger = this.options.runLedger + const handle = await startRun(session, command, terminal.currentCwd, terminal.env, { + ...(ledger + ? { + beforeStart: (run: RecordedRun) => + ledger.record({ ...run, callId: toolCallId, delivered: false }), + abandon: (runId: string) => ledger.forget(runId), + } + : {}), + }) if ('error' in handle) throw new TerminalError('SPAWN_FAILED', handle.error) // Tracked from the moment its window exists, so sign-out can stop it even mid-wait. const pending = this.pendingRuns.get(terminal.terminalId) @@ -1316,10 +1350,12 @@ export class TerminalService { if (outcome.done) { await closeRunPane(handle, terminal.env) this.untrackRun(terminal.terminalId, handle) + this.forgetRun(handle) handle.dispose() } // Still going, it stays tracked, and nothing polls the status file again: `read` captures - // the pane instead. + // the pane instead. Its record is marked handed back only once that result reaches the model; + // see `TerminalRegistry.markRunDelivered`. const { text, truncated } = elideOutput(outcome.output) return { diff --git a/apps/desktop/src/main/terminal/registry-run-ledger.test.ts b/apps/desktop/src/main/terminal/registry-run-ledger.test.ts new file mode 100644 index 00000000000..4e3cc036c5f --- /dev/null +++ b/apps/desktop/src/main/terminal/registry-run-ledger.test.ts @@ -0,0 +1,136 @@ +import { mkdtempSync, rmSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' + +vi.mock('electron', () => import('@/test/electron-mock')) + +/** + * Each recorded run's pane as tmux has it: running, stopped by Sim, already gone, or one tmux + * cannot answer for. A pane not listed is running. + */ +const { panes } = vi.hoisted(() => ({ + panes: new Map(), +})) + +vi.mock('@/main/terminal/tmux', async () => { + const actual = + await vi.importActual('@/main/terminal/tmux') + const state = (runId: string) => panes.get(runId) ?? 'running' + return { + ...actual, + recordedRunState: async (run: { runId: string }) => { + const pane = state(run.runId) + return pane === 'running' ? 'ours' : pane === 'unknown' ? 'unknown' : 'gone' + }, + stopRecordedRun: async (run: { runId: string }) => { + if (state(run.runId) === 'unknown') return 'unknown' + if (state(run.runId) === 'running') panes.set(run.runId, 'stopped') + return 'gone' + }, + } +}) + +import { TerminalRegistry } from '@/main/terminal/registry' +import { createRunLedger } from '@/main/terminal/run-ledger' + +const dirs: string[] = [] + +function ledgerDir(): string { + const dir = mkdtempSync(join(tmpdir(), 'sim-registry-ledger-')) + dirs.push(dir) + return join(dir, 'terminal-runs') +} + +function run(runId: string, pane: string, delivered = false) { + return { runId, pane, socket: '/tmp/tmux-501/default', callId: `call-${runId}`, delivered } +} + +afterEach(() => { + panes.clear() + for (const dir of dirs.splice(0)) rmSync(dir, { recursive: true, force: true }) +}) + +describe('stopping recorded tmux runs', () => { + it('forgets each run once nothing of it is left, and keeps one tmux could not answer for', async () => { + const dir = ledgerDir() + const previous = createRunLedger(dir) + previous.record(run('stopped', '%1')) + previous.record(run('unanswered', '%2')) + panes.set('unanswered', 'unknown') + // Signing out ends a run whose result the model has, too. + const ledger = createRunLedger(dir) + + await new TerminalRegistry(undefined, undefined, ledger).stopAgentCommands() + + expect(ledger.list().map((record) => record.runId)).toEqual(['unanswered']) + }) + + it("at launch, leaves the same user's runs the model can come back to, and stops the rest", async () => { + const dir = ledgerDir() + const previous = createRunLedger(dir) + // Sim took its result: the model has the pane. + previous.record(run('handed-back', '%1', true)) + previous.record(run('never-handed-back', '%2')) + // Its result is in the journal, unacknowledged: recovery will hand the pane to the model. + previous.record(run('on-its-way', '%3')) + previous.record(run('handed-back-and-gone', '%4', true)) + panes.set('handed-back-and-gone', 'gone') + const ledger = createRunLedger(dir) + + await new TerminalRegistry(undefined, undefined, ledger).stopUncollectableRuns( + new Set(['call-on-its-way']) + ) + + expect(Object.fromEntries(panes)).toEqual({ + 'never-handed-back': 'stopped', + 'handed-back-and-gone': 'gone', + }) + expect( + ledger + .list() + .map((record) => record.runId) + .sort() + ).toEqual(['handed-back', 'on-its-way']) + }) + + it('keeps meaning to stop a run a stop for everything could not confirm', async () => { + const dir = ledgerDir() + createRunLedger(dir).record(run('unconfirmed', '%1', true)) + panes.set('unconfirmed', 'unknown') + const ledger = createRunLedger(dir) + // Sign-out could not confirm the run ended. + await new TerminalRegistry(undefined, undefined, ledger).stopAgentCommands() + panes.set('unconfirmed', 'running') + + // A later launch for the same user, with the journal cleared, still stops it. + await new TerminalRegistry(undefined, undefined, ledger).stopUncollectableRuns(new Set()) + + expect(panes.get('unconfirmed')).toBe('stopped') + expect(ledger.list()).toEqual([]) + }) + + it('notes a run as handed back once its call result is durable', () => { + const dir = ledgerDir() + const ledger = createRunLedger(dir) + ledger.record(run('watched', '%1')) + ledger.record(run('other', '%2')) + + new TerminalRegistry(undefined, undefined, ledger).markRunDelivered('call-watched') + + expect( + Object.fromEntries(ledger.list().map((record) => [record.runId, record.delivered])) + ).toEqual({ watched: true, other: false }) + }) + + it("at launch, stops the previous process's runs and none this one has started", async () => { + const dir = ledgerDir() + createRunLedger(dir).record(run('previous', '%1')) + const ledger = createRunLedger(dir) + ledger.record(run('current', '%2')) + + await new TerminalRegistry(undefined, undefined, ledger).stopRecordedRuns({ excludeLive: true }) + + expect(ledger.list().map((record) => record.runId)).toEqual(['current']) + }) +}) diff --git a/apps/desktop/src/main/terminal/registry.test.ts b/apps/desktop/src/main/terminal/registry.test.ts index 3926b3cfd79..5f723173378 100644 --- a/apps/desktop/src/main/terminal/registry.test.ts +++ b/apps/desktop/src/main/terminal/registry.test.ts @@ -1,4 +1,5 @@ import { tmpdir } from 'node:os' +import { join } from 'node:path' import type { TerminalCommandEvent } from '@sim/terminal-protocol' import { beforeEach, describe, expect, it, vi } from 'vitest' @@ -89,11 +90,13 @@ vi.mock('@/main/terminal/session', () => ({ }, })) +import { TerminalService, type TerminalServiceOptions } from '@/main/terminal' import { type ScopedTerminalSink, TerminalRegistry, type TerminalScopePersistence, } from '@/main/terminal/registry' +import { createRunLedger } from '@/main/terminal/run-ledger' function registry(): TerminalRegistry { return new TerminalRegistry() @@ -186,6 +189,32 @@ describe('TerminalRegistry', () => { terminals.dispose() }) + it('gives a service rebuilt after a failed restore the run ledger too', () => { + const ledger = createRunLedger(join(tmpdir(), `sim-registry-ledger-${process.pid}`)) + const built: Array = [] + const persistence: TerminalScopePersistence = { + load: () => ({ v: 1 as const, tabs: [{ cwd: tmpdir() }, { cwd: tmpdir() }], activeIndex: 0 }), + save: () => true, + migrate: () => true, + disposeScope: () => {}, + } + const terminals = new TerminalRegistry( + persistence, + (_scope, options) => { + built.push(options.runLedger) + return new TerminalService(options) + }, + ledger + ) + createControl.failAt = 2 + + expect(() => terminals.restoreScope('chat-A')).toThrow('PTY spawn failed') + + // The service that failed to restore, and the one built in its place. + expect(built).toHaveLength(2) + expect(built.every((runLedger) => runLedger === ledger)).toBe(true) + }) + it('rolls back a partial restore before retrying the complete descriptor', () => { const persistedTabs = [{ cwd: tmpdir() }, { cwd: process.cwd() }, { cwd: tmpdir() }] const persistence: TerminalScopePersistence = { diff --git a/apps/desktop/src/main/terminal/registry.ts b/apps/desktop/src/main/terminal/registry.ts index 0bdf7548537..6ecc7a8a38a 100644 --- a/apps/desktop/src/main/terminal/registry.ts +++ b/apps/desktop/src/main/terminal/registry.ts @@ -19,6 +19,11 @@ import { type TerminalServiceOptions, type TerminalSink, } from '@/main/terminal' +import type { RunLedger, RunRecord } from '@/main/terminal/run-ledger' +import { recordedRunState, stopRecordedRun } from '@/main/terminal/tmux' + +/** How long a recorded run gets to end on Ctrl-C before its pane is closed. */ +const RECORDED_RUN_GRACE_MS = 2_000 /** Native PTYs and their headless xterm buffers are process-wide resources. */ export const MAX_TERMINALS_PER_PROCESS = 48 @@ -106,7 +111,8 @@ export class TerminalRegistry { constructor( private readonly persistence?: TerminalScopePersistence, - private readonly serviceFactory: TerminalServiceFactory = createTerminalService + private readonly serviceFactory: TerminalServiceFactory = createTerminalService, + private readonly runLedger?: RunLedger ) {} setSink(sink: ScopedTerminalSink | null): void { @@ -369,11 +375,70 @@ export class TerminalRegistry { return true } - /** Stops every command the agent started in any chat's terminals; the user's own are untouched. */ + /** + * Stops every command the agent started in any chat's terminals; the user's own are untouched. + * That includes tmux runs no live terminal holds any more: a chat put away, or a previous + * process that quit or crashed while they ran. + */ async stopAgentCommands(): Promise { await Promise.allSettled( [...this.entries.values()].map((entry) => entry.service.stopAgentCommands()) ) + await this.stopRecordedRuns() + } + + /** + * At launch, for the same user: leaves a previous process's tmux run going only when the model + * has, or will get, the result that handed it back as still going: Sim acknowledged it + * (`delivered`), or the executor's journal holds it for recovery to send (`pendingResults`). + * Every other run is stopped, as is one a stop for everything could not confirm. + */ + stopUncollectableRuns(pendingResults: ReadonlySet): Promise { + return this.stopRecordedRuns({ + excludeLive: true, + keep: (run) => !run.mustStop && (run.delivered || pendingResults.has(run.callId)), + }) + } + + /** + * Notes that a call's result reached the model, so a tmux run it handed back as still going may + * be left to the model across a restart. + */ + markRunDelivered(callId: string): void { + const ledger = this.runLedger + if (!ledger) return + for (const run of ledger.list()) { + if (run.callId === callId) ledger.markDelivered(run.runId) + } + } + + /** + * Stops the recorded tmux runs, each only while its pane still carries its tag, and drops the + * records with nothing left to stop. `excludeLive` skips the runs this process has started; + * `keep` names runs to leave going, such as a previous process's runs whose results the model + * already has and may come back to. + */ + async stopRecordedRuns( + options: { + excludeLive?: boolean + /** Runs to leave going; their records are only dropped once their panes are gone. */ + keep?: (run: RunRecord) => boolean + } = {} + ): Promise { + const ledger = this.runLedger + if (!ledger) return + await Promise.allSettled( + ledger.list({ excludeLive: options.excludeLive }).map(async (run) => { + const keeping = options.keep?.(run) ?? false + const state = keeping + ? await recordedRunState(run, process.env) + : await stopRecordedRun(run, process.env, RECORDED_RUN_GRACE_MS) + if (state === 'gone') ledger.forget(run.runId) + // A run this sweep meant to stop but could not confirm stays meant to stop, so no later + // sweep keeps it. + else if (!keeping) ledger.markMustStop(run.runId) + }) + ) } /** Tears down every shell owned by every chat scope. */ @@ -405,6 +470,7 @@ export class TerminalRegistry { service: this.serviceFactory(scope, { loadCwd: () => this.entries.get(scope)?.persisted?.tabs[0]?.cwd, canSpawn: () => this.liveTerminalCount() < MAX_TERMINALS_PER_PROCESS, + runLedger: this.runLedger, }), persisted, restoreApplied: false, @@ -457,6 +523,7 @@ export class TerminalRegistry { replacement = this.serviceFactory(entry.scope, { loadCwd: () => entry.persisted?.tabs[0]?.cwd, canSpawn: () => this.liveTerminalCount() < MAX_TERMINALS_PER_PROCESS, + runLedger: this.runLedger, }) } catch { this.entries.delete(entry.scope) diff --git a/apps/desktop/src/main/terminal/run-ledger.test.ts b/apps/desktop/src/main/terminal/run-ledger.test.ts new file mode 100644 index 00000000000..9baf932b211 --- /dev/null +++ b/apps/desktop/src/main/terminal/run-ledger.test.ts @@ -0,0 +1,115 @@ +import { mkdirSync, mkdtempSync, readdirSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import { createRunLedger } from '@/main/terminal/run-ledger' + +const dirs: string[] = [] + +function scratch(): string { + const dir = mkdtempSync(join(tmpdir(), 'sim-run-ledger-')) + dirs.push(dir) + return join(dir, 'terminal-runs') +} + +const RUN = { + runId: 'run-1', + pane: '%3', + socket: '/tmp/tmux-501/default', + callId: 'call-1', + delivered: false, +} + +afterEach(() => { + for (const dir of dirs.splice(0)) rmSync(dir, { recursive: true, force: true }) +}) + +describe('the tmux run ledger', () => { + it('keeps a run recorded until it is forgotten, for the next process too', () => { + const dir = scratch() + createRunLedger(dir).record(RUN) + + const nextProcess = createRunLedger(dir) + expect(nextProcess.list()).toEqual([RUN]) + + nextProcess.forget(RUN.runId) + expect(createRunLedger(dir).list()).toEqual([]) + }) + + it("leaves out this process's own runs when asked", () => { + const dir = scratch() + createRunLedger(dir).record(RUN) + const ledger = createRunLedger(dir) + ledger.record({ ...RUN, runId: 'run-2', pane: '%4' }) + + expect(ledger.list({ excludeLive: true })).toEqual([RUN]) + expect(ledger.list()).toHaveLength(2) + }) + + it('drops files nothing could act on safely', () => { + const dir = scratch() + const ledger = createRunLedger(dir) + ledger.record(RUN) + writeFileSync(join(dir, 'garbled.json'), '{not json') + writeFileSync(join(dir, 'run-9.json'), JSON.stringify({ ...RUN, socket: 'relative.sock' })) + // A record under another run's name could stop the wrong run. + writeFileSync(join(dir, 'run-8.json'), JSON.stringify({ ...RUN, runId: 'run-7' })) + // Only a pane id names one pane; a target like this one names whatever pane is active there. + writeFileSync( + join(dir, 'run-6.json'), + JSON.stringify({ ...RUN, runId: 'run-6', pane: 'work:0.0' }) + ) + + expect(ledger.list()).toEqual([RUN]) + expect(readdirSync(dir)).toEqual(['run-1.json']) + }) + + it('cleans up after a write that never finished, keeping the saved record', () => { + const dir = scratch() + const ledger = createRunLedger(dir) + ledger.record(RUN) + writeFileSync(join(dir, 'run-2.json.4242.1.tmp'), '{"runId":"run-2","pa') + + expect(ledger.list()).toEqual([RUN]) + expect(readdirSync(dir)).toEqual(['run-1.json']) + }) + + it('keeps sweeping, and lets a call finish, when an entry cannot be removed', () => { + const dir = scratch() + const ledger = createRunLedger(dir) + ledger.record(RUN) + // Not a file, so removing it fails. + mkdirSync(join(dir, 'run-3.json')) + mkdirSync(join(dir, 'run-4.json')) + + expect(ledger.list()).toEqual([RUN]) + expect(() => ledger.forget('run-4')).not.toThrow() + }) + + it('reports a record it could not save, so the run is not started', () => { + const dir = scratch() + // A file where the directory should be: nothing can be saved under it. + writeFileSync(join(dir, '..', 'blocked'), '') + const ledger = createRunLedger(join(dir, '..', 'blocked')) + + expect(ledger.record(RUN)).toBe(false) + }) + + it('notes a run handed back as still going, for the next process too', () => { + const dir = scratch() + const ledger = createRunLedger(dir) + ledger.record(RUN) + + ledger.markDelivered(RUN.runId) + + expect(createRunLedger(dir).list()).toEqual([{ ...RUN, delivered: true }]) + }) + + it('records nothing for a run tag that is not a plain id', () => { + const dir = scratch() + const ledger = createRunLedger(dir) + ledger.record({ ...RUN, runId: '../escape' }) + + expect(ledger.list()).toEqual([]) + }) +}) diff --git a/apps/desktop/src/main/terminal/run-ledger.ts b/apps/desktop/src/main/terminal/run-ledger.ts new file mode 100644 index 00000000000..badf2806feb --- /dev/null +++ b/apps/desktop/src/main/terminal/run-ledger.ts @@ -0,0 +1,164 @@ +/** + * A durable record of the agent's tagged tmux runs, so a later process can still stop them. + * + * A tmux run outlives the app: quitting, crashing or signing out mid-teardown leaves its command + * going in the user's tmux server, and the next launch would otherwise know nothing about it. + * Each record names the run's tag, its pane and its tmux server's socket, and nothing else: no + * command line and no output, since it lives outside the account's encrypted data. A record is + * saved before the run's command may start, and removed once the run has ended, its pane is gone, + * or Sim has stopped it. + */ + +import { readdirSync, readFileSync, rmSync } from 'node:fs' +import { join } from 'node:path' +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { writeJsonFileAtomicallySync } from '@/main/atomic-json-file' +import type { RecordedRun } from '@/main/terminal/tmux' + +const logger = createLogger('DesktopTerminalRunLedger') + +/** A recorded run, with the call it belongs to and whether that call's result went back. */ +export interface RunRecord extends RecordedRun { + /** The tool call that started the run, to match it against the executor's journal. */ + callId: string + /** + * True once the run's call handed back its result while the run went on (`running`, with its + * pane): from then on the model can come back to the pane, so a restart must leave it be. + */ + delivered: boolean + /** A stop for everything (sign-out, Terminal off) could not confirm this run ended. */ + mustStop?: boolean +} + +export interface RunLedger { + /** Saves a run's record; false when it could not be saved, so the run must not start. */ + record(run: RunRecord): boolean + /** Notes that the run's call has handed back its result, with the run still going. */ + markDelivered(runId: string): void + /** Notes that the run must be stopped, whatever a later launch would otherwise decide. */ + markMustStop(runId: string): void + forget(runId: string): void + /** Every recorded run; `excludeLive` leaves out runs this process recorded. */ + list(options?: { excludeLive?: boolean }): RunRecord[] +} + +/** Run tags are generated ids; anything else in the directory is not a record. */ +const RUN_ID = /^[A-Za-z0-9_-]{1,128}$/ + +function parseRecord(text: string): RunRecord | null { + try { + const parsed = JSON.parse(text) as Partial + if ( + typeof parsed.runId === 'string' && + RUN_ID.test(parsed.runId) && + typeof parsed.pane === 'string' && + /^%\d+$/.test(parsed.pane) && + typeof parsed.callId === 'string' && + typeof parsed.delivered === 'boolean' && + typeof parsed.socket === 'string' && + parsed.socket.startsWith('/') + ) { + return { + runId: parsed.runId, + pane: parsed.pane, + socket: parsed.socket, + callId: parsed.callId, + delivered: parsed.delivered, + ...(parsed.mustStop === true ? { mustStop: true } : {}), + } + } + } catch { + // Unreadable: treated as no record below. + } + return null +} + +/** Removes a file the ledger no longer needs; a failure is logged, never thrown at a caller. */ +function remove(path: string): void { + try { + rmSync(path, { force: true }) + } catch (error) { + logger.warn('Could not remove a tmux run record', { error: getErrorMessage(error) }) + } +} + +export function createRunLedger(dir: string): RunLedger { + /** Runs recorded by this process, still going as far as it knows. */ + const live = new Set() + const pathFor = (runId: string) => join(dir, `${runId}.json`) + + /** Rewrites a saved record; `change` returns null to leave it as it is. */ + const update = (runId: string, change: (record: RunRecord) => RunRecord | null): void => { + if (!RUN_ID.test(runId)) return + let record: RunRecord | null = null + try { + record = parseRecord(readFileSync(pathFor(runId), 'utf8')) + } catch { + record = null + } + const changed = record ? change(record) : null + if (!changed) return + try { + writeJsonFileAtomicallySync(pathFor(runId), changed) + } catch (error) { + logger.warn('Could not update a tmux run record', { error: getErrorMessage(error) }) + } + } + + return { + record(run) { + if (!RUN_ID.test(run.runId)) return false + try { + writeJsonFileAtomicallySync(pathFor(run.runId), run) + live.add(run.runId) + return true + } catch (error) { + logger.warn('Could not record a tmux run', { error: getErrorMessage(error) }) + return false + } + }, + markDelivered(runId) { + // Left undelivered, a restart stops the run: the conservative side. + update(runId, (record) => (record.delivered ? null : { ...record, delivered: true })) + }, + markMustStop(runId) { + update(runId, (record) => (record.mustStop ? null : { ...record, mustStop: true })) + }, + forget(runId) { + live.delete(runId) + if (RUN_ID.test(runId)) remove(pathFor(runId)) + }, + list(options = {}) { + let names: string[] + try { + names = readdirSync(dir) + } catch { + return [] + } + const runs: RunRecord[] = [] + for (const name of names) { + // A write that never finished leaves only its temporary file behind. + if (name.endsWith('.tmp')) { + remove(join(dir, name)) + continue + } + if (!name.endsWith('.json')) continue + let record: RunRecord | null = null + try { + record = parseRecord(readFileSync(join(dir, name), 'utf8')) + } catch { + record = null + } + if (!record || `${record.runId}.json` !== name) { + // Nothing could act on it safely; it only takes up space. + remove(join(dir, name)) + continue + } + if (options.excludeLive && live.has(record.runId)) continue + runs.push(record) + } + return runs + }, + } +} diff --git a/apps/desktop/src/main/terminal/service.test.ts b/apps/desktop/src/main/terminal/service.test.ts index 55204d4ce8e..b2f24c9cec8 100644 --- a/apps/desktop/src/main/terminal/service.test.ts +++ b/apps/desktop/src/main/terminal/service.test.ts @@ -1,8 +1,10 @@ import { existsSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' import { tmpdir } from 'node:os' import { join } from 'node:path' +import { sleep } from '@sim/utils/helpers' import { describe, expect, it, vi } from 'vitest' import { TerminalService } from '@/main/terminal' +import { createRunLedger } from '@/main/terminal/run-ledger' /** * A tmux attachment the service sees only when a test turns it on. Runs get real status files; @@ -20,6 +22,8 @@ const tmuxFake = vi.hoisted(() => ({ untracked: false, /** Run panes tmux still shows, kept open after their command ends (`remain-on-exit`). */ open: new Set(), + /** tmux stops answering: no pane can be confirmed, so none is closed. */ + unanswered: false, statusPaths: new Map(), })) @@ -36,8 +40,20 @@ vi.mock('@/main/terminal/tmux', async () => { : actual.resolveAttachment(pid, env), startRun: async (...args: Parameters) => { if (!tmuxFake.on) return actual.startRun(...args) - const dir = mkdtempSync(join(tmpdir(), 'sim-tmux-fake-')) + const options = args[4] const pane = `%${nextPane++}` + const runId = tmuxFake.untracked ? null : `run-${pane.slice(1)}` + const socket = tmuxFake.untracked ? null : '/tmp/tmux-fake/default' + // Like the real one, a tagged run starts only once its record is saved. + if ( + runId && + socket && + options?.beforeStart && + !options.beforeStart({ runId, pane, socket }) + ) { + return { error: 'The command could not be recorded for a later stop, so it was not run.' } + } + const dir = mkdtempSync(join(tmpdir(), 'sim-tmux-fake-')) const statusPath = join(dir, 'status') writeFileSync(join(dir, 'out'), 'partial output') tmuxFake.statusPaths.set(pane, statusPath) @@ -45,7 +61,8 @@ vi.mock('@/main/terminal/tmux', async () => { return { window: `@${pane.slice(1)}`, pane, - runId: tmuxFake.untracked ? null : `run-${pane}`, + runId, + socket, outPath: join(dir, 'out'), statusPath, dispose: () => rmSync(dir, { recursive: true, force: true }), @@ -53,6 +70,7 @@ vi.mock('@/main/terminal/tmux', async () => { }, runPaneState: async (...args: Parameters) => { if (!tmuxFake.on) return actual.runPaneState(...args) + if (tmuxFake.unanswered) return 'unknown' if (tmuxFake.gone.has(args[0].pane)) return 'gone' return args[0].runId === null ? 'unknown' : 'ours' }, @@ -65,7 +83,13 @@ vi.mock('@/main/terminal/tmux', async () => { }, closeRunPane: async (...args: Parameters) => { if (!tmuxFake.on) return actual.closeRunPane(...args) + if (tmuxFake.unanswered) return + // An untracked run's pane is proven its own only from the run's files, read after tmux + // has answered, as the real check does. + await sleep(10) + if (args[0].runId === null && !existsSync(args[0].statusPath)) return tmuxFake.open.delete(args[0].pane) + tmuxFake.gone.add(args[0].pane) }, } }) @@ -623,6 +647,189 @@ describe('agent commands in tmux', () => { } }) + it('keeps a record of a tagged run exactly as long as the run goes on', async () => { + tmuxFake.on = true + tmuxFake.statusPaths.clear() + const scratch = mkdtempSync(join(tmpdir(), 'sim-ledger-')) + const ledgerDir = join(scratch, 'terminal-runs') + const ledger = createRunLedger(ledgerDir) + try { + const terminal = new TerminalService({ loadCwd: () => '/tmp', runLedger: ledger }) + terminal.start({ cols: 80, rows: 24 }) + await terminal.executeTool('call-long', 'run', { command: 'make build', waitSeconds: 1 }) + const [[pane = '', statusPath = ''] = []] = [...tmuxFake.statusPaths] + + // Still going after its call returned: a later process must be able to find it. + expect(createRunLedger(ledgerDir).list()).toEqual([ + { + runId: `run-${pane.slice(1)}`, + pane, + socket: '/tmp/tmux-fake/default', + callId: 'call-long', + // Handed back, but only the executor's journal can make that durable. + delivered: false, + }, + ]) + + writeFileSync(statusPath, '0') + await terminal.executeTool('call-next', 'run', { command: 'ls', waitSeconds: 1 }) + + // The first finished and is forgotten; the one still going is recorded in its place. + expect( + createRunLedger(ledgerDir) + .list() + .map((run) => run.pane) + ).toEqual([[...tmuxFake.statusPaths.keys()][1]]) + } finally { + tmuxFake.on = false + rmSync(scratch, { recursive: true, force: true }) + } + }) + + it('closes and then forgets a run that finished before its terminal closed', async () => { + tmuxFake.on = true + tmuxFake.statusPaths.clear() + tmuxFake.open.clear() + const scratch = mkdtempSync(join(tmpdir(), 'sim-ledger-')) + const ledgerDir = join(scratch, 'terminal-runs') + try { + const terminal = new TerminalService({ + loadCwd: () => '/tmp', + runLedger: createRunLedger(ledgerDir), + }) + const { activeTerminalId } = terminal.start({ cols: 80, rows: 24 }) + await terminal.executeTool('call-long', 'run', { command: 'make build', waitSeconds: 1 }) + const [[, statusPath = ''] = []] = [...tmuxFake.statusPaths] + writeFileSync(statusPath, '0') + + const [pane = ''] = [...tmuxFake.statusPaths.keys()] + // Its dead pane is still open, as with `remain-on-exit`. + expect(tmuxFake.open.has(pane)).toBe(true) + + terminal.closeTerminal(activeTerminalId as string) + + await vi.waitFor(() => expect(createRunLedger(ledgerDir).list()).toEqual([])) + expect(tmuxFake.open.has(pane)).toBe(false) + } finally { + tmuxFake.on = false + rmSync(scratch, { recursive: true, force: true }) + } + }) + + it("forgets a run that finishes within its call, and a closed tab's run once its pane is gone", async () => { + tmuxFake.on = true + tmuxFake.statusPaths.clear() + const scratch = mkdtempSync(join(tmpdir(), 'sim-ledger-')) + const ledgerDir = join(scratch, 'terminal-runs') + try { + const terminal = new TerminalService({ + loadCwd: () => '/tmp', + runLedger: createRunLedger(ledgerDir), + }) + const { activeTerminalId } = terminal.start({ cols: 80, rows: 24 }) + // Finishes while its call still waits on it. + const quick = terminal.executeTool('call-quick', 'run', { command: 'ls', waitSeconds: 30 }) + await vi.waitFor(() => expect(tmuxFake.statusPaths.size).toBe(1)) + writeFileSync([...tmuxFake.statusPaths.values()][0] ?? '', '0') + await quick + expect(createRunLedger(ledgerDir).list()).toEqual([]) + + // Still going when its tab closes; its pane goes later, and the next run's bookkeeping sees. + await terminal.executeTool('call-long', 'run', { command: 'make build', waitSeconds: 1 }) + const [, longPane = ''] = [...tmuxFake.statusPaths.keys()] + terminal.closeTerminal(activeTerminalId as string) + expect( + createRunLedger(ledgerDir) + .list() + .map((run) => run.pane) + ).toEqual([longPane]) + tmuxFake.gone.add(longPane) + await terminal.executeTool('call-new', 'new', {}) + await terminal.executeTool('call-next', 'run', { command: 'pwd', waitSeconds: 1 }) + + expect( + createRunLedger(ledgerDir) + .list() + .map((run) => run.pane) + ).not.toContain(longPane) + } finally { + tmuxFake.on = false + tmuxFake.gone.clear() + rmSync(scratch, { recursive: true, force: true }) + } + }) + + it("closes a finished untracked run's pane when its terminal closes", async () => { + tmuxFake.on = true + tmuxFake.untracked = true + tmuxFake.statusPaths.clear() + tmuxFake.open.clear() + try { + const terminal = new TerminalService({ loadCwd: () => '/tmp' }) + const { activeTerminalId } = terminal.start({ cols: 80, rows: 24 }) + await terminal.executeTool('call-old-tmux', 'run', { command: 'make build', waitSeconds: 1 }) + const [[pane = '', statusPath = ''] = []] = [...tmuxFake.statusPaths] + writeFileSync(statusPath, '0') + + terminal.closeTerminal(activeTerminalId as string) + + await vi.waitFor(() => expect(tmuxFake.open.has(pane)).toBe(false)) + } finally { + tmuxFake.on = false + tmuxFake.untracked = false + } + }) + + it('keeps the record of a finished run whose pane tmux could not confirm closing', async () => { + tmuxFake.on = true + tmuxFake.statusPaths.clear() + const scratch = mkdtempSync(join(tmpdir(), 'sim-ledger-')) + const ledgerDir = join(scratch, 'terminal-runs') + try { + const terminal = new TerminalService({ + loadCwd: () => '/tmp', + runLedger: createRunLedger(ledgerDir), + }) + const { activeTerminalId } = terminal.start({ cols: 80, rows: 24 }) + await terminal.executeTool('call-long', 'run', { command: 'make build', waitSeconds: 1 }) + writeFileSync([...tmuxFake.statusPaths.values()][0] ?? '', '0') + tmuxFake.unanswered = true + + terminal.closeTerminal(activeTerminalId as string) + await sleep(200) + + // The next sweep will close its pane. + expect(createRunLedger(ledgerDir).list()).toHaveLength(1) + } finally { + tmuxFake.on = false + tmuxFake.unanswered = false + rmSync(scratch, { recursive: true, force: true }) + } + }) + + it('never starts a tagged run it could not record', async () => { + tmuxFake.on = true + tmuxFake.statusPaths.clear() + const scratch = mkdtempSync(join(tmpdir(), 'sim-ledger-')) + // A file where the ledger's directory should be: no record can be saved. + writeFileSync(join(scratch, 'terminal-runs'), '') + try { + const terminal = new TerminalService({ + loadCwd: () => '/tmp', + runLedger: createRunLedger(join(scratch, 'terminal-runs')), + }) + terminal.start({ cols: 80, rows: 24 }) + + await expect( + terminal.executeTool('call-unrecorded', 'run', { command: 'make build', waitSeconds: 1 }) + ).resolves.toMatchObject({ ok: false, code: 'SPAWN_FAILED' }) + expect(tmuxFake.statusPaths.size).toBe(0) + } finally { + tmuxFake.on = false + rmSync(scratch, { recursive: true, force: true }) + } + }) + it("closes a run's pane when a later run reaps it after it finished", async () => { tmuxFake.on = true tmuxFake.statusPaths.clear() diff --git a/apps/desktop/src/main/terminal/tmux.test.ts b/apps/desktop/src/main/terminal/tmux.test.ts index 88c1ccb1942..adb0cfdf139 100644 --- a/apps/desktop/src/main/terminal/tmux.test.ts +++ b/apps/desktop/src/main/terminal/tmux.test.ts @@ -12,6 +12,7 @@ import { resolveAttachment, runPaneState, startRun, + stopRecordedRun, stopRun, type TmuxRunHandle, } from '@/main/terminal/tmux' @@ -71,6 +72,7 @@ describe('run status files', () => { window: '@1', pane: '%1', runId: 'run-1', + socket: null, outPath: join(dir, 'out'), statusPath: join(dir, 'status'), dispose: () => {}, @@ -119,6 +121,8 @@ describe('run status files', () => { * reshapes the server directly (a split, a closed window, a restart) by rewriting that file. */ interface FakeTmuxState { + /** The server's socket; `-S` naming any other reaches no server. */ + socket: string nextWindow: number nextPane: number panes: Record; command?: string }> @@ -126,6 +130,14 @@ interface FakeTmuxState { log: string[] /** Commands the fake fails, with the error tmux would print. */ fail?: Record + /** Once a key is sent, tmux stops answering: every later command fails like a dying server. */ + dieAfterKeys?: boolean + /** On this many-th display-message, the pane is retagged as another run's (a restart race). */ + retagAtCheck?: number + /** display-message calls so far. */ + checks?: number + /** tmux restarts and the user's pane takes the id just before the next guarded action. */ + retagBeforeAction?: boolean /** Attached clients, as `list-clients` reports them. */ clients?: Array<{ pid: string; tty: string; session: string }> /** Commands the fake holds until the file named here exists, like a busy tmux server. */ @@ -138,8 +150,16 @@ const FAKE_TMUX = ` const fs = require('node:fs') const file = process.env.FAKE_TMUX_STATE const state = JSON.parse(fs.readFileSync(file, 'utf8')) -const args = process.argv.slice(2) +let args = process.argv.slice(2) const save = () => fs.writeFileSync(file, JSON.stringify(state)) +// \`-S socket\` names the server; any server but this one does not exist. +if (args[0] === '-S') { + if (args[1] !== state.socket) { + process.stderr.write('no server running on ' + args[1]) + process.exit(1) + } + args = args.slice(2) +} const target = () => args[args.indexOf('-t') + 1] const fail = (message) => { process.stderr.write(message); process.exit(1) } // Prints a format's output as tmux 3.4 and 3.5 do: a backslash doubled, and every other control @@ -156,6 +176,18 @@ if (state.hold && state.hold[args[0]]) { } state.held = state.held.filter((command) => command !== args[0]) } +// \`if-shell -F -t pane '#{==:#{option},value}' command\`: the check and the action in one command. +if (args[0] === 'if-shell') { + const pane = state.panes[target()] + if (pane && state.retagBeforeAction) { + pane.options['@sim-run-id'] = 'someone-else' + state.retagBeforeAction = false + save() + } + const check = /^#\\{==:#\\{([^}]+)\\},(.*)\\}$/.exec(args[args.length - 2]) + if (!pane || !check || (pane.options[check[1]] ?? '') !== check[2]) process.exit(0) + args = args[args.length - 1].split(' ') +} if (state.fail && state.fail[args[0]]) fail(state.fail[args[0]]) switch (args[0]) { case 'new-window': { @@ -180,6 +212,11 @@ switch (args[0]) { break } case 'display-message': { + state.checks = (state.checks ?? 0) + 1 + if (state.retagAtCheck === state.checks && state.panes[target()]) { + state.panes[target()].options['@sim-run-id'] = 'someone-else' + } + save() // Like tmux 3.x, a pane that is gone answers with an empty line rather than an error. const pane = state.panes[target()] const name = args[args.length - 1].slice(2, -1) @@ -187,6 +224,8 @@ switch (args[0]) { ? '' : name === 'pane_id' ? target() + : name === 'socket_path' + ? state.socket : name === 'pane_start_command' ? (pane.command ?? '') : (pane.options[name] ?? '') @@ -208,6 +247,9 @@ switch (args[0]) { case 'kill-pane': { if (!state.panes[target()]) fail("can't find pane") state.log.push(args[0] + ' ' + target() + (args[0] === 'send-keys' ? ' ' + args[args.length - 1] : '')) + if (args[0] === 'send-keys' && state.dieAfterKeys) { + state.fail = { 'display-message': 'server exited unexpectedly', 'kill-pane': 'server exited unexpectedly' } + } if (args[0] === 'kill-pane') delete state.panes[target()] save() break @@ -224,7 +266,7 @@ function fakeTmux(options: { exec?: boolean } = {}) { writeFileSync(binary, `#!${process.execPath}\n${FAKE_TMUX}`) chmodSync(binary, 0o755) const write = (state: FakeTmuxState) => writeFileSync(stateFile, JSON.stringify(state)) - write({ nextWindow: 0, nextPane: 0, panes: {}, log: [] }) + write({ nextWindow: 0, nextPane: 0, panes: {}, log: [], socket: join(dir, 'server.sock') }) const read = (): FakeTmuxState => JSON.parse(readFileSync(stateFile, 'utf8')) return { dir, @@ -274,6 +316,177 @@ describe('finding the tmux session a shell runs', () => { }) }) +describe('stopping a run another process started, from its record', () => { + const dirs: string[] = [] + + afterEach(() => { + for (const dir of dirs.splice(0)) rmSync(dir, { recursive: true, force: true }) + }) + + async function recorded(tmux: ReturnType) { + dirs.push(tmux.dir) + const run = await startRun('agent', 'sleep 600', null, tmux.env) + if ('error' in run) throw new Error(run.error) + return { run, record: { runId: run.runId ?? '', pane: run.pane, socket: run.socket ?? '' } } + } + + it('records the server a tagged run runs on', async () => { + const tmux = fakeTmux() + const { run } = await recorded(tmux) + + expect(run.socket).toBe(tmux.read().socket) + }) + + it('interrupts and then closes the pane while it still carries the run tag', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + + expect(await stopRecordedRun(record, tmux.env, 0)).toBe('gone') + expect(tmux.read().log).toEqual([`send-keys ${record.pane} C-c`, `kill-pane ${record.pane}`]) + expect(tmux.read().panes).toEqual({}) + }) + + it('never touches a pane that took the recorded id after tmux restarted', async () => { + const tmux = fakeTmux() + const { run, record } = await recorded(tmux) + tmux.restart() + const state = tmux.read() + state.panes[run.pane] = { window: run.window, options: {}, command: 'zsh' } + tmux.write(state) + + expect(await stopRecordedRun(record, tmux.env, 0)).toBe('gone') + expect(tmux.read().log).toEqual([]) + expect(Object.keys(tmux.read().panes)).toEqual([run.pane]) + }) + + it('never asks a different tmux server about the pane', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + + const elsewhere = { ...record, socket: join(tmux.dir, 'another.sock') } + expect(await stopRecordedRun(elsewhere, tmux.env, 0)).toBe('gone') + expect(tmux.read().log).toEqual([]) + expect(Object.keys(tmux.read().panes)).toEqual([record.pane]) + }) + + it('checks the pane is still the run once more right before closing it', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + // The first check (before Ctrl-C) sees the run; the next, right before closing, does not. + tmux.write({ ...tmux.read(), checks: 0, retagAtCheck: 2 }) + + expect(await stopRecordedRun(record, tmux.env, 0)).toBe('gone') + expect(tmux.read().log).toEqual([`send-keys ${record.pane} C-c`]) + expect(Object.keys(tmux.read().panes)).toEqual([record.pane]) + }) + + it('never starts a tagged run whose server it cannot name for a record', async () => { + const tmux = fakeTmux({ exec: true }) + dirs.push(tmux.dir) + tmux.write({ ...tmux.read(), socket: '' }) + const marker = join(tmux.dir, 'ran') + + const result = await startRun('agent', `touch ${JSON.stringify(marker)}`, null, tmux.env, { + beforeStart: () => true, + }) + + expect(result).toMatchObject({ error: expect.stringContaining('was not run') }) + // The pane it opened and tagged is closed with it. + expect(tmux.read().panes).toEqual({}) + await sleep(1_500) + expect(existsSync(marker)).toBe(false) + }) + + it('forgets the saved record of a run that then could not start', async () => { + const tmux = fakeTmux() + dirs.push(tmux.dir) + const saved = new Set() + let runDir = '' + + const result = await startRun('agent', 'sleep 600', null, tmux.env, { + beforeStart: (run) => { + // The run's directory turns read-only, so its go file cannot be written. + const command = tmux.read().panes[run.pane]?.command ?? '' + const script = /"([^"]+)\/run\.sh"/.exec(command)?.[1] ?? '' + runDir = script + chmodSync(runDir, 0o500) + saved.add(run.runId) + return true + }, + abandon: (runId) => saved.delete(runId), + }) + + if (runDir) chmodSync(runDir, 0o700) + if (runDir) dirs.push(runDir) + expect(result).toMatchObject({ error: expect.stringContaining('could not be started') }) + expect([...saved]).toEqual([]) + }) + + it('acts on no pane that took the recorded id between the last check and the action', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + tmux.write({ ...tmux.read(), retagBeforeAction: true }) + + await stopRecordedRun(record, tmux.env, 0) + + expect(tmux.read().log).toEqual([]) + expect(Object.keys(tmux.read().panes)).toEqual([record.pane]) + }) + + it('keeps the record when tmux stops answering part-way through the stop', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + tmux.write({ ...tmux.read(), dieAfterKeys: true }) + + expect(await stopRecordedRun(record, tmux.env, 0)).toBe('unknown') + // Its pane could not be confirmed as the run's, so it was not closed. + expect(tmux.read().log).toEqual([`send-keys ${record.pane} C-c`]) + }) + + it('saves the record before the command may start, and never starts one it could not save', async () => { + const tmux = fakeTmux({ exec: true }) + dirs.push(tmux.dir) + const marker = join(tmux.dir, 'ran') + let ranBeforeRecord = true + const run = await startRun('agent', `touch ${JSON.stringify(marker)}`, null, tmux.env, { + beforeStart: () => { + // Time enough for a command already released to have run. + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 1_000) + ranBeforeRecord = existsSync(marker) + return true + }, + }) + if ('error' in run) throw new Error(run.error) + expect(ranBeforeRecord).toBe(false) + await expect.poll(() => existsSync(marker), { timeout: 10_000 }).toBe(true) + + const unsaved = fakeTmux({ exec: true }) + dirs.push(unsaved.dir) + const unsavedMarker = join(unsaved.dir, 'ran') + const refused = await startRun( + 'agent', + `touch ${JSON.stringify(unsavedMarker)}`, + null, + unsaved.env, + { + beforeStart: () => false, + } + ) + expect(refused).toMatchObject({ error: expect.stringContaining('was not run') }) + await sleep(1_500) + expect(existsSync(unsavedMarker)).toBe(false) + }, 20_000) + + it('leaves the record for later while tmux cannot be asked', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + tmux.write({ ...tmux.read(), fail: { 'display-message': 'server exited unexpectedly' } }) + + expect(await stopRecordedRun(record, tmux.env, 0)).toBe('unknown') + expect(tmux.read().log).toEqual([]) + }) +}) + describe('stopping a tmux run touches only its own pane', () => { const dirs: string[] = [] @@ -299,6 +512,17 @@ describe('stopping a tmux run touches only its own pane', () => { expect(Object.keys(tmux.read().panes)).toEqual([users]) }) + it('sends a live run nothing once its pane is retagged between the check and the action', async () => { + const tmux = fakeTmux() + const run = await started(tmux) + tmux.write({ ...tmux.read(), retagBeforeAction: true }) + + await stopRun(run, tmux.env, 0) + + expect(tmux.read().log).toEqual([]) + expect(Object.keys(tmux.read().panes)).toEqual([run.pane]) + }) + it('sends nothing to a pane that reused the run pane id after tmux restarted', async () => { const tmux = fakeTmux() const run = await started(tmux) diff --git a/apps/desktop/src/main/terminal/tmux.ts b/apps/desktop/src/main/terminal/tmux.ts index 6fb40d82ba3..37c205c9dad 100644 --- a/apps/desktop/src/main/terminal/tmux.ts +++ b/apps/desktop/src/main/terminal/tmux.ts @@ -317,6 +317,8 @@ export interface TmuxRunHandle { * nothing ever stops it, since nothing could tell its pane from one of the user's. */ runId: string | null + /** The tmux server's socket, for a record a later process can act on; null if unknown. */ + socket: string | null outPath: string statusPath: string dispose(): void @@ -355,7 +357,16 @@ export async function startRun( session: string, command: string, cwd: string | null, - env: NodeJS.ProcessEnv + env: NodeJS.ProcessEnv, + options: { + /** + * Records a tagged run before its command may start; false keeps the command from starting, + * since a run no later process could find must not outlive this one. + */ + beforeStart?: (run: RecordedRun) => boolean + /** Undoes `beforeStart` for a run that then could not start after all. */ + abandon?: (runId: string) => void + } = {} ): Promise { const dir = mkdtempSync(join(tmpdir(), 'sim-tmux-run-')) const outPath = join(dir, 'out') @@ -439,14 +450,27 @@ export async function startRun( }) } const runId = tagged.ok ? tag : null + // Which server the pane lives on, so a later process stops it there and nowhere else. + const shown = runId + ? await runTmux(['display-message', '-p', '-t', pane, '#{socket_path}'], env) + : null + const socket = shown?.ok && shown.stdout.trim().startsWith('/') ? shown.stdout.trim() : null + if (runId && options.beforeStart && !(socket && options.beforeStart({ runId, pane, socket }))) { + // The pane was tagged a moment ago, so it is closed as the run's; without the go file its + // command never starts in any case. + await runTmux(ifTagged(pane, runId, `kill-pane -t ${pane}`), env) + dispose() + return { error: 'The command could not be recorded for a later stop, so it was not run.' } + } try { writeFileSync(goPath, '') } catch (error) { + if (runId) options.abandon?.(runId) dispose() return { error: `The command could not be started: ${getErrorMessage(error)}` } } - return { window, pane, runId, outPath, statusPath, dispose } + return { window, pane, runId, socket, outPath, statusPath, dispose } } /** @@ -477,18 +501,94 @@ export async function runPaneState( * step first checks the pane is still the run's, and only that pane is ever closed, so a pane the * user split off beside it, or a window that reused its ids, is never touched. */ +/** + * Arguments for one tmux command that runs `command` on a pane only while the pane carries the + * run's tag. tmux checks the tag and acts within the one command, so a restart between a check + * and an action can never hand the action to a pane that took the id. + */ +function ifTagged(pane: string, runId: string, command: string, socket?: string): string[] { + return [ + ...(socket ? ['-S', socket] : []), + 'if-shell', + '-F', + '-t', + pane, + `#{==:#{${RUN_ID_OPTION}},${runId}}`, + command, + ] +} + export async function stopRun( handle: TmuxRunHandle, env: NodeJS.ProcessEnv, graceMs: number ): Promise { - if ((await runPaneState(handle, env)) !== 'ours') return - await sendKey(handle.pane, 'C-c', env) + if ((await runPaneState(handle, env)) !== 'ours' || !handle.runId) return + await runTmux(ifTagged(handle.pane, handle.runId, `send-keys -t ${handle.pane} C-c`), env) const deadline = Date.now() + graceMs while (!isRunComplete(handle) && Date.now() < deadline) await sleep(100) if (!isRunComplete(handle)) await closeRunPane(handle, env) } +/** A recorded run as a later process finds it: what it takes to find its pane again. */ +export interface RecordedRun { + /** The tag on the run's pane; only a pane carrying it is ever acted on. */ + runId: string + pane: string + /** The tmux server's socket, so a different server is never asked about this pane. */ + socket: string +} + +/** + * Whether a recorded run's pane still carries its tag, asked of the run's own server only: `gone` + * when that server or pane no longer exists or the pane is not tagged as this run's, `unknown` + * when tmux could not be asked. + */ +export async function recordedRunState( + run: RecordedRun, + env: NodeJS.ProcessEnv +): Promise<'ours' | 'gone' | 'unknown'> { + const shown = await runTmux( + ['-S', run.socket, 'display-message', '-p', '-t', run.pane, `#{${RUN_ID_OPTION}}`], + env + ) + if (shown.ok) return shown.stdout.trim() === run.runId ? 'ours' : 'gone' + return /can't find|no server running|error connecting|no such file/i.test(shown.stderr) + ? 'gone' + : 'unknown' +} + +/** + * Stops a run another process started, from its record: Ctrl-C in its pane, then closing that pane + * if the command ignored it. Every step first checks, on the run's own server, that the pane still + * carries the run's tag, so a pane that took its id after a tmux restart is never touched. + * Resolves `gone` once nothing of the run is left to stop, and `unknown` when tmux could not say. + */ +export async function stopRecordedRun( + run: RecordedRun, + env: NodeJS.ProcessEnv, + graceMs: number +): Promise<'gone' | 'unknown'> { + const before = await recordedRunState(run, env) + if (before !== 'ours') return before + await runTmux(ifTagged(run.pane, run.runId, `send-keys -t ${run.pane} C-c`, run.socket), env) + const deadline = Date.now() + graceMs + let state: 'ours' | 'gone' | 'unknown' = 'ours' + while (Date.now() < deadline) { + await sleep(100) + state = await recordedRunState(run, env) + if (state !== 'ours') break + } + // The pane is closed only while it is confirmed the run's; a pane tmux could not answer for is + // left alone, and its record kept for the next sweep. + if (state === 'ours') state = await recordedRunState(run, env) + if (state === 'ours') { + await runTmux(ifTagged(run.pane, run.runId, `kill-pane -t ${run.pane}`, run.socket), env) + state = await recordedRunState(run, env) + } + return state === 'gone' ? 'gone' : 'unknown' +} + function readIfPresent(path: string): string | null { try { return readFileSync(path, 'utf8') @@ -576,7 +676,12 @@ export async function closeRunPane(handle: TmuxRunHandle, env: NodeJS.ProcessEnv isRunComplete(handle) && (await startedByRun(handle, env)) if (state !== 'ours' && !finishedUntracked) return - const killed = await runTmux(['kill-pane', '-t', handle.pane], env) + const killed = await runTmux( + handle.runId + ? ifTagged(handle.pane, handle.runId, `kill-pane -t ${handle.pane}`) + : ['kill-pane', '-t', handle.pane], + env + ) if (!killed.ok) { logger.warn('Could not close the tmux run pane', { error: killed.stderr.trim() }) }