Skip to content

Commit d8356a1

Browse files
committed
feat(desktop): device registry, inbox and leased claims for a background executor
Adds the server half of the desktop background executor's device protocol, behind the mothership-desktop-background-executor flag (off by default): - desktop_devices binds an install id to a user and the Better Auth session that registered it; copilot_runs.desktop_device_id records which device a turn's desktop tools run on (nothing sets it yet). - POST /api/desktop/devices registers the device and returns the timing contract; GET /api/desktop/inbox is the reconciling pull (calls offered to the device, calls awaiting approval, and claimed calls to cancel); GET /api/desktop/inbox/stream is an SSE doorbell over Redis pub/sub that keeps a 45 s presence key while it is open. - POST /api/desktop/tool/claim, /lease and /complete claim an offered call with an execution token and a 60 s lease, renew it (410 once stopped or lapsed), and record the sealed result idempotently: a retry is a duplicate, a result for a call Sim settled first is superseded. Every operation requires the device to be bound to the caller's user and session, the run to be bound to that device, and the run's tool admission to be open.
1 parent 2613543 commit d8356a1

26 files changed

Lines changed: 32997 additions & 0 deletions

File tree

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
import { registerDesktopDeviceContract } from '@/lib/api/contracts/desktop-executor'
2+
import {
3+
defineInternalJsonRoute,
4+
internalRateLimits,
5+
internalSessionAuth,
6+
} from '@/lib/api/server/routes'
7+
import { desktopExecutorErrorPolicy } from '@/lib/api/server/routes/desktop-executor'
8+
import { registerDesktopDevice } from '@/lib/desktop/application/executor'
9+
10+
export const POST = defineInternalJsonRoute({
11+
contract: registerDesktopDeviceContract,
12+
auth: internalSessionAuth,
13+
operation: registerDesktopDevice.operation,
14+
rateLimit: internalRateLimits.user({ bucketName: 'desktop-device-register' }),
15+
errorPolicy: desktopExecutorErrorPolicy,
16+
mapInput: ({ body }) => body,
17+
useCase: registerDesktopDevice,
18+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
19+
})
Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
import { listDesktopInboxContract } from '@/lib/api/contracts/desktop-executor'
2+
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
3+
import {
4+
desktopExecutorErrorPolicy,
5+
desktopExecutorRateLimit,
6+
} from '@/lib/api/server/routes/desktop-executor'
7+
import { listDesktopInbox } from '@/lib/desktop/application/executor'
8+
9+
export const dynamic = 'force-dynamic'
10+
11+
export const GET = defineInternalJsonRoute({
12+
contract: listDesktopInboxContract,
13+
auth: internalSessionAuth,
14+
operation: listDesktopInbox.operation,
15+
rateLimit: desktopExecutorRateLimit,
16+
errorPolicy: desktopExecutorErrorPolicy,
17+
mapInput: ({ query }) => ({ deviceId: query.deviceId }),
18+
useCase: listDesktopInbox,
19+
present: ({ items, hasActiveRun }) => ({
20+
hasActiveRun,
21+
items: items.map((item) =>
22+
item.kind === 'call' ? { ...item, createdAt: item.createdAt.toISOString() } : item
23+
),
24+
}),
25+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
26+
})
Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
import { createLogger } from '@sim/logger'
2+
import type { NextRequest } from 'next/server'
3+
import { desktopInboxStreamContract } from '@/lib/api/contracts/desktop-executor'
4+
import { parseRequest } from '@/lib/api/server'
5+
import {
6+
InternalUnauthenticatedError,
7+
internalSessionAuth,
8+
} from '@/lib/api/server/routes/internal-json-route'
9+
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
10+
import { openDesktopInboxStream } from '@/lib/desktop/application/executor'
11+
import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors'
12+
import { createSSEStream } from '@/lib/events/sse-endpoint'
13+
14+
export const dynamic = 'force-dynamic'
15+
16+
const logger = createLogger('DesktopInboxStream')
17+
18+
/**
19+
* The background executor's doorbell. A raw route because it streams: authentication, the device
20+
* check, and presence all run through the `openDesktopInboxStream` use case.
21+
*/
22+
export const GET = withRouteHandler(async (request: NextRequest) => {
23+
try {
24+
const principal = await internalSessionAuth.authenticate()
25+
const parsed = await parseRequest(desktopInboxStreamContract, request, {})
26+
if (!parsed.success) return parsed.response
27+
const { deviceId } = parsed.data.query
28+
const inbox = await openDesktopInboxStream.execute({ principal, input: { deviceId } })
29+
return createSSEStream(request, {
30+
label: 'desktop-inbox',
31+
revalidate: inbox.revalidate,
32+
subscriptions: [{ subscribe: inbox.subscribe }],
33+
})
34+
} catch (error) {
35+
if (error instanceof InternalUnauthenticatedError)
36+
return new Response('Unauthorized', { status: 401 })
37+
if (error instanceof DesktopDeviceUnrecognizedError)
38+
return new Response(error.message, { status: 401 })
39+
logger.error('Failed to open the desktop inbox stream', error)
40+
return new Response('Unable to open the desktop inbox', { status: 500 })
41+
}
42+
})
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
import { claimDesktopToolContract } from '@/lib/api/contracts/desktop-executor'
2+
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
3+
import {
4+
desktopExecutorErrorPolicy,
5+
desktopExecutorRateLimit,
6+
} from '@/lib/api/server/routes/desktop-executor'
7+
import { claimDesktopTool } from '@/lib/desktop/application/executor'
8+
9+
export const POST = defineInternalJsonRoute({
10+
contract: claimDesktopToolContract,
11+
auth: internalSessionAuth,
12+
operation: claimDesktopTool.operation,
13+
rateLimit: desktopExecutorRateLimit,
14+
errorPolicy: desktopExecutorErrorPolicy,
15+
mapInput: ({ body }) => body,
16+
useCase: claimDesktopTool,
17+
present: (claim) => ({ ...claim, leaseExpiresAt: claim.leaseExpiresAt.toISOString() }),
18+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
19+
})
Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
import { completeDesktopToolContract } from '@/lib/api/contracts/desktop-executor'
2+
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
3+
import {
4+
desktopExecutorErrorPolicy,
5+
desktopExecutorRateLimit,
6+
} from '@/lib/api/server/routes/desktop-executor'
7+
import { completeDesktopTool } from '@/lib/desktop/application/executor'
8+
9+
export const POST = defineInternalJsonRoute({
10+
contract: completeDesktopToolContract,
11+
auth: internalSessionAuth,
12+
operation: completeDesktopTool.operation,
13+
rateLimit: desktopExecutorRateLimit,
14+
errorPolicy: desktopExecutorErrorPolicy,
15+
mapInput: ({ body }) => body,
16+
useCase: completeDesktopTool,
17+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
18+
})
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
import { renewDesktopToolLeaseContract } from '@/lib/api/contracts/desktop-executor'
2+
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
3+
import {
4+
desktopExecutorErrorPolicy,
5+
desktopExecutorRateLimit,
6+
} from '@/lib/api/server/routes/desktop-executor'
7+
import { renewDesktopToolLease } from '@/lib/desktop/application/executor'
8+
9+
export const POST = defineInternalJsonRoute({
10+
contract: renewDesktopToolLeaseContract,
11+
auth: internalSessionAuth,
12+
operation: renewDesktopToolLease.operation,
13+
rateLimit: desktopExecutorRateLimit,
14+
errorPolicy: desktopExecutorErrorPolicy,
15+
mapInput: ({ body }) => body,
16+
useCase: renewDesktopToolLease,
17+
present: ({ leaseExpiresAt }) => ({ leaseExpiresAt: leaseExpiresAt.toISOString() }),
18+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
19+
})
Lines changed: 183 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,183 @@
1+
import { z } from 'zod'
2+
import { desktopToolCallIdSchema } from '@/lib/api/contracts/desktop-tool-authorization'
3+
import { defineRouteContract } from '@/lib/api/contracts/types'
4+
5+
/**
6+
* The desktop background executor's wire protocol. The Electron main process calls every route
7+
* here with the cookie of the Better Auth session it registered under; Sim refuses any other
8+
* session for that device.
9+
*/
10+
11+
/** The install id the desktop app generates once and keeps in its user data. */
12+
const desktopDeviceIdSchema = z.string().uuid('Device ID must be a UUID')
13+
14+
const desktopDeviceCapabilitiesSchema = z.object({
15+
executor: z.number().int().min(1).max(1000),
16+
browser: z.boolean(),
17+
terminal: z.boolean(),
18+
localFiles: z.boolean(),
19+
})
20+
21+
const registerDesktopDeviceBodySchema = z.object({
22+
deviceId: desktopDeviceIdSchema,
23+
name: z.string().trim().min(1, 'Device name is required').max(128, 'Device name is too long'),
24+
appVersion: z.string().trim().min(1, 'App version is required').max(64),
25+
platform: z.string().trim().min(1, 'Platform is required').max(64),
26+
capabilities: desktopDeviceCapabilitiesSchema,
27+
})
28+
export type RegisterDesktopDeviceBody = z.input<typeof registerDesktopDeviceBodySchema>
29+
30+
/**
31+
* `enabled: false` is the kill switch: the device keeps finishing calls on runs already bound to
32+
* it, but no new turn binds to it.
33+
*/
34+
export const registerDesktopDeviceResponseSchema = z.object({
35+
enabled: z.boolean(),
36+
protocolVersion: z.number().int().min(1),
37+
leaseMs: z.number().int().positive(),
38+
leaseRenewMs: z.number().int().positive(),
39+
pickupGraceMs: z.number().int().positive(),
40+
activeReconcileMs: z.number().int().positive(),
41+
idleReconcileMs: z.number().int().positive(),
42+
})
43+
44+
export const registerDesktopDeviceContract = defineRouteContract({
45+
method: 'POST',
46+
path: '/api/desktop/devices',
47+
body: registerDesktopDeviceBodySchema,
48+
response: { mode: 'json', schema: registerDesktopDeviceResponseSchema },
49+
error: z.object({ error: z.string() }),
50+
})
51+
52+
const desktopInboxQuerySchema = z.object({ deviceId: desktopDeviceIdSchema })
53+
54+
const desktopInboxCallSchema = z.object({
55+
kind: z.literal('call'),
56+
toolCallId: z.string().min(1),
57+
toolName: z.string().min(1),
58+
chatId: z.string().min(1),
59+
workspaceId: z.string().nullable(),
60+
createdAt: z.string(),
61+
})
62+
63+
const desktopInboxApprovalSchema = z.object({
64+
kind: z.literal('approval_needed'),
65+
toolCallId: z.string().min(1),
66+
toolName: z.string().min(1),
67+
chatId: z.string().min(1),
68+
chatTitle: z.string().nullable(),
69+
workspaceId: z.string().nullable(),
70+
/** A short description of what the call would do, such as the command to run. */
71+
summary: z.string().nullable(),
72+
})
73+
74+
const desktopInboxCancelSchema = z.object({
75+
kind: z.literal('cancel'),
76+
toolCallId: z.string().min(1),
77+
})
78+
79+
/** Ordered by when each call was persisted, which is the order the executor claims them in. */
80+
const desktopInboxItemSchema = z.discriminatedUnion('kind', [
81+
desktopInboxCallSchema,
82+
desktopInboxApprovalSchema,
83+
desktopInboxCancelSchema,
84+
])
85+
export type DesktopInboxItem = z.output<typeof desktopInboxItemSchema>
86+
87+
export const desktopInboxResponseSchema = z.object({
88+
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(),
91+
})
92+
93+
export const listDesktopInboxContract = defineRouteContract({
94+
method: 'GET',
95+
path: '/api/desktop/inbox',
96+
query: desktopInboxQuerySchema,
97+
response: { mode: 'json', schema: desktopInboxResponseSchema },
98+
error: z.object({ error: z.string() }),
99+
})
100+
101+
/**
102+
* Server-sent events: `inbox_changed` with `{ reason: 'call' | 'approval' | 'cancel' }`, plus
103+
* heartbeat comments and a `rotate` event before the server closes a long-lived stream. Events
104+
* are hints; the device re-reads `GET /api/desktop/inbox` on each one. While the stream is open
105+
* the device counts as online.
106+
*/
107+
export const desktopInboxStreamContract = defineRouteContract({
108+
method: 'GET',
109+
path: '/api/desktop/inbox/stream',
110+
query: desktopInboxQuerySchema,
111+
response: { mode: 'stream' },
112+
})
113+
114+
const claimDesktopToolBodySchema = z.object({
115+
deviceId: desktopDeviceIdSchema,
116+
toolCallId: desktopToolCallIdSchema,
117+
})
118+
export type ClaimDesktopToolBody = z.input<typeof claimDesktopToolBodySchema>
119+
120+
export const claimDesktopToolResponseSchema = z.object({
121+
toolName: z.string().min(1),
122+
args: z.record(z.string(), z.unknown()),
123+
chatId: z.string().min(1),
124+
workspaceId: z.string().nullable(),
125+
/** Presented on every lease renewal and on completion; it fences out stale owners. */
126+
executionToken: z.string().min(1),
127+
leaseExpiresAt: z.string(),
128+
})
129+
130+
export const claimDesktopToolContract = defineRouteContract({
131+
method: 'POST',
132+
path: '/api/desktop/tool/claim',
133+
body: claimDesktopToolBodySchema,
134+
response: { mode: 'json', schema: claimDesktopToolResponseSchema },
135+
error: z.object({ error: z.string() }),
136+
})
137+
138+
const renewDesktopToolLeaseBodySchema = z.object({
139+
deviceId: desktopDeviceIdSchema,
140+
toolCallId: desktopToolCallIdSchema,
141+
executionToken: z.string().min(1).max(128),
142+
})
143+
export type RenewDesktopToolLeaseBody = z.input<typeof renewDesktopToolLeaseBodySchema>
144+
145+
export const renewDesktopToolLeaseResponseSchema = z.object({ leaseExpiresAt: z.string() })
146+
147+
/** A 410 means the call was stopped, settled, or its lease lapsed: cancel the local action. */
148+
export const renewDesktopToolLeaseContract = defineRouteContract({
149+
method: 'POST',
150+
path: '/api/desktop/tool/lease',
151+
body: renewDesktopToolLeaseBodySchema,
152+
response: { mode: 'json', schema: renewDesktopToolLeaseResponseSchema },
153+
error: z.object({ error: z.string() }),
154+
})
155+
156+
const completeDesktopToolBodySchema = z.object({
157+
deviceId: desktopDeviceIdSchema,
158+
toolCallId: desktopToolCallIdSchema,
159+
executionToken: z.string().min(1).max(128),
160+
status: z.enum(['success', 'error', 'cancelled']),
161+
message: z.string().max(10_000).optional(),
162+
data: z.unknown().optional(),
163+
})
164+
export type CompleteDesktopToolBody = z.input<typeof completeDesktopToolBodySchema>
165+
166+
/**
167+
* Every 200 is an acknowledgement the device's outbox can drop the result on. `recorded`: this
168+
* request settled the call. `duplicate`: an earlier request with the same token already did.
169+
* `superseded`: Sim settled the call first (Stop, or a lapsed lease), and `status` is what the
170+
* model received.
171+
*/
172+
export const completeDesktopToolResponseSchema = z.object({
173+
outcome: z.enum(['recorded', 'duplicate', 'superseded']),
174+
status: z.enum(['completed', 'failed', 'cancelled']),
175+
})
176+
177+
export const completeDesktopToolContract = defineRouteContract({
178+
method: 'POST',
179+
path: '/api/desktop/tool/complete',
180+
body: completeDesktopToolBodySchema,
181+
response: { mode: 'json', schema: completeDesktopToolResponseSchema },
182+
error: z.object({ error: z.string() }),
183+
})
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
import {
2+
extendInternalErrorPolicy,
3+
internalErrorResponse,
4+
internalOrchestrationErrorPolicy,
5+
internalRateLimits,
6+
} from '@/lib/api/server/routes/internal-json-route'
7+
import {
8+
DesktopCallRevokedError,
9+
DesktopDeviceUnrecognizedError,
10+
} from '@/lib/desktop/executor/errors'
11+
12+
/**
13+
* A device's unrecognized binding is a 401, so it registers again; a call taken away from its
14+
* token is a 410, so it stops the local action.
15+
*/
16+
export const desktopExecutorErrorPolicy = extendInternalErrorPolicy(
17+
internalOrchestrationErrorPolicy,
18+
(error) => {
19+
if (error instanceof DesktopDeviceUnrecognizedError)
20+
return internalErrorResponse(401, { error: error.message })
21+
if (error instanceof DesktopCallRevokedError)
22+
return internalErrorResponse(410, { error: error.message })
23+
return null
24+
}
25+
)
26+
27+
/**
28+
* One busy device renews a lease per running call every 20 s and pulls its inbox on every
29+
* doorbell, so the bucket allows a sustained 10 requests a second per user.
30+
*/
31+
export const desktopExecutorRateLimit = internalRateLimits.user({
32+
bucketName: 'desktop-executor',
33+
config: { maxTokens: 600, refillRate: 600, refillIntervalMs: 60_000 },
34+
})

