diff --git a/apps/desktop/src/main/desktop-executor/executor.test.ts b/apps/desktop/src/main/desktop-executor/executor.test.ts index 795d4df54de..509f1c5e271 100644 --- a/apps/desktop/src/main/desktop-executor/executor.test.ts +++ b/apps/desktop/src/main/desktop-executor/executor.test.ts @@ -224,6 +224,27 @@ describe('claiming', () => { await vi.waitFor(() => expect(delivered).toEqual(['call-recorded'])) }) + it('reports no result that stands in for what the action produced as reaching the model', async () => { + const { sim, runner, executor, delivered } = setup() + const standIns: DesktopToolCompletion[] = [ + { status: 'cancelled', message: 'Stopped.' }, + { status: 'error', message: 'Too large.', data: { resultOmitted: true } }, + { status: 'error', message: 'Not started.', data: { notStarted: true } }, + { status: 'error', message: 'Unknown.', data: { outcomeUnknown: true } }, + ] + for (const [index, completion] of standIns.entries()) { + runner.immediate = completion + sim.inbox = [callItem(`stand-in-${index}`, 'chat-a')] + await executor.reconcile() + await vi.waitFor(() => expect(sim.completions).toHaveLength(index + 1)) + } + + runner.immediate = DONE + sim.inbox = [callItem('real', 'chat-a')] + await executor.reconcile() + await vi.waitFor(() => expect(delivered).toEqual(['real'])) + }) + 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 6ca794182cd..d41998bc84d 100644 --- a/apps/desktop/src/main/desktop-executor/executor.ts +++ b/apps/desktop/src/main/desktop-executor/executor.ts @@ -56,6 +56,20 @@ export interface DesktopToolRunner { export type DesktopApprovalItem = Extract +/** + * Whether a result hands the model what the action produced: not a stop, and not one standing in + * for an action that did not start, whose outcome is unknown, or whose output was too large to send. + */ +export function isDeliveredResult(completion: DesktopToolCompletion): boolean { + const data = completion.data + return ( + completion.status !== 'cancelled' && + data?.notStarted !== true && + data?.outcomeUnknown !== true && + data?.resultOmitted !== true + ) +} + export interface DesktopExecutorOptions { client: DesktopExecutorClient journal: ExecutorJournal @@ -68,10 +82,11 @@ export interface DesktopExecutorOptions { /** 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. + * Called once Sim has taken a call's real result ({@link isDeliveredResult}) 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 + onResultDelivered?: (toolCallId: string) => void maxHeldCalls?: number /** First delivery retry delay; tests shorten it. */ retryBaseMs?: number @@ -407,7 +422,9 @@ export class DesktopExecutor { }) 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) + if (outcome !== 'superseded' && isDeliveredResult(pending)) { + this.options.onResultDelivered?.(toolCallId) + } 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 924ad2b551b..363e418fedf 100644 --- a/apps/desktop/src/main/desktop-executor/service.test.ts +++ b/apps/desktop/src/main/desktop-executor/service.test.ts @@ -64,6 +64,7 @@ function fakeSim(protocolVersion = 1) { }) } if (path === '/api/desktop/tool/lease') return Response.json({ renewed: true }) + if (path === '/api/desktop/tool/complete') return Response.json({ outcome: 'recorded' }) return new Promise((_resolve, reject) => init.signal?.addEventListener('abort', () => reject(new Error('aborted'))) ) @@ -114,6 +115,18 @@ describe('results recovery will hand to the model', () => { executionToken: 't4', completion: { status: 'error', message: 'x', data: { notStarted: true } }, }) + await journal.put({ + toolCallId: 'stopped', + state: 'result', + executionToken: 't6', + completion: { status: 'cancelled', message: 'Stopped.' }, + }) + await journal.put({ + toolCallId: 'too-large', + state: 'result', + executionToken: 't7', + completion: { status: 'error', message: 'x', data: { resultOmitted: true } }, + }) await journal.put({ toolCallId: 'handed-back', state: 'result', @@ -124,6 +137,39 @@ describe('results recovery will hand to the model', () => { expect([...(await desktopExecutor.pendingResults())]).toEqual(['handed-back']) }) + + it('counts an unreadable journal as holding none', async () => { + const userData = await mkdtemp(join(tmpdir(), 'sim-executor-service-')) + // A directory where the journal file should be: every read of it fails. + await mkdir(join(userData, 'desktop-executor-journal.json')) + const { desktopExecutor } = await service(1, userData) + + expect([...(await desktopExecutor.pendingResults())]).toEqual([]) + }) + + it('names what the journal held before recovery sent it, even when asked after', async () => { + const userData = await mkdtemp(join(tmpdir(), 'sim-executor-service-')) + const path = join(userData, 'desktop-executor-journal.json') + await createExecutorJournal(path).put({ + toolCallId: 'handed-back', + state: 'result', + executionToken: 't1', + completion: { status: 'success', message: 'running', data: { status: 'running' } }, + }) + const { sim, desktopExecutor } = await service(1, userData) + desktopExecutor.start() + await vi.waitFor(() => expect(sim.registrations).toHaveLength(1)) + sim.registrations[0]?.(true) + + // Recovery hands the result to Sim and drops it from the journal. + await vi.waitFor(async () => { + expect(sim.requests).toContain('POST /api/desktop/tool/complete') + expect(await createExecutorJournal(path).load()).toEqual([]) + }) + + expect([...(await desktopExecutor.pendingResults())]).toEqual(['handed-back']) + await desktopExecutor.signOut() + }) }) describe('desktop executor registration', () => { diff --git a/apps/desktop/src/main/desktop-executor/service.ts b/apps/desktop/src/main/desktop-executor/service.ts index 7611dfd5c2f..a46ffc29377 100644 --- a/apps/desktop/src/main/desktop-executor/service.ts +++ b/apps/desktop/src/main/desktop-executor/service.ts @@ -7,7 +7,6 @@ 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' @@ -32,6 +31,7 @@ import { type DesktopApprovalItem, DesktopExecutor, type DesktopToolRunner, + isDeliveredResult, } from '@/main/desktop-executor/executor' import { createExecutorJournal } from '@/main/desktop-executor/journal' import { @@ -62,8 +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 + /** Sim has taken a call's real result as the call's own, so the model has it. */ + onResultDelivered?: (toolCallId: string) => void } export interface DesktopExecutorService { @@ -72,9 +72,9 @@ export interface DesktopExecutorService { 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. + * The calls whose real result ({@link isDeliveredResult}) is in the journal and not yet + * acknowledged, so recovery will hand it to the model. Read once, before recovery can change the + * journal, whenever it is asked; empty when the journal cannot be read. */ pendingResults(): Promise> /** Stores one entry of a claimed import, as this device's registered session. */ @@ -138,6 +138,24 @@ export function createDesktopExecutorService( let registrationFailedOffline = false let suspended = false let started = false + /** The journal's pending real results as they stood before this process's recovery. */ + let pendingSnapshot: Promise> | null = null + const pendingResultsSnapshot = (): Promise> => { + // An unreadable journal loads as empty, so it holds none. + pendingSnapshot ??= journal + .load() + .then( + (entries) => + new Set( + entries.flatMap((entry) => + entry.state === 'result' && isDeliveredResult(entry.completion) + ? [entry.toolCallId] + : [] + ) + ) + ) + return pendingSnapshot + } /** Bumped on sign-out, so work started for the previous session cannot resume it. */ let generation = 0 let executorDeviceId: string | null = null @@ -250,6 +268,8 @@ export function createDesktopExecutorService( ...(deps.onBusyChange ? { onBusyChange: deps.onBusyChange } : {}), ...(deps.onResultDelivered ? { onResultDelivered: deps.onResultDelivered } : {}), }) + // Recovery rewrites the journal, so what it held before is read first. + await pendingResultsSnapshot() await executor.recover() // Signed out while recovering: sign-out already disposed this executor. if (registrationGeneration !== generation || !executor) return @@ -453,21 +473,8 @@ 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 + pendingResults() { + return pendingResultsSnapshot() }, importEntry(request, signal) { if (!client) throw new Error('The Sim desktop app is not signed in to Sim.') diff --git a/apps/desktop/src/main/index.ts b/apps/desktop/src/main/index.ts index d777615929f..1cd08b432bc 100644 --- a/apps/desktop/src/main/index.ts +++ b/apps/desktop/src/main/index.ts @@ -607,13 +607,9 @@ 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) - } - }, + // A result the model has makes a tmux run it handed back as still going collectable across a + // restart. + onResultDelivered: (toolCallId) => terminal.markRunDelivered(toolCallId), runner: createDesktopToolRunner({ preferences: () => desktopSettings.getPreferences(), accountDataAvailable, @@ -817,13 +813,10 @@ 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)) + // The same user's tmux runs from a previous process: one whose pane the model has, or will get + // from recovery, is left to the model, which may come back to it; any other has nothing left + // to collect what it does, so it is stopped, while its pane still carries its tag. + void desktopExecutor.pendingResults().then((pending) => terminal.stopUncollectableRuns(pending)) if (!accountDataAvailable()) { logger.warn( @@ -974,7 +967,6 @@ 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 0150434ef0a..6c118f29254 100644 --- a/apps/desktop/src/main/terminal/index.ts +++ b/apps/desktop/src/main/terminal/index.ts @@ -1310,8 +1310,7 @@ export class TerminalService { ...(ledger ? { beforeStart: (run: RecordedRun) => - ledger.record({ ...run, callId: toolCallId, delivered: false }), - abandon: (runId: string) => ledger.forget(runId), + ledger.record({ ...run, callId: toolCallId, state: 'started' }), } : {}), }) diff --git a/apps/desktop/src/main/terminal/registry-run-ledger.test.ts b/apps/desktop/src/main/terminal/registry-run-ledger.test.ts index 4e3cc036c5f..a505e633b51 100644 --- a/apps/desktop/src/main/terminal/registry-run-ledger.test.ts +++ b/apps/desktop/src/main/terminal/registry-run-ledger.test.ts @@ -6,11 +6,12 @@ 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. + * Each recorded run's pane as tmux has it: running, dead after its command ended + * (`remain-on-exit`), stopped by Sim, already gone, or one tmux cannot answer for. A pane not + * listed is running. */ const { panes } = vi.hoisted(() => ({ - panes: new Map(), + panes: new Map(), })) vi.mock('@/main/terminal/tmux', async () => { @@ -21,18 +22,22 @@ vi.mock('@/main/terminal/tmux', async () => { ...actual, recordedRunState: async (run: { runId: string }) => { const pane = state(run.runId) - return pane === 'running' ? 'ours' : pane === 'unknown' ? 'unknown' : 'gone' + if (pane === 'running') return 'ours' + if (pane === 'dead') return 'finished' + return 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') + if (state(run.runId) === 'running' || state(run.runId) === 'dead') { + panes.set(run.runId, 'stopped') + } return 'gone' }, } }) import { TerminalRegistry } from '@/main/terminal/registry' -import { createRunLedger } from '@/main/terminal/run-ledger' +import { createRunLedger, type RunState } from '@/main/terminal/run-ledger' const dirs: string[] = [] @@ -42,8 +47,8 @@ function ledgerDir(): string { 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 } +function run(runId: string, pane: string, state: RunState = 'started') { + return { runId, pane, socket: '/tmp/tmux-501/default', callId: `call-${runId}`, state } } afterEach(() => { @@ -70,12 +75,15 @@ describe('stopping recorded tmux runs', () => { 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('handed-back', '%1', 'delivered')) 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)) + previous.record(run('handed-back-and-gone', '%4', 'delivered')) panes.set('handed-back-and-gone', 'gone') + // Handed back, but its command has since ended: only its dead pane is left. + previous.record(run('handed-back-and-done', '%5', 'delivered')) + panes.set('handed-back-and-done', 'dead') const ledger = createRunLedger(dir) await new TerminalRegistry(undefined, undefined, ledger).stopUncollectableRuns( @@ -85,6 +93,7 @@ describe('stopping recorded tmux runs', () => { expect(Object.fromEntries(panes)).toEqual({ 'never-handed-back': 'stopped', 'handed-back-and-gone': 'gone', + 'handed-back-and-done': 'stopped', }) expect( ledger @@ -96,7 +105,7 @@ describe('stopping recorded tmux runs', () => { 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)) + createRunLedger(dir).record(run('unconfirmed', '%1', 'delivered')) panes.set('unconfirmed', 'unknown') const ledger = createRunLedger(dir) // Sign-out could not confirm the run ended. @@ -110,7 +119,7 @@ describe('stopping recorded tmux runs', () => { expect(ledger.list()).toEqual([]) }) - it('notes a run as handed back once its call result is durable', () => { + it('notes the run of a call whose result reached the model as handed back', () => { const dir = ledgerDir() const ledger = createRunLedger(dir) ledger.record(run('watched', '%1')) @@ -118,9 +127,9 @@ describe('stopping recorded tmux runs', () => { new TerminalRegistry(undefined, undefined, ledger).markRunDelivered('call-watched') - expect( - Object.fromEntries(ledger.list().map((record) => [record.runId, record.delivered])) - ).toEqual({ watched: true, other: false }) + expect(Object.fromEntries(ledger.list().map((record) => [record.runId, record.state]))).toEqual( + { watched: 'delivered', other: 'started' } + ) }) it("at launch, stops the previous process's runs and none this one has started", async () => { diff --git a/apps/desktop/src/main/terminal/registry.ts b/apps/desktop/src/main/terminal/registry.ts index 6ecc7a8a38a..39e41edbaaa 100644 --- a/apps/desktop/src/main/terminal/registry.ts +++ b/apps/desktop/src/main/terminal/registry.ts @@ -391,12 +391,13 @@ export class TerminalRegistry { * 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. + * Every other run is stopped, as is one a stop for everything could not confirm (`stop`). */ stopUncollectableRuns(pendingResults: ReadonlySet): Promise { return this.stopRecordedRuns({ excludeLive: true, - keep: (run) => !run.mustStop && (run.delivered || pendingResults.has(run.callId)), + keep: (run) => + run.state === 'delivered' || (run.state === 'started' && pendingResults.has(run.callId)), }) } @@ -405,11 +406,8 @@ export class TerminalRegistry { * 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) - } + const runId = this.runLedger?.runOf(callId) + if (runId) this.runLedger?.advance(runId, 'delivered') } /** @@ -429,14 +427,15 @@ export class TerminalRegistry { 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) + let keeping = options.keep?.(run) ?? false + let state = keeping ? await recordedRunState(run, process.env) : 'ours' + // A kept run whose command already ended leaves only its dead pane (`remain-on-exit`). + if (state === 'finished') keeping = false + if (!keeping) state = 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) + else if (!keeping) ledger.advance(run.runId, 'stop') }) ) } diff --git a/apps/desktop/src/main/terminal/run-ledger.test.ts b/apps/desktop/src/main/terminal/run-ledger.test.ts index 9baf932b211..69220826737 100644 --- a/apps/desktop/src/main/terminal/run-ledger.test.ts +++ b/apps/desktop/src/main/terminal/run-ledger.test.ts @@ -17,7 +17,7 @@ const RUN = { pane: '%3', socket: '/tmp/tmux-501/default', callId: 'call-1', - delivered: false, + state: 'started' as const, } afterEach(() => { @@ -54,6 +54,11 @@ describe('the tmux run ledger', () => { 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' })) + // A state no version wrote: nothing could tell what to do with the run. + writeFileSync( + join(dir, 'run-5.json'), + JSON.stringify({ ...RUN, runId: 'run-5', state: 'done' }) + ) // Only a pane id names one pane; a target like this one names whatever pane is active there. writeFileSync( join(dir, 'run-6.json'), @@ -95,14 +100,59 @@ describe('the tmux run ledger', () => { expect(ledger.record(RUN)).toBe(false) }) - it('notes a run handed back as still going, for the next process too', () => { + it('moves a run forward only, for the next process too', () => { const dir = scratch() const ledger = createRunLedger(dir) ledger.record(RUN) - ledger.markDelivered(RUN.runId) + ledger.advance(RUN.runId, 'delivered') + expect(createRunLedger(dir).list()).toEqual([{ ...RUN, state: 'delivered' }]) + + ledger.advance(RUN.runId, 'stop') + // A late acknowledgement never undoes a stop. + ledger.advance(RUN.runId, 'delivered') + expect(createRunLedger(dir).list()).toEqual([{ ...RUN, state: 'stop' }]) + }) + + it('reads a record saved before runs had a state as the state it meant', () => { + const dir = scratch() + mkdirSync(dir, { recursive: true }) + const { state: _state, ...saved } = RUN + const legacy = { + 'run-started': { delivered: false }, + 'run-delivered': { delivered: true }, + 'run-stop': { delivered: false, mustStop: true }, + // A stop wins over an acknowledgement. + 'run-stop-delivered': { delivered: true, mustStop: true }, + } + for (const [runId, fields] of Object.entries(legacy)) { + writeFileSync(join(dir, `${runId}.json`), JSON.stringify({ ...saved, runId, ...fields })) + } + + const states = Object.fromEntries( + createRunLedger(dir) + .list() + .map((record) => [record.runId, record.state]) + ) + expect(states).toEqual({ + 'run-started': 'started', + 'run-delivered': 'delivered', + 'run-stop': 'stop', + 'run-stop-delivered': 'stop', + }) + }) + + it('finds the run a call started, a previous process recorded it or this one', () => { + const dir = scratch() + // Recovery may hand the model the result of a previous process's call. + createRunLedger(dir).record(RUN) + const ledger = createRunLedger(dir) + ledger.record({ ...RUN, runId: 'run-2', callId: 'call-2' }) - expect(createRunLedger(dir).list()).toEqual([{ ...RUN, delivered: true }]) + expect(ledger.runOf(RUN.callId)).toBe(RUN.runId) + expect(ledger.runOf('call-2')).toBe('run-2') + ledger.forget(RUN.runId) + expect(ledger.runOf(RUN.callId)).toBeUndefined() }) it('records nothing for a run tag that is not a plain id', () => { diff --git a/apps/desktop/src/main/terminal/run-ledger.ts b/apps/desktop/src/main/terminal/run-ledger.ts index badf2806feb..8a1e2cb38ca 100644 --- a/apps/desktop/src/main/terminal/run-ledger.ts +++ b/apps/desktop/src/main/terminal/run-ledger.ts @@ -18,26 +18,32 @@ 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. */ +/** + * Where a recorded run stands, moving only forward: + * - `started`: its call has not handed its result to the model; + * - `delivered`: the model has the result that handed the run back as still going (with its + * pane), so a restart must leave the run be; + * - `stop`: a stop for everything (sign-out, Terminal off) could not confirm the run ended, so + * it is stopped whatever comes after, a late `delivered` included. + */ +export type RunState = 'started' | 'delivered' | 'stop' + +const RUN_STATES: readonly RunState[] = ['started', 'delivered', 'stop'] + +/** A recorded run, with the call it belongs to and where it stands. */ 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 + state: RunState } 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 + /** Moves a run forward to `state`; a run already there or past it stays as it is. */ + advance(runId: string, state: RunState): void + /** The recorded run a tool call started, if the ledger still holds one. */ + runOf(callId: string): string | undefined forget(runId: string): void /** Every recorded run; `excludeLive` leaves out runs this process recorded. */ list(options?: { excludeLive?: boolean }): RunRecord[] @@ -46,16 +52,31 @@ export interface RunLedger { /** Run tags are generated ids; anything else in the directory is not a record. */ const RUN_ID = /^[A-Za-z0-9_-]{1,128}$/ +/** A record as saved: the current shape, or the earlier one with `delivered` and `mustStop`. */ +type SavedRecord = Partial & { delivered?: unknown; mustStop?: unknown } + +/** + * The record's state. A record from before `state` existed (a released build wrote them) is read + * as the state it meant: `mustStop` as `stop`, which wins, then `delivered` as `delivered`. + */ +function stateOf(saved: SavedRecord): RunState | null { + if (RUN_STATES.includes(saved.state as RunState)) return saved.state as RunState + if (typeof saved.delivered !== 'boolean') return null + if (saved.mustStop === true) return 'stop' + return saved.delivered ? 'delivered' : 'started' +} + function parseRecord(text: string): RunRecord | null { try { - const parsed = JSON.parse(text) as Partial + const parsed = JSON.parse(text) as SavedRecord + const state = stateOf(parsed) 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' && + state && typeof parsed.socket === 'string' && parsed.socket.startsWith('/') ) { @@ -64,8 +85,7 @@ function parseRecord(text: string): RunRecord | null { pane: parsed.pane, socket: parsed.socket, callId: parsed.callId, - delivered: parsed.delivered, - ...(parsed.mustStop === true ? { mustStop: true } : {}), + state, } } } catch { @@ -86,24 +106,45 @@ function remove(path: string): void { export function createRunLedger(dir: string): RunLedger { /** Runs recorded by this process, still going as far as it knows. */ const live = new Set() + /** Each recorded run by the call that started it: read from the directory once, then kept. */ + let byCall: Map | null = null 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 + /** Every saved record; anything else in the directory is removed. */ + const readAll = (): RunRecord[] => { + let names: string[] try { - record = parseRecord(readFileSync(pathFor(runId), 'utf8')) + names = readdirSync(dir) } catch { - record = null + return [] } - 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) }) + 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 + } + runs.push(record) } + return runs + } + + const index = (): Map => { + byCall ??= new Map(readAll().map((run) => [run.callId, run.runId])) + return byCall } return { @@ -112,53 +153,40 @@ export function createRunLedger(dir: string): RunLedger { try { writeJsonFileAtomicallySync(pathFor(run.runId), run) live.add(run.runId) + index().set(run.callId, 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 })) + advance(runId, state) { + if (!RUN_ID.test(runId)) return + let record: RunRecord | null = null + try { + record = parseRecord(readFileSync(pathFor(runId), 'utf8')) + } catch { + record = null + } + if (!record || RUN_STATES.indexOf(state) <= RUN_STATES.indexOf(record.state)) return + try { + writeJsonFileAtomicallySync(pathFor(runId), { ...record, state }) + } catch (error) { + // Left behind, a restart stops the run: the conservative side. + logger.warn('Could not update a tmux run record', { error: getErrorMessage(error) }) + } }, - markMustStop(runId) { - update(runId, (record) => (record.mustStop ? null : { ...record, mustStop: true })) + runOf(callId) { + return index().get(callId) }, forget(runId) { live.delete(runId) + for (const [callId, indexed] of byCall ?? []) if (indexed === runId) byCall?.delete(callId) 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 + const runs = readAll() + return options.excludeLive ? runs.filter((run) => !live.has(run.runId)) : runs }, } } diff --git a/apps/desktop/src/main/terminal/service.test.ts b/apps/desktop/src/main/terminal/service.test.ts index b2f24c9cec8..df6aa112d85 100644 --- a/apps/desktop/src/main/terminal/service.test.ts +++ b/apps/desktop/src/main/terminal/service.test.ts @@ -666,8 +666,8 @@ describe('agent commands in tmux', () => { pane, socket: '/tmp/tmux-fake/default', callId: 'call-long', - // Handed back, but only the executor's journal can make that durable. - delivered: false, + // Handed back, but the model has it only once Sim acknowledges the result. + state: 'started', }, ]) diff --git a/apps/desktop/src/main/terminal/tmux.test.ts b/apps/desktop/src/main/terminal/tmux.test.ts index 378c52e29e2..db2fe9f5288 100644 --- a/apps/desktop/src/main/terminal/tmux.test.ts +++ b/apps/desktop/src/main/terminal/tmux.test.ts @@ -10,6 +10,7 @@ import { listPanes, parseFormatLines, pollRun, + recordedRunState, resolveAttachment, runPaneState, startRun, @@ -188,7 +189,16 @@ interface FakeTmuxState { socket: string nextWindow: number nextPane: number - panes: Record; command?: string }> + panes: Record< + string, + { + window: string + options: Record + command?: string + /** Its command has ended, and `remain-on-exit` keeps the pane open. */ + dead?: boolean + } + > /** Every command that reached a pane: `send-keys %1 C-c`, `kill-pane %1`. */ log: string[] /** Commands the fake fails, with the error tmux would print. */ @@ -205,6 +215,8 @@ interface FakeTmuxState { clients?: Array<{ pid: string; tty: string; session: string }> /** A session's panes as `list-panes` reports them, with fields others can set. */ listed?: Array<{ windowName: string; command: string; cwd: string }> + /** list-panes writes its output in two pieces, split inside a multi-byte character. */ + splitWrites?: boolean /** Every `-F` format the fake was asked for, in order. */ formats?: string[] /** Commands the fake holds until the file named here exists, like a busy tmux server. */ @@ -286,16 +298,17 @@ switch (args[0]) { 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) - const value = !pane - ? '' - : name === 'pane_id' + const field = (name) => + name === 'pane_id' ? target() : name === 'socket_path' ? state.socket : name === 'pane_start_command' ? (pane.command ?? '') + : name === 'pane_dead' + ? (pane.dead ? '1' : '0') : (pane.options[name] ?? '') + const value = pane ? args[args.length - 1].replace(/#\\{([^}]+)\\}/g, (_, name) => field(name)) : '' process.stdout.write(value + '\\n') break } @@ -305,6 +318,7 @@ switch (args[0]) { const format = args[args.indexOf('-F') + 1] state.formats = [...(state.formats ?? []), format] save() + let output = '' ;(state.listed ?? []).forEach((pane, index) => { const fields = { session_name: 'work', @@ -320,8 +334,18 @@ switch (args[0]) { fields[name].split(from).join(to) ) .replace(/#\\{([a-z_]+)\\}/g, (_, name) => fields[name]) - process.stdout.write(line + '\\n') + output += line + '\\n' }) + const bytes = Buffer.from(output) + if (!state.splitWrites) { + process.stdout.write(bytes) + break + } + // The first piece ends inside a character, and arrives well before the rest. + const cut = bytes.findIndex((byte) => byte >= 0x80) + 1 + fs.writeSync(1, bytes.subarray(0, cut)) + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 100) + fs.writeSync(1, bytes.subarray(cut)) break } case 'list-clients': { @@ -398,7 +422,7 @@ describe("listing a session's panes", () => { tmux.write({ ...tmux.read(), listed: [ - { windowName: 'build\nlogs', command: 'make', cwd: '/tmp/a\nuser:0.0' }, + { windowName: 'build\nlogs', command: 'make\nuser:0.9', cwd: '/tmp/a\nuser:0.0' }, { windowName: 'zsh', command: 'zsh', cwd: '/tmp' }, ], }) @@ -407,7 +431,7 @@ describe("listing a session's panes", () => { { target: 'work:0.0', windowName: 'buildlogs', - command: 'make', + command: 'makeuser:0.9', cwd: '/tmp/auser:0.0', active: true, }, @@ -415,6 +439,18 @@ describe("listing a session's panes", () => { ]) }) + it('reads a character tmux wrote in two pieces as one', async () => { + const tmux = fakeTmux() + dirs.push(tmux.dir) + tmux.write({ + ...tmux.read(), + splitWrites: true, + listed: [{ windowName: 'café ✓', command: 'zsh', cwd: '/tmp' }], + }) + + expect((await listPanes('work', tmux.env))[0]?.windowName).toBe('café ✓') + }) + it('frames each call with a marker of its own', async () => { const tmux = fakeTmux() dirs.push(tmux.dir) @@ -481,6 +517,19 @@ describe('stopping a run another process started, from its record', () => { expect(tmux.read().panes).toEqual({}) }) + it('reads a pane kept open after its command ended as finished, and closes it without keys', async () => { + const tmux = fakeTmux() + const { record } = await recorded(tmux) + const state = tmux.read() + const pane = state.panes[record.pane] + if (pane) pane.dead = true + tmux.write(state) + + expect(await recordedRunState(record, tmux.env)).toBe('finished') + expect(await stopRecordedRun(record, tmux.env, 0)).toBe('gone') + expect(tmux.read().log).toEqual([`kill-pane ${record.pane}`]) + }) + it('never touches a pane that took the recorded id after tmux restarted', async () => { const tmux = fakeTmux() const { run, record } = await recorded(tmux) @@ -532,29 +581,25 @@ describe('stopping a run another process started, from its record', () => { expect(existsSync(marker)).toBe(false) }) - it('forgets the saved record of a run that then could not start', async () => { + it('closes the tagged pane 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 + runDir = /"([^"]+)\/run\.sh"/.exec(command)?.[1] ?? '' 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([]) + expect(tmux.read().panes).toEqual({}) }) it('acts on no pane that took the recorded id between the last check and the action', async () => { diff --git a/apps/desktop/src/main/terminal/tmux.ts b/apps/desktop/src/main/terminal/tmux.ts index b239bb9c66e..782b1a84990 100644 --- a/apps/desktop/src/main/terminal/tmux.ts +++ b/apps/desktop/src/main/terminal/tmux.ts @@ -391,8 +391,6 @@ export async function startRun( * 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-')) @@ -492,7 +490,8 @@ export async function startRun( try { writeFileSync(goPath, '') } catch (error) { - if (runId) options.abandon?.(runId) + // Like a refused run: its tagged pane closes now, and the next sweep drops a saved record. + if (runId) await runTmux(ifTagged(pane, runId, `kill-pane -t ${pane}`), env) dispose() return { error: `The command could not be started: ${getErrorMessage(error)}` } } @@ -523,11 +522,6 @@ export async function runPaneState( return /can't find|no server running/i.test(shown.stderr) ? 'gone' : 'unknown' } -/** - * Stops a run: Ctrl-C in its own pane, then closing that pane if the command ignored it. Every - * 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 @@ -545,6 +539,11 @@ function ifTagged(pane: string, runId: string, command: string, socket?: string) ] } +/** + * Stops a run: Ctrl-C in its own pane, then closing that pane if the command ignored it. Every + * 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. + */ export async function stopRun( handle: TmuxRunHandle, env: NodeJS.ProcessEnv, @@ -574,12 +573,17 @@ export interface RecordedRun { export async function recordedRunState( run: RecordedRun, env: NodeJS.ProcessEnv -): Promise<'ours' | 'gone' | 'unknown'> { +): Promise<'ours' | 'finished' | 'gone' | 'unknown'> { const shown = await runTmux( - ['-S', run.socket, 'display-message', '-p', '-t', run.pane, `#{${RUN_ID_OPTION}}`], + ['-S', run.socket, 'display-message', '-p', '-t', run.pane, `#{${RUN_ID_OPTION}} #{pane_dead}`], env ) - if (shown.ok) return shown.stdout.trim() === run.runId ? 'ours' : 'gone' + if (shown.ok) { + const [tag, dead] = shown.stdout.trim().split(' ') + if (tag !== run.runId) return 'gone' + // Kept open after its command ended (`remain-on-exit`): nothing in it is left to run. + return dead === '1' ? 'finished' : 'ours' + } return /can't find|no server running|error connecting|no such file/i.test(shown.stderr) ? 'gone' : 'unknown' @@ -597,19 +601,21 @@ export async function stopRecordedRun( 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 + if (before === 'gone' || before === 'unknown') return before + let state: 'ours' | 'finished' | 'gone' | 'unknown' = before + if (before === 'ours') { + await runTmux(ifTagged(run.pane, run.runId, `send-keys -t ${run.pane} C-c`, run.socket), env) + const deadline = Date.now() + graceMs + 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') { + if (state === 'ours' || state === 'finished') state = await recordedRunState(run, env) + if (state === 'ours' || state === 'finished') { await runTmux(ifTagged(run.pane, run.runId, `kill-pane -t ${run.pane}`, run.socket), env) state = await recordedRunState(run, env) }