Skip to content

Commit c3a81b0

Browse files
committed
fix(desktop): make the executor's journal, sign-out and terminal stops hold under races
- A journal transition that cannot be made durable now fails: an unrecorded claim is not made, and an action whose start could not be recorded is not run and is reported as not run. The journal keeps result data within what a restart can read back. - Sign-out fences registrations in flight, and a rotated install id never reuses the previous id's client. - A stopped call's result carries no content, only the acknowledgement. - Terminal stops: a Stop before the command starts runs nothing; escalation signals only the stopped command's own process group, rechecked before each signal; a stopped tmux run ends its wait at once. - A chat's next terminal operation waits for one that outlived its deadline. - The doorbell bounds its handshake and reconnects at once on wake. - Approval notifications name neither the chat nor the command. - A user-local glob that stopped early says so.
1 parent 5e79e72 commit c3a81b0

19 files changed

Lines changed: 676 additions & 148 deletions

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,7 @@ function harness(options: { enabled?: boolean; focusedChatId?: string | null } =
6161
}
6262

6363
describe('approval notifications', () => {
64-
it('notifies once per waiting call and opens its chat, never naming the command', () => {
64+
it('notifies once per waiting call and opens its chat, naming neither chat nor command', () => {
6565
const { notifier, notifications, openRoute } = harness()
6666

6767
notifier.update([approval('call-1')])
@@ -70,6 +70,7 @@ describe('approval notifications', () => {
7070
expect(notifications).toHaveLength(1)
7171
expect(notifications[0]?.show).toHaveBeenCalledOnce()
7272
expect(notifications[0]?.options.body).not.toContain('rm -rf')
73+
expect(notifications[0]?.options.body).not.toContain('Fix CI')
7374
notifications[0]?.click()
7475
expect(openRoute).toHaveBeenCalledWith('/workspace/ws-1/chat/chat-b')
7576
})

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

Lines changed: 3 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
/**
22
* Tells the user when a chat they are not looking at waits for their approval. The notification
33
* only opens the chat at its approval card; approving happens there, where the user sees what
4-
* they are approving. It closes itself once the call is decided, and never names the command,
5-
* since a notification can show on a locked screen.
4+
* they are approving. It closes itself once the call is decided, and names neither the chat nor
5+
* the command, since a notification can show on a locked screen.
66
*/
77
import type { DesktopApprovalItem } from '@/main/desktop-executor/executor'
88

@@ -25,12 +25,6 @@ export interface ApprovalNotifierDeps {
2525
}) => ApprovalNotification | null
2626
}
2727

28-
function approvalBody(item: DesktopApprovalItem): string {
29-
return item.chatTitle
30-
? `“${item.chatTitle}” is waiting for your approval.`
31-
: 'A chat is waiting for your approval.'
32-
}
33-
3428
export function createApprovalNotifier(deps: ApprovalNotifierDeps) {
3529
/** Calls already brought to the user's attention, by notification or by being on screen. */
3630
const shown = new Map<string, ApprovalNotification | null>()
@@ -54,7 +48,7 @@ export function createApprovalNotifier(deps: ApprovalNotifierDeps) {
5448
}
5549
const notification = deps.createNotification({
5650
title: 'Approval needed',
57-
body: approvalBody(item),
51+
body: 'A chat is waiting for your approval.',
5852
silent: !preferences.notificationSounds,
5953
})
6054
if (!notification) return

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

Lines changed: 29 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -6,25 +6,25 @@ const encoder = new TextEncoder()
66

77
/** A stream the test writes SSE text into, closing when told to or when its reader aborts. */
88
function controllableStream(signal: AbortSignal) {
9-
let controller!: ReadableStreamDefaultController<Uint8Array>
9+
let controller: ReadableStreamDefaultController<Uint8Array> | undefined
1010
const stream = new ReadableStream<Uint8Array>({
1111
start(c) {
1212
controller = c
1313
},
1414
})
1515
signal.addEventListener('abort', () => {
1616
try {
17-
controller.error(new Error('aborted'))
17+
controller?.error(new Error('aborted'))
1818
} catch {}
1919
})
2020
return {
2121
stream,
22-
write: (text: string) => controller.enqueue(encoder.encode(text)),
23-
end: () => controller.close(),
22+
write: (text: string) => controller?.enqueue(encoder.encode(text)),
23+
end: () => controller?.close(),
2424
}
2525
}
2626

27-
function harness(options: { staleAfterMs?: number } = {}) {
27+
function harness(options: { staleAfterMs?: number; retryBaseMs?: number } = {}) {
2828
const connections: ReturnType<typeof controllableStream>[] = []
2929
const failures: DeviceRequestError[] = []
3030
const onRing = vi.fn()
@@ -106,4 +106,28 @@ describe('InboxDoorbell', () => {
106106
await vi.waitFor(() => expect(onUnregistered).toHaveBeenCalled())
107107
doorbell.stop()
108108
})
109+
110+
it('reconnects at once on wake, without waiting out a backoff', async () => {
111+
const { doorbell, connections, onRing } = harness({ retryBaseMs: 60_000 })
112+
doorbell.start()
113+
await vi.waitFor(() => expect(connections).toHaveLength(1))
114+
115+
doorbell.wake()
116+
117+
await vi.waitFor(() => expect(connections).toHaveLength(2))
118+
await vi.waitFor(() => expect(onRing).toHaveBeenCalledTimes(2))
119+
doorbell.stop()
120+
})
121+
122+
it('starts on wake after sleep stopped it', async () => {
123+
const { doorbell, connections } = harness({ retryBaseMs: 60_000 })
124+
doorbell.start()
125+
await vi.waitFor(() => expect(connections).toHaveLength(1))
126+
doorbell.stop()
127+
128+
doorbell.wake()
129+
130+
await vi.waitFor(() => expect(connections).toHaveLength(2))
131+
doorbell.stop()
132+
})
109133
})

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

Lines changed: 32 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ const logger = createLogger('DesktopExecutorDoorbell')
1616
/** Sim heartbeats every 30 s; two missed ones mean the connection is gone. */
1717
const STALE_STREAM_MS = 75_000
1818
const RECONNECT_MAX_MS = 30_000
19+
const HANDSHAKE_TIMEOUT_MS = 15_000
1920

2021
interface DoorbellOptions {
2122
client: Pick<DesktopExecutorClient, 'openInboxStream'>
@@ -54,6 +55,10 @@ export function parseServerSentEvents(buffer: string): {
5455

5556
export class InboxDoorbell {
5657
private running = false
58+
/** Which loop is current; a loop left over from before a stop ends at its next turn. */
59+
private loopGeneration = 0
60+
/** The connection was dropped on purpose; open the next one without backing off. */
61+
private reconnectNow = false
5762
private connection: AbortController | null = null
5863
private sleeper: AbortController | null = null
5964

@@ -62,30 +67,46 @@ export class InboxDoorbell {
6267
start(): void {
6368
if (this.running) return
6469
this.running = true
65-
void this.loop()
70+
this.loopGeneration += 1
71+
void this.loop(this.loopGeneration)
6672
}
6773

6874
stop(): void {
6975
this.running = false
76+
this.loopGeneration += 1
7077
this.connection?.abort()
7178
this.sleeper?.abort()
7279
}
7380

74-
/** Drops the current connection and reconnects now, as after waking or coming back online. */
75-
reconnect(): void {
81+
/**
82+
* After waking or coming back online: a stopped doorbell starts, and a running one drops its
83+
* connection, which may be dead without knowing it, and reconnects at once.
84+
*/
85+
wake(): void {
86+
if (!this.running) {
87+
this.start()
88+
return
89+
}
90+
this.reconnectNow = true
7691
this.connection?.abort()
7792
this.sleeper?.abort()
7893
}
7994

80-
private async loop(): Promise<void> {
95+
private async loop(generation: number): Promise<void> {
8196
let attempt = 0
82-
while (this.running) {
97+
const current = () => this.loopGeneration === generation
98+
while (current()) {
8399
const rotated = await this.connectOnce().then(
84100
(result) => {
85101
attempt = 0
86102
return result === 'rotated'
87103
},
88104
(error: unknown) => {
105+
if (this.reconnectNow) {
106+
this.reconnectNow = false
107+
attempt = 0
108+
return true
109+
}
89110
attempt += 1
90111
if (error instanceof DeviceRequestError && error.unregistered) {
91112
this.options.onUnregistered()
@@ -97,7 +118,7 @@ export class InboxDoorbell {
97118
return false
98119
}
99120
)
100-
if (!this.running || rotated) continue
121+
if (!current() || rotated) continue
101122
this.sleeper = new AbortController()
102123
await interruptibleSleep(
103124
attempt === 0
@@ -123,7 +144,11 @@ export class InboxDoorbell {
123144
staleTimer = setTimeout(() => connection.abort(), staleAfterMs)
124145
}
125146
try {
126-
const stream = await this.options.client.openInboxStream(connection.signal)
147+
// A stream whose response never starts is as dead as one that went quiet.
148+
const handshake = setTimeout(() => connection.abort(), HANDSHAKE_TIMEOUT_MS)
149+
const stream = await this.options.client
150+
.openInboxStream(connection.signal)
151+
.finally(() => clearTimeout(handshake))
127152
touch()
128153
this.options.onRing()
129154
const reader = stream.getReader()

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

Lines changed: 74 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -95,10 +95,13 @@ class FakeSim {
9595
class MemoryJournal implements ExecutorJournal {
9696
readonly entries = new Map<string, JournalEntry>()
9797
readonly history: JournalEntry[] = []
98+
/** A transition the disk refuses, as a full disk or a failed encryption would. */
99+
failOn: JournalEntry['state'] | null = null
98100
async load() {
99101
return [...this.entries.values()]
100102
}
101103
async put(entry: JournalEntry) {
104+
if (entry.state === this.failOn) throw new Error('disk full')
102105
this.entries.set(entry.toolCallId, entry)
103106
this.history.push(entry)
104107
}
@@ -117,16 +120,19 @@ class FakeRunner implements DesktopToolRunner {
117120
/** When set, every call finishes at once with this completion. */
118121
immediate: DesktopToolCompletion | null = null
119122
onStart: ((call: ClaimedDesktopCall) => void) | null = null
123+
/** An action that finishes with its full result even after being told to stop. */
124+
ignoresAbort = false
120125

121126
async run(call: ClaimedDesktopCall, signal: AbortSignal) {
122127
this.started.push(call.toolCallId)
123128
this.onStart?.(call)
124129
if (this.immediate) return this.immediate
125130
const result = deferred<DesktopToolCompletion>()
126131
this.pending.set(call.toolCallId, result)
127-
signal.addEventListener('abort', () =>
128-
result.resolve({ status: 'error', message: 'This browser action was cancelled.' })
129-
)
132+
signal.addEventListener('abort', () => {
133+
if (!this.ignoresAbort)
134+
result.resolve({ status: 'error', message: 'This browser action was cancelled.' })
135+
})
130136
return result.promise
131137
}
132138

@@ -313,7 +319,7 @@ describe('reporting', () => {
313319
sim.inbox = [callItem('call-1', 'chat-a')]
314320

315321
await executor.reconcile()
316-
await vi.waitFor(() => expect(journal.entries.get('call-1')?.state).toBe('result'))
322+
await vi.waitFor(() => expect(journal.history.map((entry) => entry.state)).toContain('result'))
317323

318324
await vi.waitFor(() => expect(sim.completions).toHaveLength(1))
319325
expect(runner.started).toEqual(['call-1'])
@@ -367,6 +373,23 @@ describe('stopping', () => {
367373
expect(runner.started).toEqual(['a-1'])
368374
})
369375

376+
it('sends nothing the stopped action produced, only the acknowledgement', async () => {
377+
const { sim, runner, executor } = setup()
378+
runner.ignoresAbort = true
379+
sim.inbox = [callItem('call-1', 'chat-a', 'read_local_file')]
380+
await executor.reconcile()
381+
await vi.waitFor(() => expect(runner.started).toEqual(['call-1']))
382+
383+
sim.inbox = [{ kind: 'cancel', toolCallId: 'call-1' }]
384+
await executor.reconcile()
385+
await vi.waitFor(() => expect(runner.cancelled).toEqual(['call-1']))
386+
runner.finish('call-1', { status: 'success', message: 'read', data: { text: 'secret' } })
387+
388+
await vi.waitFor(() => expect(sim.completions).toHaveLength(1))
389+
expect(sim.completions[0]?.completion.status).toBe('cancelled')
390+
expect(JSON.stringify(sim.completions[0])).not.toContain('secret')
391+
})
392+
370393
it('ignores a cancel for a call it never held', async () => {
371394
const { sim, runner, executor } = setup()
372395
sim.inbox = [{ kind: 'cancel', toolCallId: 'someone-else' }]
@@ -431,4 +454,51 @@ describe('registration', () => {
431454
expect(journal.entries.size).toBe(0)
432455
expect(executor.heldCallCount()).toBe(0)
433456
})
457+
458+
it('does not keep a claim that Sim answers after sign-out', async () => {
459+
const { sim, journal, runner, executor } = setup({ leaseRenewMs: 10 })
460+
const answer = deferred<void>()
461+
const claim = sim.client.claim
462+
sim.client.claim = async (toolCallId) => {
463+
await answer.promise
464+
return claim(toolCallId)
465+
}
466+
sim.inbox = [callItem('call-1', 'chat-a')]
467+
const reading = executor.reconcile()
468+
await vi.waitFor(() => expect(journal.entries.get('call-1')?.state).toBe('claiming'))
469+
470+
const signingOut = executor.dispose()
471+
answer.resolve()
472+
await Promise.all([reading, signingOut])
473+
await sleep(40)
474+
475+
expect(executor.heldCallCount()).toBe(0)
476+
expect(runner.started).toEqual([])
477+
expect(sim.renewals).toEqual([])
478+
expect(journal.entries.size).toBe(0)
479+
})
480+
})
481+
482+
describe('recording', () => {
483+
it('leaves a call on offer when it cannot record the claim', async () => {
484+
const { sim, journal, executor } = setup()
485+
journal.failOn = 'claiming'
486+
sim.inbox = [callItem('call-1', 'chat-a')]
487+
488+
await executor.reconcile()
489+
490+
expect(sim.claims).toEqual([])
491+
})
492+
493+
it('never starts an action it could not record, and reports it as not run', async () => {
494+
const { sim, journal, runner, executor } = setup()
495+
journal.failOn = 'started'
496+
sim.inbox = [callItem('call-1', 'chat-a')]
497+
498+
await executor.reconcile()
499+
500+
await vi.waitFor(() => expect(sim.completions).toHaveLength(1))
501+
expect(sim.completions[0]?.completion.data).toMatchObject({ notStarted: true })
502+
expect(runner.started).toEqual([])
503+
})
434504
})

0 commit comments

Comments
 (0)