Skip to content

Commit c8ef18d

Browse files
authored
fix(desktop): inbox persistence order, desktop-only overdue sweep, and ungated pickup windows (#8685)
* fix(desktop): scope the overdue sweep to desktop calls and keep ungated pickup windows - The overdue listing and the per-call deadline read only consider desktop tool calls, and the settlement checks the call is one the desktop runs, so a bound run's workflow call or Sim-files VFS read is never failed with the desktop's not-started result. - The cron scans only runs inside the inbox's horizon, oldest calls first, within its batch limit. - Recording a decision clears the pickup deadline only for a call that was gated; a call that was never gated keeps the window it was offered with. * fix(desktop): limit the overdue sweep over desktop calls only, and bound it by deadline The sweep's tool-name filter admitted read, grep and glob calls on Sim's own files, which the settlement then skipped, so enough of them could fill every batch and starve real desktop calls. Queries now use the SQL form of isDesktopToolCall (a desktop tool by name, or a VFS read of a granted local folder) before their limit. The sweep's horizon is now on the deadline that lapsed, not on the run's start: Sim does not enforce a run's wall clock, so a long-lived run's recently overdue call is still settled. * fix(desktop): hand a device its calls in the order they were persisted The inbox ordered calls by created_at, then tool_call_id. Calls of one turn can be persisted in the same millisecond, and then the tie broke on random ids: a device could run a click before the type the model emitted first. Calls now carry persist_seq, a strictly increasing number assigned on insert (pre-persist writes them in emission order), and the inbox and the overdue sweep order by it. The migration adds the column without a default and then sets the default, so existing rows are not rewritten; rows persisted before it have no position and sort first, as the oldest. * fix(desktop): reach every overdue desktop call, and declare the desktop tool names as const A 24 h horizon on the lapsed deadline meant a call the backstop missed for a day was never settled. The sweep starts from the few unsettled (pending or running) calls in persistence order, so it needs no horizon to stay small. The desktop tool names are a literal array declared as const, and the lookup set is derived from it. * test(mothership): give the hand-built tool call table the persistence sequence column
1 parent 5fccf71 commit c8ef18d

13 files changed

Lines changed: 30470 additions & 19 deletions

File tree

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

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -865,6 +865,28 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
865865
])
866866
})
867867

