Skip to content

Commit 4cebfa2

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 claimDesktopToolCall, 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 9af89fc commit 4cebfa2

33 files changed

Lines changed: 33215 additions & 85 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
@@ -18,18 +18,13 @@ import {
1818
isWorkflowToolExecutionClaimable,
1919
} from '@/lib/mothership/async-runs/lifecycle'
2020
import {
21-
completeAsyncToolCall,
22-
completeClaimedAsyncToolCall,
23-
completePendingAsyncToolCall,
24-
detachAsyncToolCall,
2521
getAsyncToolCall,
2622
getClaimedWorkflowExecutionId,
2723
getRunSegment,
2824
} from '@/lib/mothership/async-runs/repository'
2925
import { CopilotConfirmOutcome } from '@/lib/mothership/generated/trace-attribute-values-v1'
3026
import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1'
3127
import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1'
32-
import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confirm'
3328
import {
3429
authenticateCopilotRequestSessionOnly,
3530
createInternalServerErrorResponse,
@@ -39,6 +34,11 @@ import {
3934
} from '@/lib/mothership/request/http'
4035
import { withIncomingGoSpan } from '@/lib/mothership/request/otel'
4136
import { sealClientToolSettlement } from '@/lib/mothership/request/tools/client-completion-seal.server'
37+
import {
38+
type ClientToolSettlementGuard,
39+
clientToolCompletionMessage,
40+
settleClientToolCall,
41+
} from '@/lib/mothership/request/tools/client-settlement.server'
4242
import { isWorkflowToolName } from '@/lib/mothership/tools/client-executed-tools'
4343
import { getDesktopToolClaimOwner, isNativeDesktopTool } from '@/lib/mothership/tools/desktop-tools'
4444
import {
@@ -58,20 +58,6 @@ const NATIVE_HANDOFF_INTERRUPTED_MESSAGE =
5858

5959
type ToolCallStatusUpdateOutcome = 'updated' | 'conflict' | 'failed'
6060

61-
interface UpdateToolCallStatusOptions {
62-
executionId?: string
63-
completionGuard?:
64-
| { status: typeof ASYNC_TOOL_STATUS.pending }
65-
| { status: typeof ASYNC_TOOL_STATUS.running; claimedBy: string }
66-
}
67-
68-
function getClientToolCompletionMessage(status: AsyncConfirmationStatus): string {
69-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.success) return 'Tool completed'
70-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) return 'Tool is running in background'
71-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled) return 'Tool cancelled'
72-
return 'Tool failed'
73-
}
74-
7561
function createConfirmationResponse(
7662
toolCallId: string,
7763
status: AsyncConfirmationStatus,
@@ -104,53 +90,18 @@ async function updateToolCallStatus(
10490
status: AsyncConfirmationStatus,
10591
message?: string,
10692
data?: AsyncCompletionData,
107-
options: UpdateToolCallStatusOptions = {}
93+
options: { executionId?: string; guard?: ClientToolSettlementGuard } = {}
10894
): Promise<ToolCallStatusUpdateOutcome> {
10995
const toolCallId = existing.toolCallId
11096
try {
111-
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) {
112-
const detached = options.executionId
113-
? await detachAsyncToolCall(toolCallId, { preserveClaim: true })
114-
: await detachAsyncToolCall(toolCallId)
115-
if (!detached) return 'conflict'
116-
publishToolConfirmation({
117-
toolCallId,
118-
status,
119-
message: message || undefined,
120-
timestamp: new Date().toISOString(),
121-
data,
122-
...(options.executionId ? { executionId: options.executionId } : {}),
123-
})
124-
return 'updated'
125-
}
126-
const durableStatus =
127-
status === ASYNC_TOOL_CONFIRMATION_STATUS.success
128-
? ASYNC_TOOL_STATUS.completed
129-
: status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled
130-
? ASYNC_TOOL_STATUS.cancelled
131-
: ASYNC_TOOL_STATUS.failed
132-
const completionInput = {
133-
toolCallId,
134-
status: durableStatus,
135-
result: data ?? null,
136-
error: status === 'success' ? null : message || status,
137-
}
138-
const completed =
139-
options.completionGuard?.status === ASYNC_TOOL_STATUS.pending
140-
? await completePendingAsyncToolCall(completionInput)
141-
: options.completionGuard?.status === ASYNC_TOOL_STATUS.running
142-
? await completeClaimedAsyncToolCall(completionInput, options.completionGuard.claimedBy)
143-
: await completeAsyncToolCall(completionInput)
144-
if (!completed) return 'conflict'
145-
publishToolConfirmation({
97+
return await settleClientToolCall({
14698
toolCallId,
14799
status,
148-
message: message || undefined,
149-
timestamp: new Date().toISOString(),
100+
message: message ?? '',
150101
data,
151-
...(options.executionId ? { executionId: options.executionId } : {}),
102+
executionId: options.executionId,
103+
guard: options.guard ?? { kind: 'open' },
152104
})
153-
return 'updated'
154105
} catch (error) {
155106
logger.error('Failed to update tool call status', {
156107
toolCallId,
@@ -419,7 +370,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
419370
),
420371
}
421372
: {
422-
message: getClientToolCompletionMessage(status),
373+
message: clientToolCompletionMessage(status),
423374
data: await sealClientToolSettlement(existing.result, {
424375
toolCallId,
425376
runId: existing.runId,
@@ -447,9 +398,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
447398
projected.data,
448399
{
449400
...(isWorkflowTool && executionId ? { executionId } : {}),
450-
...(isPreclaimNativeTerminalOutcome
451-
? { completionGuard: { status: ASYNC_TOOL_STATUS.pending } as const }
452-
: {}),
401+
...(isPreclaimNativeTerminalOutcome ? { guard: { kind: 'pending' } as const } : {}),
453402
}
454403
)
455404

@@ -460,12 +409,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
460409
ASYNC_TOOL_CONFIRMATION_STATUS.error,
461410
projected.message,
462411
projected.data,
463-
{
464-
completionGuard: {
465-
status: ASYNC_TOOL_STATUS.running,
466-
claimedBy: nativeClaimOwner,
467-
},
468-
}
412+
{ guard: { kind: 'claimed', claimedBy: nativeClaimOwner } }
469413
)
470414
: updateOutcome
471415

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)