Skip to content

Commit 9a1aca0

Browse files
authored
fix(desktop): stop a run the model never got the result of (#8749)
* fix(desktop): stop a run the model never got the result of * fix(desktop): never stop a run the model may hold * fix(desktop): only a recovered result leaves a duplicate open * fix(desktop): keep a recovered result's mark while it is parked
1 parent a7aaa2c commit 9a1aca0

8 files changed

Lines changed: 219 additions & 11 deletions

File tree

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

Lines changed: 114 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -170,6 +170,8 @@ function setup(
170170
const busy: boolean[] = []
171171
/** Calls whose result the executor reported as reaching the model, in order. */
172172
const delivered: string[] = []
173+
/** Calls Sim finished with without their real result reaching the model, in order. */
174+
const notDelivered: string[] = []
173175
const executor = new DesktopExecutor({
174176
client: sim.client,
175177
journal,
@@ -180,12 +182,23 @@ function setup(
180182
onBusyChange: (value) => busy.push(value),
181183
onApprovals: (items) => approvals.push(items),
182184
onResultDelivered: (toolCallId) => delivered.push(toolCallId),
185+
onResultNotDelivered: (toolCallId) => notDelivered.push(toolCallId),
183186
...(options.maxHeldCalls ? { maxHeldCalls: options.maxHeldCalls } : {}),
184187
...(options.deliveryAwakeLimitMs !== undefined
185188
? { deliveryAwakeLimitMs: options.deliveryAwakeLimitMs }
186189
: {}),
187190
})
188-
return { sim, journal, runner, executor, onUnregistered, busy, approvals, delivered }
191+
return {
192+
sim,
193+
journal,
194+
runner,
195+
executor,
196+
onUnregistered,
197+
busy,
198+
approvals,
199+
delivered,
200+
notDelivered,
201+
}
189202
}
190203

191204
describe('claiming', () => {
@@ -207,13 +220,13 @@ describe('claiming', () => {
207220
})
208221

209222
it('reports a result as reaching the model only once Sim takes it as the call own', async () => {
210-
const { sim, journal, runner, executor, delivered } = setup()
223+
const { sim, journal, runner, executor, delivered, notDelivered } = setup()
211224
runner.immediate = DONE
212225
// Sim settled this call first: the result never reached the model.
213226
sim.completionOutcome = 'superseded'
214227
sim.inbox = [callItem('call-superseded', 'chat-a')]
215228
await executor.reconcile()
216-
await vi.waitFor(() => expect(sim.completions).toHaveLength(1))
229+
await vi.waitFor(() => expect(notDelivered).toEqual(['call-superseded']))
217230
expect(delivered).toEqual([])
218231

219232
// Taken by Sim, even with nothing written locally (no OS encryption, say).
@@ -222,10 +235,57 @@ describe('claiming', () => {
222235
sim.inbox = [callItem('call-recorded', 'chat-a')]
223236
await executor.reconcile()
224237
await vi.waitFor(() => expect(delivered).toEqual(['call-recorded']))
238+
expect(notDelivered).toEqual(['call-superseded'])
239+
})
240+
241+
it('reports a result Sim settled first, while a failed send waited to retry, as never reaching the model', async () => {
242+
const { sim, runner, executor, delivered, notDelivered } = setup()
243+
// The terminal handed back a command still going, and the first send of that result fails.
244+
runner.immediate = { status: 'success', message: 'running', data: { status: 'running' } }
245+
sim.completeErrors = [new DeviceRequestError(503, 'unavailable')]
246+
sim.inbox = [callItem('call-1', 'chat-a', 'terminal')]
247+
await executor.reconcile()
248+
await vi.waitFor(() => expect(runner.started).toEqual(['call-1']))
249+
250+
// The user pressed Stop meanwhile, so Sim settled the call before the retry reached it.
251+
sim.completionOutcome = 'superseded'
252+
sim.inbox = [{ kind: 'cancel', toolCallId: 'call-1' }]
253+
await executor.reconcile()
254+
255+
await vi.waitFor(() => expect(notDelivered).toEqual(['call-1']))
256+
expect(delivered).toEqual([])
257+
})
258+
259+
it('reports a stand-in Sim holds as a duplicate as never reaching the model, within one app run', async () => {
260+
const { sim, runner, executor, delivered, notDelivered } = setup()
261+
// Too large to send, so Sim never took the real result; the stand-in's first answer was lost.
262+
runner.immediate = DONE
263+
sim.completeErrors = [
264+
new DeviceRequestError(413, 'too large'),
265+
new DeviceRequestError(0, 'connection reset'),
266+
]
267+
sim.completionOutcome = 'duplicate'
268+
sim.inbox = [callItem('call-1', 'chat-a', 'terminal')]
269+
270+
await executor.reconcile()
271+
272+
await vi.waitFor(() => expect(notDelivered).toEqual(['call-1']))
273+
expect(delivered).toEqual([])
274+
})
275+
276+
it('reports a result Sim refused as never reaching the model', async () => {
277+
const { sim, runner, executor, notDelivered } = setup()
278+
runner.immediate = DONE
279+
sim.completeErrors = [new DeviceRequestError(403, 'forbidden')]
280+
sim.inbox = [callItem('call-refused', 'chat-a')]
281+
282+
await executor.reconcile()
283+
284+
await vi.waitFor(() => expect(notDelivered).toEqual(['call-refused']))
225285
})
226286

227287
it('reports no result that stands in for what the action produced as reaching the model', async () => {
228-
const { sim, runner, executor, delivered } = setup()
288+
const { sim, runner, executor, delivered, notDelivered } = setup()
229289
const standIns: DesktopToolCompletion[] = [
230290
{ status: 'cancelled', message: 'Stopped.' },
231291
{ status: 'error', message: 'Too large.', data: { resultOmitted: true } },
@@ -243,6 +303,7 @@ describe('claiming', () => {
243303
sim.inbox = [callItem('real', 'chat-a')]
244304
await executor.reconcile()
245305
await vi.waitFor(() => expect(delivered).toEqual(['real']))
306+
expect(notDelivered).toEqual(['stand-in-0', 'stand-in-1', 'stand-in-2', 'stand-in-3'])
246307
})
247308

248309
it('claims a whole backlog at once, before any of it runs', async () => {
@@ -492,6 +553,55 @@ describe('restarting', () => {
492553
expect(runner.started).toEqual([])
493554
await vi.waitFor(() => expect(journal.entries.size).toBe(0))
494555
})
556+
557+
it('leaves open whether the model has a result when Sim already held one for the call', async () => {
558+
const { sim, journal, executor, delivered, notDelivered } = setup()
559+
// The previous run's send of the real result may have landed just before the app went down.
560+
sim.completionOutcome = 'duplicate'
561+
await journal.put({ toolCallId: 'started-1', state: 'started', executionToken: 't-started' })
562+
563+
await executor.recover()
564+
565+
await vi.waitFor(() => expect(journal.entries.size).toBe(0))
566+
expect(sim.completions).toHaveLength(1)
567+
expect(delivered).toEqual([])
568+
expect(notDelivered).toEqual([])
569+
})
570+
571+
it('still leaves a recovered result open when Sim held one, after parking it', async () => {
572+
const { sim, journal, executor, onUnregistered, delivered, notDelivered } = setup()
573+
sim.completeErrors = [new DeviceRequestError(401, 'unregistered')]
574+
sim.completionOutcome = 'duplicate'
575+
await journal.put({ toolCallId: 'started-1', state: 'started', executionToken: 't-started' })
576+
await executor.recover()
577+
await vi.waitFor(() => expect(onUnregistered).toHaveBeenCalledTimes(1))
578+
579+
// Registered again.
580+
executor.resumeParked()
581+
582+
await vi.waitFor(() => expect(sim.completions).toHaveLength(1))
583+
await vi.waitFor(() => expect(journal.entries.size).toBe(0))
584+
expect(delivered).toEqual([])
585+
expect(notDelivered).toEqual([])
586+
})
587+
588+
it('reports a recovered result too large to send, then held as a duplicate, as never reaching the model', async () => {
589+
const { sim, journal, executor, delivered, notDelivered } = setup()
590+
// No run could have sent this result: it is too large every time.
591+
sim.completeErrors = [new DeviceRequestError(413, 'too large')]
592+
sim.completionOutcome = 'duplicate'
593+
await journal.put({
594+
toolCallId: 'result-1',
595+
state: 'result',
596+
executionToken: 't',
597+
completion: DONE,
598+
})
599+
600+
await executor.recover()
601+
602+
await vi.waitFor(() => expect(notDelivered).toEqual(['result-1']))
603+
expect(delivered).toEqual([])
604+
})
495605
})
496606

497607
describe('delivery', () => {

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

Lines changed: 24 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,12 @@ export interface DesktopExecutorOptions {
8787
* pane still running) is in use.
8888
*/
8989
onResultDelivered?: (toolCallId: string) => void
90+
/**
91+
* Called once Sim is done with a call without taking its real result: it settled the call
92+
* first, refused the result, or recorded one that stands in for it. The model never learns of
93+
* anything the action handed back as still going, so nothing will come back to it.
94+
*/
95+
onResultNotDelivered?: (toolCallId: string) => void
9096
maxHeldCalls?: number
9197
/** First delivery retry delay; tests shorten it. */
9298
retryBaseMs?: number
@@ -123,6 +129,11 @@ export class DesktopExecutor {
123129
private recovering: Promise<void> | null = null
124130
/** Results a previous app run left, still on their way to Sim. */
125131
private readonly recoveringIds = new Set<string>()
132+
/**
133+
* Results a previous app run left, until Sim answers for them, parked ones included: that run's
134+
* send of the real result may already have landed.
135+
*/
136+
private readonly recoveredResults = new Set<string>()
126137
/** Results that have failed to reach Sim for longer than {@link DELIVERY_AWAKE_LIMIT_MS}. */
127138
private readonly stalledDeliveries = new Set<string>()
128139
/** Results Sim refused because it no longer recognized this device; sent once it registers again. */
@@ -184,6 +195,7 @@ export class DesktopExecutor {
184195
})
185196
// Held awake like a running call: the result exists only on this machine until Sim has it.
186197
this.recoveringIds.add(entry.toolCallId)
198+
this.recoveredResults.add(entry.toolCallId)
187199
this.updateBusy()
188200
void this.deliver(entry.toolCallId, entry.executionToken, completion).finally(() => {
189201
this.recoveringIds.delete(entry.toolCallId)
@@ -413,6 +425,7 @@ export class DesktopExecutor {
413425
sendingSince: number
414426
): Promise<void> {
415427
let pending = completion
428+
let notDelivered = true
416429
for (let attempt = 1; !this.disposed; attempt++) {
417430
try {
418431
const outcome = await this.options.client.complete({
@@ -422,9 +435,13 @@ export class DesktopExecutor {
422435
})
423436
logger.info('Desktop call result acknowledged', { toolCallId, outcome })
424437
// Superseded: Sim settled the call first, so this result never reached the model.
425-
if (outcome !== 'superseded' && isDeliveredResult(pending)) {
426-
this.options.onResultDelivered?.(toolCallId)
427-
}
438+
const delivered = outcome !== 'superseded' && isDeliveredResult(pending)
439+
if (delivered) this.options.onResultDelivered?.(toolCallId)
440+
// A duplicate holds what this app run sent, unless the result was recovered unchanged from
441+
// an earlier run, whose send of the real result may have landed before the restart.
442+
const takenEarlier =
443+
outcome === 'duplicate' && pending === completion && this.recoveredResults.has(toolCallId)
444+
notDelivered = !delivered && !takenEarlier
428445
break
429446
} catch (error) {
430447
// Encoding failed on this machine, so nothing was sent; the same data would fail again.
@@ -481,7 +498,10 @@ export class DesktopExecutor {
481498
)
482499
}
483500
}
484-
if (!this.disposed) await this.forget(toolCallId)
501+
this.recoveredResults.delete(toolCallId)
502+
if (this.disposed) return
503+
if (notDelivered) this.options.onResultNotDelivered?.(toolCallId)
504+
await this.forget(toolCallId)
485505
}
486506

487507
private async renew(entry: HeldCall): Promise<void> {

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,8 @@ export interface DesktopExecutorServiceDeps {
6464
onBusyChange?: (busy: boolean) => void
6565
/** Sim has taken a call's real result as the call's own, so the model has it. */
6666
onResultDelivered?: (toolCallId: string) => void
67+
/** Sim is done with a call whose real result the model never got. */
68+
onResultNotDelivered?: (toolCallId: string) => void
6769
}
6870

6971
export interface DesktopExecutorService {
@@ -275,6 +277,7 @@ export function createDesktopExecutorService(
275277
...(deps.onApprovals ? { onApprovals: deps.onApprovals } : {}),
276278
...(deps.onBusyChange ? { onBusyChange: deps.onBusyChange } : {}),
277279
...(deps.onResultDelivered ? { onResultDelivered: deps.onResultDelivered } : {}),
280+
...(deps.onResultNotDelivered ? { onResultNotDelivered: deps.onResultNotDelivered } : {}),
278281
})
279282
// Recovery rewrites the journal, so what it held before is read first.
280283
await pendingResultsSnapshot()

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -610,6 +610,9 @@ function main(): void {
610610
// A result the model has makes a tmux run it handed back as still going collectable across a
611611
// restart.
612612
onResultDelivered: (toolCallId) => terminal.markRunDelivered(toolCallId),
613+
// A result the model never got leaves such a run with no one to come back to it, so it is
614+
// stopped, as the launch sweep stops a previous process's.
615+
onResultNotDelivered: (toolCallId) => void terminal.stopUndeliveredRun(toolCallId),
613616
runner: createDesktopToolRunner({
614617
preferences: () => desktopSettings.getPreferences(),
615618
accountDataAvailable,

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -938,6 +938,13 @@ export class TerminalService {
938938
await Promise.allSettled(stops)
939939
}
940940

941+
/** Stops the command a plain shell's `run` handed back as still going, if it still is. */
942+
async stopAgentCommand(toolCallId: string): Promise<void> {
943+
await Promise.allSettled(
944+
[...this.sessions.values()].map((session) => this.stopCommand(session, toolCallId))
945+
)
946+
}
947+
941948
/** Waits for the command a run started to end, up to `ms`. */
942949
private async commandEnds(
943950
session: TerminalSession,

‎apps/desktop/src/main/terminal/registry-run-ledger.test.ts‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,32 @@ describe('stopping recorded tmux runs', () => {
119119
expect(ledger.list()).toEqual([])
120120
})
121121

122+
it("stops only the run of a call whose result never reached the model, this process's or a previous one's", async () => {
123+
const dir = ledgerDir()
124+
createRunLedger(dir).record(run('previous', '%1'))
125+
const ledger = createRunLedger(dir)
126+
ledger.record(run('current', '%2'))
127+
ledger.record(run('other', '%3'))
128+
// Recorded as handed back: the model has its pane.
129+
ledger.record(run('handed-back', '%4', 'delivered'))
130+
const registry = new TerminalRegistry(undefined, undefined, ledger)
131+
132+
await registry.stopUndeliveredRun('call-current')
133+
await registry.stopUndeliveredRun('call-previous')
134+
await registry.stopUndeliveredRun('call-handed-back')
135+
136+
expect(panes.get('current')).toBe('stopped')
137+
expect(panes.get('previous')).toBe('stopped')
138+
expect(panes.get('other')).toBeUndefined()
139+
expect(panes.get('handed-back')).toBeUndefined()
140+
expect(
141+
ledger
142+
.list()
143+
.map((record) => record.runId)
144+
.sort()
145+
).toEqual(['handed-back', 'other'])
146+
})
147+
122148
it('notes the run of a call whose result reached the model as handed back', () => {
123149
const dir = ledgerDir()
124150
const ledger = createRunLedger(dir)

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

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -410,23 +410,39 @@ export class TerminalRegistry {
410410
if (runId) this.runLedger?.advance(runId, 'delivered')
411411
}
412412

413+
/**
414+
* Stops what a call's `run` handed back as still going once Sim is done with the call without
415+
* its result reaching the model: a plain shell's command, or a tmux run, this process's or a
416+
* previous one's. A run recorded as handed back is the model's, and is left going.
417+
*/
418+
async stopUndeliveredRun(callId: string): Promise<void> {
419+
await Promise.allSettled(
420+
[...this.entries.values()].map((entry) => entry.service.stopAgentCommand(callId))
421+
)
422+
await this.stopRecordedRuns({ callId, keep: (run) => run.state === 'delivered' })
423+
}
424+
413425
/**
414426
* Stops the recorded tmux runs, each only while its pane still carries its tag, and drops the
415427
* records with nothing left to stop. `excludeLive` skips the runs this process has started;
416-
* `keep` names runs to leave going, such as a previous process's runs whose results the model
417-
* already has and may come back to.
428+
* `callId` limits it to the run that call started; `keep` names runs to leave going, such as a
429+
* previous process's runs whose results the model already has and may come back to.
418430
*/
419431
async stopRecordedRuns(
420432
options: {
421433
excludeLive?: boolean
434+
callId?: string
422435
/** Runs to leave going; their records are only dropped once their panes are gone. */
423436
keep?: (run: RunRecord) => boolean
424437
} = {}
425438
): Promise<void> {
426439
const ledger = this.runLedger
427440
if (!ledger) return
441+
const runs = ledger
442+
.list({ excludeLive: options.excludeLive })
443+
.filter((run) => options.callId === undefined || run.callId === options.callId)
428444
await Promise.allSettled(
429-
ledger.list({ excludeLive: options.excludeLive }).map(async (run) => {
445+
runs.map(async (run) => {
430446
let keeping = options.keep?.(run) ?? false
431447
let state = keeping ? await recordedRunState(run, process.env) : 'ours'
432448
// A kept run whose command already ended leaves only its dead pane (`remain-on-exit`).

0 commit comments

Comments
 (0)