Skip to content

Commit 574d70f

Browse files
committed
fix(desktop): address review on the background executor protocol
- Presence is a per-connection set, so overlapping streams never erase each other's presence. - A result posted after Stop settles the call as stopped (superseded) instead of recording a success for a stopped turn. - Nothing reads as awaiting approval while tool permissions are off. - One reconcile cadence (10 s), shorter than the pickup grace, so a lost doorbell cannot fail the first call of a turn the device has not seen. - The inbox query skips unclaimed calls on stopped runs and non-desktop tools, so they cannot crowd actionable rows past its limit. - The doorbell stream applies the executor rate limit; inbox ids use the claim route's id schema; the index build recovers from an interrupted concurrent build; the E2E waits for the cancel doorbell specifically.
1 parent d8356a1 commit 574d70f

11 files changed

Lines changed: 160 additions & 98 deletions

File tree

‎apps/sim/app/api/desktop/inbox/route.ts‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,7 @@ export const GET = defineInternalJsonRoute({
1616
errorPolicy: desktopExecutorErrorPolicy,
1717
mapInput: ({ query }) => ({ deviceId: query.deviceId }),
1818
useCase: listDesktopInbox,
19-
present: ({ items, hasActiveRun }) => ({
20-
hasActiveRun,
19+
present: ({ items }) => ({
2120
items: items.map((item) =>
2221
item.kind === 'call' ? { ...item, createdAt: item.createdAt.toISOString() } : item
2322
),

‎apps/sim/app/api/desktop/inbox/stream/route.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { createLogger } from '@sim/logger'
22
import type { NextRequest } from 'next/server'
33
import { desktopInboxStreamContract } from '@/lib/api/contracts/desktop-executor'
44
import { parseRequest } from '@/lib/api/server'
5+
import { desktopExecutorRateLimit } from '@/lib/api/server/routes/desktop-executor'
56
import {
67
InternalUnauthenticatedError,
78
internalSessionAuth,
@@ -22,6 +23,8 @@ const logger = createLogger('DesktopInboxStream')
2223
export const GET = withRouteHandler(async (request: NextRequest) => {
2324
try {
2425
const principal = await internalSessionAuth.authenticate()
26+
const limited = await desktopExecutorRateLimit.enforce(request, principal)
27+
if (limited) return limited
2528
const parsed = await parseRequest(desktopInboxStreamContract, request, {})
2629
if (!parsed.success) return parsed.response
2730
const { deviceId } = parsed.data.query

‎apps/sim/lib/api/contracts/desktop-executor.ts‎

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -37,8 +37,7 @@ export const registerDesktopDeviceResponseSchema = z.object({
3737
leaseMs: z.number().int().positive(),
3838
leaseRenewMs: z.number().int().positive(),
3939
pickupGraceMs: z.number().int().positive(),
40-
activeReconcileMs: z.number().int().positive(),
41-
idleReconcileMs: z.number().int().positive(),
40+
reconcileMs: z.number().int().positive(),
4241
})
4342

4443
export const registerDesktopDeviceContract = defineRouteContract({
@@ -53,7 +52,7 @@ const desktopInboxQuerySchema = z.object({ deviceId: desktopDeviceIdSchema })
5352

5453
const desktopInboxCallSchema = z.object({
5554
kind: z.literal('call'),
56-
toolCallId: z.string().min(1),
55+
toolCallId: desktopToolCallIdSchema,
5756
toolName: z.string().min(1),
5857
chatId: z.string().min(1),
5958
workspaceId: z.string().nullable(),
@@ -62,7 +61,7 @@ const desktopInboxCallSchema = z.object({
6261

6362
const desktopInboxApprovalSchema = z.object({
6463
kind: z.literal('approval_needed'),
65-
toolCallId: z.string().min(1),
64+
toolCallId: desktopToolCallIdSchema,
6665
toolName: z.string().min(1),
6766
chatId: z.string().min(1),
6867
chatTitle: z.string().nullable(),
@@ -73,7 +72,7 @@ const desktopInboxApprovalSchema = z.object({
7372

7473
const desktopInboxCancelSchema = z.object({
7574
kind: z.literal('cancel'),
76-
toolCallId: z.string().min(1),
75+
toolCallId: desktopToolCallIdSchema,
7776
})
7877

7978
/** Ordered by when each call was persisted, which is the order the executor claims them in. */
@@ -86,8 +85,6 @@ export type DesktopInboxItem = z.output<typeof desktopInboxItemSchema>
8685

8786
export const desktopInboxResponseSchema = z.object({
8887
items: z.array(desktopInboxItemSchema),
89-
/** True while any run bound to this device is active, so the device pulls at the active cadence. */
90-
hasActiveRun: z.boolean(),
9188
})
9289

9390
export const listDesktopInboxContract = defineRouteContract({

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

Lines changed: 39 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -9,8 +9,9 @@ import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest'
99
const { redisUrl } = await vi.hoisted(async () => {
1010
const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure')
1111
const url = readTestRedisUrl()
12-
/** The real Redis module reads this at import. */
12+
/** The real Redis module and the tool-permission switch read these at import. */
1313
if (url) process.env.REDIS_URL = url
14+
process.env.COPILOT_TOOL_PERMISSIONS_ENABLED = 'true'
1415
return { redisUrl: url }
1516
})
1617

@@ -28,6 +29,7 @@ import {
2829
workspace,
2930
} from '@sim/db/schema'
3031
import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock'
32+
import { sleep } from '@sim/utils/helpers'
3133
import { generateId } from '@sim/utils/id'
3234
import { eq, inArray, sql } from 'drizzle-orm'
3335
import { closeRedisConnection } from '@/lib/core/config/redis'
@@ -278,7 +280,7 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
278280
).rejects.toBeInstanceOf(DesktopDeviceUnrecognizedError)
279281
await expect(
280282
listDesktopInbox.execute({ principal: newer, input: { deviceId: desktop.deviceId } })
281-
).resolves.toEqual({ items: [], hasActiveRun: false })
283+
).resolves.toEqual({ items: [] })
282284
})
283285

284286
it('disconnects the device when its session is signed out', async () => {
@@ -519,18 +521,20 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
519521
input: { deviceId: desktop.deviceId },
520522
})
521523
expect(inbox.items).toEqual([{ kind: 'cancel', toolCallId: running }])
522-
expect(inbox.hasActiveRun).toBe(false)
523524

524-
/** The device's cancelled result acknowledges the cancellation. */
525-
await completeDesktopTool.execute({
526-
principal: desktop.principal,
527-
input: {
528-
deviceId: desktop.deviceId,
529-
toolCallId: running,
530-
executionToken,
531-
status: 'cancelled',
532-
},
533-
})
525+
/** Stop already answered the turn: even a success after it settles the call as stopped. */
526+
await expect(
527+
completeDesktopTool.execute({
528+
principal: desktop.principal,
529+
input: {
530+
deviceId: desktop.deviceId,
531+
toolCallId: running,
532+
executionToken,
533+
status: 'success',
534+
},
535+
})
536+
).resolves.toEqual({ outcome: 'superseded', status: 'cancelled' })
537+
expect((await row(running)).status).toBe('cancelled')
534538
const after = await listDesktopInbox.execute({
535539
principal: desktop.principal,
536540
input: { deviceId: desktop.deviceId },
@@ -654,7 +658,6 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
654658
principal: desktop.principal,
655659
input: { deviceId: desktop.deviceId },
656660
})
657-
expect(inbox.hasActiveRun).toBe(true)
658661
expect(inbox.items.map((item) => item.kind === 'call' && item.toolCallId)).toEqual([
659662
first,
660663
second,
@@ -682,6 +685,28 @@ describe.runIf(Boolean(redisUrl))('desktop background executor protocol', () =>
682685
await expect.poll(() => isDesktopPresent(desktop.deviceId)).toBe(false)
683686
})
684687

688+
it('keeps the device online while a replacement stream overlaps the one it replaces', async () => {
689+
const desktop = await signedInDesktop()
690+
const open = () =>
691+
openDesktopInboxStream.execute({
692+
principal: desktop.principal,
693+
input: { deviceId: desktop.deviceId },
694+
})
695+
/** The stream that wrote presence last is the one that closes, while the other stays open. */
696+
const closeStaying = (await open()).subscribe(() => {})
697+
await expect.poll(() => isDesktopPresent(desktop.deviceId)).toBe(true)
698+
const closeLeaving = (await open()).subscribe(() => {})
699+
try {
700+
await sleep(200)
701+
closeLeaving()
702+
await sleep(200)
703+
expect(await isDesktopPresent(desktop.deviceId)).toBe(true)
704+
} finally {
705+
closeStaying()
706+
}
707+
await expect.poll(() => isDesktopPresent(desktop.deviceId)).toBe(false)
708+
})
709+
685710
it('ends a stream whose session was signed out at its next revalidation', async () => {
686711
const desktop = await signedInDesktop()
687712
const stream = await openDesktopInboxStream.execute({

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

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -12,8 +12,7 @@ import {
1212
DESKTOP_CALL_LEASE_SECONDS,
1313
DESKTOP_CALL_PICKUP_GRACE_MS,
1414
DESKTOP_EXECUTOR_PROTOCOL_VERSION,
15-
DESKTOP_INBOX_ACTIVE_RECONCILE_MS,
16-
DESKTOP_INBOX_IDLE_RECONCILE_MS,
15+
DESKTOP_INBOX_RECONCILE_MS,
1716
DESKTOP_PRESENCE_REFRESH_MS,
1817
} from '@/lib/desktop/executor/constants'
1918
import {
@@ -35,7 +34,6 @@ import {
3534
claimOfferedDesktopCall,
3635
getBoundDesktopCall,
3736
getBoundDesktopDevice,
38-
hasActiveDesktopRun,
3937
listDesktopInboxRows,
4038
recordDesktopCallResult,
4139
renewDesktopCallLease,
@@ -126,8 +124,7 @@ export const registerDesktopDevice = defineAuthorizedCredentialUserUseCase({
126124
leaseMs: DESKTOP_CALL_LEASE_SECONDS * 1000,
127125
leaseRenewMs: DESKTOP_CALL_LEASE_RENEW_MS,
128126
pickupGraceMs: DESKTOP_CALL_PICKUP_GRACE_MS,
129-
activeReconcileMs: DESKTOP_INBOX_ACTIVE_RECONCILE_MS,
130-
idleReconcileMs: DESKTOP_INBOX_IDLE_RECONCILE_MS,
127+
reconcileMs: DESKTOP_INBOX_RECONCILE_MS,
131128
}
132129
},
133130
})
@@ -146,15 +143,14 @@ export const listDesktopInbox = defineAuthorizedCredentialUserUseCase({
146143
}: {
147144
principal: SessionPrincipal
148145
input: DeviceInput
149-
}): Promise<{ items: DesktopInboxEntry[]; hasActiveRun: boolean }> {
146+
}): Promise<{ items: DesktopInboxEntry[] }> {
150147
await requireBoundDevice(principal, input.deviceId)
151148
const identity = { deviceId: input.deviceId, userId: principal.userId }
152-
const [rows, hasActiveRun] = await Promise.all([
149+
const [rows] = await Promise.all([
153150
listDesktopInboxRows(identity),
154-
hasActiveDesktopRun(identity),
155151
touchDesktopDevice(input.deviceId),
156152
])
157-
return { items: classifyDesktopInbox(rows), hasActiveRun }
153+
return { items: classifyDesktopInbox(rows) }
158154
},
159155
})
160156

@@ -320,6 +316,9 @@ export interface CompleteDesktopToolInput extends DesktopCallTokenInput {
320316
data?: unknown
321317
}
322318

319+
const STOPPED_WHILE_RUNNING_MESSAGE =
320+
'Stopped by the user while the Sim desktop app was running this action. It may already have taken effect; inspect the current state before repeating it.'
321+
323322
const COMPLETION_STATUS = {
324323
success: { durable: ASYNC_TOOL_STATUS.completed, message: 'Tool completed' },
325324
error: { durable: ASYNC_TOOL_STATUS.failed, message: 'Tool failed' },
@@ -396,6 +395,7 @@ export const completeDesktopTool = defineAuthorizedCredentialUserUseCase({
396395
toolCallId: call.toolCallId,
397396
runId: call.runId,
398397
ownerToken: input.executionToken,
398+
stoppedMessage: STOPPED_WHILE_RUNNING_MESSAGE,
399399
})
400400
if (acknowledged.outcome === 'unknown')
401401
throw new OrchestrationError('not_found', 'Desktop tool call not found')

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

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,10 @@ export const DESKTOP_CALL_LEASE_RENEW_MS = 20_000
1616
export const DESKTOP_CALL_PICKUP_GRACE_MS = 15_000
1717

1818
/**
19-
* The device's safety-net inbox pull while it has an active bound run. Shorter than the pickup
20-
* grace, so a lost doorbell still lets the device claim a call before it fails.
19+
* The device's safety-net inbox pull. Shorter than the pickup grace, so a lost doorbell still lets
20+
* the device claim a call before it fails, including the first call of a turn it has not seen.
2121
*/
22-
export const DESKTOP_INBOX_ACTIVE_RECONCILE_MS = 10_000
23-
24-
/** The device's safety-net inbox pull while none of its runs are active. */
25-
export const DESKTOP_INBOX_IDLE_RECONCILE_MS = 60_000
22+
export const DESKTOP_INBOX_RECONCILE_MS = 10_000
2623

2724
/** Presence outlives two missed refreshes, so a rolling deploy does not read as offline. */
2825
export const DESKTOP_PRESENCE_TTL_SECONDS = 45

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

Lines changed: 14 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -2,17 +2,10 @@ import { getRedisClient } from '@/lib/core/config/redis'
22
import { DESKTOP_PRESENCE_TTL_SECONDS } from '@/lib/desktop/executor/constants'
33

44
/**
5-
* A device is present while one of its inbox streams is open. Each stream writes its own
6-
* connection id, so a stream that closes after its replacement opened cannot erase the newer
7-
* stream's presence.
5+
* A device is present while one of its inbox streams is open. Each stream holds its own entry in
6+
* the device's set, scored by when that entry expires, so streams that overlap during a reconnect
7+
* or rotation never erase each other's presence. The key itself expires with its newest entry.
88
*/
9-
const RELEASE_OWN_PRESENCE = `
10-
if redis.call('GET', KEYS[1]) == ARGV[1] then
11-
return redis.call('DEL', KEYS[1])
12-
end
13-
return 0
14-
`
15-
169
function presenceKey(deviceId: string): string {
1710
return `desktop:presence:${deviceId}`
1811
}
@@ -26,22 +19,28 @@ export function isDesktopPresenceAvailable(): boolean {
2619
export async function markDesktopPresent(deviceId: string, connectionId: string): Promise<void> {
2720
const redis = getRedisClient()
2821
if (!redis) throw new Error('Desktop presence requires Redis')
29-
await redis.set(presenceKey(deviceId), connectionId, 'EX', DESKTOP_PRESENCE_TTL_SECONDS)
22+
const key = presenceKey(deviceId)
23+
await redis
24+
.multi()
25+
.zremrangebyscore(key, '-inf', Date.now())
26+
.zadd(key, Date.now() + DESKTOP_PRESENCE_TTL_SECONDS * 1000, connectionId)
27+
.expire(key, DESKTOP_PRESENCE_TTL_SECONDS)
28+
.exec()
3029
}
3130

32-
/** Clears presence only when this connection is still the one recorded. */
31+
/** Removes only this connection's entry; any other open stream keeps the device present. */
3332
export async function releaseDesktopPresence(
3433
deviceId: string,
3534
connectionId: string
3635
): Promise<void> {
3736
const redis = getRedisClient()
3837
if (!redis) return
39-
await redis.eval(RELEASE_OWN_PRESENCE, 1, presenceKey(deviceId), connectionId)
38+
await redis.zrem(presenceKey(deviceId), connectionId)
4039
}
4140

42-
/** A device with no open inbox stream cannot pick up a call; a missing Redis reads as absent. */
41+
/** A device with no unexpired stream entry cannot pick up a call; a missing Redis reads as absent. */
4342
export async function isDesktopPresent(deviceId: string): Promise<boolean> {
4443
const redis = getRedisClient()
4544
if (!redis) return false
46-
return (await redis.exists(presenceKey(deviceId))) === 1
45+
return (await redis.zcount(presenceKey(deviceId), Date.now(), '+inf')) > 0
4746
}

0 commit comments

Comments
 (0)