Skip to content

Commit ea73898

Browse files
committed
fix(desktop): resume supervising a desktop call from its row after a restart
Supervision now resumes from the stored pickup deadline or lease when the call is waited on again (the worker's handoff re-dispatches it), instead of giving up because the call was already offered. The existing-claim waiter uses the same bound-run budget as the resume gate, the desktop composer's behaviour without a background executor is pinned by a test, and the supervisor tests wait for presence before supervising.
1 parent b41f394 commit ea73898

7 files changed

Lines changed: 185 additions & 60 deletions

File tree

‎apps/sim/lib/desktop/executor/supervisor.integration.ts‎

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ import {
3939
renewDesktopToolLease,
4040
resolveTurnDesktopDevice,
4141
} from '@/lib/desktop/application/executor'
42+
import { isDesktopPresent } from '@/lib/desktop/executor/presence'
4243
import { failLapsedDesktopCall } from '@/lib/desktop/executor/repository'
4344
import {
4445
DESKTOP_LEASE_LOST_MESSAGE,
@@ -185,6 +186,7 @@ describe.runIf(Boolean(redisUrl))('desktop call supervision and turn binding', (
185186
const heard: string[] = []
186187
const close = stream.subscribe((_event, data) => heard.push(String(data.reason)))
187188
try {
189+
await expect.poll(() => isDesktopPresent(desktop.deviceId)).toBe(true)
188190
return await run(heard)
189191
} finally {
190192
close()
@@ -445,6 +447,80 @@ describe.runIf(Boolean(redisUrl))('desktop call supervision and turn binding', (
445447
expect((await row(toolCallId)).status).toBe('running')
446448
})
447449

450+
it('resumes supervising an offered call after the process that offered it restarts', async () => {
451+
const desktop = await signedInDesktop()
452+
const run = await boundRun(desktop.userId, desktop.deviceId)
453+
const toolCallId = await pendingCall(run.runId)
454+
const supervise = (signal?: AbortSignal) =>
455+
superviseDesktopCall({
456+
toolCallId,
457+
runId: run.runId,
458+
userId: desktop.userId,
459+
deviceId: desktop.deviceId,
460+
pickupGraceMs: 800,
461+
signal,
462+
})
463+
464+
await online(desktop, async () => {
465+
/** The first supervisor offers the call, then its process dies. */
466+
const died = new AbortController()
467+
const first = supervise(died.signal)
468+
await expect
469+
.poll(async () => (await row(toolCallId)).executionLeaseExpiresAt)
470+
.not.toBeNull()
471+
died.abort()
472+
await expect(first).resolves.toBe('aborted')
473+
474+
/** The worker's handoff waits on the call again; supervision resumes from the row. */
475+
const started = Date.now()
476+
await expect(supervise()).resolves.toBe('not_responding')
477+
expect(Date.now() - started).toBeLessThan(1_500)
478+
})
479+
await expect(unsealed(toolCallId, run.runId, desktop.userId)).resolves.toMatchObject({
480+
data: { notStarted: true, reason: 'not_responding' },
481+
})
482+
})
483+
484+
it('resumes supervising a claimed call after a restart and settles its lapsed lease', async () => {
485+
const desktop = await signedInDesktop()
486+
const run = await boundRun(desktop.userId, desktop.deviceId)
487+
const toolCallId = await pendingCall(run.runId)
488+
489+
await online(desktop, async () => {
490+
const died = new AbortController()
491+
const first = superviseDesktopCall({
492+
toolCallId,
493+
runId: run.runId,
494+
userId: desktop.userId,
495+
deviceId: desktop.deviceId,
496+
pickupGraceMs: 2_000,
497+
signal: died.signal,
498+
})
499+
await expect
500+
.poll(async () => (await row(toolCallId)).executionLeaseExpiresAt)
501+
.not.toBeNull()
502+
await claimDesktopTool.execute({
503+
principal: desktop.principal,
504+
input: { deviceId: desktop.deviceId, toolCallId },
505+
})
506+
died.abort()
507+
await expect(first).resolves.toBe('aborted')
508+
await db
509+
.update(copilotAsyncToolCalls)
510+
.set({ executionLeaseExpiresAt: sql`clock_timestamp() + interval '300 milliseconds'` })
511+
.where(eq(copilotAsyncToolCalls.toolCallId, toolCallId))
512+
513+
await expect(
514+
superviseDesktopCall({
515+
toolCallId,
516+
runId: run.runId,
517+
userId: desktop.userId,
518+
deviceId: desktop.deviceId,
519+
})
520+
).resolves.toBe('lease_lost')
521+
})
522+
})
523+
448524
it('never offers a call on a run the chat view serves', async () => {
449525
const desktop = await signedInDesktop()
450526
const run = await boundRun(desktop.userId, null)

‎apps/sim/lib/desktop/executor/supervisor.ts‎

Lines changed: 68 additions & 52 deletions
Original file line numberDiff line numberDiff line change
@@ -81,13 +81,67 @@ function wake(input: SuperviseDesktopCallInput, message: string, data: unknown)
8181
})
8282
}
8383

84+
/** Fails a call nobody claimed as not started; returns whether this settlement won. */
85+
async function failUnclaimed(
86+
input: SuperviseDesktopCallInput,
87+
reason: 'offline' | 'not_responding',
88+
message: string,
89+
deadlinePassed: boolean
90+
): Promise<boolean> {
91+
const state = await getDesktopCallState(input.toolCallId)
92+
const result = await sealFailure(input, state?.result, { reason, message, notStarted: true })
93+
const failed = await failUnclaimedDesktopCall({
94+
toolCallId: input.toolCallId,
95+
runId: input.runId,
96+
result,
97+
error: message,
98+
deadlinePassed,
99+
})
100+
if (failed) wake(input, message, result)
101+
return failed
102+
}
103+
104+
/** Fails a claimed call whose lease lapsed as outcome unknown; returns whether this settlement won. */
105+
async function failLapsed(
106+
input: SuperviseDesktopCallInput,
107+
ownerToken: string,
108+
existing: unknown
109+
): Promise<boolean> {
110+
const result = await sealFailure(input, existing, {
111+
reason: 'lease_lost',
112+
message: DESKTOP_LEASE_LOST_MESSAGE,
113+
outcomeUnknown: true,
114+
doNotRetry: true,
115+
})
116+
const failed = await failLapsedDesktopCall({
117+
toolCallId: input.toolCallId,
118+
runId: input.runId,
119+
ownerToken,
120+
result,
121+
error: DESKTOP_LEASE_LOST_MESSAGE,
122+
})
123+
if (failed) wake(input, DESKTOP_LEASE_LOST_MESSAGE, result)
124+
return failed
125+
}
126+
127+
/** A call some supervisor offered (pickup window open or closed) or the device claimed. */
128+
function isOfferedOrClaimed(state: Awaited<ReturnType<typeof getDesktopCallState>>): boolean {
129+
if (!state) return false
130+
if (state.status === ASYNC_TOOL_STATUS.pending)
131+
return !state.ownerToken && state.msUntilLeaseEnd !== null
132+
return state.status === ASYNC_TOOL_STATUS.running && Boolean(state.ownerToken)
133+
}
134+
84135
/**
85136
* Owns a desktop call on a device-bound run from the moment it may run until it settles: offers
86137
* it to the device and rings the doorbell, fails it at once as not started when the device is
87138
* offline or after the pickup window when nobody claimed it, and fails it as outcome unknown when
88139
* a claimed call's lease lapses. A running call stays alive for as long as the device renews it.
89140
* Every settlement is a CAS against the device's own transitions, so exactly one wins, and every
90141
* failure is published so the run's waiter reads it at once.
142+
*
143+
* All of its state is on the row: when the process supervising a call restarts, the next wait on
144+
* the call (the worker's handoff re-dispatches it) resumes supervision from the stored deadline.
91145
*/
92146
export async function superviseDesktopCall(
93147
input: SuperviseDesktopCallInput
@@ -97,65 +151,27 @@ export async function superviseDesktopCall(
97151
runId: input.runId,
98152
pickupGraceMs: input.pickupGraceMs ?? DESKTOP_CALL_PICKUP_GRACE_MS,
99153
})
100-
if (!offered) return 'not_offered'
101-
ringDesktopInbox(input.deviceId, 'call')
102-
103-
const failUnclaimed = async (
104-
reason: 'offline' | 'not_responding',
105-
message: string,
106-
deadlinePassed: boolean
107-
) => {
108-
const state = await getDesktopCallState(input.toolCallId)
109-
const result = await sealFailure(input, state?.result, { reason, message, notStarted: true })
110-
const failed = await failUnclaimedDesktopCall({
111-
toolCallId: input.toolCallId,
112-
runId: input.runId,
113-
result,
114-
error: message,
115-
deadlinePassed,
116-
})
117-
if (failed) wake(input, message, result)
118-
return failed
119-
}
120-
121-
if (!(await isDesktopPresent(input.deviceId).catch(() => false))) {
122-
if (await failUnclaimed('offline', DESKTOP_OFFLINE_MESSAGE, false)) return 'offline'
154+
if (offered) {
155+
ringDesktopInbox(input.deviceId, 'call')
156+
if (!(await isDesktopPresent(input.deviceId).catch(() => false))) {
157+
if (await failUnclaimed(input, 'offline', DESKTOP_OFFLINE_MESSAGE, false)) return 'offline'
158+
}
159+
} else if (!isOfferedOrClaimed(await getDesktopCallState(input.toolCallId))) {
160+
return 'not_offered'
123161
}
124162

125163
for (;;) {
126164
if (input.signal?.aborted) return 'aborted'
127165
const state = await getDesktopCallState(input.toolCallId)
128166
if (!state || isTerminalAsyncStatus(state.status)) return 'settled'
167+
if (!isOfferedOrClaimed(state)) return 'settled'
129168
const remainingMs = state.msUntilLeaseEnd ?? 0
130-
if (state.status === ASYNC_TOOL_STATUS.pending && !state.ownerToken) {
131-
if (remainingMs <= 0) {
132-
if (await failUnclaimed('not_responding', DESKTOP_NOT_RESPONDING_MESSAGE, true))
133-
return 'not_responding'
134-
continue
135-
}
136-
} else if (state.status === ASYNC_TOOL_STATUS.running && state.ownerToken) {
137-
if (remainingMs <= 0) {
138-
const result = await sealFailure(input, state.result, {
139-
reason: 'lease_lost',
140-
message: DESKTOP_LEASE_LOST_MESSAGE,
141-
outcomeUnknown: true,
142-
doNotRetry: true,
143-
})
144-
const failed = await failLapsedDesktopCall({
145-
toolCallId: input.toolCallId,
146-
runId: input.runId,
147-
ownerToken: state.ownerToken,
148-
result,
149-
error: DESKTOP_LEASE_LOST_MESSAGE,
150-
})
151-
if (failed) {
152-
wake(input, DESKTOP_LEASE_LOST_MESSAGE, result)
153-
return 'lease_lost'
154-
}
155-
continue
156-
}
157-
} else {
158-
return 'settled'
169+
if (remainingMs <= 0) {
170+
const won = state.ownerToken
171+
? await failLapsed(input, state.ownerToken, state.result)
172+
: await failUnclaimed(input, 'not_responding', DESKTOP_NOT_RESPONDING_MESSAGE, true)
173+
if (won) return state.ownerToken ? 'lease_lost' : 'not_responding'
174+
continue
159175
}
160176
await interruptibleSleep(Math.min(remainingMs + DEADLINE_SLACK_MS, MAX_WAIT_MS), input.signal)
161177
}

‎apps/sim/lib/desktop/index.test.ts‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,37 @@ describe('desktop surface availability', () => {
5757
).toBeUndefined()
5858
})
5959

60+
it('offers a turn to the background executor only when the shell has a registered one', async () => {
61+
setDesktopPreferencesSnapshot({
62+
...ENABLED_PREFERENCES,
63+
browserEnabled: false,
64+
terminalEnabled: false,
65+
})
66+
const offerFrom = async (bridge: Record<string, unknown>) => {
67+
installBridge({ localFiles: vi.fn(), ...bridge })
68+
const { desktopCapabilities } = await getDesktopChatCapabilities('chat-1')
69+
return { deviceId: desktopCapabilities?.deviceId, executor: desktopCapabilities?.executor }
70+
}
71+
const unset = { deviceId: undefined, executor: undefined }
72+
73+
expect(await offerFrom({})).toEqual(unset)
74+
expect(await offerFrom({ desktopExecutor: { getDevice: vi.fn(async () => null) } })).toEqual(
75+
unset
76+
)
77+
expect(
78+
await offerFrom({
79+
desktopExecutor: { getDevice: vi.fn(async () => Promise.reject(new Error('IPC closed'))) },
80+
})
81+
).toEqual(unset)
82+
expect(
83+
await offerFrom({
84+
desktopExecutor: {
85+
getDevice: vi.fn(async () => ({ deviceId: 'device-1', protocolVersion: 1 })),
86+
},
87+
})
88+
).toEqual({ deviceId: 'device-1', executor: 1 })
89+
})
90+
6091
it('bounds terminal hints before adding them to a chat request', async () => {
6192
const oversizedValue = 'x'.repeat(1100)
6293
installBridge({

‎apps/sim/lib/desktop/index.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,7 @@ export async function getDesktopChatCapabilities(
215215
.then((state) => state.sessions)
216216
.catch(() => [])
217217
: []
218+
// Shells without a background executor (every build before it ships) leave the turn unbound.
218219
const executorDevice = bridge?.desktopExecutor
219220
? await bridge.desktopExecutor.getDevice().catch(() => null)
220221
: null

‎apps/sim/lib/mothership/chat/post.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -377,6 +377,7 @@ const ChatMessageSchema = z
377377
)
378378

379379
type UnifiedChatRequest = z.infer<typeof ChatMessageSchema>
380+
380381
/**
381382
* The desktop a turn asks its background executor to run on. Only a desktop composer that speaks
382383
* the executor protocol and switched on at least one desktop surface asks; Assistant turns have

‎apps/sim/lib/mothership/request/handlers/tool.ts‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -194,6 +194,13 @@ function rebindResolvedIntegrationCall(
194194
return true
195195
}
196196

197+
/** Reads the run's device binding once per streaming context. */
198+
async function resolveRunDesktopDevice(context: StreamingContext): Promise<string | null> {
199+
if (context.desktopDeviceId === undefined)
200+
context.desktopDeviceId = context.runId ? await getRunDesktopDeviceId(context.runId) : null
201+
return context.desktopDeviceId
202+
}
203+
197204
/**
198205
* Upsert the durable `async_tool_calls` row before the authoritative tool-call
199206
* SSE frame is forwarded to the client, so `/api/copilot/confirm` and
@@ -208,13 +215,6 @@ function rebindResolvedIntegrationCall(
208215
* Also stamps `awaiting_approval` onto the outgoing frame so the browser and
209216
* the persisted content block both record that the call is gated.
210217
*/
211-
/** Reads the run's device binding once per streaming context. */
212-
async function resolveRunDesktopDevice(context: StreamingContext): Promise<string | null> {
213-
if (context.desktopDeviceId === undefined)
214-
context.desktopDeviceId = context.runId ? await getRunDesktopDeviceId(context.runId) : null
215-
return context.desktopDeviceId
216-
}
217-
218218
export async function prePersistClientExecutableToolCall(
219219
event: StreamEvent,
220220
context: StreamingContext,

‎apps/sim/lib/mothership/request/tools/executor.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -668,7 +668,7 @@ async function executeToolAndReportInner(
668668
/** The winning controller owns execution; this promise observes its durable result. */
669669
const completion = await waitForToolConfirmation(
670670
toolCall.id,
671-
pendingToolWaitBudgetMs(toolCall),
671+
pendingToolWaitBudgetMs(toolCall, context.desktopDeviceId),
672672
options?.abortSignal ?? execContext.abortSignal,
673673
{
674674
executionScope: { runId: context.runId, userId: execContext.userId },

0 commit comments

Comments
 (0)