Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 22 additions & 1 deletion apps/desktop/src/main/desktop-executor/executor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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', () => {
Expand All @@ -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')]
Expand Down
7 changes: 7 additions & 0 deletions apps/desktop/src/main/desktop-executor/executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
31 changes: 31 additions & 0 deletions apps/desktop/src/main/desktop-executor/service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand Down Expand Up @@ -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()
Expand Down
26 changes: 26 additions & 0 deletions apps/desktop/src/main/desktop-executor/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -61,13 +62,21 @@ 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 {
start(): void
/** 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<Set<string>>
/** Stores one entry of a claimed import, as this device's registered session. */
importEntry(
request: DesktopImportEntryRequest,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -443,6 +453,22 @@ export function createDesktopExecutorService(
getDevice() {
return device
},
async pendingResults() {
const pending = new Set<string>()
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)
Expand Down
38 changes: 30 additions & 8 deletions apps/desktop/src/main/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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<BrowserWindow>()
Expand Down Expand Up @@ -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)
Comment thread
waleedlatif1 marked this conversation as resolved.
}
},
runner: createDesktopToolRunner({
preferences: () => desktopSettings.getPreferences(),
accountDataAvailable,
Expand Down Expand Up @@ -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'
Expand Down Expand Up @@ -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()
Expand Down
46 changes: 41 additions & 5 deletions apps/desktop/src/main/terminal/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -49,6 +50,7 @@ import {
killPane,
listPanes,
pollRun,
type RecordedRun,
resolveAttachment,
runPaneState,
sendKey,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -490,7 +494,10 @@ export class TerminalService {
private async reapFinishedRuns(terminalId: string, env: NodeJS.ProcessEnv): Promise<void> {
// 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
Expand All @@ -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)
Expand All @@ -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)
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -1269,6 +1293,7 @@ export class TerminalService {
* see through tmux.
*/
private async runInTmux(
toolCallId: string,
terminal: TerminalSession,
session: string,
args: TerminalToolArgs,
Expand All @@ -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)
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading