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
21 changes: 21 additions & 0 deletions apps/desktop/src/main/desktop-executor/executor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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')]
Expand Down
25 changes: 21 additions & 4 deletions apps/desktop/src/main/desktop-executor/executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,20 @@ export interface DesktopToolRunner {

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

/**
* 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
Expand All @@ -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
Expand Down Expand Up @@ -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.
Expand Down
46 changes: 46 additions & 0 deletions apps/desktop/src/main/desktop-executor/service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Response>((_resolve, reject) =>
init.signal?.addEventListener('abort', () => reject(new Error('aborted')))
)
Expand Down Expand Up @@ -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',
Expand All @@ -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', () => {
Expand Down
49 changes: 28 additions & 21 deletions apps/desktop/src/main/desktop-executor/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -32,6 +31,7 @@ import {
type DesktopApprovalItem,
DesktopExecutor,
type DesktopToolRunner,
isDeliveredResult,
} from '@/main/desktop-executor/executor'
import { createExecutorJournal } from '@/main/desktop-executor/journal'
import {
Expand Down Expand Up @@ -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 {
Expand All @@ -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<Set<string>>
/** Stores one entry of a claimed import, as this device's registered session. */
Expand Down Expand Up @@ -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<Set<string>> | null = null
const pendingResultsSnapshot = (): Promise<Set<string>> => {
// 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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -453,21 +473,8 @@ 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
pendingResults() {
return pendingResultsSnapshot()
},
importEntry(request, signal) {
if (!client) throw new Error('The Sim desktop app is not signed in to Sim.')
Expand Down
22 changes: 7 additions & 15 deletions apps/desktop/src/main/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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()
Expand Down
3 changes: 1 addition & 2 deletions apps/desktop/src/main/terminal/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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' }),
}
: {}),
})
Expand Down
Loading
Loading