‎apps/sim/lib/core/config/env.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -601,6 +601,7 @@ export const env = createEnv({
601601
DASHBOARDS: z.boolean().optional(),
602602
PROJECT_API_ENABLED: z.boolean().optional(), // Fallback for the `projects` feature flag off AppConfig
603603
MSHIP_MODEL_SELECTOR: z.boolean().optional(),
604+
MSHIP_DESKTOP_BACKGROUND_EXECUTOR: z.boolean().optional(), // Fallback for the `mothership-desktop-background-executor` feature flag off AppConfig
604605
INBOX_ENABLED: z.boolean().optional(), // Enable inbox (Sim Mailer) on self-hosted (bypasses hosted requirements)
605606
SANDBOXES_ENABLED: z.boolean().optional(), // Enable custom sandboxes on self-hosted (bypasses hosted requirements)
606607

‎apps/sim/lib/core/config/feature-flags.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ setEnv({
4444
TABLE_ROW_TTL: undefined,
4545
MSHIP_MODEL_SELECTOR: undefined,
4646
MSHIP_PLAN_MODE: undefined,
47+
MSHIP_DESKTOP_BACKGROUND_EXECUTOR: undefined,
4748
AGENT_MEMORY_HISTORY: undefined,
4849
CREDENTIAL_GROUPS: undefined,
4950
KNOWLEDGE_MEMBER_ACCESS: undefined,

0 commit comments

Comments
 (0)