Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
82 changes: 13 additions & 69 deletions apps/sim/app/api/copilot/confirm/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,18 +18,13 @@ import {
isWorkflowToolExecutionClaimable,
} from '@/lib/mothership/async-runs/lifecycle'
import {
completeAsyncToolCall,
completeClaimedAsyncToolCall,
completePendingAsyncToolCall,
detachAsyncToolCall,
getAsyncToolCall,
getClaimedWorkflowExecutionId,
getRunSegment,
} from '@/lib/mothership/async-runs/repository'
import { CopilotConfirmOutcome } from '@/lib/mothership/generated/trace-attribute-values-v1'
import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1'
import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1'
import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confirm'
import {
authenticateCopilotRequestSessionOnly,
createInternalServerErrorResponse,
Expand All @@ -39,6 +34,11 @@ import {
} from '@/lib/mothership/request/http'
import { withIncomingGoSpan } from '@/lib/mothership/request/otel'
import { sealClientToolSettlement } from '@/lib/mothership/request/tools/client-completion-seal.server'
import {
type ClientToolSettlementGuard,
clientToolCompletionMessage,
settleClientToolCall,
} from '@/lib/mothership/request/tools/client-settlement.server'
import { isWorkflowToolName } from '@/lib/mothership/tools/client-executed-tools'
import { getDesktopToolClaimOwner, isNativeDesktopTool } from '@/lib/mothership/tools/desktop-tools'
import {
Expand All @@ -58,20 +58,6 @@ const NATIVE_HANDOFF_INTERRUPTED_MESSAGE =

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

interface UpdateToolCallStatusOptions {
executionId?: string
completionGuard?:
| { status: typeof ASYNC_TOOL_STATUS.pending }
| { status: typeof ASYNC_TOOL_STATUS.running; claimedBy: string }
}

function getClientToolCompletionMessage(status: AsyncConfirmationStatus): string {
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.success) return 'Tool completed'
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) return 'Tool is running in background'
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled) return 'Tool cancelled'
return 'Tool failed'
}

