Skip to content

Commit 54db79e

Browse files
committed
refactor(desktop): one forward-only state per recorded tmux run, and one check for a delivered result
A recorded run is started, delivered or stop, and only moves forward, so a late acknowledgement can never undo a stop. The executor alone decides which results reach the model (not a stop, nor one that stands in for a run that did not start, has an unknown outcome or was too large to send), and snapshots the journal before its own recovery. A delivery looks up its run by call instead of reading every record. A run that cannot start after tagging closes its pane like a refused one. At launch, a kept run whose pane outlived its command is closed and forgotten. Covers a newline in a pane's running command and a character split across tmux's output.
1 parent 566cea6 commit 54db79e

13 files changed

Lines changed: 348 additions & 172 deletions

File tree

‎apps/desktop/src/main/desktop-executor/executor.test.ts‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -224,6 +224,27 @@ describe('claiming', () => {
224224
await vi.waitFor(() => expect(delivered).toEqual(['call-recorded']))
225225
})
226226

227+
it('reports no result that stands in for what the action produced as reaching the model', async () => {
228+
const { sim, runner, executor, delivered } = setup()
229+
const standIns: DesktopToolCompletion[] = [
230+
{ status: 'cancelled', message: 'Stopped.' },
231+
{ status: 'error', message: 'Too large.', data: { resultOmitted: true } },
232+
{ status: 'error', message: 'Not started.', data: { notStarted: true } },
233+
{ status: 'error', message: 'Unknown.', data: { outcomeUnknown: true } },
234+
]
235+
for (const [index, completion] of standIns.entries()) {
236+
runner.immediate = completion
237+
sim.inbox = [callItem(`stand-in-${index}`, 'chat-a')]
238+
await executor.reconcile()
239+
await vi.waitFor(() => expect(sim.completions).toHaveLength(index + 1))
240+
}
241+
242+
runner.immediate = DONE
243+
sim.inbox = [callItem('real', 'chat-a')]
244+
await executor.reconcile()
245+
await vi.waitFor(() => expect(delivered).toEqual(['real']))
246+
})
247+
227248
it('claims a whole backlog at once, before any of it runs', async () => {
228249
const { sim, runner, executor } = setup()
229250
sim.inbox = [callItem('a-1', 'chat-a'), callItem('a-2', 'chat-a'), callItem('a-3', 'chat-a')]

‎apps/desktop/src/main/desktop-executor/executor.ts‎

Lines changed: 21 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,20 @@ export interface DesktopToolRunner {
5656

5757
export type DesktopApprovalItem = Extract<DesktopInboxItem, { kind: 'approval_needed' }>
5858

59+
/**
60+
* Whether a result hands the model what the action produced: not a stop, and not one standing in
61+
* for an action that did not start, whose outcome is unknown, or whose output was too large to send.
62+
*/
63+
export function isDeliveredResult(completion: DesktopToolCompletion): boolean {
64+
const data = completion.data
65+
return (
66+
completion.status !== 'cancelled' &&
67+
data?.notStarted !== true &&
68+
data?.outcomeUnknown !== true &&
69+
data?.resultOmitted !== true
70+
)
71+
}
72+
5973
export interface DesktopExecutorOptions {
6074
client: DesktopExecutorClient
6175
journal: ExecutorJournal
@@ -68,10 +82,11 @@ export interface DesktopExecutorOptions {
6882
/** Called whenever the number of held calls changes between zero and more. */
6983
onBusyChange?: (busy: boolean) => void
7084
/**
71-
* Called once Sim has taken a call's result as the call's own (recorded, or a duplicate of one
72-
* it recorded): the model has it, so anything it hands back (a pane still running) is in use.
85+
* Called once Sim has taken a call's real result ({@link isDeliveredResult}) as the call's own
86+
* (recorded, or a duplicate of one it recorded): the model has it, so anything it hands back (a
87+
* pane still running) is in use.
7388
*/
74-
onResultDelivered?: (toolCallId: string, completion: DesktopToolCompletion) => void
89+
onResultDelivered?: (toolCallId: string) => void
7590
maxHeldCalls?: number
7691
/** First delivery retry delay; tests shorten it. */
7792
retryBaseMs?: number
@@ -407,7 +422,9 @@ export class DesktopExecutor {
407422
})
408423
logger.info('Desktop call result acknowledged', { toolCallId, outcome })
409424
// Superseded: Sim settled the call first, so this result never reached the model.
410-
if (outcome !== 'superseded') this.options.onResultDelivered?.(toolCallId, pending)
425+
if (outcome !== 'superseded' && isDeliveredResult(pending)) {
426+
this.options.onResultDelivered?.(toolCallId)
427+
}
411428
break
412429
} catch (error) {
413430
// Encoding failed on this machine, so nothing was sent; the same data would fail again.

‎apps/desktop/src/main/desktop-executor/service.test.ts‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ function fakeSim(protocolVersion = 1) {
6464
})
6565
}
6666
if (path === '/api/desktop/tool/lease') return Response.json({ renewed: true })
67+
if (path === '/api/desktop/tool/complete') return Response.json({ outcome: 'recorded' })
6768
return new Promise<Response>((_resolve, reject) =>
6869
init.signal?.addEventListener('abort', () => reject(new Error('aborted')))
6970
)
@@ -114,6 +115,18 @@ describe('results recovery will hand to the model', () => {
114115
executionToken: 't4',
115116
completion: { status: 'error', message: 'x', data: { notStarted: true } },
116117
})
118+
await journal.put({
119+
toolCallId: 'stopped',
120+
state: 'result',
121+
executionToken: 't6',
122+
completion: { status: 'cancelled', message: 'Stopped.' },
123+
})
124+
await journal.put({
125+
toolCallId: 'too-large',
126+
state: 'result',
127+
executionToken: 't7',
128+
completion: { status: 'error', message: 'x', data: { resultOmitted: true } },
129+
})
117130
await journal.put({
118131
toolCallId: 'handed-back',
119132
state: 'result',
@@ -124,6 +137,39 @@ describe('results recovery will hand to the model', () => {
124137

125138
expect([...(await desktopExecutor.pendingResults())]).toEqual(['handed-back'])
126139
})
140+
141+
it('counts an unreadable journal as holding none', async () => {
142+
const userData = await mkdtemp(join(tmpdir(), 'sim-executor-service-'))
143+
// A directory where the journal file should be: every read of it fails.
144+
await mkdir(join(userData, 'desktop-executor-journal.json'))
145+
const { desktopExecutor } = await service(1, userData)
146+
147+
expect([...(await desktopExecutor.pendingResults())]).toEqual([])
148+
})
149+
150+
it('names what the journal held before recovery sent it, even when asked after', async () => {
151+
const userData = await mkdtemp(join(tmpdir(), 'sim-executor-service-'))
152+
const path = join(userData, 'desktop-executor-journal.json')
153+
await createExecutorJournal(path).put({
154+
toolCallId: 'handed-back',
155+
state: 'result',
156+
executionToken: 't1',
157+
completion: { status: 'success', message: 'running', data: { status: 'running' } },
158+
})
159+
const { sim, desktopExecutor } = await service(1, userData)
160+
desktopExecutor.start()
161+
await vi.waitFor(() => expect(sim.registrations).toHaveLength(1))
162+
sim.registrations[0]?.(true)
163+
164+
// Recovery hands the result to Sim and drops it from the journal.
165+
await vi.waitFor(async () => {
166+
expect(sim.requests).toContain('POST /api/desktop/tool/complete')
167+
expect(await createExecutorJournal(path).load()).toEqual([])
168+
})
169+
170+
expect([...(await desktopExecutor.pendingResults())]).toEqual(['handed-back'])
171+
await desktopExecutor.signOut()
172+
})
127173
})
128174

129175
describe('desktop executor registration', () => {

‎apps/desktop/src/main/desktop-executor/service.ts‎

Lines changed: 28 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@
77
import { hostname } from 'node:os'
88
import { join } from 'node:path'
99
import type { DesktopExecutorDevice } from '@sim/desktop-bridge'
10-
import type { DesktopToolCompletion } from '@sim/desktop-bridge/tool-results'
1110
import { createLogger } from '@sim/logger'
1211
import { getErrorMessage } from '@sim/utils/errors'
1312
import { generateId } from '@sim/utils/id'
@@ -32,6 +31,7 @@ import {
3231
type DesktopApprovalItem,
3332
DesktopExecutor,
3433
type DesktopToolRunner,
34+
isDeliveredResult,
3535
} from '@/main/desktop-executor/executor'
3636
import { createExecutorJournal } from '@/main/desktop-executor/journal'
3737
import {
@@ -62,8 +62,8 @@ export interface DesktopExecutorServiceDeps {
6262
onApprovals?: (items: DesktopApprovalItem[]) => void
6363
/** Whether any chat has desktop work claimed on this machine changed. */
6464
onBusyChange?: (busy: boolean) => void
65-
/** Sim has taken a call's result as the call's own, so the model has it. */
66-
onResultDelivered?: (toolCallId: string, completion: DesktopToolCompletion) => void
65+
/** Sim has taken a call's real result as the call's own, so the model has it. */
66+
onResultDelivered?: (toolCallId: string) => void
6767
}
6868

6969
export interface DesktopExecutorService {
@@ -72,9 +72,9 @@ export interface DesktopExecutorService {
7272
refreshRegistration(): void
7373
getDevice(): DesktopExecutorDevice | null
7474
/**
75-
* The calls whose real result (not one reported as not started or outcome unknown) is in the
76-
* journal and not yet acknowledged, so recovery will hand it to the model. Read before recovery
77-
* changes the journal; empty when it cannot be read.
75+
* The calls whose real result ({@link isDeliveredResult}) is in the journal and not yet
76+
* acknowledged, so recovery will hand it to the model. Read once, before recovery can change the
77+
* journal, whenever it is asked; empty when the journal cannot be read.
7878
*/
7979
pendingResults(): Promise<Set<string>>
8080
/** Stores one entry of a claimed import, as this device's registered session. */
@@ -138,6 +138,24 @@ export function createDesktopExecutorService(
138138
let registrationFailedOffline = false
139139
let suspended = false
140140
let started = false
141+
/** The journal's pending real results as they stood before this process's recovery. */
142+
let pendingSnapshot: Promise<Set<string>> | null = null
143+
const pendingResultsSnapshot = (): Promise<Set<string>> => {
144+
// An unreadable journal loads as empty, so it holds none.
145+
pendingSnapshot ??= journal
146+
.load()
147+
.then(
148+
(entries) =>
149+
new Set(
150+
entries.flatMap((entry) =>
151+
entry.state === 'result' && isDeliveredResult(entry.completion)
152+
? [entry.toolCallId]
153+
: []
154+
)
155+
)
156+
)
157+
return pendingSnapshot
158+
}
141159
/** Bumped on sign-out, so work started for the previous session cannot resume it. */
142160
let generation = 0
143161
let executorDeviceId: string | null = null
@@ -250,6 +268,8 @@ export function createDesktopExecutorService(
250268
...(deps.onBusyChange ? { onBusyChange: deps.onBusyChange } : {}),
251269
...(deps.onResultDelivered ? { onResultDelivered: deps.onResultDelivered } : {}),
252270
})
271+
// Recovery rewrites the journal, so what it held before is read first.
272+
await pendingResultsSnapshot()
253273
await executor.recover()
254274
// Signed out while recovering: sign-out already disposed this executor.
255275
if (registrationGeneration !== generation || !executor) return
@@ -453,21 +473,8 @@ export function createDesktopExecutorService(
453473
getDevice() {
454474
return device
455475
},
456-
async pendingResults() {
457-
const pending = new Set<string>()
458-
try {
459-
for (const entry of await journal.load()) {
460-
if (entry.state !== 'result') continue
461-
const data = entry.completion.data
462-
if (data?.outcomeUnknown === true || data?.notStarted === true) continue
463-
pending.add(entry.toolCallId)
464-
}
465-
} catch (error) {
466-
logger.warn('Could not read the executor journal for pending results', {
467-
error: getErrorMessage(error),
468-
})
469-
}
470-
return pending
476+
pendingResults() {
477+
return pendingResultsSnapshot()
471478
},
472479
importEntry(request, signal) {
473480
if (!client) throw new Error('The Sim desktop app is not signed in to Sim.')

‎apps/desktop/src/main/index.ts‎

Lines changed: 7 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -607,13 +607,9 @@ function main(): void {
607607
accountDataAvailable,
608608
onApprovals: (items) => approvalNotifier.update(items),
609609
onBusyChange: (busy) => sleepBlocker.setBusy(busy),
610-
// A result the model has (not one reported as not started or outcome unknown) makes a tmux
611-
// run it handed back as still going collectable across a restart.
612-
onResultDelivered: (toolCallId, completion) => {
613-
if (completion.data?.outcomeUnknown !== true && completion.data?.notStarted !== true) {
614-
terminal.markRunDelivered(toolCallId)
615-
}
616-
},
610+
// A result the model has makes a tmux run it handed back as still going collectable across a
611+
// restart.
612+
onResultDelivered: (toolCallId) => terminal.markRunDelivered(toolCallId),
617613
runner: createDesktopToolRunner({
618614
preferences: () => desktopSettings.getPreferences(),
619615
accountDataAvailable,
@@ -817,13 +813,10 @@ function main(): void {
817813
}
818814
}
819815

820-
// The same user's tmux runs from a previous process: a run whose call never handed back its
821-
// result (or whose result the journal will report as unknown) has nothing left to collect what
822-
// it does, so it is stopped, while its pane still carries its tag. A run already handed back as
823-
// still going, with its pane, is left to the model, which may come back to it. Read before the
824-
// executor starts, since its recovery rewrites the journal.
825-
const pendingResults = desktopExecutor.pendingResults()
826-
void pendingResults.then((pending) => terminal.stopUncollectableRuns(pending))
816+
// The same user's tmux runs from a previous process: one whose pane the model has, or will get
817+
// from recovery, is left to the model, which may come back to it; any other has nothing left
818+
// to collect what it does, so it is stopped, while its pane still carries its tag.
819+
void desktopExecutor.pendingResults().then((pending) => terminal.stopUncollectableRuns(pending))
827820

828821
if (!accountDataAvailable()) {
829822
logger.warn(
@@ -974,7 +967,6 @@ function main(): void {
974967
ensureAppSession().cookies.on('changed', (_event, cookie, _cause, removed) => {
975968
if (!removed && isSessionCookieName(cookie.name)) desktopExecutor.refreshRegistration()
976969
})
977-
await pendingResults
978970
desktopExecutor.start()
979971
}
980972
await ensureMainWindow()

‎apps/desktop/src/main/terminal/index.ts‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1310,8 +1310,7 @@ export class TerminalService {
13101310
...(ledger
13111311
? {
13121312
beforeStart: (run: RecordedRun) =>
1313-
ledger.record({ ...run, callId: toolCallId, delivered: false }),
1314-
abandon: (runId: string) => ledger.forget(runId),
1313+
ledger.record({ ...run, callId: toolCallId, state: 'started' }),
13151314
}
13161315
: {}),
13171316
})

0 commit comments

Comments
 (0)