868+
it('lists calls persisted in the same millisecond in the order they were persisted', async () => {
869+
const desktop = await signedInDesktop()
870+
const run = await boundRun(desktop)
871+
const sameMillisecond = new Date()
872+
/** Ids that sort in reverse, so neither the clock nor the id can produce this order. */
873+
const persisted = ['type-zz', 'click-mm', 'submit-aa'].map(
874+
(name) => `${name}-${generateId()}`
875+
)
876+
for (const toolCallId of persisted) {
877+
await db.insert(copilotAsyncToolCalls).values({
878+
runId: run.runId,
879+
toolCallId,
880+
toolName: 'browser_click',
881+
args: { ref: 'e1' },
882+
createdAt: sameMillisecond,
883+
pickupDeadlineAt: sql`now() + interval '1 minute'`,
884+
})
885+
}
886+
887+
expect((await inbox(desktop)).items.map((item) => item.toolCallId)).toEqual(persisted)
888+
})
889+
868890
it('lists pending calls in persistence order even when their doorbell was never heard', async () => {
869891
const desktop = await signedInDesktop()
870892
const run = await boundRun(desktop)

‎apps/sim/lib/desktop/executor/bound-turn.integration.ts‎

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -731,6 +731,138 @@ describe.runIf(Boolean(redisUrl))("a turn bound to a desktop's background execut
731731
TURN_WAIT_MS
732732
)
733733

734+
it(
735+
'leaves a bound run’s calls the desktop never runs to their own path',
736+
async () => {
737+
const desktop = await signedInDesktop()
738+
const run = await boundRun(desktop)
739+
const longAgo = new Date(Date.now() - 600_000)
740+
const workflowCall = generateId()
741+
const simFileRead = generateId()
742+
await db.insert(copilotAsyncToolCalls).values([
743+
{
744+
runId: run.runId,
745+
toolCallId: workflowCall,
746+
toolName: 'run_workflow',
747+
args: {},
748+
createdAt: longAgo,
749+
},
750+
{
751+
runId: run.runId,
752+
toolCallId: simFileRead,
753+
toolName: 'read',
754+
args: { path: 'workspace/notes.md' },
755+
createdAt: longAgo,
756+
},
757+
])
758+
759+
await runCleanupStaleExecutions()
760+
761+
expect((await storedCall(workflowCall)).status).toBe('pending')
762+
expect((await storedCall(simFileRead)).status).toBe('pending')
763+
},
764+
TURN_WAIT_MS
765+
)
766+
767+
it(
768+
'reaches an abandoned desktop call however many older Sim-file reads the run left pending',
769+
async () => {
770+
const desktop = await signedInDesktop()
771+
const run = await boundRun(desktop)
772+
await db.insert(copilotAsyncToolCalls).values(
773+
Array.from({ length: 201 }, () => ({
774+
runId: run.runId,
775+
toolCallId: generateId(),
776+
toolName: 'read',
777+
args: { path: 'workspace/notes.md' },
778+
createdAt: new Date(Date.now() - 900_000),
779+
}))
780+
)
781+
const abandoned = generateId()
782+
await db.insert(copilotAsyncToolCalls).values({
783+
runId: run.runId,
784+
toolCallId: abandoned,
785+
toolName: 'browser_click',
786+
args: { ref: 'e1' },
787+
createdAt: new Date(Date.now() - 600_000),
788+
})
789+
790+
await runCleanupStaleExecutions()
791+
792+
expect((await storedCall(abandoned)).status).toBe('failed')
793+
},
794+
TURN_WAIT_MS
795+
)
796+
797+
it(
798+
'still settles a call that stayed overdue for days, however long the backstop missed it',
799+
async () => {
800+
const desktop = await signedInDesktop()
801+
const run = await boundRun(desktop)
802+
const toolCallId = generateId()
803+
await db.insert(copilotAsyncToolCalls).values({
804+
runId: run.runId,
805+
toolCallId,
806+
toolName: 'browser_click',
807+
args: { ref: 'e1' },
808+
createdAt: new Date(Date.now() - 2 * 24 * 3_600_000),
809+
})
810+
811+
await runCleanupStaleExecutions()
812+
813+
expect((await storedCall(toolCallId)).status).toBe('failed')
814+
},
815+
TURN_WAIT_MS
816+
)
817+
818+
it(
819+
'settles a call whose window lapsed recently on a run that started long ago',
820+
async () => {
821+
const desktop = await signedInDesktop()
822+
const run = await boundRun(desktop)
823+
await db
824+
.update(copilotRuns)
825+
.set({ startedAt: new Date(Date.now() - 2 * 24 * 3_600_000) })
826+
.where(eq(copilotRuns.id, run.runId))
827+
const toolCallId = generateId()
828+
await db.insert(copilotAsyncToolCalls).values({
829+
runId: run.runId,
830+
toolCallId,
831+
toolName: 'browser_click',
832+
args: { ref: 'e1' },
833+
createdAt: new Date(Date.now() - 600_000),
834+
})
835+
836+
await runCleanupStaleExecutions()
837+
838+
expect((await storedCall(toolCallId)).status).toBe('failed')
839+
},
840+
TURN_WAIT_MS
841+
)
842+
843+
it(
844+
'keeps the pickup window of a call that was never gated when a decision is posted for it',
845+
async () => {
846+
const desktop = await signedInDesktop()
847+
const run = await boundRun(desktop)
848+
await desktop.pull()
849+
850+
const { toolCallId, answer, context } = await agentCalls(run, 'browser_click', { ref: 'e1' })
851+
const offeredRow = await offered(toolCallId)
852+
const decided = await post(toolPermissionPOST, '/api/copilot/tool-permission', {
853+
decisions: [{ toolCallId, decision: 'allow' }],
854+
})
855+
expect(decided.status).toBe(200)
856+
857+
expect((await storedCall(toolCallId)).pickupDeadlineAt).toEqual(offeredRow.pickupDeadlineAt)
858+
const { executionToken } = await desktop.claim(toolCallId)
859+
await desktop.complete(toolCallId, executionToken, { clicked: true })
860+
await answer
861+
expect(resultOf(context, toolCallId)).toMatchObject({ success: true })
862+
},
863+
TURN_WAIT_MS
864+
)
865+
734866
it(
735867
'settles a call Sim never offered once its pickup window would have closed',
736868
async () => {

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

Lines changed: 63 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,8 @@ import {
3131
isTerminalAsyncStatus,
3232
} from '@/lib/mothership/async-runs/lifecycle'
3333
import { DESKTOP_TOOL_PICKUP_GRACE_MS } from '@/lib/mothership/constants'
34-
import { DESKTOP_TOOL_CALL_NAMES } from '@/lib/mothership/tools/desktop-tools'
34+
import { NAMED_DESKTOP_TOOL_NAMES } from '@/lib/mothership/tools/desktop-tools'
35+
import { USER_LOCAL_VFS_ROOT } from '@/lib/mothership/tools/local-filesystem'
3536

3637
const LIVE_RUN_STATUSES: CopilotRunStatus[] = ['active', 'paused_waiting_for_tool', 'resuming']
3738

@@ -61,8 +62,51 @@ function pickupOverdueAt(at: SQL) {
6162
sql`${pickupDeadline} <= ${at}`
6263
)
6364
}
65+
66+
function isUserLocalVfsPath(path: SQL) {
67+
return sql`(${path} = ${USER_LOCAL_VFS_ROOT} OR ${path} LIKE ${`${USER_LOCAL_VFS_ROOT}/%`})`
68+
}
69+
70+
/**
71+
* The SQL form of `isDesktopToolCall`, so a query limits only over calls the desktop runs: a
72+
* desktop tool by name, or a VFS read of a granted local folder (not a read of Sim's own files).
73+
*/
74+
const isDesktopToolCallRow = or(
75+
inArray(copilotAsyncToolCalls.toolName, [...NAMED_DESKTOP_TOOL_NAMES]),
76+
and(
77+
inArray(copilotAsyncToolCalls.toolName, ['read', 'grep']),
78+
isUserLocalVfsPath(sql`${copilotAsyncToolCalls.args}->>'path'`)
79+
),
80+
and(
81+
eq(copilotAsyncToolCalls.toolName, 'glob'),
82+
isUserLocalVfsPath(sql`${copilotAsyncToolCalls.args}->>'pattern'`)
83+
)
84+
)
6485
const INBOX_ROW_LIMIT = 500
6586

87+
/**
88+
* Persistence order, which follows the order the model emitted the calls in: calls of one turn
89+
* can share a millisecond, so the timestamp alone cannot order them. Rows persisted before the
90+
* sequence existed have none and come first, as they are the oldest.
91+
*/
92+
const persistOrder = [
93+
sql`${copilotAsyncToolCalls.persistSeq} ASC NULLS FIRST`,
94+
asc(copilotAsyncToolCalls.createdAt),
95+
asc(copilotAsyncToolCalls.toolCallId),
96+
]
97+
98+
function comparePersistOrder(
99+
a: { persistSeq: number | null; createdAt: Date; toolCallId: string },
100+
b: { persistSeq: number | null; createdAt: Date; toolCallId: string }
101+
): number {
102+
if (a.persistSeq !== b.persistSeq) {
103+
if (a.persistSeq === null) return -1
104+
if (b.persistSeq === null) return 1
105+
return a.persistSeq - b.persistSeq
106+
}
107+
return a.createdAt.getTime() - b.createdAt.getTime() || a.toolCallId.localeCompare(b.toolCallId)
108+
}
109+
66110
export interface DesktopDeviceRegistration {
67111
id: string
68112
userId: string
@@ -157,7 +201,7 @@ export async function touchDesktopDevice(deviceId: string): Promise<void> {
157201
* out the other: unclaimed calls on its recent open runs that are offered or waiting for the
158202
* user's decision, and calls it claimed that Sim settled without its result and it has not yet
159203
* acknowledged (cancel items). A call the device is still running is never listed: it already
160-
* holds it. Ordered by persistence time, the order the device claims in.
204+
* holds it. Ordered by persistence, the order the device claims in.
161205
*/
162206
export async function listDesktopInboxRows(identity: Omit<DesktopDeviceIdentity, 'sessionId'>) {
163207
const rowsWhere = (state: SQL | undefined) =>
@@ -171,6 +215,7 @@ export async function listDesktopInboxRows(identity: Omit<DesktopDeviceIdentity,
171215
permissionDecision: copilotAsyncToolCalls.permissionDecision,
172216
claimed: sql<boolean>`${copilotAsyncToolCalls.executionOwnerToken} IS NOT NULL`,
173217
createdAt: copilotAsyncToolCalls.createdAt,
218+
persistSeq: copilotAsyncToolCalls.persistSeq,
174219
chatId: copilotRuns.chatId,
175220
chatTitle: copilotChats.title,
176221
workspaceId: copilotRuns.workspaceId,
@@ -183,11 +228,11 @@ export async function listDesktopInboxRows(identity: Omit<DesktopDeviceIdentity,
183228
eq(copilotRuns.desktopDeviceId, identity.deviceId),
184229
eq(copilotRuns.userId, identity.userId),
185230
sql`${copilotRuns.startedAt} > now() - make_interval(hours => ${DESKTOP_INBOX_HORIZON_HOURS})`,
186-
inArray(copilotAsyncToolCalls.toolName, [...DESKTOP_TOOL_CALL_NAMES]),
231+
isDesktopToolCallRow,
187232
state
188233
)
189234
)
190-
.orderBy(asc(copilotAsyncToolCalls.createdAt), asc(copilotAsyncToolCalls.toolCallId))
235+
.orderBy(...persistOrder)
191236
.limit(INBOX_ROW_LIMIT)
192237
const [waiting, cancelled] = await Promise.all([
193238
rowsWhere(
@@ -207,10 +252,7 @@ export async function listDesktopInboxRows(identity: Omit<DesktopDeviceIdentity,
207252
)
208253
),
209254
])
210-
return [...waiting, ...cancelled].sort(
211-
(a, b) =>
212-
a.createdAt.getTime() - b.createdAt.getTime() || a.toolCallId.localeCompare(b.toolCallId)
213-
)
255+
return [...waiting, ...cancelled].sort(comparePersistOrder)
214256
}
215257

216258
export type DesktopInboxRow = Awaited<ReturnType<typeof listDesktopInboxRows>>[number]
@@ -336,12 +378,14 @@ export async function offerDesktopToolCall(input: {
336378
/**
337379
* A bound desktop call with the deadlines only Sim enforces, read on the database clock every CAS
338380
* uses: whether its pickup window closed while it is unclaimed, and whether its lease lapsed while
339-
* it runs. Null for a call on a run no device is bound to.
381+
* it runs. Null for a call on a run no device is bound to, and for a tool the desktop never runs.
340382
*/
341383
export async function getDesktopToolCallDeadlines(toolCallId: string) {
342384
const [row] = await db
343385
.select({
344386
toolCallId: copilotAsyncToolCalls.toolCallId,
387+
toolName: copilotAsyncToolCalls.toolName,
388+
args: copilotAsyncToolCalls.args,
345389
runId: copilotRuns.id,
346390
userId: copilotRuns.userId,
347391
deviceId: copilotRuns.desktopDeviceId,
@@ -361,7 +405,11 @@ export async function getDesktopToolCallDeadlines(toolCallId: string) {
361405
.innerJoin(copilotRuns, eq(copilotRuns.id, copilotAsyncToolCalls.runId))
362406
.leftJoin(desktopDevices, eq(desktopDevices.id, copilotRuns.desktopDeviceId))
363407
.where(
364-
and(eq(copilotAsyncToolCalls.toolCallId, toolCallId), isNotNull(copilotRuns.desktopDeviceId))
408+
and(
409+
eq(copilotAsyncToolCalls.toolCallId, toolCallId),
410+
isNotNull(copilotRuns.desktopDeviceId),
411+
isDesktopToolCallRow
412+
)
365413
)
366414
.limit(1)
367415
return row?.deviceId ? { ...row, deviceId: row.deviceId } : null
@@ -374,7 +422,9 @@ export type DesktopToolCallDeadlines = NonNullable<
374422
/**
375423
* Bound desktop calls a deadline passed for at least `slackMs` ago: unclaimed past their pickup
376424
* deadline (offered or not), or claimed by the executor with a lapsed lease. A live waiter settles
377-
* these within its 5 s poll, so anything this finds lost its waiter.
425+
* these within its 5 s poll, so anything this finds lost its waiter. The scan starts from the
426+
* few unsettled calls (pending or running), in persistence order, so however long a call stayed
427+
* overdue it is still reached.
378428
*/
379429
export async function listOverdueDesktopToolCalls(input: { slackMs: number; limit: number }) {
380430
const overdue = sql`clock_timestamp() - ${input.slackMs} * interval '1 millisecond'`
@@ -385,6 +435,7 @@ export async function listOverdueDesktopToolCalls(input: { slackMs: number; limi
385435
.where(
386436
and(
387437
isNotNull(copilotRuns.desktopDeviceId),
438+
isDesktopToolCallRow,
388439
or(
389440
pickupOverdueAt(overdue),
390441
and(
@@ -397,6 +448,7 @@ export async function listOverdueDesktopToolCalls(input: { slackMs: number; limi
397448
)
398449
)
399450
)
451+
.orderBy(...persistOrder)
400452
.limit(input.limit)
401453
return rows.map((row) => row.toolCallId)
402454
}

‎apps/sim/lib/mothership/async-runs/browser-download-claim.integration.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ describe('browser download admission with PostgreSQL', () => {
4848
result jsonb, error text, permission_decision text, permission_decided_at timestamp,
4949
permission_requested_at timestamp,
5050
pickup_deadline_at timestamp with time zone,
51+
persist_seq bigint,
5152
claimed_at timestamp, claimed_by text, execution_started_at timestamp,
5253
execution_settled_at timestamp, execution_owner_token text,
5354
execution_lease_expires_at timestamptz, execution_revoked_at timestamptz,

‎apps/sim/lib/mothership/async-runs/repository.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1738,8 +1738,9 @@ export async function recordToolPermissionDecision(
17381738
.set({
17391739
permissionDecision: decision,
17401740
permissionDecidedAt: now,
1741-
// No pickup window runs while the user decides: an allowed call is offered afresh.
1742-
pickupDeadlineAt: null,
1741+
// No pickup window runs while the user decides: a gated call is offered afresh once
1742+
// allowed. A call that was never gated keeps the window it was offered with.
1743+
pickupDeadlineAt: sql`CASE WHEN ${copilotAsyncToolCalls.permissionRequestedAt} IS NULL THEN ${copilotAsyncToolCalls.pickupDeadlineAt} END`,
17431744
updatedAt: now,
17441745
})
17451746
.where(

‎apps/sim/lib/mothership/request/tools/desktop-wait.test.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,8 @@ describe('settleAbandonedDesktopToolCalls', () => {
9898
.mockRejectedValueOnce(new Error('connection reset'))
9999
.mockResolvedValueOnce({
100100
toolCallId: 'overdue',
101+
toolName: 'browser_click',
102+
args: { ref: 'e1' },
101103
runId: 'run-1',
102104
userId: 'user-1',
103105
deviceId: 'device-1',

‎apps/sim/lib/mothership/request/tools/desktop-wait.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { createLogger } from '@sim/logger'
22
import { toError } from '@sim/utils/errors'
33
import { interruptibleSleep } from '@sim/utils/helpers'
4+
import { toRecordOrNull } from '@sim/utils/object'
45
import { ringDesktopInbox } from '@/lib/desktop/executor/doorbell'
56
import { isDesktopPresent } from '@/lib/desktop/executor/presence'
67
import {
@@ -23,6 +24,7 @@ import {
2324
type ClientToolSettlementGuard,
2425
settleClientToolCall,
2526
} from '@/lib/mothership/request/tools/client-settlement.server'
27+
import { isDesktopToolCall } from '@/lib/mothership/tools/desktop-tools'
2628
import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
2729

2830
const logger = createLogger('CopilotDesktopToolWait')
@@ -230,7 +232,9 @@ async function readDesktopPresence(deviceId: string): Promise<boolean | null> {
230232
*/
231233
async function settleOverdueDesktopToolCall(toolCallId: string): Promise<boolean> {
232234
const call = await getDesktopToolCallDeadlines(toolCallId)
233-
if (!call) return false
235+
// A VFS read of Sim's own files shares its tool name with a local read, but no desktop runs it.
236+
if (!call || !isDesktopToolCall(call.toolName, toRecordOrNull(call.args) ?? undefined))
237+
return false
234238
if (call.status === ASYNC_TOOL_STATUS.pending) {
235239
const present = await readDesktopPresence(call.deviceId)
236240
// Before its pickup window closes, a call fails early only when its device is known to be away:

0 commit comments

Comments
 (0)