Skip to content

Commit e0e3fb7

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), built on the existing tool-execution primitives rather than new ones: - 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). - Claims, renewals and results reuse claimSimToolExecution, renewSimToolExecutionLease and the execution-owner fence, extended for a device: only a pending call on a run bound to that device, that the user is not still being asked about, with the per-surface claim owner; a result is accepted after its lease lapsed until something else settles the call. - /api/copilot/confirm and the device's complete share one sealing and settlement function, so a device result is sealed, settled and published exactly like a chat view's. - POST /api/desktop/devices registers; GET /api/desktop/inbox lists calls, approval_needed items (from permission_requested_at) and cancel items; GET /api/desktop/inbox/stream is an SSE doorbell over Redis pub/sub. - Presence is driven by the device: every pull, lease renewal and stream open refreshes a 45 s key, and only its TTL removes it.
1 parent 8e796d1 commit e0e3fb7

31 files changed

Lines changed: 33085 additions & 118 deletions

File tree

‎apps/sim/app/api/copilot/confirm/route.ts‎

Lines changed: 29 additions & 93 deletions
Original file line numberDiff line numberDiff line change
@@ -20,18 +20,13 @@ import {
2020
isWorkflowToolExecutionClaimable,
2121
} from '@/lib/mothership/async-runs/lifecycle'
2222
import {
23-
completeAsyncToolCall,
24-
completeClaimedAsyncToolCall,
25-
completePendingAsyncToolCall,
26-
detachAsyncToolCall,
2723
getAsyncToolCall,
2824
getClaimedWorkflowExecutionId,
2925
getRunSegment,
3026
} from '@/lib/mothership/async-runs/repository'
3127
import { CopilotConfirmOutcome } from '@/lib/mothership/generated/trace-attribute-values-v1'
3228
import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1'
3329
import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1'
34-
import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confirm'
3530
import {
3631
authenticateCopilotRequestSessionOnly,
3732
createInternalServerErrorResponse,
@@ -41,9 +36,11 @@ import {
4136
} from '@/lib/mothership/request/http'
4237
import { withIncomingGoSpan } from '@/lib/mothership/request/otel'
4338
import {
44-
retainSealedClientToolContext,
45-
sealClientToolCompletion,
46-
} from '@/lib/mothership/request/tools/client-completion-seal.server'
39+
type ClientToolSettlementGuard,
40+
clientToolCompletionMessage,
41+
sealClientToolResult,
42+
settleClientToolCall,
43+
} from '@/lib/mothership/request/tools/client-settlement.server'
4744
import { isWorkflowToolName } from '@/lib/mothership/tools/client-executed-tools'
4845
import {
4946
createStructuralWorkflowToolCompletionData,
@@ -63,20 +60,6 @@ const NATIVE_HANDOFF_INTERRUPTED_MESSAGE =
6360

6461
type ToolCallStatusUpdateOutcome = 'updated' | 'conflict' | 'failed'
6562

66-
interface UpdateToolCallStatusOptions {
67-
executionId?: string
68-
completionGuard?:
69-
| { status: typeof ASYNC_TOOL_STATUS.pending }
70-
| { status: typeof ASYNC_TOOL_STATUS.running; claimedBy: string }
71-
}
72-
73-
function getClientToolCompletionMessage(status: AsyncConfirmationStatus): string {
74-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.success) return 'Tool completed'
75-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) return 'Tool is running in background'
76-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled) return 'Tool cancelled'
77-
return 'Tool failed'
78-
}
79-
8063
function createConfirmationResponse(
8164
toolCallId: string,
8265
status: AsyncConfirmationStatus,
@@ -109,53 +92,18 @@ async function updateToolCallStatus(
10992
status: AsyncConfirmationStatus,
11093
message?: string,
11194
data?: AsyncCompletionData,
112-
options: UpdateToolCallStatusOptions = {}
95+
options: { executionId?: string; guard?: ClientToolSettlementGuard } = {}
11396
): Promise<ToolCallStatusUpdateOutcome> {
11497
const toolCallId = existing.toolCallId
11598
try {
116-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) {
117-
const detached = options.executionId
118-
? await detachAsyncToolCall(toolCallId, { preserveClaim: true })
119-
: await detachAsyncToolCall(toolCallId)
120-
if (!detached) return 'conflict'
121-
publishToolConfirmation({
122-
toolCallId,
123-
status,
124-
message: message || undefined,
125-
timestamp: new Date().toISOString(),
126-
data,
127-
...(options.executionId ? { executionId: options.executionId } : {}),
128-
})
129-
return 'updated'
130-
}
131-
const durableStatus =
132-
status === ASYNC_TOOL_CONFIRMATION_STATUS.success
133-
? ASYNC_TOOL_STATUS.completed
134-
: status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled
135-
? ASYNC_TOOL_STATUS.cancelled
136-
: ASYNC_TOOL_STATUS.failed
137-
const completionInput = {
138-
toolCallId,
139-
status: durableStatus,
140-
result: data ?? null,
141-
error: status === 'success' ? null : message || status,
142-
}
143-
const completed =
144-
options.completionGuard?.status === ASYNC_TOOL_STATUS.pending
145-
? await completePendingAsyncToolCall(completionInput)
146-
: options.completionGuard?.status === ASYNC_TOOL_STATUS.running
147-
? await completeClaimedAsyncToolCall(completionInput, options.completionGuard.claimedBy)
148-
: await completeAsyncToolCall(completionInput)
149-
if (!completed) return 'conflict'
150-
publishToolConfirmation({
99+
return await settleClientToolCall({
151100
toolCallId,
152101
status,
153-
message: message || undefined,
154-
timestamp: new Date().toISOString(),
102+
message: message ?? '',
155103
data,
156-
...(options.executionId ? { executionId: options.executionId } : {}),
104+
executionId: options.executionId,
105+
guard: options.guard ?? { kind: 'open' },
157106
})
158-
return 'updated'
159107
} catch (error) {
160108
logger.error('Failed to update tool call status', {
161109
toolCallId,
@@ -416,28 +364,23 @@ export const POST = withRouteHandler((req: NextRequest) => {
416364
),
417365
}
418366
: {
419-
message: getClientToolCompletionMessage(status),
420-
data: {
421-
...retainSealedClientToolContext(existing.result),
422-
...(await sealClientToolCompletion({
423-
toolCallId,
424-
runId: existing.runId,
425-
userId: authenticatedUserId,
426-
...(isIndeterminateNativeExit
427-
? {
428-
message: NATIVE_HANDOFF_INTERRUPTED_MESSAGE,
429-
data: {
430-
error: NATIVE_HANDOFF_INTERRUPTED_MESSAGE,
431-
outcomeUnknown: true,
432-
doNotRetry: true,
433-
},
434-
}
435-
: {
436-
...(message !== undefined ? { message } : {}),
437-
...(data !== undefined ? { data } : {}),
438-
}),
439-
})),
440-
},
367+
message: clientToolCompletionMessage(status),
368+
data: await sealClientToolResult({
369+
toolCallId,
370+
runId: existing.runId,
371+
userId: authenticatedUserId,
372+
storedResult: existing.result,
373+
...(isIndeterminateNativeExit
374+
? {
375+
message: NATIVE_HANDOFF_INTERRUPTED_MESSAGE,
376+
data: {
377+
error: NATIVE_HANDOFF_INTERRUPTED_MESSAGE,
378+
outcomeUnknown: true,
379+
doNotRetry: true,
380+
},
381+
}
382+
: { message, data }),
383+
}),
441384
}
442385

443386
const updateOutcome = await updateToolCallStatus(
@@ -447,9 +390,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
447390
projected.data,
448391
{
449392
...(isWorkflowTool && executionId ? { executionId } : {}),
450-
...(isPreclaimNativeTerminalOutcome
451-
? { completionGuard: { status: ASYNC_TOOL_STATUS.pending } as const }
452-
: {}),
393+
...(isPreclaimNativeTerminalOutcome ? { guard: { kind: 'pending' } as const } : {}),
453394
}
454395
)
455396

@@ -460,12 +401,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
460401
ASYNC_TOOL_CONFIRMATION_STATUS.error,
461402
projected.message,
462403
projected.data,
463-
{
464-
completionGuard: {
465-
status: ASYNC_TOOL_STATUS.running,
466-
claimedBy: nativeClaimOwner,
467-
},
468-
}
404+
{ guard: { kind: 'claimed', claimedBy: nativeClaimOwner } }
469405
)
470406
: updateOutcome
471407

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: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
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 }) => ({
20+
items: items.map((item) =>
21+
item.kind === 'call' ? { ...item, createdAt: item.createdAt.toISOString() } : item
22+
),
23+
}),
24+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
25+
})
Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
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 { desktopExecutorRateLimit } from '@/lib/api/server/routes/desktop-executor'
6+
import {
7+
InternalUnauthenticatedError,
8+
internalSessionAuth,
9+
} from '@/lib/api/server/routes/internal-json-route'
10+
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
11+
import { openDesktopInboxStream } from '@/lib/desktop/application/executor'
12+
import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors'
13+
import { createSSEStream } from '@/lib/events/sse-endpoint'
14+
15+
export const dynamic = 'force-dynamic'
16+
17+
const logger = createLogger('DesktopInboxStream')
18+
19+
/**
20+
* The background executor's doorbell. A raw route because it streams: authentication, the device
21+
* check, and presence all run through the `openDesktopInboxStream` use case.
22+
*/
23+
export const GET = withRouteHandler(async (request: NextRequest) => {
24+
try {
25+
const principal = await internalSessionAuth.authenticate()
26+
const limited = await desktopExecutorRateLimit.enforce(request, principal)
27+
if (limited) return limited
28+
const parsed = await parseRequest(desktopInboxStreamContract, request, {})
29+
if (!parsed.success) return parsed.response
30+
const { deviceId } = parsed.data.query
31+
const inbox = await openDesktopInboxStream.execute({ principal, input: { deviceId } })
32+
return createSSEStream(request, {
33+
label: 'desktop-inbox',
34+
revalidate: inbox.revalidate,
35+
subscriptions: [{ subscribe: inbox.subscribe }],
36+
})
37+
} catch (error) {
38+
if (error instanceof InternalUnauthenticatedError)
39+
return new Response('Unauthorized', { status: 401 })
40+
if (error instanceof DesktopDeviceUnrecognizedError)
41+
return new Response(error.message, { status: 401 })
42+
logger.error('Failed to open the desktop inbox stream', error)
43+
return new Response('Unable to open the desktop inbox', { status: 500 })
44+
}
45+
})
Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
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+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
18+
})
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: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
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+
staticResponseHeaders: { 'Cache-Control': 'no-store' },
18+
})

0 commit comments

Comments
 (0)