function createConfirmationResponse(
toolCallId: string,
status: AsyncConfirmationStatus,
Expand Down Expand Up @@ -104,53 +90,18 @@ async function updateToolCallStatus(
status: AsyncConfirmationStatus,
message?: string,
data?: AsyncCompletionData,
options: UpdateToolCallStatusOptions = {}
options: { executionId?: string; guard?: ClientToolSettlementGuard } = {}
): Promise<ToolCallStatusUpdateOutcome> {
const toolCallId = existing.toolCallId
try {
if (status === ASYNC_TOOL_CONFIRMATION_STATUS.background) {
const detached = options.executionId
? await detachAsyncToolCall(toolCallId, { preserveClaim: true })
: await detachAsyncToolCall(toolCallId)
if (!detached) return 'conflict'
publishToolConfirmation({
toolCallId,
status,
message: message || undefined,
timestamp: new Date().toISOString(),
data,
...(options.executionId ? { executionId: options.executionId } : {}),
})
return 'updated'
}
const durableStatus =
status === ASYNC_TOOL_CONFIRMATION_STATUS.success
? ASYNC_TOOL_STATUS.completed
: status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled
? ASYNC_TOOL_STATUS.cancelled
: ASYNC_TOOL_STATUS.failed
const completionInput = {
toolCallId,
status: durableStatus,
result: data ?? null,
error: status === 'success' ? null : message || status,
}
const completed =
options.completionGuard?.status === ASYNC_TOOL_STATUS.pending
? await completePendingAsyncToolCall(completionInput)
: options.completionGuard?.status === ASYNC_TOOL_STATUS.running
? await completeClaimedAsyncToolCall(completionInput, options.completionGuard.claimedBy)
: await completeAsyncToolCall(completionInput)
if (!completed) return 'conflict'
publishToolConfirmation({
return await settleClientToolCall({
toolCallId,
status,
message: message || undefined,
timestamp: new Date().toISOString(),
message: message ?? '',
data,
...(options.executionId ? { executionId: options.executionId } : {}),
executionId: options.executionId,
guard: options.guard ?? { kind: 'open' },
})
return 'updated'
} catch (error) {
logger.error('Failed to update tool call status', {
toolCallId,
Expand Down Expand Up @@ -419,7 +370,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
),
}
: {
message: getClientToolCompletionMessage(status),
message: clientToolCompletionMessage(status),
data: await sealClientToolSettlement(existing.result, {
toolCallId,
runId: existing.runId,
Expand Down Expand Up @@ -447,9 +398,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
projected.data,
{
...(isWorkflowTool && executionId ? { executionId } : {}),
...(isPreclaimNativeTerminalOutcome
? { completionGuard: { status: ASYNC_TOOL_STATUS.pending } as const }
: {}),
...(isPreclaimNativeTerminalOutcome ? { guard: { kind: 'pending' } as const } : {}),
}
)

Expand All @@ -460,12 +409,7 @@ export const POST = withRouteHandler((req: NextRequest) => {
ASYNC_TOOL_CONFIRMATION_STATUS.error,
projected.message,
projected.data,
{
completionGuard: {
status: ASYNC_TOOL_STATUS.running,
claimedBy: nativeClaimOwner,
},
}
{ guard: { kind: 'claimed', claimedBy: nativeClaimOwner } }
)
: updateOutcome

Expand Down
19 changes: 19 additions & 0 deletions apps/sim/app/api/desktop/devices/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
import { registerDesktopDeviceContract } from '@/lib/api/contracts/desktop-executor'
import {
defineInternalJsonRoute,
internalRateLimits,
internalSessionAuth,
} from '@/lib/api/server/routes'
import { desktopExecutorErrorPolicy } from '@/lib/api/server/routes/desktop-executor'
import { registerDesktopDevice } from '@/lib/desktop/application/executor'

export const POST = defineInternalJsonRoute({
contract: registerDesktopDeviceContract,
auth: internalSessionAuth,
operation: registerDesktopDevice.operation,
rateLimit: internalRateLimits.user({ bucketName: 'desktop-device-register' }),
errorPolicy: desktopExecutorErrorPolicy,
mapInput: ({ body }) => body,
useCase: registerDesktopDevice,
staticResponseHeaders: { 'Cache-Control': 'no-store' },
})
25 changes: 25 additions & 0 deletions apps/sim/app/api/desktop/inbox/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
import { listDesktopInboxContract } from '@/lib/api/contracts/desktop-executor'
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
import {
desktopExecutorErrorPolicy,
desktopExecutorRateLimit,
} from '@/lib/api/server/routes/desktop-executor'
import { listDesktopInbox } from '@/lib/desktop/application/executor'

export const dynamic = 'force-dynamic'

export const GET = defineInternalJsonRoute({
contract: listDesktopInboxContract,
auth: internalSessionAuth,
operation: listDesktopInbox.operation,
rateLimit: desktopExecutorRateLimit,
errorPolicy: desktopExecutorErrorPolicy,
mapInput: ({ query }) => ({ deviceId: query.deviceId }),
useCase: listDesktopInbox,
present: ({ items }) => ({
items: items.map((item) =>
item.kind === 'call' ? { ...item, createdAt: item.createdAt.toISOString() } : item
),
}),
staticResponseHeaders: { 'Cache-Control': 'no-store' },
})
45 changes: 45 additions & 0 deletions apps/sim/app/api/desktop/inbox/stream/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
import { createLogger } from '@sim/logger'
import type { NextRequest } from 'next/server'
import { desktopInboxStreamContract } from '@/lib/api/contracts/desktop-executor'
import { parseRequest } from '@/lib/api/server'
import { desktopExecutorRateLimit } from '@/lib/api/server/routes/desktop-executor'
import {
InternalUnauthenticatedError,
internalSessionAuth,
} from '@/lib/api/server/routes/internal-json-route'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { openDesktopInboxStream } from '@/lib/desktop/application/executor'
import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors'
import { createSSEStream } from '@/lib/events/sse-endpoint'

export const dynamic = 'force-dynamic'

const logger = createLogger('DesktopInboxStream')

/**
* The background executor's doorbell. A raw route because it streams: authentication, the device
* check, and presence all run through the `openDesktopInboxStream` use case.
*/
export const GET = withRouteHandler(async (request: NextRequest) => {
Comment thread
waleedlatif1 marked this conversation as resolved.
try {
const principal = await internalSessionAuth.authenticate()
const limited = await desktopExecutorRateLimit.enforce(request, principal)
if (limited) return limited
const parsed = await parseRequest(desktopInboxStreamContract, request, {})
if (!parsed.success) return parsed.response
const { deviceId } = parsed.data.query
const inbox = await openDesktopInboxStream.execute({ principal, input: { deviceId } })
return createSSEStream(request, {
Comment thread
waleedlatif1 marked this conversation as resolved.
label: 'desktop-inbox',
revalidate: inbox.revalidate,
subscriptions: [{ subscribe: inbox.subscribe }],
})
} catch (error) {
if (error instanceof InternalUnauthenticatedError)
return new Response('Unauthorized', { status: 401 })
if (error instanceof DesktopDeviceUnrecognizedError)
return new Response(error.message, { status: 401 })
logger.error('Failed to open the desktop inbox stream', error)
return new Response('Unable to open the desktop inbox', { status: 500 })
}
})
18 changes: 18 additions & 0 deletions apps/sim/app/api/desktop/tool/claim/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
import { claimDesktopToolContract } from '@/lib/api/contracts/desktop-executor'
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
import {
desktopExecutorErrorPolicy,
desktopExecutorRateLimit,
} from '@/lib/api/server/routes/desktop-executor'
import { claimDesktopTool } from '@/lib/desktop/application/executor'

export const POST = defineInternalJsonRoute({
contract: claimDesktopToolContract,
auth: internalSessionAuth,
operation: claimDesktopTool.operation,
rateLimit: desktopExecutorRateLimit,
errorPolicy: desktopExecutorErrorPolicy,
mapInput: ({ body }) => body,
useCase: claimDesktopTool,
staticResponseHeaders: { 'Cache-Control': 'no-store' },
})
19 changes: 19 additions & 0 deletions apps/sim/app/api/desktop/tool/complete/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
import { completeDesktopToolContract } from '@/lib/api/contracts/desktop-executor'
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
import {
desktopExecutorErrorPolicy,
desktopExecutorRateLimit,
} from '@/lib/api/server/routes/desktop-executor'
import { completeDesktopTool } from '@/lib/desktop/application/executor'

export const POST = defineInternalJsonRoute({
contract: completeDesktopToolContract,
auth: internalSessionAuth,
operation: completeDesktopTool.operation,
rateLimit: desktopExecutorRateLimit,
errorPolicy: desktopExecutorErrorPolicy,
mapInput: ({ body }) => body,
useCase: completeDesktopTool,
present: ({ outcome, status }) => ({ outcome, status }),
staticResponseHeaders: { 'Cache-Control': 'no-store' },
})
18 changes: 18 additions & 0 deletions apps/sim/app/api/desktop/tool/lease/route.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
import { renewDesktopToolLeaseContract } from '@/lib/api/contracts/desktop-executor'
import { defineInternalJsonRoute, internalSessionAuth } from '@/lib/api/server/routes'
import {
desktopExecutorErrorPolicy,
desktopExecutorRateLimit,
} from '@/lib/api/server/routes/desktop-executor'
import { renewDesktopToolLease } from '@/lib/desktop/application/executor'

export const POST = defineInternalJsonRoute({
contract: renewDesktopToolLeaseContract,
auth: internalSessionAuth,
operation: renewDesktopToolLease.operation,
rateLimit: desktopExecutorRateLimit,
errorPolicy: desktopExecutorErrorPolicy,
mapInput: ({ body }) => body,
useCase: renewDesktopToolLease,
staticResponseHeaders: { 'Cache-Control': 'no-store' },
})
Loading
Loading