Skip to content

Commit 4d32e76

Browse files
committed
fix(desktop): leave a stopped run's desktop calls to Stop
The supervisor's not-started and outcome-unknown settlements now hold the run row against Stop, like the claim and a device result, and refuse once admission is closed; a supervisor whose run was stopped exits. An unreadable presence falls back to the pickup window instead of failing the call as offline.
1 parent ea73898 commit 4d32e76

3 files changed

Lines changed: 155 additions & 48 deletions

File tree

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

Lines changed: 75 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,8 @@ import {
88
desktopDevices,
99
} from '@sim/db/schema'
1010
import { sanitizeValueForJsonb } from '@sim/utils/string'
11-
import { and, asc, eq, inArray, isNotNull, isNull, ne, or, sql } from 'drizzle-orm'
11+
import { and, asc, eq, inArray, isNotNull, isNull, ne, notInArray, or, sql } from 'drizzle-orm'
12+
import type { DbTransaction } from '@/lib/db/types'
1213
import {
1314
DESKTOP_CALL_LEASE_SECONDS,
1415
DESKTOP_INBOX_HORIZON_HOURS,
@@ -485,8 +486,10 @@ export async function getDesktopCallState(toolCallId: string) {
485486
msUntilLeaseEnd: sql<
486487
number | null
487488
>`(extract(epoch from (${copilotAsyncToolCalls.executionLeaseExpiresAt} - clock_timestamp())) * 1000)::float8`,
489+
runOpen: sql<boolean>`${copilotRuns.toolAdmissionClosedAt} IS NULL AND ${copilotRuns.status} NOT IN ('complete', 'error', 'cancelled')`,
488490
})
489491
.from(copilotAsyncToolCalls)
492+
.innerJoin(copilotRuns, eq(copilotRuns.id, copilotAsyncToolCalls.runId))
490493
.where(eq(copilotAsyncToolCalls.toolCallId, toolCallId))
491494
.limit(1)
492495
return row ?? null
@@ -499,35 +502,61 @@ interface DesktopCallFailure {
499502
error: string
500503
}
501504

505+
/**
506+
* Runs a settlement while holding the run row against Stop, the same serialization the claim and a
507+
* device result use. Nothing is settled on a stopped or ended run: Stop answers for its calls.
508+
*/
509+
async function settleOnOpenRun(
510+
runId: string,
511+
settle: (tx: DbTransaction) => Promise<boolean>
512+
): Promise<boolean> {
513+
return db.transaction(async (tx) => {
514+
const [run] = await tx
515+
.select({ id: copilotRuns.id })
516+
.from(copilotRuns)
517+
.where(
518+
and(
519+
eq(copilotRuns.id, runId),
520+
isNull(copilotRuns.toolAdmissionClosedAt),
521+
notInArray(copilotRuns.status, TERMINAL_RUN_STATUSES)
522+
)
523+
)
524+
.for('share')
525+
return run ? settle(tx) : false
526+
})
527+
}
528+
502529
/**
503530
* Fails a call nobody claimed: the inverse CAS of the claim, so exactly one of the two wins.
504531
* `deadlinePassed` limits it to a call whose pickup window has closed.
505532
*/
506533
export async function failUnclaimedDesktopCall(
507534
input: DesktopCallFailure & { deadlinePassed: boolean }
508535
): Promise<boolean> {
509-
const [row] = await db
510-
.update(copilotAsyncToolCalls)
511-
.set({
512-
status: ASYNC_TOOL_STATUS.failed,
513-
result: sanitizeValueForJsonb(input.result),
514-
error: input.error,
515-
completedAt: sql`now()`,
516-
updatedAt: sql`now()`,
517-
})
518-
.where(
519-
and(
520-
eq(copilotAsyncToolCalls.toolCallId, input.toolCallId),
521-
eq(copilotAsyncToolCalls.runId, input.runId),
522-
eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.pending),
523-
isNull(copilotAsyncToolCalls.executionOwnerToken),
524-
input.deadlinePassed
525-
? sql`${copilotAsyncToolCalls.executionLeaseExpiresAt} <= clock_timestamp()`
526-
: undefined
536+
return settleOnOpenRun(input.runId, async (tx) => {
537+
const [row] = await tx
538+
.update(copilotAsyncToolCalls)
539+
.set({
540+
status: ASYNC_TOOL_STATUS.failed,
541+
result: sanitizeValueForJsonb(input.result),
542+
error: input.error,
543+
completedAt: sql`now()`,
544+
updatedAt: sql`now()`,
545+
})
546+
.where(
547+
and(
548+
eq(copilotAsyncToolCalls.toolCallId, input.toolCallId),
549+
eq(copilotAsyncToolCalls.runId, input.runId),
550+
eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.pending),
551+
isNull(copilotAsyncToolCalls.executionOwnerToken),
552+
input.deadlinePassed
553+
? sql`${copilotAsyncToolCalls.executionLeaseExpiresAt} <= clock_timestamp()`
554+
: undefined
555+
)
527556
)
528-
)
529-
.returning({ toolCallId: copilotAsyncToolCalls.toolCallId })
530-
return Boolean(row)
557+
.returning({ toolCallId: copilotAsyncToolCalls.toolCallId })
558+
return Boolean(row)
559+
})
531560
}
532561

