Skip to content

Commit 9e92d98

Browse files
committed
fix(desktop): serialize a device result with Stop through the run row
A result now locks the run row before checking tool admission, the same serialization the claim and Stop use, so it can no longer commit after a concurrent Stop. The schema declares the run-device index concurrent, as the migration builds it.
1 parent 574d70f commit 9e92d98

4 files changed

Lines changed: 77 additions & 29 deletions

File tree

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

Lines changed: 43 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ import {
2828
user,
2929
workspace,
3030
} from '@sim/db/schema'
31+
import { createDeferred } from '@sim/testing/helpers/deferred'
3132
import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock'
3233
import { sleep } from '@sim/utils/helpers'
3334
import { generateId } from '@sim/utils/id'
@@ -544,6 +545,37 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
544545
})
545546

546547
describe('completion', () => {
548+
it('serializes a result with a concurrent Stop, so it never lands after the Stop commits', async () => {
549+
const desktop = await signedInDesktop()
550+
const run = await boundRun(desktop)
551+
const toolCallId = await pendingCall(run.runId)
552+
await offer(toolCallId)
553+
const { executionToken } = await claim(desktop.principal, desktop.deviceId, toolCallId)
554+
555+
/** Stop closes admission on the run row; hold its transaction open while the result arrives. */
556+
const stopHeld = createDeferred<void>()
557+
const releaseStop = createDeferred<void>()
558+
const stopping = db.transaction(async (tx) => {
559+
await tx
560+
.update(copilotRuns)
561+
.set({ toolAdmissionClosedAt: sql`now()` })
562+
.where(eq(copilotRuns.id, run.runId))
563+
stopHeld.resolve()
564+
await releaseStop.promise
565+
})
566+
await stopHeld.promise
567+
const completing = completeDesktopTool.execute({
568+
principal: desktop.principal,
569+
input: { deviceId: desktop.deviceId, toolCallId, executionToken, status: 'success' },
570+
})
571+
await sleep(300)
572+
releaseStop.resolve()
573+
await stopping
574+
575+
await expect(completing).resolves.toEqual({ outcome: 'superseded', status: 'cancelled' })
576+
expect((await row(toolCallId)).status).toBe('cancelled')
577+
})
578+
547579
it('records a result once, wakes the waiting run, and answers a retry as a duplicate', async () => {
548580
const desktop = await signedInDesktop()
549581
const run = await boundRun(desktop)
@@ -713,9 +745,17 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
713745
principal: desktop.principal,
714746
input: { deviceId: desktop.deviceId },
715747
})
716-
await expect(stream.revalidate()).resolves.toBeUndefined()
717-
await db.delete(session).where(eq(session.id, desktop.sessionId))
718-
await expect(stream.revalidate()).rejects.toBeInstanceOf(DesktopDeviceUnrecognizedError)
748+
const close = stream.subscribe(() => {})
749+
try {
750+
await expect.poll(() => isDesktopPresent(desktop.deviceId)).toBe(true)
751+
await expect(stream.revalidate()).resolves.toBeUndefined()
752+
await db.delete(session).where(eq(session.id, desktop.sessionId))
753+
/** The SSE transport closes a stream whose revalidation rejects; closing releases presence. */
754+
await expect(stream.revalidate()).rejects.toBeInstanceOf(DesktopDeviceUnrecognizedError)
755+
} finally {
756+
close()
757+
}
758+
await expect.poll(() => isDesktopPresent(desktop.deviceId)).toBe(false)
719759
})
720760
})
721761
})

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

Lines changed: 31 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -301,31 +301,38 @@ export async function recordDesktopCallResult(input: {
301301
result: AsyncCompletionData
302302
error: string | null
303303
}): Promise<boolean> {
304-
const [row] = await db
305-
.update(copilotAsyncToolCalls)
306-
.set({
307-
status: input.status,
308-
result: sanitizeValueForJsonb(input.result),
309-
error: input.error,
310-
claimedBy: null,
311-
claimedAt: null,
312-
completedAt: sql`now()`,
313-
executionSettledAt: sql`now()`,
314-
updatedAt: sql`now()`,
315-
})
316-
.where(
317-
and(
318-
eq(copilotAsyncToolCalls.toolCallId, input.toolCallId),
319-
eq(copilotAsyncToolCalls.runId, input.runId),
320-
eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.running),
321-
eq(copilotAsyncToolCalls.executionOwnerToken, input.ownerToken),
322-
isNull(copilotAsyncToolCalls.executionRevokedAt),
323-
sql`EXISTS (SELECT 1 FROM ${copilotRuns} r WHERE r.id = ${copilotAsyncToolCalls.runId}
324-
AND r.tool_admission_closed_at IS NULL)`
304+
return db.transaction(async (tx) => {
305+
/** Serialized with Stop through the run row, exactly like the claim. */
306+
const [run] = await tx
307+
.select({ toolAdmissionClosedAt: copilotRuns.toolAdmissionClosedAt })
308+
.from(copilotRuns)
309+
.where(eq(copilotRuns.id, input.runId))
310+
.for('update')
311+
if (!run || run.toolAdmissionClosedAt) return false
312+
const [row] = await tx
313+
.update(copilotAsyncToolCalls)
314+
.set({
315+
status: input.status,
316+
result: sanitizeValueForJsonb(input.result),
317+
error: input.error,
318+
claimedBy: null,
319+
claimedAt: null,
320+
completedAt: sql`now()`,
321+
executionSettledAt: sql`now()`,
322+
updatedAt: sql`now()`,
323+
})
324+
.where(
325+
and(
326+
eq(copilotAsyncToolCalls.toolCallId, input.toolCallId),
327+
eq(copilotAsyncToolCalls.runId, input.runId),
328+
eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.running),
329+
eq(copilotAsyncToolCalls.executionOwnerToken, input.ownerToken),
330+
isNull(copilotAsyncToolCalls.executionRevokedAt)
331+
)
325332
)
326-
)
327-
.returning({ toolCallId: copilotAsyncToolCalls.toolCallId })
328-
return Boolean(row)
333+
.returning({ toolCallId: copilotAsyncToolCalls.toolCallId })
334+
return Boolean(row)
335+
})
329336
}
330337

331338
export type DesktopCallAcknowledgement =

‎packages/db/migrations/meta/0395_snapshot.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3165,7 +3165,7 @@
31653165
],
31663166
"isUnique": false,
31673167
"where": "\"copilot_runs\".\"desktop_device_id\" IS NOT NULL",
3168-
"concurrently": false,
3168+
"concurrently": true,
31693169
"method": "btree",
31703170
"with": {}
31713171
}

‎packages/db/schema.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4420,7 +4420,8 @@ export const copilotRuns = pgTable(
44204420
streamIdUnique: uniqueIndex('copilot_runs_stream_id_unique').on(table.streamId),
44214421
desktopDeviceStartedAtIdx: index('copilot_runs_desktop_device_started_at_idx')
44224422
.on(table.desktopDeviceId, table.startedAt)
4423-
.where(sql`${table.desktopDeviceId} IS NOT NULL`),
4423+
.where(sql`${table.desktopDeviceId} IS NOT NULL`)
4424+
.concurrently(),
44244425
})
44254426
)
44264427

0 commit comments

Comments
 (0)