Skip to content

Commit 48b2051

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 claimToolExecution, renewSimToolExecutionLease and the execution-owner fence. A background executor's claim is the desktop claim plus an execution lease, only on a run bound to that device; 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 46cd710 commit 48b2051

33 files changed

Lines changed: 33144 additions & 96 deletions

File tree

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

Lines changed: 13 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -19,18 +19,13 @@ import {
1919
isWorkflowToolExecutionClaimable,
2020
} from '@/lib/mothership/async-runs/lifecycle'
2121
import {
22-
completeAsyncToolCall,
23-
completeClaimedAsyncToolCall,
24-
completePendingAsyncToolCall,
25-
detachAsyncToolCall,
2622
getAsyncToolCall,
2723
getClaimedWorkflowExecutionId,
2824
getRunSegment,
2925
} from '@/lib/mothership/async-runs/repository'
3026
import { CopilotConfirmOutcome } from '@/lib/mothership/generated/trace-attribute-values-v1'
3127
import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1'
3228
import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1'
33-
import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confirm'
3429
import {
3530
authenticateCopilotRequestSessionOnly,
3631
createInternalServerErrorResponse,
@@ -40,6 +35,11 @@ import {
4035
} from '@/lib/mothership/request/http'
4136
import { withIncomingGoSpan } from '@/lib/mothership/request/otel'
4237
import { sealClientToolSettlement } from '@/lib/mothership/request/tools/client-completion-seal.server'
38+
import {
39+
type ClientToolSettlementGuard,
40+
clientToolCompletionMessage,
41+
settleClientToolCall,
42+
} from '@/lib/mothership/request/tools/client-settlement.server'
4343
import { isWorkflowToolName } from '@/lib/mothership/tools/client-executed-tools'
4444
import { getDesktopToolClaimOwner } from '@/lib/mothership/tools/desktop-tools'
4545
import {
@@ -60,20 +60,6 @@ const NATIVE_HANDOFF_INTERRUPTED_MESSAGE =
6060

6161
type ToolCallStatusUpdateOutcome = 'updated' | 'conflict' | 'failed'
6262

63-
interface UpdateToolCallStatusOptions {
64-
executionId?: string
65-
completionGuard?:
66-
| { status: typeof ASYNC_TOOL_STATUS.pending }
67-
| { status: typeof ASYNC_TOOL_STATUS.running; claimedBy: string }
68-
}
69-
70-
function getClientToolCompletionMessage(status: AsyncConfirmationStatus): string {
71-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.success) return 'Tool completed'
72-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) return 'Tool is running in background'
73-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled) return 'Tool cancelled'
74-
return 'Tool failed'
75-
}
76-
7763
function createConfirmationResponse(
7864
toolCallId: string,
7965
status: AsyncConfirmationStatus,
@@ -106,53 +92,18 @@ async function updateToolCallStatus(
10692
status: AsyncConfirmationStatus,
10793
message?: string,
10894
data?: AsyncCompletionData,
109-
options: UpdateToolCallStatusOptions = {}
95+
options: { executionId?: string; guard?: ClientToolSettlementGuard } = {}
11096
): Promise<ToolCallStatusUpdateOutcome> {
11197
const toolCallId = existing.toolCallId
11298
try {
113-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) {
114-
const detached = options.executionId
115-
? await detachAsyncToolCall(toolCallId, { preserveClaim: true })
116-
: await detachAsyncToolCall(toolCallId)
117-
if (!detached) return 'conflict'
118-
publishToolConfirmation({
119-
toolCallId,
120-
status,
121-
message: message || undefined,
122-
timestamp: new Date().toISOString(),
123-
data,
124-
...(options.executionId ? { executionId: options.executionId } : {}),
125-
})
126-
return 'updated'
127-
}
128-
const durableStatus =
129-
status === ASYNC_TOOL_CONFIRMATION_STATUS.success
130-
? ASYNC_TOOL_STATUS.completed
131-
: status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled
132-
? ASYNC_TOOL_STATUS.cancelled
133-
: ASYNC_TOOL_STATUS.failed
134-
const completionInput = {
135-
toolCallId,
136-
status: durableStatus,
137-
result: data ?? null,
138-
error: status === 'success' ? null : message || status,
139-
}
140-
const completed =
141-
options.completionGuard?.status === ASYNC_TOOL_STATUS.pending
142-
? await completePendingAsyncToolCall(completionInput)
143-
: options.completionGuard?.status === ASYNC_TOOL_STATUS.running
144-
? await completeClaimedAsyncToolCall(completionInput, options.completionGuard.claimedBy)
145-
: await completeAsyncToolCall(completionInput)
146-
if (!completed) return 'conflict'
147-
publishToolConfirmation({
99+
return await settleClientToolCall({
148100
toolCallId,
149101
status,
150-
message: message || undefined,
151-
timestamp: new Date().toISOString(),
102+
message: message ?? '',
152103
data,
153-
...(options.executionId ? { executionId: options.executionId } : {}),
104+
executionId: options.executionId,
105+
guard: options.guard ?? { kind: 'open' },
154106
})
155-
return 'updated'
156107
} catch (error) {
157108
logger.error('Failed to update tool call status', {
158109
toolCallId,
@@ -413,7 +364,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
413364
),
414365
}
415366
: {
416-
message: getClientToolCompletionMessage(status),
367+
message: clientToolCompletionMessage(status),
417368
data: await sealClientToolSettlement(existing.result, {
418369
toolCallId,
419370
runId: existing.runId,
@@ -441,9 +392,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
441392
projected.data,
442393
{
443394
...(isWorkflowTool && executionId ? { executionId } : {}),
444-
...(isPreclaimNativeTerminalOutcome
445-
? { completionGuard: { status: ASYNC_TOOL_STATUS.pending } as const }
446-
: {}),
395+
...(isPreclaimNativeTerminalOutcome ? { guard: { kind: 'pending' } as const } : {}),
447396
}
448397
)
449398

@@ -454,12 +403,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
454403
ASYNC_TOOL_CONFIRMATION_STATUS.error,
455404
projected.message,
456405
projected.data,
457-
{
458-
completionGuard: {
459-
status: ASYNC_TOOL_STATUS.running,
460-
claimedBy: nativeClaimOwner,
461-
},
462-
}
406+
{ guard: { kind: 'claimed', claimedBy: nativeClaimOwner } }
463407
)
464408
: updateOutcome
465409

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: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
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+
present: ({ outcome, status }) => ({ outcome, status }),
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 { 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)