533562
/**
@@ -537,27 +566,29 @@ export async function failUnclaimedDesktopCall(
537566
export async function failLapsedDesktopCall(
538567
input: DesktopCallFailure & { ownerToken: string }
539568
): Promise<boolean> {
540-
const [row] = await db
541-
.update(copilotAsyncToolCalls)
542-
.set({
543-
status: ASYNC_TOOL_STATUS.failed,
544-
result: sanitizeValueForJsonb(input.result),
545-
error: input.error,
546-
claimedBy: null,
547-
claimedAt: null,
548-
completedAt: sql`now()`,
549-
updatedAt: sql`now()`,
550-
})
551-
.where(
552-
and(
553-
eq(copilotAsyncToolCalls.toolCallId, input.toolCallId),
554-
eq(copilotAsyncToolCalls.runId, input.runId),
555-
eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.running),
556-
eq(copilotAsyncToolCalls.executionOwnerToken, input.ownerToken),
557-
isNull(copilotAsyncToolCalls.executionRevokedAt),
558-
sql`${copilotAsyncToolCalls.executionLeaseExpiresAt} <= clock_timestamp()`
569+
return settleOnOpenRun(input.runId, async (tx) => {
570+
const [row] = await tx
571+
.update(copilotAsyncToolCalls)
572+
.set({
573+
status: ASYNC_TOOL_STATUS.failed,
574+
result: sanitizeValueForJsonb(input.result),
575+
error: input.error,
576+
claimedBy: null,
577+
claimedAt: null,
578+
completedAt: sql`now()`,
579+
updatedAt: sql`now()`,
580+
})
581+
.where(
582+
and(
583+
eq(copilotAsyncToolCalls.toolCallId, input.toolCallId),
584+
eq(copilotAsyncToolCalls.runId, input.runId),
585+
eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.running),
586+
eq(copilotAsyncToolCalls.executionOwnerToken, input.ownerToken),
587+
isNull(copilotAsyncToolCalls.executionRevokedAt),
588+
sql`${copilotAsyncToolCalls.executionLeaseExpiresAt} <= clock_timestamp()`
589+
)
559590
)
560-
)
561-
.returning({ toolCallId: copilotAsyncToolCalls.toolCallId })
562-
return Boolean(row)
591+
.returning({ toolCallId: copilotAsyncToolCalls.toolCallId })
592+
return Boolean(row)
593+
})
563594
}

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

Lines changed: 68 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,11 +20,13 @@ import { db } from '@sim/db'
2020
import {
2121
copilotAsyncToolCalls,
2222
copilotChats,
23+
copilotRuns,
2324
desktopDevices,
2425
session,
2526
user,
2627
workspace,
2728
} from '@sim/db/schema'
29+
import { createDeferred } from '@sim/testing/helpers/deferred'
2830
import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock'
2931
import { sleep } from '@sim/utils/helpers'
3032
import { generateId } from '@sim/utils/id'
@@ -40,7 +42,7 @@ import {
4042
resolveTurnDesktopDevice,
4143
} from '@/lib/desktop/application/executor'
4244
import { isDesktopPresent } from '@/lib/desktop/executor/presence'
43-
import { failLapsedDesktopCall } from '@/lib/desktop/executor/repository'
45+
import { failLapsedDesktopCall, failUnclaimedDesktopCall } from '@/lib/desktop/executor/repository'
4446
import {
4547
DESKTOP_LEASE_LOST_MESSAGE,
4648
DESKTOP_NOT_RESPONDING_MESSAGE,
@@ -521,6 +523,71 @@ describe.runIf(Boolean(redisUrl))('desktop call supervision and turn binding', (
521523
})
522524
})
523525

526+
it('serializes a not-started settlement with a concurrent Stop', async () => {
527+
const desktop = await signedInDesktop()
528+
const run = await boundRun(desktop.userId, desktop.deviceId)
529+
const toolCallId = await pendingCall(run.runId)
530+
531+
/** Stop's run-row update, held open until the settlement is queued behind it. */
532+
const stopHeld = createDeferred<number>()
533+
const releaseStop = createDeferred<void>()
534+
const stopping = db.transaction(async (tx) => {
535+
await tx
536+
.update(copilotRuns)
537+
.set({ toolAdmissionClosedAt: sql`now()` })
538+
.where(eq(copilotRuns.id, run.runId))
539+
const [backend] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`)
540+
stopHeld.resolve(backend.pid)
541+
await releaseStop.promise
542+
})
543+
const stopPid = await stopHeld.promise
544+
const settling = failUnclaimedDesktopCall({
545+
toolCallId,
546+
runId: run.runId,
547+
result: {},
548+
error: 'offline',
549+
deadlinePassed: false,
550+
})
551+
await expect
552+
.poll(async () => {
553+
const [waiter] = await db.execute<{ pid: number }>(sql`
554+
SELECT pid FROM pg_stat_activity WHERE wait_event_type = 'Lock'
555+
AND ${stopPid}::int = ANY(pg_blocking_pids(pid)) LIMIT 1`)
556+
return Boolean(waiter)
557+
})
558+
.toBe(true)
559+
releaseStop.resolve()
560+
await stopping
561+
562+
await expect(settling).resolves.toBe(false)
563+
expect((await row(toolCallId)).status).toBe('pending')
564+
})
565+
566+
it('leaves a call on a stopped run to Stop rather than failing it as not responding', async () => {
567+
const desktop = await signedInDesktop()
568+
const run = await boundRun(desktop.userId, desktop.deviceId)
569+
const toolCallId = await pendingCall(run.runId)
570+
571+
await online(desktop, async () => {
572+
const supervision = superviseDesktopCall({
573+
toolCallId,
574+
runId: run.runId,
575+
userId: desktop.userId,
576+
deviceId: desktop.deviceId,
577+
pickupGraceMs: 400,
578+
})
579+
await expect
580+
.poll(async () => (await row(toolCallId)).executionLeaseExpiresAt)
581+
.not.toBeNull()
582+
await db
583+
.update(copilotRuns)
584+
.set({ toolAdmissionClosedAt: sql`now()` })
585+
.where(eq(copilotRuns.id, run.runId))
586+
await expect(supervision).resolves.toBe('settled')
587+
})
588+
expect((await row(toolCallId)).status).toBe('pending')
589+
})
590+
524591
it('never offers a call on a run the chat view serves', async () => {
525592
const desktop = await signedInDesktop()
526593
const run = await boundRun(desktop.userId, null)

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

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ export type DesktopCallSupervision =
3939
| 'offline'
4040
| 'not_responding'
4141
| 'lease_lost'
42-
/** The device's result, Stop, or another settlement ended it. */
42+
/** The device's result, Stop, or another settlement ended it, or its run was stopped. */
4343
| 'settled'
4444
| 'aborted'
4545

@@ -153,7 +153,15 @@ export async function superviseDesktopCall(
153153
})
154154
if (offered) {
155155
ringDesktopInbox(input.deviceId, 'call')
156-
if (!(await isDesktopPresent(input.deviceId).catch(() => false))) {
156+
/** Only a confirmed absence fails fast; an unreadable presence falls back to the pickup window. */
157+
const present = await isDesktopPresent(input.deviceId).catch((error) => {
158+
logger.warn('Desktop presence could not be read; waiting out the pickup window', {
159+
deviceId: input.deviceId,
160+
error: getErrorMessage(error),
161+
})
162+
return null
163+
})
164+
if (present === false) {
157165
if (await failUnclaimed(input, 'offline', DESKTOP_OFFLINE_MESSAGE, false)) return 'offline'
158166
}
159167
} else if (!isOfferedOrClaimed(await getDesktopCallState(input.toolCallId))) {
@@ -163,7 +171,8 @@ export async function superviseDesktopCall(
163171
for (;;) {
164172
if (input.signal?.aborted) return 'aborted'
165173
const state = await getDesktopCallState(input.toolCallId)
166-
if (!state || isTerminalAsyncStatus(state.status)) return 'settled'
174+
/** Settled, held by a legacy claim, or its run was stopped: Stop answers for its calls. */
175+
if (!state || isTerminalAsyncStatus(state.status) || !state.runOpen) return 'settled'
167176
if (!isOfferedOrClaimed(state)) return 'settled'
168177
const remainingMs = state.msUntilLeaseEnd ?? 0
169178
if (remainingMs <= 0) {

0 commit comments

Comments
 (0)