From bd54c7c0e903b490d183e8cc212364f16fb9e405 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 5 Oct 2026 18:53:05 -0700 Subject: [PATCH 1/4] fix(mothership): fail unclaimed desktop calls fast and stop dropping them silently A desktop call reaches the user's machine only through the chat view showing that chat, so one issued while the user is elsewhere was never claimed and the turn waited out a 90-160 s watchdog that called it hung. One desktop wait now fails a call still unclaimed after a 15 s pickup grace as never started, with the inverse of the desktop's claim (pending -> failed) and a result sealed like any client completion. A call claimed in time keeps waiting, and a result that lands as the grace runs out is returned, not discarded. Local reads (read_local_file, user-local VFS reads) join them for a desktop that advertises localReadClaims: the desktop claims each read through authorize, as it already does for imports. Older desktops keep today's behaviour. The single "hung and was abandoned" result becomes two: notStarted for an unclaimed call, and outcomeUnknown/doNotRetry for one that started and lost its result. Both are settled through one sealed failure helper. A terminal run's server budget now follows the wait the desktop holds it for. A not-started report can settle only an unclaimed call. In the renderer, a running browser action is cancelled only by Stop: recovering the stream or leaving the chat view lets it finish and report. A terminal call delivered too late is reported as not started instead of being skipped. --- apps/desktop/src/main/ipc.test.ts | 4 +- apps/desktop/src/main/ipc.ts | 7 +- apps/desktop/src/main/terminal/index.ts | 15 +- apps/desktop/src/preload/index.ts | 1 + apps/sim/app/api/copilot/confirm/route.ts | 11 + .../app/api/desktop/tool/authorize/route.ts | 18 + .../home/hooks/use-chat.dom.test.tsx | 113 ++++++- .../[workspaceId]/home/hooks/use-chat.ts | 27 +- apps/sim/lib/desktop/index.ts | 2 + apps/sim/lib/mothership/chat/post.ts | 2 + apps/sim/lib/mothership/constants.ts | 12 + .../request/handlers/handlers.test.ts | 21 +- .../lib/mothership/request/handlers/tool.ts | 30 +- .../request/tools/desktop-wait.test.ts | 55 +++ .../mothership/request/tools/desktop-wait.ts | 95 ++++++ .../mothership/request/tools/executor.test.ts | 56 ++- .../lib/mothership/request/tools/executor.ts | 145 ++++---- .../request/tools/tool-call-failure.ts | 87 +++++ .../tools/workflow-client-fallback.test.ts | 14 +- apps/sim/lib/mothership/request/types.ts | 7 + .../client/desktop-tool-pickup.integration.ts | 319 ++++++++++++++++++ .../client/terminal-tool-execution.test.ts | 19 ++ .../tools/client/terminal-tool-execution.ts | 20 +- .../sim/lib/mothership/tools/desktop-tools.ts | 22 ++ packages/desktop-bridge/src/index.ts | 6 + packages/terminal-protocol/src/index.ts | 12 +- packages/testing/src/mocks/index.ts | 4 + .../mothership-client-tool-waiter.mock.ts | 31 ++ 28 files changed, 1041 insertions(+), 114 deletions(-) create mode 100644 apps/sim/lib/mothership/request/tools/desktop-wait.test.ts create mode 100644 apps/sim/lib/mothership/request/tools/desktop-wait.ts create mode 100644 apps/sim/lib/mothership/request/tools/tool-call-failure.ts create mode 100644 apps/sim/lib/mothership/tools/client/desktop-tool-pickup.integration.ts create mode 100644 packages/testing/src/mocks/mothership-client-tool-waiter.mock.ts diff --git a/apps/desktop/src/main/ipc.test.ts b/apps/desktop/src/main/ipc.test.ts index 20769df0376..f322f03e454 100644 --- a/apps/desktop/src/main/ipc.test.ts +++ b/apps/desktop/src/main/ipc.test.ts @@ -485,7 +485,7 @@ describe('registerIpcHandlers', () => { expect(mounts).not.toHaveBeenCalled() expect(fetchAuthorization).toHaveBeenCalledWith( `${APP}/api/desktop/tool/authorize`, - expect.objectContaining({ body: JSON.stringify({ toolCallId: 'tool-native' }) }) + expect.objectContaining({ body: JSON.stringify({ toolCallId: 'tool-native', claim: true }) }) ) expect( await handler?.(evilEvent, { operation: 'read', toolCallId: 'tool-native' }) @@ -569,7 +569,7 @@ describe('registerIpcHandlers', () => { expect.objectContaining({ method: 'POST', credentials: 'include', - body: JSON.stringify({ toolCallId: 'tool-1' }), + body: JSON.stringify({ toolCallId: 'tool-1', claim: true }), }) ) }) diff --git a/apps/desktop/src/main/ipc.ts b/apps/desktop/src/main/ipc.ts index b7a16b1d278..423dca0b7e6 100644 --- a/apps/desktop/src/main/ipc.ts +++ b/apps/desktop/src/main/ipc.ts @@ -578,10 +578,13 @@ async function authorizeLocalFilesystemTool( request: unknown ): Promise { if (typeof request !== 'object' || request === null) return false + // Claimed like an import: the server sees the read picked up, and refuses one it already + // failed as never started. const authorization = await fetchDesktopToolAuthorization( event, deps, - (request as { requestId?: unknown }).requestId + (request as { requestId?: unknown }).requestId, + true ) return authorization ? deps.localFilesystem.isAuthorizedClientToolRequest(request, authorization) @@ -2113,7 +2116,7 @@ export function registerIpcHandlers(deps: IpcDeps): void { event, deps, request.toolCallId, - request.operation === 'manifest', + request.operation === 'manifest' || request.operation === 'read', (status) => { failureStatus = status } diff --git a/apps/desktop/src/main/terminal/index.ts b/apps/desktop/src/main/terminal/index.ts index d20e68a57ec..925330f72fb 100644 --- a/apps/desktop/src/main/terminal/index.ts +++ b/apps/desktop/src/main/terminal/index.ts @@ -14,11 +14,10 @@ import { homedir } from 'node:os' import type { TerminalShortcutCommand } from '@sim/desktop-bridge' import { createLogger } from '@sim/logger' import { - DEFAULT_RUN_WAIT_MS, isTerminalControlKey, MAX_INPUT_KEYS, - MAX_RUN_WAIT_MS, MAX_TOOL_OUTPUT_CHARS, + resolveRunWaitMs, type TerminalCommandEvent, type TerminalControlKey, type TerminalCwdResult, @@ -109,14 +108,6 @@ const HANDOFF_MAX_MS = 12 * 60 * 60 * 1000 */ const HANDOFF_SETTLE_MS = 5_000 -/** How long to hold the turn before handing a still-running command back. */ -function resolveWaitMs(waitSeconds: number | undefined): number { - const requested = Number(waitSeconds) - return Number.isFinite(requested) && requested > 0 - ? Math.min(requested * 1000, MAX_RUN_WAIT_MS) - : DEFAULT_RUN_WAIT_MS -} - function elideOutput(value: string): { text: string; truncated: boolean } { return elide(value, MAX_TOOL_OUTPUT_CHARS) } @@ -1097,7 +1088,7 @@ export class TerminalService { const handle = await startRun(session, command, terminal.currentCwd, terminal.env) if ('error' in handle) throw new TerminalError('SPAWN_FAILED', handle.error) - const waitMs = resolveWaitMs(args.waitSeconds) + const waitMs = resolveRunWaitMs(args.waitSeconds) const outcome = await awaitRun(handle, waitMs) if (outcome.done) { await closeRunWindow(handle, terminal.env) @@ -1149,7 +1140,7 @@ export class TerminalService { ) } - return session.runCommand(command, toolCallId, resolveWaitMs(args.waitSeconds)) + return session.runCommand(command, toolCallId, resolveRunWaitMs(args.waitSeconds)) } private spawn( diff --git a/apps/desktop/src/preload/index.ts b/apps/desktop/src/preload/index.ts index 00fb6180a9b..4122466b9eb 100644 --- a/apps/desktop/src/preload/index.ts +++ b/apps/desktop/src/preload/index.ts @@ -156,6 +156,7 @@ const api: SimDesktopApi = { ipcRenderer.invoke('desktop:local-filesystem', request), localFiles: (request: DesktopLocalFileRequest): Promise => ipcRenderer.invoke('desktop:local-files', request), + localReadClaims: true, onCommand: (callback: (command: DesktopCommand) => void): (() => void) => { const listener = (_event: unknown, command: DesktopCommand) => callback(command) ipcRenderer.on('desktop:command', listener) diff --git a/apps/sim/app/api/copilot/confirm/route.ts b/apps/sim/app/api/copilot/confirm/route.ts index 71c12bd1953..d389a63b839 100644 --- a/apps/sim/app/api/copilot/confirm/route.ts +++ b/apps/sim/app/api/copilot/confirm/route.ts @@ -330,6 +330,17 @@ export const POST = withRouteHandler((req: NextRequest) => { span.setAttribute(TraceAttr.CopilotConfirmOutcome, CopilotConfirmOutcome.ToolCallNotFound) return createNotFoundResponse('Running client tool call not found') } + // A reporter that says the call never started (a stale replay, a closed view) cannot speak + // for a call the desktop claimed: only the claim's own result may settle it. + if ( + isNativeClientTool && + isPlainRecord(data) && + data.notStarted === true && + existing.status !== ASYNC_TOOL_STATUS.pending + ) { + span.setAttribute(TraceAttr.CopilotConfirmOutcome, CopilotConfirmOutcome.ToolCallNotFound) + return createNotFoundResponse('Pending client tool call not found') + } let effectiveStatus = status let executionId = submittedExecutionId diff --git a/apps/sim/app/api/desktop/tool/authorize/route.ts b/apps/sim/app/api/desktop/tool/authorize/route.ts index d3dd0a3b7f0..37d2c5d41e9 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.ts @@ -19,6 +19,7 @@ import { createNotFoundResponse, createUnauthorizedResponse, } from '@/lib/mothership/request/http' +import { isLocalReadToolCall } from '@/lib/mothership/tools/desktop-tools' import { isUserLocalVfsToolCall } from '@/lib/mothership/tools/local-filesystem' const admissionClosedResponse = () => @@ -129,6 +130,23 @@ export const POST = withRouteHandler(async (request: NextRequest) => { return createNotFoundResponse('The import must be started before reading file bytes') } + // A desktop that claims local reads claims each one before its first read; its later reads of + // the same call ride on that claim. A call persisted running (an older desktop's turn) is read + // as before. + if (parsed.data.body.claim && isLocalReadToolCall(toolCall.toolName, args)) { + const notPending = () => createNotFoundResponse('Pending client tool call not found') + if (toolCall.status === 'pending') { + const { outcome } = await claimToolExecution({ + toolCallId: toolCall.toolCallId, + runId: toolCall.runId, + userId, + claimedBy: DESKTOP_TOOL_CLAIM_OWNER.files, + }) + if (outcome !== 'claimed') return refusedClaimResponse(outcome, notPending) + } else if (toolCall.claimedBy !== null && toolCall.claimedBy !== DESKTOP_TOOL_CLAIM_OWNER.files) + return notPending() + } + // Browser and terminal actions are one-shot side effects on the user's own // machine, so the pending call is claimed here, atomically, before crossing // the Electron boundary — a replayed renderer event must not run a command diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 1781894479c..9fcb5c7d609 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -21,13 +21,15 @@ import { act, type ReactNode, StrictMode, useEffect, useState } from 'react' import { authClientMock, authClientMockFns } from '@sim/testing/mocks/auth-client.mock' +import { libDesktopMock, libDesktopMockFns } from '@sim/testing/mocks/lib-desktop.mock' import { nextNavigationMock, nextNavigationMockFns } from '@sim/testing/mocks/next-navigation.mock' import { sleep } from '@sim/utils/helpers' import { QueryClient, QueryClientProvider } from '@tanstack/react-query' import { createRoot, type Root } from 'react-dom/client' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -const { mockRequestJson, mockExecuteWorkflow } = vi.hoisted(() => ({ +const { mockRequestJson, mockExecuteWorkflow, mockExecuteBrowserToolOnClient } = vi.hoisted(() => ({ + mockExecuteBrowserToolOnClient: vi.fn(), mockRequestJson: vi.fn(), mockExecuteWorkflow: vi.fn< @@ -44,6 +46,10 @@ vi.mock('@/app/workspace/[workspaceId]/providers/feature-flags-provider', () => })) vi.mock('next/navigation', () => nextNavigationMock) +vi.mock('@/lib/desktop', () => libDesktopMock) +vi.mock('@/lib/mothership/tools/client/browser-tool-execution', () => ({ + executeBrowserToolOnClient: mockExecuteBrowserToolOnClient, +})) vi.mock('@/lib/auth/auth-client', () => authClientMock) vi.unmock('@/stores/execution/store') vi.unmock('@/stores/terminal') @@ -2039,4 +2045,109 @@ describe('useChat remount send recovery', () => { ?.messages.map((message) => message.id) ).toEqual(['saved-user', 'saved-assistant']) }) + describe('a desktop browser action in flight', () => { + const chatId = 'chat-browser-action' + const history: MothershipChatHistory = { + id: chatId, + mode: 'agent', + title: 'Browser action', + messages: [], + activeStreamId: null, + resources: [], + } + + /** Opens a turn whose stream delivers one desktop browser call and stays open. */ + async function startBrowserAction() { + let streamId: string | undefined + const replays: string[] = [] + mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => + Promise.resolve( + contract.path === '/api/mothership/chats/[chatId]' + ? { chat: { ...history, activeStreamId: streamId ?? null } } + : { chats: [] } + ) + ) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + const url = String(input) + if (url.includes('/api/mothership/chat/stream')) replays.push(url) + if (url !== '/api/mothership/chat' || init?.method !== 'POST') { + return fetchStub(input, init) + } + streamId = JSON.parse(String(init.body)).userMessageId + const call: MothershipStreamV1EventEnvelope = { + v: 1, + seq: 1, + ts: new Date().toISOString(), + type: 'tool', + stream: { streamId: streamId ?? '' }, + payload: { + phase: 'call', + executor: 'client', + mode: 'async', + toolName: 'browser_list_tabs', + toolCallId: 'browser-call', + arguments: {}, + }, + } + return new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(call)}\n\n`)) + }, + }), + { headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } } + ) + }) + const chat = renderUseChatInChat(chatId, history) + await act(async () => { + void chat.getResult().sendMessage('List my tabs') + }) + await waitFor(() => mockExecuteBrowserToolOnClient.mock.calls.length === 1) + const toolSignal = mockExecuteBrowserToolOnClient.mock.calls[0]?.[5] + if (!(toolSignal instanceof AbortSignal)) + throw new Error('The browser action has no lifetime') + return { ...chat, toolSignal, replays } + } + + beforeEach(() => { + libDesktopMockFns.mockIsDesktopApp.mockReturnValue(true) + }) + + afterEach(() => { + libDesktopMockFns.mockIsDesktopApp.mockReset() + }) + + it('keeps running when the window returns to view and the stream is recovered', async () => { + const { toolSignal, replays } = await startBrowserAction() + + Object.defineProperty(document, 'visibilityState', { + configurable: true, + get: () => 'visible', + }) + await act(async () => { + document.dispatchEvent(new Event('visibilitychange')) + }) + await waitFor(() => replays.length > 0) + + expect(toolSignal.aborted).toBe(false) + }) + + it('keeps running when the chat view unmounts, so it finishes and reports its result', async () => { + const { toolSignal, unmount } = await startBrowserAction() + + unmount() + + expect(toolSignal.aborted).toBe(false) + }) + + it('is cancelled when the user stops the chat', async () => { + const { toolSignal, getResult } = await startBrowserAction() + + await act(async () => { + await getResult().stopGeneration() + }) + + expect(toolSignal.aborted).toBe(true) + }) + }) }) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index 0a85a441003..8bed508318c 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -452,6 +452,25 @@ export async function waitForDetachedChatResolution( } } +/** The abort reason of the user's Stop. */ +const USER_STOP_ABORT_REASON = 'user_stop:client_stopGeneration' + +/** + * The lifetime a browser action started from one stream observes: only the user's Stop cancels + * it. Replacing the stream reader (the window returning to view, a history reconnect) or leaving + * the chat view leaves it running, so it finishes and reports its own result. + */ +function browserToolLifetime(streamSignal: AbortSignal | undefined): AbortSignal | undefined { + if (!streamSignal) return undefined + const lifetime = new AbortController() + const followStop = () => { + if (streamSignal.reason === USER_STOP_ABORT_REASON) lifetime.abort(USER_STOP_ABORT_REASON) + } + if (streamSignal.aborted) followStop() + else streamSignal.addEventListener('abort', followStop, { once: true }) + return lifetime.signal +} + /** * Runs a browser tool on the desktop client. The agent's tab reaches the * resource strip through the desktop tab list, so nothing is opened here. @@ -2151,7 +2170,7 @@ export function useChat( shouldContinue?: () => boolean } ) => { - const streamAbortSignal = abortControllerRef.current?.signal + const browserToolSignal = browserToolLifetime(abortControllerRef.current?.signal) const activityTracker = getResourceActivityTracker( expectedGen ?? streamGenRef.current, options?.targetChatId @@ -2172,7 +2191,7 @@ export function useChat( eventTs?: string ) => { const scopeId = activityScopeId() - startClientBrowserTool(toolCallId, toolName, toolArgs, scopeId, eventTs, streamAbortSignal) + startClientBrowserTool(toolCallId, toolName, toolArgs, scopeId, eventTs, browserToolSignal) } const startClientTerminalToolForStream = ( toolCallId: string, @@ -4424,7 +4443,7 @@ export function useChat( streamReaderRef.current = null const stoppedController = abortControllerRef.current if (stoppedController !== pendingAdmission?.controller) { - stoppedController?.abort('user_stop:client_stopGeneration') + stoppedController?.abort(USER_STOP_ABORT_REASON) } abortControllerRef.current = null setTransportIdle() @@ -4520,7 +4539,7 @@ export function useChat( if (pendingAdmission && pendingAdmission.userMessageId === sid) { const admittedChatId = await pendingAdmission.settled resolvedChatId ??= admittedChatId - pendingAdmission.controller.abort('user_stop:client_stopGeneration') + pendingAdmission.controller.abort(USER_STOP_ABORT_REASON) } if (!resolvedChatId && sid) { resolvedChatId = await resolveChatIdForStream(sid, { preferExistingChatId: false }) diff --git a/apps/sim/lib/desktop/index.ts b/apps/sim/lib/desktop/index.ts index 7ab60648924..230c4b76ddc 100644 --- a/apps/sim/lib/desktop/index.ts +++ b/apps/sim/lib/desktop/index.ts @@ -156,6 +156,7 @@ export interface DesktopChatCapabilities { desktopCapabilities?: { localFiles?: true localFilesystem?: true + localReadClaims?: true browser?: true terminal?: true browserSessions?: BrowserKnownSession[] @@ -217,6 +218,7 @@ export async function getDesktopChatCapabilities( ? { desktopCapabilities: { ...(localFiles ? { localFiles: true as const } : {}), + ...(bridge?.localReadClaims === true ? { localReadClaims: true as const } : {}), ...(localFilesystem ? { localFilesystem: true as const } : {}), ...(browser ? { browser: true as const } : {}), ...(terminal ? { terminal: true as const } : {}), diff --git a/apps/sim/lib/mothership/chat/post.ts b/apps/sim/lib/mothership/chat/post.ts index 76cc073961a..9979bc05d4b 100644 --- a/apps/sim/lib/mothership/chat/post.ts +++ b/apps/sim/lib/mothership/chat/post.ts @@ -312,6 +312,7 @@ const ChatMessageSchema = z .object({ localFilesystem: z.boolean().optional(), localFiles: z.boolean().optional(), + localReadClaims: z.boolean().optional(), browser: z.boolean().optional(), terminal: z.boolean().optional(), terminals: z @@ -1517,6 +1518,7 @@ export async function handleUnifiedChatPost(req: NextRequest) { // Executor routing is decided HERE, once per turn, from the caller's declared // capabilities — dispatch never discovers client absence by burning a grace timer. clientToolPickupExpected, + desktopClaimsLocalReads: body.desktopCapabilities?.localReadClaims === true, executionContext, billingAttribution: executionContext.billingAttribution, onComplete: buildOnComplete({ diff --git a/apps/sim/lib/mothership/constants.ts b/apps/sim/lib/mothership/constants.ts index f5b6a7b95f5..90cdb8037a7 100644 --- a/apps/sim/lib/mothership/constants.ts +++ b/apps/sim/lib/mothership/constants.ts @@ -66,6 +66,18 @@ export const CHAT_RUN_DEADLINE_MS = 3_600_000 */ export const COPILOT_WORKFLOW_TOOL_CLIENT_GRACE_MS = 30_000 +/** + * How long a desktop tool call (browser, terminal, local import) waits for the desktop app to + * claim it before it fails as never started. + * + * Same cause as the workflow grace: only the chat view showing this chat starts the call, so a + * call issued while the user is on another chat or page is claimed by nobody. A live view claims + * within a second or two (stream frame -> IPC -> authorize), and there is no server fallback to + * run instead, so the call fails with a "not started, safe to retry" result rather than parking + * until a watchdog calls it hung. + */ +export const DESKTOP_TOOL_PICKUP_GRACE_MS = 15_000 + /** SessionStorage key for persisting active stream metadata across page reloads. */ export const STREAM_STORAGE_KEY = 'copilot_active_stream' diff --git a/apps/sim/lib/mothership/request/handlers/handlers.test.ts b/apps/sim/lib/mothership/request/handlers/handlers.test.ts index 92c889615e5..f1080512fb6 100644 --- a/apps/sim/lib/mothership/request/handlers/handlers.test.ts +++ b/apps/sim/lib/mothership/request/handlers/handlers.test.ts @@ -2,6 +2,10 @@ import { mothershipAsyncRunsMock, mothershipAsyncRunsMockFns, } from '@sim/testing/mocks/mothership-async-runs.mock' +import { + mothershipClientToolWaiterMock, + mothershipClientToolWaiterMockFns, +} from '@sim/testing/mocks/mothership-client-tool-waiter.mock' import { sleep } from '@sim/utils/helpers' import { beforeEach, describe, expect, it, vi } from 'vitest' import { AsyncToolCallOwnershipError } from '@/lib/mothership/async-runs/errors' @@ -17,12 +21,11 @@ const { isSimExecuted, executeTool, ensureHandlersRegistered, toolRequiresApprov }) ) -const { waitForClientToolCompletion, waitForToolCompletion, waitForWorkflowToolCompletion } = - vi.hoisted(() => ({ - waitForClientToolCompletion: vi.fn(), - waitForToolCompletion: vi.fn(), - waitForWorkflowToolCompletion: vi.fn(), - })) +const { + mockWaitForClientToolCompletion: waitForClientToolCompletion, + mockWaitForToolCompletion: waitForToolCompletion, + mockWaitForWorkflowToolCompletion: waitForWorkflowToolCompletion, +} = mothershipClientToolWaiterMockFns const { sealClientToolContext } = vi.hoisted(() => ({ sealClientToolContext: vi.fn(), @@ -44,11 +47,7 @@ vi.mock('@/lib/mothership/request/tools/tables', () => ({ maybeWriteReadCsvToTable: vi.fn(async (_toolName, _params, result) => result), })) -vi.mock('@/lib/mothership/request/tools/client', () => ({ - waitForClientToolCompletion, - waitForToolCompletion, - waitForWorkflowToolCompletion, -})) +vi.mock('@/lib/mothership/request/tools/client', () => mothershipClientToolWaiterMock) vi.mock('@/lib/mothership/request/tools/client-completion-seal.server', () => ({ sealClientToolContext, diff --git a/apps/sim/lib/mothership/request/handlers/tool.ts b/apps/sim/lib/mothership/request/handlers/tool.ts index f98031289a7..5cb2fa59d04 100644 --- a/apps/sim/lib/mothership/request/handlers/tool.ts +++ b/apps/sim/lib/mothership/request/handlers/tool.ts @@ -9,6 +9,7 @@ import { upsertAsyncToolCall } from '@/lib/mothership/async-runs/repository' import { CLIENT_TOOL_RESULT_TIMEOUT_MS, COPILOT_WORKFLOW_TOOL_CLIENT_GRACE_MS, + DESKTOP_TOOL_PICKUP_GRACE_MS, } from '@/lib/mothership/constants' import { MothershipStreamV1AsyncToolRecordStatus, @@ -30,6 +31,7 @@ import { markToolResultSeen, wasToolResultSeen } from '@/lib/mothership/request/ import { setTerminalToolCallState } from '@/lib/mothership/request/tool-call-state' import { waitForClientToolCompletion } from '@/lib/mothership/request/tools/client' import { sealClientToolContext } from '@/lib/mothership/request/tools/client-completion-seal.server' +import { waitForDesktopToolCall } from '@/lib/mothership/request/tools/desktop-wait' import { executeToolAndReport } from '@/lib/mothership/request/tools/executor' import { runGatedToolExecution, @@ -47,7 +49,7 @@ import type { import { getToolEntry, isSimExecuted } from '@/lib/mothership/tool-executor' import { isToolHiddenInUi } from '@/lib/mothership/tools/client/hidden-tools' import { isWorkflowToolName } from '@/lib/mothership/tools/client-executed-tools' -import { getDesktopToolClaimOwner } from '@/lib/mothership/tools/desktop-tools' +import { isClaimedOnPickup } from '@/lib/mothership/tools/desktop-tools' import { isUserLocalVfsToolCall } from '@/lib/mothership/tools/local-filesystem' import { extractStreamingStringArgument } from '@/lib/mothership/tools/streaming-args' import { readToolActivity } from '@/lib/mothership/tools/tool-activity' @@ -304,14 +306,16 @@ export async function prePersistClientExecutableToolCall( toolName: data.toolName, args: data.arguments, sealedContext, - // Browser and terminal actions cross a second, native authorization - // boundary. Leave those rows pending until Electron atomically claims - // them — the authorize endpoint only hands over a pending call, so a row - // that arrives already running can never be executed natively. All other - // client tools retain the established "already dispatched" running state. - // A gated tool is likewise pending: nothing has been dispatched yet. + // Desktop calls the desktop claims on pickup (browser, terminal, import, + // and local reads when this turn's desktop claims them) cross a second, + // native authorization boundary. Leave those rows pending until Electron + // atomically claims them, so a replayed event cannot act twice and a call + // nobody picks up can fail fast. All other client tools retain the + // established "already dispatched" running state. A gated tool is + // likewise pending: nothing has been dispatched yet. status: - gated || getDesktopToolClaimOwner(data.toolName) + gated || + isClaimedOnPickup(data.toolName, data.arguments, options?.desktopClaimsLocalReads === true) ? MothershipStreamV1AsyncToolRecordStatus.pending : MothershipStreamV1AsyncToolRecordStatus.running, permissionRequested: gated, @@ -935,6 +939,16 @@ async function dispatchToolExecution( return race.signal ?? errorCompletion('Tool completion missing') } completion = race.completion ?? null + } else if (isClaimedOnPickup(toolName, args, options.desktopClaimsLocalReads === true)) { + completion = await waitForDesktopToolCall({ + toolCallId, + runId: context.runId, + userId: execContext.userId, + timeoutMs, + pickupGraceMs: DESKTOP_TOOL_PICKUP_GRACE_MS, + abortSignal: options.abortSignal, + registry: execContext.resolvedSecretTraceRegistry, + }) } else { completion = await waitForClientToolCompletion({ toolCallId, diff --git a/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts b/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts new file mode 100644 index 00000000000..5358721b65f --- /dev/null +++ b/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts @@ -0,0 +1,55 @@ +import { + mothershipAsyncRunsMock, + mothershipAsyncRunsMockFns, +} from '@sim/testing/mocks/mothership-async-runs.mock' +import { + mothershipClientToolWaiterMock, + mothershipClientToolWaiterMockFns, +} from '@sim/testing/mocks/mothership-client-tool-waiter.mock' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +vi.mock('@/lib/mothership/async-runs/repository', () => mothershipAsyncRunsMock) +vi.mock('@/lib/mothership/request/tools/client', () => mothershipClientToolWaiterMock) + +import { waitForDesktopToolCall } from '@/lib/mothership/request/tools/desktop-wait' + +const waitForClientToolCompletion = + mothershipClientToolWaiterMockFns.mockWaitForClientToolCompletion +const completePendingAsyncToolCall = mothershipAsyncRunsMockFns.mockCompletePendingAsyncToolCall + +const GRACE_MS = 15_000 +const params = { + toolCallId: 'browser-call', + runId: 'run-1', + userId: 'user-1', + timeoutMs: 3_600_000, + pickupGraceMs: GRACE_MS, +} +const desktopResult = { status: 'success' as const, message: 'Tool completed', data: { tabs: 2 } } + +describe('waitForDesktopToolCall', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('returns the desktop result that landed as the pickup grace ran out', async () => { + waitForClientToolCompletion.mockImplementationOnce( + ({ abortSignal }: { abortSignal: AbortSignal }) => + new Promise((resolve) => { + // The confirmation arrived at the grace boundary: the waiter is already restoring it, + // so stopping the wait no longer turns it into a missing result. + abortSignal.addEventListener('abort', () => resolve(desktopResult), { once: true }) + }) + ) + completePendingAsyncToolCall.mockResolvedValueOnce(null) + + const answer = waitForDesktopToolCall(params) + await vi.advanceTimersByTimeAsync(GRACE_MS) + + expect(await answer).toEqual(desktopResult) + }) +}) diff --git a/apps/sim/lib/mothership/request/tools/desktop-wait.ts b/apps/sim/lib/mothership/request/tools/desktop-wait.ts new file mode 100644 index 00000000000..770283d4e75 --- /dev/null +++ b/apps/sim/lib/mothership/request/tools/desktop-wait.ts @@ -0,0 +1,95 @@ +import { createLogger } from '@sim/logger' +import { interruptibleSleep } from '@sim/utils/helpers' +import type { AsyncTerminalCompletionSnapshot } from '@/lib/mothership/async-runs/lifecycle' +import { MothershipStreamV1ToolOutcome } from '@/lib/mothership/generated/mothership-stream-v1' +import { CopilotDegradedReason } from '@/lib/mothership/generated/trace-attribute-values-v1' +import { recordDegraded } from '@/lib/mothership/request/metrics' +import { waitForClientToolCompletion } from '@/lib/mothership/request/tools/client' +import { settleToolCallFailure } from '@/lib/mothership/request/tools/tool-call-failure' +import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' + +const logger = createLogger('CopilotDesktopToolWait') + +const DESKTOP_TOOL_NOT_STARTED_MESSAGE = + 'Not run: this chat is not open in the Sim desktop app, so nothing picked up this call and nothing happened on the user’s computer. It is safe to retry once the user opens this chat in the Sim desktop app.' + +/** The model-facing result of a desktop call that was never picked up. */ +export function desktopToolNotStarted(): { message: string; data: Record } { + return { + message: DESKTOP_TOOL_NOT_STARTED_MESSAGE, + data: { error: DESKTOP_TOOL_NOT_STARTED_MESSAGE, notStarted: true }, + } +} + +interface WaitForDesktopToolCallParams { + toolCallId: string + runId?: string + userId: string + timeoutMs: number + /** How long the desktop has to claim the call before it fails as never started. */ + pickupGraceMs: number + abortSignal?: AbortSignal + registry?: ResolvedSecretTraceRegistry +} + +/** + * Waits for a desktop tool call the desktop claims on pickup, and fails it fast when nothing + * picks it up. + * + * Only the chat view showing this chat starts a desktop call, so one issued while the user is + * elsewhere is never claimed. After the pickup grace the still-pending call is settled as never + * started with the inverse of the desktop's claim (`pending -> failed`), so exactly one of the two + * wins: a desktop that claims it later is refused, and a call claimed in time keeps waiting for + * its result within the original `timeoutMs`. + */ +export async function waitForDesktopToolCall( + params: WaitForDesktopToolCallParams +): Promise { + const { toolCallId, timeoutMs, pickupGraceMs, abortSignal } = params + const startedAt = Date.now() + const stopPickupWait = new AbortController() + const pickupWait = waitForClientToolCompletion({ + ...params, + abortSignal: abortSignal + ? AbortSignal.any([abortSignal, stopPickupWait.signal]) + : stopPickupWait.signal, + }) + const graceOver = new AbortController() + const first = await Promise.race([ + pickupWait.then((completion) => ({ kind: 'reported' as const, completion })), + interruptibleSleep(Math.min(pickupGraceMs, timeoutMs), graceOver.signal).then(() => ({ + kind: 'grace' as const, + })), + ]) + graceOver.abort() + if (first.kind === 'reported') return first.completion + if (abortSignal?.aborted) return pickupWait + + // A result that landed as the grace ran out is already being restored: it is the answer. + stopPickupWait.abort() + const lateResult = await pickupWait + if (lateResult) return lateResult + + const notStarted = desktopToolNotStarted() + const settlement = await settleToolCallFailure({ + toolCallId, + runId: params.runId, + userId: params.userId, + registry: params.registry, + ...notStarted, + unclaimedOnly: true, + }) + if (settlement !== 'settled') { + return waitForClientToolCompletion({ + ...params, + timeoutMs: Math.max(0, timeoutMs - (Date.now() - startedAt)), + }) + } + + recordDegraded(CopilotDegradedReason.ClientPickupTimeout) + logger.info('No desktop picked up the tool call within its grace; failed it as not started', { + toolCallId, + pickupGraceMs, + }) + return { status: MothershipStreamV1ToolOutcome.error, ...notStarted } +} diff --git a/apps/sim/lib/mothership/request/tools/executor.test.ts b/apps/sim/lib/mothership/request/tools/executor.test.ts index 72f44af3fc9..aa2461e61fc 100644 --- a/apps/sim/lib/mothership/request/tools/executor.test.ts +++ b/apps/sim/lib/mothership/request/tools/executor.test.ts @@ -87,6 +87,7 @@ vi.mock('@/lib/mothership/chat/delegation', () => ({ const { mockCompleteAsyncToolCall: completeAsyncToolCall, + mockCompletePendingAsyncToolCall: completePendingAsyncToolCall, mockMarkAsyncToolRunning: markAsyncToolRunning, mockUpsertAsyncToolCall: upsertAsyncToolCall, mockClaimToolExecution: claimToolExecution, @@ -250,6 +251,23 @@ describe('pendingToolWaitBudgetMs', () => { ).toBe(195_000) }) + it('outlasts the wait window the desktop holds a terminal run for', () => { + expect( + pendingToolWaitBudgetMs({ + name: 'terminal', + status: 'executing', + params: { operation: 'run', args: { command: 'bun install', waitSeconds: 120 } }, + }) + ).toBeGreaterThan(120_000) + expect( + pendingToolWaitBudgetMs({ + name: 'terminal', + status: 'executing', + params: { operation: 'read', args: {} }, + }) + ).toBe(TOOL_WATCHDOG_DEFAULT_MS) + }) + it('falls back to the tool\u2019s own watchdog once it is actually executing', () => { expect(pendingToolWaitBudgetMs({ name: 'terminal_run', status: 'executing' })).toBe( TOOL_WATCHDOG_DEFAULT_MS @@ -874,7 +892,7 @@ describe('watchdog completion provenance', () => { status: 'error', message: expect.stringContaining('outcome is unknown'), data: { - error: expect.stringContaining('hung'), + error: expect.stringContaining('never came back'), outcomeUnknown: true, doNotRetry: true, }, @@ -892,6 +910,36 @@ describe('watchdog completion provenance', () => { ) }) + it('tells the model a desktop call nobody picked up never started', async () => { + const { toolCall, context, execContext } = createHungClient() + toolCall.name = 'browser_click' + completePendingAsyncToolCall.mockImplementationOnce(async (input) => ({ ...input })) + + await failPendingToolCall(toolCall.id, context, execContext) + + expect(toolCall.result).toEqual({ + success: false, + output: { error: expect.stringContaining('safe to retry'), notStarted: true }, + }) + }) + + it('tells the model a desktop call it picked up lost its result and may have acted', async () => { + const { toolCall, context, execContext } = createHungClient() + toolCall.name = 'browser_click' + completePendingAsyncToolCall.mockResolvedValueOnce(null) + + await failPendingToolCall(toolCall.id, context, execContext) + + expect(toolCall.result).toEqual({ + success: false, + output: { + error: expect.stringContaining('desktop app started this action'), + outcomeUnknown: true, + doNotRetry: true, + }, + }) + }) + it('preserves an actual completion that settles while encryption is pending', async () => { const { toolCall, context, execContext } = createHungClient() let finishEncryption: (value: { encrypted: string }) => void = () => {} @@ -1042,7 +1090,7 @@ describe('watchdog completion provenance', () => { const completion = await execution expect(completion.status).toBe('error') - expect(completion.message).toContain('hung') + expect(completion.message).toContain('never came back') expect(toolCall.status).toBe('error') expect(completeAsyncToolCall).toHaveBeenCalledTimes(1) expect(publishToolConfirmation).toHaveBeenCalledTimes(1) @@ -1112,7 +1160,7 @@ describe('watchdog completion provenance', () => { const completion = await execution expect(completion.status).toBe('error') - expect(completion.message).toContain('hung') + expect(completion.message).toContain('never came back') expect(toolCall.status).toBe('error') expect(publishToolConfirmation).toHaveBeenCalledTimes(1) expect(onEvent).not.toHaveBeenCalled() @@ -1146,7 +1194,7 @@ describe('watchdog completion provenance', () => { const completion = await execution expect(completion.status).toBe('error') - expect(completion.message).toContain('hung') + expect(completion.message).toContain('never came back') expect(completeAsyncToolCall).toHaveBeenCalledTimes(1) expect(publishToolConfirmation).toHaveBeenCalledTimes(1) expect(JSON.stringify(completion)).not.toContain('late rejected secret output') diff --git a/apps/sim/lib/mothership/request/tools/executor.ts b/apps/sim/lib/mothership/request/tools/executor.ts index 4262b930af0..2a95b47b384 100644 --- a/apps/sim/lib/mothership/request/tools/executor.ts +++ b/apps/sim/lib/mothership/request/tools/executor.ts @@ -1,7 +1,8 @@ import { browserToolRendererTimeoutMs, isCurrentBrowserToolName } from '@sim/browser-protocol' import { createLogger } from '@sim/logger' +import { isTerminalToolName, resolveRunWaitMs } from '@sim/terminal-protocol' import { toError } from '@sim/utils/errors' -import { isRecordLike } from '@sim/utils/object' +import { isRecordLike, toRecord } from '@sim/utils/object' import { AsyncToolCallOwnershipError } from '@/lib/mothership/async-runs/errors' import type { AsyncCompletionEnvelope, @@ -9,7 +10,6 @@ import type { } from '@/lib/mothership/async-runs/lifecycle' import { type CompleteAsyncToolCallInput, - completeAsyncToolCall, markAsyncToolRunning, upsertAsyncToolCall, } from '@/lib/mothership/async-runs/repository' @@ -59,10 +59,7 @@ import { requireToolCallError, setTerminalToolCallState, } from '@/lib/mothership/request/tool-call-state' -import { - sealClientToolCompletion, - sealClientToolContext, -} from '@/lib/mothership/request/tools/client-completion-seal.server' +import { desktopToolNotStarted } from '@/lib/mothership/request/tools/desktop-wait' import { type ToolExecutionLifetime, withToolExecutionLifetime, @@ -78,6 +75,7 @@ import { maybeWriteOutputToTable, maybeWriteReadCsvToTable, } from '@/lib/mothership/request/tools/tables' +import { settleToolCallFailure } from '@/lib/mothership/request/tools/tool-call-failure' import { applyCreateWorkflowOutputToContext } from '@/lib/mothership/request/tools/workflow-context' import { type ExecutionContext, @@ -88,6 +86,7 @@ import { type ToolCallState, } from '@/lib/mothership/request/types' import { ensureHandlersRegistered, executeTool } from '@/lib/mothership/tool-executor' +import { isDesktopToolCall } from '@/lib/mothership/tools/desktop-tools' import { withSandboxResourceScope } from '@/lib/mothership/tools/sandbox-resources' import { isMcpTool } from '@/executor/constants' @@ -268,9 +267,23 @@ export function pendingToolWaitBudgetMs( if (executableName && isCurrentBrowserToolName(executableName)) { return browserToolRendererTimeoutMs(executableName, toolCall?.params) } + if ( + executableName && + isTerminalToolName(executableName) && + toolCall?.params?.operation === 'run' + ) { + const runWaitMs = resolveRunWaitMs(toRecord(toolCall.params.args).waitSeconds) + return Math.max(TOOL_WATCHDOG_DEFAULT_MS, runWaitMs + TERMINAL_RUN_REPORT_HEADROOM_MS) + } return toolWatchdogTimeoutMs(executableName) } +/** + * Past a `terminal_run`'s wait window the desktop still captures the output and reports it. The + * desktop holds a run for up to its requested wait, so the resume gate must not call it hung first. + */ +const TERMINAL_RUN_REPORT_HEADROOM_MS = 30_000 + /** * Bare timeout/abort messages (AbortSignal.timeout's "The operation timed out.", * a DOMException's "This operation was aborted") strip everything the model @@ -397,8 +410,10 @@ async function executeToolWithWatchdog( } } -const HUNG_TOOL_MESSAGE = - 'Tool execution hung and was abandoned so the conversation could continue. Its outcome is unknown; do not retry it automatically.' +const TOOL_RESULT_LOST_MESSAGE = + 'This tool started, but its result never came back, so it was abandoned to let the conversation continue. Its outcome is unknown: it may already have taken effect, so inspect the current state before repeating it, and do not retry it automatically.' +const DESKTOP_TOOL_RESULT_LOST_MESSAGE = + 'The Sim desktop app started this action, but its result never came back (the chat view closed or the app stopped responding). Its outcome is unknown: it may already have taken effect, so inspect the current state before repeating it, and do not retry it automatically.' const UNAVAILABLE_TOOL_SETTLEMENT_MESSAGE = 'The tool result could not be restored before the conversation resumed. Its outcome is unknown; do not retry it automatically.' @@ -406,81 +421,87 @@ const UNAVAILABLE_TOOL_SETTLEMENT_MESSAGE = * Settles an abandoned tool with a fixed server-owned failure. Client waiters consume the same * sealed transport as ordinary client completions; no abandoned tool content is certified. * Execution ownership remains held while retained work cleans up. + * + * Without an explicit `failureMessage` the model learns which of two things happened: a desktop + * call nothing claimed never started (`notStarted`, safe to retry), and anything else started and + * lost its result (`outcomeUnknown`, `doNotRetry`). */ export async function failPendingToolCall( toolCallId: string, context: StreamingContext, execContext: ExecutionContext, - failureMessage: string = HUNG_TOOL_MESSAGE + failureMessage?: string ): Promise { const toolCall = context.toolCalls.get(toolCallId) if (!toolCall || toolCall.endTime || isTerminalToolCallStatus(toolCall.status)) return - const failure = { error: failureMessage, outcomeUnknown: true, doNotRetry: true } - let durableData: unknown = failure - let completed = false - let lostSettlementRace = false - try { - if (context.runId && execContext.resolvedSecretTraceRegistry) { - const binding = { toolCallId, runId: context.runId, userId: execContext.userId } - const [completion, provenance] = await Promise.all([ - sealClientToolCompletion({ ...binding, message: failureMessage, data: failure }), - sealClientToolContext({ - ...binding, - registry: execContext.resolvedSecretTraceRegistry, - /** The fixed failure contains no output or arguments from the abandoned tool. */ - toolInput: undefined, - }), - ]) - durableData = { ...completion, ...provenance } - } - if (toolCall.endTime || isTerminalToolCallStatus(toolCall.status)) return - - completed = Boolean( - await completeAsyncToolCall({ - toolCallId, - status: MothershipStreamV1AsyncToolRecordStatus.failed, - result: durableData, - error: failureMessage, - }) - ) - if (!completed) { - if (toolCall.endTime || isTerminalToolCallStatus(toolCall.status)) return - lostSettlementRace = true - } - } catch (error) { - logger.warn('Failed to persist force-failed async tool status', { - toolCallId, - error: toError(error).message, + const desktopCall = isDesktopToolCall(toolCall.execName ?? toolCall.name, toolCall.params) + if (failureMessage === undefined && desktopCall) { + const settled = await settleAbandonedToolCall(toolCall, context, execContext, { + ...desktopToolNotStarted(), + unclaimedOnly: true, }) - if (toolCall.endTime || isTerminalToolCallStatus(toolCall.status)) return + if (settled) return + } + const message = + failureMessage ?? (desktopCall ? DESKTOP_TOOL_RESULT_LOST_MESSAGE : TOOL_RESULT_LOST_MESSAGE) + await settleAbandonedToolCall(toolCall, context, execContext, { + message, + data: { error: message, outcomeUnknown: true, doNotRetry: true }, + unclaimedOnly: false, + resultLost: failureMessage === undefined, + }) +} + +/** + * Settles one abandoned tool locally once its durable failure is decided. With `unclaimedOnly` + * it applies only while the call is still unclaimed and reports whether it did. + */ +async function settleAbandonedToolCall( + toolCall: ToolCallState, + context: StreamingContext, + execContext: ExecutionContext, + failure: { + message: string + data: Record + unclaimedOnly: boolean + resultLost?: boolean } +): Promise { + const settledLocally = () => + Boolean(toolCall.endTime || isTerminalToolCallStatus(toolCall.status)) + const settlement = await settleToolCallFailure({ + toolCallId: toolCall.id, + runId: context.runId, + userId: execContext.userId, + registry: execContext.resolvedSecretTraceRegistry, + message: failure.message, + data: failure.data, + unclaimedOnly: failure.unclaimedOnly, + settledLocally, + }) + if (settlement === 'superseded') return true + if (failure.unclaimedOnly && settlement !== 'settled') return false + if (settlement !== 'settled' && settledLocally()) return true /** A durable winner whose waiter is still hung must not become a fabricated local success. */ + const lostSettlementRace = settlement === 'lost' const message = - lostSettlementRace && failureMessage === HUNG_TOOL_MESSAGE - ? UNAVAILABLE_TOOL_SETTLEMENT_MESSAGE - : failureMessage + lostSettlementRace && failure.resultLost ? UNAVAILABLE_TOOL_SETTLEMENT_MESSAGE : failure.message setTerminalToolCallState(toolCall, { status: MothershipStreamV1ToolOutcome.error, - output: { ...failure, error: message }, + output: { ...failure.data, error: message }, error: message, }) logger.error('Tool call failed', { - toolCallId, + toolCallId: toolCall.id, toolName: toolCall.name, - persisted: completed, + persisted: settlement === 'settled', lostSettlementRace, + notStarted: failure.unclaimedOnly, }) - markToolResultSeen(context, toolCallId) - if (completed) { - publishTerminalToolConfirmation({ - toolCallId, - status: MothershipStreamV1ToolOutcome.error, - message: failureMessage, - data: durableData, - }) - } + markToolResultSeen(context, toolCall.id) + return true } function cancelledCompletion(message: string): AsyncCompletionSignal { @@ -604,7 +625,7 @@ export async function executeToolAndReport( toolCall.id, context, execContext, - admissionFailure ? message : HUNG_TOOL_MESSAGE + admissionFailure ? message : undefined ) return terminalCompletionFromToolCall(toolCall) } diff --git a/apps/sim/lib/mothership/request/tools/tool-call-failure.ts b/apps/sim/lib/mothership/request/tools/tool-call-failure.ts new file mode 100644 index 00000000000..e7ed181a2d6 --- /dev/null +++ b/apps/sim/lib/mothership/request/tools/tool-call-failure.ts @@ -0,0 +1,87 @@ +import { createLogger } from '@sim/logger' +import { toError } from '@sim/utils/errors' +import { + completeAsyncToolCall, + completePendingAsyncToolCall, +} from '@/lib/mothership/async-runs/repository' +import { + MothershipStreamV1AsyncToolRecordStatus, + MothershipStreamV1ToolOutcome, +} from '@/lib/mothership/generated/mothership-stream-v1' +import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confirm' +import { + sealClientToolCompletion, + sealClientToolContext, +} from '@/lib/mothership/request/tools/client-completion-seal.server' +import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' + +const logger = createLogger('CopilotToolCallFailure') + +/** + * How a server-owned failure ended: + * - `settled`: this write won and its result was published; + * - `lost`: another transition (a claim or a result) already moved the row; + * - `superseded`: the call settled locally while the failure was being sealed; + * - `failed`: the failure could not be persisted. + */ +type ToolCallFailureSettlement = 'settled' | 'lost' | 'superseded' | 'failed' + +/** + * Settles a call Sim gives up on with a fixed server-owned failure, sealed like a client + * completion so its waiter (in this process or another) restores it. No content of the abandoned + * call is certified. With `unclaimedOnly` the failure applies only while nothing has claimed the + * call, the inverse of the desktop's claim, so exactly one of the two wins. + */ +export async function settleToolCallFailure(input: { + toolCallId: string + runId?: string + userId: string + registry?: ResolvedSecretTraceRegistry + message: string + data: Record + unclaimedOnly: boolean + settledLocally?: () => boolean +}): Promise { + const { toolCallId, message, data } = input + try { + let durableData: unknown = data + if (input.runId && input.registry) { + const binding = { toolCallId, runId: input.runId, userId: input.userId } + const [completion, provenance] = await Promise.all([ + sealClientToolCompletion({ ...binding, message, data }), + sealClientToolContext({ + ...binding, + registry: input.registry, + /** The fixed failure contains no output or arguments from the abandoned call. */ + toolInput: undefined, + }), + ]) + durableData = { ...completion, ...provenance } + } + if (input.settledLocally?.()) return 'superseded' + const failure = { + toolCallId, + status: MothershipStreamV1AsyncToolRecordStatus.failed, + result: durableData, + error: message, + } + const completed = input.unclaimedOnly + ? await completePendingAsyncToolCall(failure) + : await completeAsyncToolCall(failure) + if (!completed) return 'lost' + publishToolConfirmation({ + toolCallId, + status: MothershipStreamV1ToolOutcome.error, + message, + data: durableData, + timestamp: new Date().toISOString(), + }) + return 'settled' + } catch (error) { + logger.warn('Failed to persist a server-owned tool failure', { + toolCallId, + error: toError(error).message, + }) + return 'failed' + } +} diff --git a/apps/sim/lib/mothership/request/tools/workflow-client-fallback.test.ts b/apps/sim/lib/mothership/request/tools/workflow-client-fallback.test.ts index 08373313ba2..cea0caa9dee 100644 --- a/apps/sim/lib/mothership/request/tools/workflow-client-fallback.test.ts +++ b/apps/sim/lib/mothership/request/tools/workflow-client-fallback.test.ts @@ -2,18 +2,22 @@ import { mothershipAsyncRunsMock, mothershipAsyncRunsMockFns, } from '@sim/testing/mocks/mothership-async-runs.mock' +import { + mothershipClientToolWaiterMock, + mothershipClientToolWaiterMockFns, +} from '@sim/testing/mocks/mothership-client-tool-waiter.mock' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' -const { waitForWorkflowToolCompletion, recordDegraded } = vi.hoisted(() => ({ - waitForWorkflowToolCompletion: vi.fn(), +const { recordDegraded } = vi.hoisted(() => ({ recordDegraded: vi.fn(), })) vi.mock('@/lib/mothership/request/metrics', () => ({ recordDegraded })) -vi.mock('@/lib/mothership/request/tools/client', () => ({ - waitForWorkflowToolCompletion, -})) +vi.mock('@/lib/mothership/request/tools/client', () => mothershipClientToolWaiterMock) + +const waitForWorkflowToolCompletion = + mothershipClientToolWaiterMockFns.mockWaitForWorkflowToolCompletion vi.mock('@/lib/mothership/async-runs/repository', () => mothershipAsyncRunsMock) diff --git a/apps/sim/lib/mothership/request/types.ts b/apps/sim/lib/mothership/request/types.ts index e09e8038afe..6bedec758ca 100644 --- a/apps/sim/lib/mothership/request/types.ts +++ b/apps/sim/lib/mothership/request/types.ts @@ -265,6 +265,13 @@ export interface OrchestratorOptions { * `interactive`, which is a trust classification, not executor routing. */ clientToolPickupExpected?: boolean + /** + * The turn's desktop claims local reads (`read_local_file`, user-local VFS reads) through + * authorize before reading, so they are persisted pending and fail fast when nothing picks them + * up. Absent for older desktops, and on a recovered leg, where local reads keep the established + * running state and the authorize route accepts them as before. + */ + desktopClaimsLocalReads?: boolean } export interface OrchestratorResult { diff --git a/apps/sim/lib/mothership/tools/client/desktop-tool-pickup.integration.ts b/apps/sim/lib/mothership/tools/client/desktop-tool-pickup.integration.ts new file mode 100644 index 00000000000..45b67358c7e --- /dev/null +++ b/apps/sim/lib/mothership/tools/client/desktop-tool-pickup.integration.ts @@ -0,0 +1,319 @@ +/** + * A desktop tool call reaches the user's machine only through the chat view that is showing that + * chat. When nothing picks the call up (the user is on another chat, another page, or the app is + * closed) the turn must learn quickly and truthfully that the call never started, instead of + * hanging until a watchdog calls it "hung". Runs against real PostgreSQL and Redis: the call is + * persisted and dispatched by the production stream handlers, the desktop claims it through the + * authorize route and reports through the confirm route. + */ +import { authMock, authMockFns } from '@sim/testing/mocks/auth.mock' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' + +const { redisUrl, inheritedEnv } = await vi.hoisted(async () => { + const { readTestRedisUrl } = await import('@sim/db/testing/test-infrastructure') + const url = readTestRedisUrl() + const inheritedEnv = { REDIS_URL: process.env.REDIS_URL } + /** The real Redis module and the confirmation channel read this at import. */ + if (url) process.env.REDIS_URL = url + return { redisUrl: url, inheritedEnv } +}) + +vi.mock('@/lib/auth', () => authMock) + +import { db } from '@sim/db' +import { copilotAsyncToolCalls, copilotChats, copilotRuns, user, workspace } from '@sim/db/schema' +import { sleep } from '@sim/utils/helpers' +import { generateId } from '@sim/utils/id' +import { eq } from 'drizzle-orm' +import { NextRequest } from 'next/server' +import { closeRedisConnection } from '@/lib/core/config/redis' +import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' +import { prePersistClientExecutableToolCall, sseHandlers } from '@/lib/mothership/request/handlers' +import { TraceCollector } from '@/lib/mothership/request/trace' +import type { StreamEvent, StreamingContext } from '@/lib/mothership/request/types' +import { POST as confirmPOST } from '@/app/api/copilot/confirm/route' +import { POST as authorizePOST } from '@/app/api/desktop/tool/authorize/route' + +const APP_ORIGIN = 'http://localhost:3000' +/** Longer than the pickup grace, far shorter than the budget the call used to wait out. */ +const TURN_WAIT_MS = 25_000 + +function post(handler: typeof authorizePOST, path: string, body: unknown): Promise { + return Promise.resolve( + handler( + new NextRequest(new URL(path, APP_ORIGIN), { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify(body), + }), + {} + ) + ) +} + +const desktopClaims = (toolCallId: string, claim?: true) => + post(authorizePOST, '/api/desktop/tool/authorize', { toolCallId, ...(claim ? { claim } : {}) }) + +afterAll(async () => { + const channels = globalThis as typeof globalThis & { + _toolConfirmationChannel?: { dispose(): void } + } + channels._toolConfirmationChannel?.dispose() + channels._toolConfirmationChannel = undefined + await closeRedisConnection() + for (const [key, value] of Object.entries(inheritedEnv)) { + if (value === undefined) delete process.env[key] + else process.env[key] = value + } +}) + +describe.runIf(Boolean(redisUrl))('a desktop tool call nobody picks up', () => { + const userId = generateId() + const workspaceId = generateId() + const chatId = generateId() + const runId = generateId() + + /** The agent issues a desktop call: the production pre-persist and dispatch path. */ + /** `desktopClaimsLocalReads`: the turn came from a desktop that claims its local reads. */ + async function agentCalls( + toolName: string, + args: Record, + desktopClaimsLocalReads = false + ) { + const toolCallId = generateId() + const context: StreamingContext = { + runId, + chatId, + messageId: generateId(), + accumulatedContent: '', + finalAssistantContent: '', + sawMainToolCall: false, + trace: new TraceCollector(), + contentBlocks: [], + toolCalls: new Map(), + pendingToolPromises: new Map(), + activeFileIntents: new Map(), + filePreviewBudget: { contentBytes: 0 }, + seenToolCalls: new Set(), + seenToolResults: new Set(), + currentThinkingBlock: null, + subagentThinkingBlocks: new Map(), + isInThinkingBlock: false, + subAgentContent: {}, + subAgentToolCalls: {}, + pendingContent: '', + streamComplete: false, + wasAborted: false, + errors: [], + toolPermissions: { enabled: false, autoAllowed: new Set(), autoAllowPermitted: true }, + } + const event: StreamEvent = { + type: 'tool', + payload: { + toolCallId, + toolName, + arguments: args, + executor: 'client', + mode: 'async', + phase: 'call', + }, + } + const options = { timeout: TURN_WAIT_MS, desktopClaimsLocalReads } + await prePersistClientExecutableToolCall(event, context, options) + await sseHandlers.tool( + event, + context, + { userId, chatId, workspaceId, workflowId: generateId() }, + options + ) + const answer = context.pendingToolPromises.get(toolCallId) + if (!answer) throw new Error('The desktop call was not dispatched to a client waiter') + return { toolCallId, answer } + } + + beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: userId, + name: 'Desktop pickup fixture', + email: `${userId}@desktop-pickup.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Desktop pickup fixture', + ownerId: userId, + billedAccountUserId: userId, + }) + await db.insert(copilotChats).values({ + id: chatId, + userId, + workspaceId, + type: 'mothership', + conversationId: generateId(), + }) + await db.insert(copilotRuns).values({ + id: runId, + executionId: generateId(), + chatId, + userId, + workspaceId, + streamId: generateId(), + toolExecutionVersion: SIM_TOOL_EXECUTION_VERSION, + status: 'paused_waiting_for_tool', + requestContext: { source: 'headless_lifecycle' }, + }) + authMockFns.mockGetSession.mockResolvedValue({ + user: { id: userId, email: `${userId}@desktop-pickup.test`, name: 'Desktop pickup' }, + session: { id: generateId(), userId }, + }) + }) + + afterAll(async () => { + await db.delete(copilotChats).where(eq(copilotChats.id, chatId)) + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(eq(user.id, userId)) + }) + + it.each([ + ['browser_snapshot', {}], + ['terminal', { operation: 'run', args: { command: 'ls' } }], + ])( + 'tells the agent a %s call never started, and refuses a late pickup', + async (toolName, args) => { + const startedAt = Date.now() + const { toolCallId, answer } = await agentCalls(toolName, args) + + const completion = await answer + + expect(Date.now() - startedAt).toBeLessThan(TURN_WAIT_MS - 5_000) + expect(completion).toMatchObject({ + status: 'error', + data: { notStarted: true, error: expect.stringContaining('safe to retry') }, + }) + const [row] = await db + .select({ status: copilotAsyncToolCalls.status }) + .from(copilotAsyncToolCalls) + .where(eq(copilotAsyncToolCalls.toolCallId, toolCallId)) + expect(row.status).toBe('failed') + expect((await desktopClaims(toolCallId)).ok).toBe(false) + }, + TURN_WAIT_MS + 10_000 + ) + + it( + 'keeps waiting for a call the desktop picked up in time, past the pickup grace', + async () => { + const { toolCallId, answer } = await agentCalls('browser_snapshot', {}) + expect((await desktopClaims(toolCallId)).status).toBe(200) + + await sleep(17_000) + const report = await post(confirmPOST, '/api/copilot/confirm', { + toolCallId, + status: 'success', + message: 'Done', + data: { text: 'page' }, + }) + + expect(report.status).toBe(200) + expect(await answer).toMatchObject({ status: 'success' }) + }, + TURN_WAIT_MS + 10_000 + ) + + it( + 'never lets a stale "not started" report end a call the desktop is running', + async () => { + const { toolCallId, answer } = await agentCalls('terminal', { + operation: 'run', + args: { command: 'bun run test' }, + }) + expect((await desktopClaims(toolCallId)).status).toBe(200) + + const staleReport = await post(confirmPOST, '/api/copilot/confirm', { + toolCallId, + status: 'error', + message: 'Not run: delivered too late', + data: { error: 'Not run: delivered too late', notStarted: true }, + }) + const realResult = await post(confirmPOST, '/api/copilot/confirm', { + toolCallId, + status: 'success', + message: 'Done', + data: { output: 'tests passed' }, + }) + + expect(staleReport.ok).toBe(false) + expect(realResult.status).toBe(200) + expect(await answer).toMatchObject({ status: 'success' }) + }, + TURN_WAIT_MS + 10_000 + ) + + it( + 'tells the agent a local read nobody picked up never started, for a desktop that claims reads', + async () => { + const startedAt = Date.now() + const { toolCallId, answer } = await agentCalls( + 'read_local_file', + { path: '/Users/me/notes.txt' }, + true + ) + + const completion = await answer + + expect(Date.now() - startedAt).toBeLessThan(TURN_WAIT_MS - 5_000) + expect(completion).toMatchObject({ status: 'error', data: { notStarted: true } }) + expect((await desktopClaims(toolCallId, true)).ok).toBe(false) + }, + TURN_WAIT_MS + 10_000 + ) + + it( + 'lets a desktop that claimed a local read keep reading it, past the pickup grace', + async () => { + const { toolCallId, answer } = await agentCalls( + 'grep', + { path: 'user-local/Project--mount-1', pattern: 'TODO' }, + true + ) + + expect((await desktopClaims(toolCallId, true)).status).toBe(200) + await sleep(17_000) + expect((await desktopClaims(toolCallId, true)).status).toBe(200) + const report = await post(confirmPOST, '/api/copilot/confirm', { + toolCallId, + status: 'success', + message: 'Done', + data: { matches: [] }, + }) + + expect(report.status).toBe(200) + expect(await answer).toMatchObject({ status: 'success' }) + }, + TURN_WAIT_MS + 10_000 + ) + + it( + 'reads as before for a desktop that does not claim local reads', + async () => { + const { toolCallId, answer } = await agentCalls('read_local_file', { + path: '/Users/me/notes.txt', + }) + + expect((await desktopClaims(toolCallId)).status).toBe(200) + const report = await post(confirmPOST, '/api/copilot/confirm', { + toolCallId, + status: 'success', + message: 'Done', + data: { text: 'notes' }, + }) + + expect(report.status).toBe(200) + expect(await answer).toMatchObject({ status: 'success' }) + }, + TURN_WAIT_MS + 10_000 + ) +}) diff --git a/apps/sim/lib/mothership/tools/client/terminal-tool-execution.test.ts b/apps/sim/lib/mothership/tools/client/terminal-tool-execution.test.ts index 46b50fb8814..95cc1f2d703 100644 --- a/apps/sim/lib/mothership/tools/client/terminal-tool-execution.test.ts +++ b/apps/sim/lib/mothership/tools/client/terminal-tool-execution.test.ts @@ -52,4 +52,23 @@ describe('terminal client execution', () => { window.dispatchEvent(new Event('pagehide')) expect(beacon).toHaveBeenCalledOnce() }) + + it('reports a call delivered too late as never started instead of dropping it', async () => { + const emittedAt = new Date(Date.now() - 5 * 60_000).toISOString() + + executeTerminalToolOnClient( + 'terminal-stale', + { operation: 'run', args: { command: 'bun run test' } }, + 'chat-1', + emittedAt + ) + await sleep(0) + + expect(reportClientToolCompletion).toHaveBeenCalledWith( + 'terminal-stale', + 'error', + expect.stringContaining('safe to retry'), + expect.objectContaining({ notStarted: true }) + ) + }) }) diff --git a/apps/sim/lib/mothership/tools/client/terminal-tool-execution.ts b/apps/sim/lib/mothership/tools/client/terminal-tool-execution.ts index 8cb295c85ca..8844c66acc1 100644 --- a/apps/sim/lib/mothership/tools/client/terminal-tool-execution.ts +++ b/apps/sim/lib/mothership/tools/client/terminal-tool-execution.ts @@ -25,6 +25,8 @@ const logger = createLogger('CopilotTerminalToolExecution') /** Tool events older than this are replays, not live instructions. */ const MAX_EVENT_AGE_MS = 120_000 +const STALE_EVENT_MESSAGE = + 'Not run: this terminal call reached the Sim desktop app too late to start safely, so nothing ran on the user’s computer. It is safe to retry.' const EXECUTED_STORAGE_PREFIX = 'sim:copilot:terminal-tool-executed:' /** @@ -120,7 +122,23 @@ export function executeTerminalToolOnClient( } const age = eventAgeMs(eventTs) if (age !== null && age > MAX_EVENT_AGE_MS) { - logger.info('Skipping stale terminal tool event', { toolCallId, operation, age }) + logger.info('Reporting stale terminal tool event as not started', { + toolCallId, + operation, + age, + }) + // Reported against the pending call only: one another window already claimed keeps its result. + void reportClientToolCompletion( + toolCallId, + ASYNC_TOOL_CONFIRMATION_STATUS.error, + STALE_EVENT_MESSAGE, + { error: STALE_EVENT_MESSAGE, notStarted: true, staleEvent: true } + ).catch((error) => { + logger.warn('Failed to report stale terminal tool event', { + toolCallId, + error: toError(error).message, + }) + }) return } markExecuted(toolCallId) diff --git a/apps/sim/lib/mothership/tools/desktop-tools.ts b/apps/sim/lib/mothership/tools/desktop-tools.ts index 9027795d190..7d56d111dd2 100644 --- a/apps/sim/lib/mothership/tools/desktop-tools.ts +++ b/apps/sim/lib/mothership/tools/desktop-tools.ts @@ -34,6 +34,28 @@ export function getDesktopToolClaimOwner(toolName: string): DesktopToolClaimOwne return undefined } +/** A read of the user's machine: `read_local_file`, or a VFS read of a granted local folder. */ +export function isLocalReadToolCall(toolName: string, args: Record | undefined) { + return toolName === 'read_local_file' || isUserLocalVfsToolCall(toolName, args) +} + +/** + * Whether the server sees this call's pickup: the desktop claims it, pending, through + * `/api/desktop/tool/authorize` before acting, so a call still pending has provably not started. + * Browser, terminal and import calls are always claimed; a local read only by a desktop that + * declared for the turn that it claims local reads. + */ +export function isClaimedOnPickup( + toolName: string, + args: Record | undefined, + desktopClaimsLocalReads: boolean +): boolean { + return ( + getDesktopToolClaimOwner(toolName) !== undefined || + (desktopClaimsLocalReads && isLocalReadToolCall(toolName, args)) + ) +} + /** What the model learns about a desktop call that Stop cancelled before the desktop picked it up. */ export const STOPPED_BEFORE_START_MESSAGE = 'Not run: the user stopped the chat before the Sim desktop app started this action. Nothing happened on their computer.' diff --git a/packages/desktop-bridge/src/index.ts b/packages/desktop-bridge/src/index.ts index 5a922d7a698..3592a6e3f30 100644 --- a/packages/desktop-bridge/src/index.ts +++ b/packages/desktop-bridge/src/index.ts @@ -1103,6 +1103,12 @@ export interface SimDesktopApi { localFilesystem(request: LocalFilesystemRequest): Promise /** Optional so older installed shells do not advertise the new native tools. */ localFiles?(request: DesktopLocalFileRequest): Promise + /** + * Present on shells that claim each local read (`read_local_file`, user-local VFS reads) through + * the server's tool authorization before reading, as they already do for imports. A turn started + * from such a shell persists those calls pending, so a read nobody picks up fails fast. + */ + localReadClaims?: true /** Subscribe to commands initiated by the native application menu. */ onCommand(callback: (command: DesktopCommand) => void): () => void windowState: SimDesktopWindowStateApi diff --git a/packages/terminal-protocol/src/index.ts b/packages/terminal-protocol/src/index.ts index 50ff3ae8095..69b57877722 100644 --- a/packages/terminal-protocol/src/index.ts +++ b/packages/terminal-protocol/src/index.ts @@ -86,9 +86,17 @@ export const MAX_CAPTURE_CHARS = 512_000 * `terminal_read`. Successive reads are also how it tells progress from a * stall — output that stops changing is a command waiting on input or wedged. */ -export const DEFAULT_RUN_WAIT_MS = 30_000 +const DEFAULT_RUN_WAIT_MS = 30_000 -export const MAX_RUN_WAIT_MS = 120_000 +const MAX_RUN_WAIT_MS = 120_000 + +/** How long one `terminal_run` holds the turn, from the `waitSeconds` the model asked for. */ +export function resolveRunWaitMs(waitSeconds: unknown): number { + const requested = Number(waitSeconds) + return Number.isFinite(requested) && requested > 0 + ? Math.min(requested * 1000, MAX_RUN_WAIT_MS) + : DEFAULT_RUN_WAIT_MS +} /** * How long output must be silent, with the cursor left mid-line, before the diff --git a/packages/testing/src/mocks/index.ts b/packages/testing/src/mocks/index.ts index 1f8cee67fec..4332a76b106 100644 --- a/packages/testing/src/mocks/index.ts +++ b/packages/testing/src/mocks/index.ts @@ -529,6 +529,10 @@ export { mothershipChatStatusMock, mothershipChatStatusMockFns, } from './mothership-chat-status.mock' +export { + mothershipClientToolWaiterMock, + mothershipClientToolWaiterMockFns, +} from './mothership-client-tool-waiter.mock' export { mothershipEnvironmentContextMock, mothershipEnvironmentContextMockFns, diff --git a/packages/testing/src/mocks/mothership-client-tool-waiter.mock.ts b/packages/testing/src/mocks/mothership-client-tool-waiter.mock.ts new file mode 100644 index 00000000000..15c73966f77 --- /dev/null +++ b/packages/testing/src/mocks/mothership-client-tool-waiter.mock.ts @@ -0,0 +1,31 @@ +import { vi } from 'vitest' + +/** + * Controllable mock functions for `@/lib/mothership/request/tools/client`, the waiters that + * resolve a client-executed tool call from its confirmation. Each defaults to resolving `null` + * (no result before the wait ended). + * + * @example + * ```ts + * import { + * mothershipClientToolWaiterMock, + * mothershipClientToolWaiterMockFns, + * } from '@sim/testing/mocks/mothership-client-tool-waiter.mock' + * + * vi.mock('@/lib/mothership/request/tools/client', () => mothershipClientToolWaiterMock) + * mothershipClientToolWaiterMockFns.mockWaitForClientToolCompletion.mockResolvedValueOnce(result) + * ``` + */ +export const mothershipClientToolWaiterMockFns = { + mockWaitForToolCompletion: vi.fn(async (..._args: unknown[]): Promise => null), + mockWaitForClientToolCompletion: vi.fn(async (..._args: unknown[]): Promise => null), + mockWaitForWorkflowToolCompletion: vi.fn(async (..._args: unknown[]): Promise => null), +} + +/** Static mock module for `@/lib/mothership/request/tools/client`. Covers every runtime export. */ +export const mothershipClientToolWaiterMock = { + waitForToolCompletion: mothershipClientToolWaiterMockFns.mockWaitForToolCompletion, + waitForClientToolCompletion: mothershipClientToolWaiterMockFns.mockWaitForClientToolCompletion, + waitForWorkflowToolCompletion: + mothershipClientToolWaiterMockFns.mockWaitForWorkflowToolCompletion, +} From 89626b1cb2ca11d13555bc62a417c74cfe15a247 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 5 Oct 2026 19:40:30 -0700 Subject: [PATCH 2/4] refactor(mothership): one desktop-tool classifier and a typed claim for each claimant The authorize and confirm routes classify desktop tools through lib/mothership/tools/desktop-tools.ts instead of inline checks. The Sim execution claim and the desktop claim share one run-admission lock but return their own outcome types, so a Sim caller can no longer receive the desktop-only awaiting_permission outcome. The terminal-status mapping moves to lifecycle as getTerminalConfirmationStatus, since confirm uses it for every client tool. --- apps/sim/app/api/copilot/confirm/route.ts | 19 +- .../api/desktop/tool/authorize/route.test.ts | 6 +- .../app/api/desktop/tool/authorize/route.ts | 42 ++- .../lib/mothership/async-runs/lifecycle.ts | 9 + .../lib/mothership/async-runs/repository.ts | 242 ++++++++++-------- .../{tool-call-failure.ts => call-failure.ts} | 0 .../lib/mothership/request/tools/client.ts | 8 +- .../mothership/request/tools/desktop-wait.ts | 2 +- .../lib/mothership/request/tools/executor.ts | 2 +- .../tools/workflow-client-settlement.ts | 4 +- .../desktop-tool-authorization.integration.ts | 4 +- .../sim/lib/mothership/tools/desktop-tools.ts | 14 +- .../lib/mothership/tools/workflow-tools.ts | 8 - .../src/mocks/mothership-async-runs.mock.ts | 2 + 14 files changed, 203 insertions(+), 159 deletions(-) rename apps/sim/lib/mothership/request/tools/{tool-call-failure.ts => call-failure.ts} (100%) diff --git a/apps/sim/app/api/copilot/confirm/route.ts b/apps/sim/app/api/copilot/confirm/route.ts index d389a63b839..f80fa0e7312 100644 --- a/apps/sim/app/api/copilot/confirm/route.ts +++ b/apps/sim/app/api/copilot/confirm/route.ts @@ -1,7 +1,5 @@ import type { Span } from '@opentelemetry/api' -import { isBrowserToolName } from '@sim/browser-protocol' import { createLogger } from '@sim/logger' -import { isTerminalToolName } from '@sim/terminal-protocol' import { getErrorMessage, toError } from '@sim/utils/errors' import { isPlainRecord } from '@sim/utils/object' import { type NextRequest, NextResponse } from 'next/server' @@ -14,6 +12,7 @@ import { type AsyncCompletionData, type AsyncConfirmationStatus, type AsyncTerminalStatus, + getTerminalConfirmationStatus, isDeliveredAsyncStatus, isTerminalAsyncStatus, isWorkflowToolExecutionClaimable, @@ -41,12 +40,11 @@ import { import { withIncomingGoSpan } from '@/lib/mothership/request/otel' import { sealClientToolSettlement } from '@/lib/mothership/request/tools/client-completion-seal.server' import { isWorkflowToolName } from '@/lib/mothership/tools/client-executed-tools' -import { getDesktopToolClaimOwner } from '@/lib/mothership/tools/desktop-tools' +import { getDesktopToolClaimOwner, isNativeDesktopTool } from '@/lib/mothership/tools/desktop-tools' import { createStructuralWorkflowToolCompletionData, getWorkflowToolCompletionExecutionId, getWorkflowToolCompletionMessage, - getWorkflowToolConfirmationStatus, getWorkflowToolLaunchError, resolveWorkflowToolTargetId, WORKFLOW_EXECUTION_BUSY, @@ -92,7 +90,7 @@ function acknowledgeSettledToolCall( toolCallId: string, storedStatus: AsyncTerminalStatus ): NextResponse { - const settledStatus = getWorkflowToolConfirmationStatus(storedStatus) + const settledStatus = getTerminalConfirmationStatus(storedStatus) span.setAttributes({ [TraceAttr.ToolConfirmationStatus]: settledStatus, [TraceAttr.CopilotConfirmOutcome]: CopilotConfirmOutcome.Delivered, @@ -264,7 +262,7 @@ export const POST = withRouteHandler((req: NextRequest) => { return createNotFoundResponse('Completed workflow execution not found') } - const terminalStatus = getWorkflowToolConfirmationStatus(existing.status) + const terminalStatus = getTerminalConfirmationStatus(existing.status) span.setAttributes({ [TraceAttr.ToolConfirmationStatus]: terminalStatus, [TraceAttr.CopilotConfirmOutcome]: CopilotConfirmOutcome.Delivered, @@ -308,10 +306,7 @@ export const POST = withRouteHandler((req: NextRequest) => { const isErrorOrCancelledOutcome = status === ASYNC_TOOL_CONFIRMATION_STATUS.error || status === ASYNC_TOOL_CONFIRMATION_STATUS.cancelled - const isNativeClientTool = - isBrowserToolName(existing.toolName) || - isTerminalToolName(existing.toolName) || - existing.toolName === 'import_local_files' + const isNativeClientTool = isNativeDesktopTool(existing.toolName) const nativeClaimOwner = getDesktopToolClaimOwner(existing.toolName) const isPreclaimNativeTerminalOutcome = nativeClaimOwner !== undefined && @@ -376,7 +371,7 @@ export const POST = withRouteHandler((req: NextRequest) => { executionId = claimedExecutionId if (status !== ASYNC_TOOL_CONFIRMATION_STATUS.background) { if (trustedExecution) { - effectiveStatus = getWorkflowToolConfirmationStatus(trustedExecution.status) + effectiveStatus = getTerminalConfirmationStatus(trustedExecution.status) } else if (!isErrorOrCancelledOutcome) { span.setAttribute( TraceAttr.CopilotConfirmOutcome, @@ -389,7 +384,7 @@ export const POST = withRouteHandler((req: NextRequest) => { executionId = submittedExecutionId } else if (trustedExecution) { executionId = trustedExecution.executionId - effectiveStatus = getWorkflowToolConfirmationStatus(trustedExecution.status) + effectiveStatus = getTerminalConfirmationStatus(trustedExecution.status) } else if (!isErrorOrCancelledOutcome) { effectiveStatus = ASYNC_TOOL_CONFIRMATION_STATUS.error executionId = undefined diff --git a/apps/sim/app/api/desktop/tool/authorize/route.test.ts b/apps/sim/app/api/desktop/tool/authorize/route.test.ts index 55ac77cb407..df6f238d96d 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.test.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.test.ts @@ -17,7 +17,7 @@ vi.mock('@/lib/mothership/async-runs/repository', () => mothershipAsyncRunsMock) import { POST } from './route' -const claimToolExecution = mothershipAsyncRunsMockFns.mockClaimToolExecution +const claimDesktopToolCall = mothershipAsyncRunsMockFns.mockClaimToolExecution const getAsyncToolCall = mothershipAsyncRunsMockFns.mockGetAsyncToolCall const getRunSegment = mothershipAsyncRunsMockFns.mockGetRunSegment const resolveInvocationWorkspace = mothershipWorkspaceTargetMockFns.mockResolveInvocationWorkspace @@ -50,7 +50,7 @@ describe('desktop tool authorization', () => { userId: 'user-1', status: 'active', }) - claimToolExecution.mockResolvedValue({ outcome: 'claimed' }) + claimDesktopToolCall.mockResolvedValue({ outcome: 'claimed' }) }) it('never returns presentation activity as an executable browser argument', async () => { @@ -193,7 +193,7 @@ describe('desktop tool authorization', () => { new OrchestrationError('not_found', 'Workspace not found') ) expect((await POST(request('import-1', true))).status).toBe(404) - claimToolExecution.mockResolvedValueOnce({ outcome: 'existing' }) + claimDesktopToolCall.mockResolvedValueOnce({ outcome: 'existing' }) expect((await POST(request('import-1', true))).status).toBe(409) }) diff --git a/apps/sim/app/api/desktop/tool/authorize/route.ts b/apps/sim/app/api/desktop/tool/authorize/route.ts index 37d2c5d41e9..7ed3bdc2f84 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.ts @@ -1,5 +1,3 @@ -import { isCurrentBrowserToolName } from '@sim/browser-protocol' -import { isTerminalToolName } from '@sim/terminal-protocol' import { isRecordLike, omit } from '@sim/utils/object' import { type NextRequest, NextResponse } from 'next/server' import { authorizeDesktopToolContract } from '@/lib/api/contracts/desktop-tool-authorization' @@ -9,18 +7,21 @@ import { withRouteHandler } from '@/lib/core/utils/with-route-handler' import { resolveInvocationWorkspace } from '@/lib/mothership/application/workspace-target' import { DESKTOP_TOOL_CLAIM_OWNER } from '@/lib/mothership/async-runs/lifecycle' import { - claimToolExecution, + claimDesktopToolCall, + type DesktopToolCallClaim, getAsyncToolCall, getRunSegment, - type ToolExecutionClaim, } from '@/lib/mothership/async-runs/repository' import { authenticateCopilotRequestSessionOnly, createNotFoundResponse, createUnauthorizedResponse, } from '@/lib/mothership/request/http' -import { isLocalReadToolCall } from '@/lib/mothership/tools/desktop-tools' -import { isUserLocalVfsToolCall } from '@/lib/mothership/tools/local-filesystem' +import { + getDesktopToolClaimOwner, + isDesktopToolCall, + isLocalReadToolCall, +} from '@/lib/mothership/tools/desktop-tools' const admissionClosedResponse = () => NextResponse.json( @@ -30,7 +31,7 @@ const admissionClosedResponse = () => /** A refused claim answers the same way for every tool, except how each reports a lost race. */ function refusedClaimResponse( - claim: Exclude, + claim: Exclude, notPending: () => NextResponse ): NextResponse { if (claim === 'closed') return admissionClosedResponse() @@ -75,16 +76,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => { } const args = isRecordLike(toolCall.args) ? (toolCall.args as Record) : {} - const isBrowserTool = isCurrentBrowserToolName(toolCall.toolName) - const isTerminalTool = isTerminalToolName(toolCall.toolName) - const isLocalFileTool = - toolCall.toolName === 'read_local_file' || toolCall.toolName === 'import_local_files' - const authorized = - isBrowserTool || - isTerminalTool || - isLocalFileTool || - isUserLocalVfsToolCall(toolCall.toolName, args) - if (!authorized) { + if (!isDesktopToolCall(toolCall.toolName, args)) { return NextResponse.json( { error: 'Tool call is not authorized for desktop execution' }, { status: 403 } @@ -116,7 +108,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => { { status: 409 } ) if (toolCall.status !== 'pending') return alreadyStarted() - const { outcome } = await claimToolExecution({ + const { outcome } = await claimDesktopToolCall({ toolCallId: toolCall.toolCallId, runId: toolCall.runId, userId, @@ -136,7 +128,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => { if (parsed.data.body.claim && isLocalReadToolCall(toolCall.toolName, args)) { const notPending = () => createNotFoundResponse('Pending client tool call not found') if (toolCall.status === 'pending') { - const { outcome } = await claimToolExecution({ + const { outcome } = await claimDesktopToolCall({ toolCallId: toolCall.toolCallId, runId: toolCall.runId, userId, @@ -151,16 +143,18 @@ export const POST = withRouteHandler(async (request: NextRequest) => { // machine, so the pending call is claimed here, atomically, before crossing // the Electron boundary — a replayed renderer event must not run a command // or click a button twice. - if (isBrowserTool || isTerminalTool) { + const actionClaimOwner = getDesktopToolClaimOwner(toolCall.toolName) + if ( + actionClaimOwner === DESKTOP_TOOL_CLAIM_OWNER.browser || + actionClaimOwner === DESKTOP_TOOL_CLAIM_OWNER.terminal + ) { const notPending = () => createNotFoundResponse('Pending client tool call not found') if (toolCall.status !== 'pending') return notPending() - const { outcome } = await claimToolExecution({ + const { outcome } = await claimDesktopToolCall({ toolCallId: toolCall.toolCallId, runId: toolCall.runId, userId, - claimedBy: isBrowserTool - ? DESKTOP_TOOL_CLAIM_OWNER.browser - : DESKTOP_TOOL_CLAIM_OWNER.terminal, + claimedBy: actionClaimOwner, }) if (outcome !== 'claimed') return refusedClaimResponse(outcome, notPending) } diff --git a/apps/sim/lib/mothership/async-runs/lifecycle.ts b/apps/sim/lib/mothership/async-runs/lifecycle.ts index 777a9656797..04792678f58 100644 --- a/apps/sim/lib/mothership/async-runs/lifecycle.ts +++ b/apps/sim/lib/mothership/async-runs/lifecycle.ts @@ -128,6 +128,15 @@ export function isWorkflowToolExecutionClaimable( ) } +/** The confirmation status a settled call reports: its durable terminal status, as the wire names it. */ +export function getTerminalConfirmationStatus( + status: AsyncTerminalStatus +): AsyncConfirmationStatus { + if (status === ASYNC_TOOL_STATUS.completed) return ASYNC_TOOL_CONFIRMATION_STATUS.success + if (status === ASYNC_TOOL_STATUS.cancelled) return ASYNC_TOOL_CONFIRMATION_STATUS.cancelled + return ASYNC_TOOL_CONFIRMATION_STATUS.error +} + export function isTerminalAsyncStatus( status: CopilotAsyncToolStatus | AsyncLifecycleStatus | string | null | undefined ): status is AsyncTerminalStatus { diff --git a/apps/sim/lib/mothership/async-runs/repository.ts b/apps/sim/lib/mothership/async-runs/repository.ts index ee01e836e59..fc68fc5cc79 100644 --- a/apps/sim/lib/mothership/async-runs/repository.ts +++ b/apps/sim/lib/mothership/async-runs/repository.ts @@ -633,46 +633,27 @@ export async function markAsyncToolRunning(toolCallId: string, claimedBy: string return markAsyncToolStatus(toolCallId, 'running', { claimedBy }) } -export type ToolExecutionClaim = - | { outcome: 'claimed' } - /** Stop, a newer turn, or the run's end closed tool admission. */ - | { outcome: 'closed' } - /** Already claimed or settled. */ - | { outcome: 'existing' } - /** Held for the user's decision, and not allowed (yet). */ - | { outcome: 'awaiting_permission' } - -/** The desktop app acting on a pending call it was handed, before any effect on the machine. */ -export interface DesktopToolExecutionClaimant { - toolCallId: string - runId: string - userId: string - claimedBy: DesktopToolClaimOwner -} - /** - * Claims a tool call exactly once, serialized with Stop by the run row lock: nothing is claimed - * once the run's tool admission has closed. Sim claims a call it is about to execute under an - * execution lease, and a terminal tool result never releases that claim. The desktop app claims a - * pending call before crossing the Electron boundary, so a replayed renderer event cannot click, - * type, or run a command twice, and a call held for the user's decision only once they allowed it. + * Locks the call's run and runs `claim` only while the run still admits tools, so a claim and Stop + * serialize on the run row: nothing is claimed once tool admission has closed. */ -export async function claimToolExecution( - input: SimToolExecutionOwner | DesktopToolExecutionClaimant -): Promise { - const simOwner = 'ownerToken' in input ? input : undefined - const claimedBy = 'ownerToken' in input ? 'sim-stream' : input.claimedBy +async function claimUnderRunAdmission( + call: { toolCallId: string; runId: string; userId: string }, + claimedBy: string, + requireCurrentVersion: boolean, + claim: (tx: RunAdmissionTransaction, thisCall: SQL | undefined) => Promise +): Promise { return await withDbSpan( TraceSpan.CopilotAsyncRunsMarkAsyncToolStatus, 'UPDATE', 'copilot_async_tool_calls', { - [TraceAttr.ToolCallId]: input.toolCallId, - [TraceAttr.RunId]: input.runId, + [TraceAttr.ToolCallId]: call.toolCallId, + [TraceAttr.RunId]: call.runId, [TraceAttr.CopilotAsyncToolClaimedBy]: claimedBy, }, () => - traceMothershipTransaction('claim_tool', async (tx) => { + traceMothershipTransaction('claim_tool', async (tx) => { const [run] = await traceMothershipQuery('SELECT FOR UPDATE', 'copilot_runs', () => tx .select({ @@ -681,86 +662,145 @@ export async function claimToolExecution( status: copilotRuns.status, }) .from(copilotRuns) - .where(and(eq(copilotRuns.id, input.runId), eq(copilotRuns.userId, input.userId))) + .where(and(eq(copilotRuns.id, call.runId), eq(copilotRuns.userId, call.userId))) .for('update') ) - if (!run || (simOwner && run.toolExecutionVersion !== SIM_TOOL_EXECUTION_VERSION)) + if ( + !run || + (requireCurrentVersion && run.toolExecutionVersion !== SIM_TOOL_EXECUTION_VERSION) + ) throw new Error('Tool execution ownership is unavailable for this run') if (run.toolAdmissionClosedAt || TERMINAL_RUN_STATUSES.includes(run.status)) return { outcome: 'closed' } - const startedAt = new Date() - const thisCall = and( - eq(copilotAsyncToolCalls.toolCallId, input.toolCallId), - eq(copilotAsyncToolCalls.runId, input.runId) - ) - const [claimed] = await traceMothershipQuery('UPDATE', 'copilot_async_tool_calls', () => - tx - .update(copilotAsyncToolCalls) - .set( - simOwner - ? { - status: ASYNC_TOOL_STATUS.running, - claimedBy, - claimedAt: startedAt, - executionStartedAt: startedAt, - executionOwnerToken: simOwner.ownerToken, - executionLeaseExpiresAt: sql`clock_timestamp() + ${SIM_TOOL_EXECUTION_LEASE_SECONDS} * interval '1 second'`, - updatedAt: startedAt, - } - : { - status: ASYNC_TOOL_STATUS.running, - claimedBy, - claimedAt: startedAt, - updatedAt: startedAt, - } - ) - .where( - simOwner - ? and( - thisCall, - isNull(copilotAsyncToolCalls.executionStartedAt), - inArray(copilotAsyncToolCalls.status, [ - ASYNC_TOOL_STATUS.pending, - ASYNC_TOOL_STATUS.running, - ]) - ) - : and( - thisCall, - eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.pending), - or( - and( - isNull(copilotAsyncToolCalls.permissionRequestedAt), - isNull(copilotAsyncToolCalls.permissionDecision) - ), - inArray(copilotAsyncToolCalls.permissionDecision, [ - ...EXECUTABLE_TOOL_PERMISSION_DECISIONS, - ]) - ) - ) - ) - .returning({ id: copilotAsyncToolCalls.id }) - ) - if (claimed) return { outcome: 'claimed' } - const [record] = await traceMothershipQuery('SELECT', 'copilot_async_tool_calls', () => - tx - .select({ - status: copilotAsyncToolCalls.status, - permissionRequestedAt: copilotAsyncToolCalls.permissionRequestedAt, - permissionDecision: copilotAsyncToolCalls.permissionDecision, - }) - .from(copilotAsyncToolCalls) - .where(thisCall) + return claim( + tx, + and( + eq(copilotAsyncToolCalls.toolCallId, call.toolCallId), + eq(copilotAsyncToolCalls.runId, call.runId) + ) ) - if (!record) throw new Error('Tool execution record is unavailable') - return !simOwner && - record.status === ASYNC_TOOL_STATUS.pending && - isAwaitingToolPermission(record) - ? { outcome: 'awaiting_permission' } - : { outcome: 'existing' } }) ) } +export type ToolExecutionClaim = + | { outcome: 'claimed' } + /** Stop, a newer turn, or the run's end closed tool admission. */ + | { outcome: 'closed' } + /** Already claimed or settled. */ + | { outcome: 'existing' } + +/** Sim claims a call it is about to execute under an execution lease; a terminal result never releases it. */ +export async function claimToolExecution( + owner: SimToolExecutionOwner +): Promise { + return await claimUnderRunAdmission(owner, 'sim-stream', true, async (tx, thisCall) => { + const startedAt = new Date() + const [claimed] = await traceMothershipQuery('UPDATE', 'copilot_async_tool_calls', () => + tx + .update(copilotAsyncToolCalls) + .set({ + status: ASYNC_TOOL_STATUS.running, + claimedBy: 'sim-stream', + claimedAt: startedAt, + executionStartedAt: startedAt, + executionOwnerToken: owner.ownerToken, + executionLeaseExpiresAt: sql`clock_timestamp() + ${SIM_TOOL_EXECUTION_LEASE_SECONDS} * interval '1 second'`, + updatedAt: startedAt, + }) + .where( + and( + thisCall, + isNull(copilotAsyncToolCalls.executionStartedAt), + inArray(copilotAsyncToolCalls.status, [ + ASYNC_TOOL_STATUS.pending, + ASYNC_TOOL_STATUS.running, + ]) + ) + ) + .returning({ id: copilotAsyncToolCalls.id }) + ) + if (claimed) return { outcome: 'claimed' } as const + const [record] = await traceMothershipQuery('SELECT', 'copilot_async_tool_calls', () => + tx.select({ id: copilotAsyncToolCalls.id }).from(copilotAsyncToolCalls).where(thisCall) + ) + if (!record) throw new Error('Tool execution record is unavailable') + return { outcome: 'existing' } as const + }) +} + +/** The desktop app acting on a pending call it was handed, before any effect on the machine. */ +export interface DesktopToolCallClaimant { + toolCallId: string + runId: string + userId: string + claimedBy: DesktopToolClaimOwner +} + +export type DesktopToolCallClaim = + | ToolExecutionClaim + /** Held for the user's decision, and not allowed (yet). */ + | { outcome: 'awaiting_permission' } + +/** + * The desktop app claims a pending call before crossing the Electron boundary, so a replayed + * renderer event cannot click, type, or run a command twice, and a call held for the user's + * decision only once they allowed it. + */ +export async function claimDesktopToolCall( + claimant: DesktopToolCallClaimant +): Promise { + return await claimUnderRunAdmission( + claimant, + claimant.claimedBy, + false, + async (tx, thisCall): Promise => { + const claimedAt = new Date() + const [claimed] = await traceMothershipQuery('UPDATE', 'copilot_async_tool_calls', () => + tx + .update(copilotAsyncToolCalls) + .set({ + status: ASYNC_TOOL_STATUS.running, + claimedBy: claimant.claimedBy, + claimedAt, + updatedAt: claimedAt, + }) + .where( + and( + thisCall, + eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.pending), + or( + and( + isNull(copilotAsyncToolCalls.permissionRequestedAt), + isNull(copilotAsyncToolCalls.permissionDecision) + ), + inArray(copilotAsyncToolCalls.permissionDecision, [ + ...EXECUTABLE_TOOL_PERMISSION_DECISIONS, + ]) + ) + ) + ) + .returning({ id: copilotAsyncToolCalls.id }) + ) + if (claimed) return { outcome: 'claimed' } + const [record] = await traceMothershipQuery('SELECT', 'copilot_async_tool_calls', () => + tx + .select({ + status: copilotAsyncToolCalls.status, + permissionRequestedAt: copilotAsyncToolCalls.permissionRequestedAt, + permissionDecision: copilotAsyncToolCalls.permissionDecision, + }) + .from(copilotAsyncToolCalls) + .where(thisCall) + ) + if (!record) throw new Error('Tool execution record is unavailable') + return record.status === ASYNC_TOOL_STATUS.pending && isAwaitingToolPermission(record) + ? { outcome: 'awaiting_permission' } + : { outcome: 'existing' } + } + ) +} + /** Expired ownership cannot be renewed, even before a follower has observed the expiry. */ export async function renewSimToolExecutionLease(owner: SimToolExecutionOwner): Promise { const [renewed] = await db @@ -1482,7 +1522,7 @@ export async function completeOwnedSimToolCall( /** * Finalizes a client tool only while it remains unclaimed. This is the inverse - * CAS of the desktop claim (`claimToolExecution`): exactly one of a renderer-side preclaim + * CAS of the desktop claim (`claimDesktopToolCall`): exactly one of a renderer-side preclaim * failure or the native authorization claim may transition the pending row. */ export async function completePendingAsyncToolCall(input: CompleteAsyncToolCallInput) { diff --git a/apps/sim/lib/mothership/request/tools/tool-call-failure.ts b/apps/sim/lib/mothership/request/tools/call-failure.ts similarity index 100% rename from apps/sim/lib/mothership/request/tools/tool-call-failure.ts rename to apps/sim/lib/mothership/request/tools/call-failure.ts diff --git a/apps/sim/lib/mothership/request/tools/client.ts b/apps/sim/lib/mothership/request/tools/client.ts index 3ae54119d1e..b725dabbc42 100644 --- a/apps/sim/lib/mothership/request/tools/client.ts +++ b/apps/sim/lib/mothership/request/tools/client.ts @@ -4,6 +4,7 @@ import { filterUndefined, isPlainRecord } from '@sim/utils/object' import { ASYNC_TOOL_CONFIRMATION_STATUS, type AsyncTerminalCompletionSnapshot, + getTerminalConfirmationStatus, isAsyncTerminalConfirmationStatus, } from '@/lib/mothership/async-runs/lifecycle' import { replaceTerminalAsyncToolCallResult } from '@/lib/mothership/async-runs/repository' @@ -23,7 +24,6 @@ import { createStructuralWorkflowToolCompletionData, getWorkflowToolCompletionExecutionId, getWorkflowToolCompletionMessage, - getWorkflowToolConfirmationStatus, getWorkflowToolLaunchError, type WorkflowToolLaunchError, } from '@/lib/mothership/tools/workflow-tools' @@ -314,7 +314,7 @@ export async function waitForWorkflowToolCompletion({ if (!trustedExecution.contentAvailable) { toolRegistry?.markIncomplete('client-tool-content-unavailable') return structuralWorkflowCompletion( - getWorkflowToolConfirmationStatus(trustedExecution.status), + getTerminalConfirmationStatus(trustedExecution.status), workflowId, executionId ) @@ -330,7 +330,7 @@ export async function waitForWorkflowToolCompletion({ origin: 'copilotToolClient.workflowExecution', }) return structuralWorkflowCompletion( - getWorkflowToolConfirmationStatus(trustedExecution.status), + getTerminalConfirmationStatus(trustedExecution.status), workflowId, executionId ) @@ -370,7 +370,7 @@ export async function waitForWorkflowToolCompletion({ if (!completion || !trustedExecution || !workflowId) return completion const executionId = trustedExecution.executionId - const status = getWorkflowToolConfirmationStatus(trustedExecution.status) + const status = getTerminalConfirmationStatus(trustedExecution.status) const genericMessage = getWorkflowToolCompletionMessage(status) const error = status !== MothershipStreamV1ToolOutcome.success diff --git a/apps/sim/lib/mothership/request/tools/desktop-wait.ts b/apps/sim/lib/mothership/request/tools/desktop-wait.ts index 770283d4e75..bbb0ec4c398 100644 --- a/apps/sim/lib/mothership/request/tools/desktop-wait.ts +++ b/apps/sim/lib/mothership/request/tools/desktop-wait.ts @@ -4,8 +4,8 @@ import type { AsyncTerminalCompletionSnapshot } from '@/lib/mothership/async-run import { MothershipStreamV1ToolOutcome } from '@/lib/mothership/generated/mothership-stream-v1' import { CopilotDegradedReason } from '@/lib/mothership/generated/trace-attribute-values-v1' import { recordDegraded } from '@/lib/mothership/request/metrics' +import { settleToolCallFailure } from '@/lib/mothership/request/tools/call-failure' import { waitForClientToolCompletion } from '@/lib/mothership/request/tools/client' -import { settleToolCallFailure } from '@/lib/mothership/request/tools/tool-call-failure' import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' const logger = createLogger('CopilotDesktopToolWait') diff --git a/apps/sim/lib/mothership/request/tools/executor.ts b/apps/sim/lib/mothership/request/tools/executor.ts index 2a95b47b384..b1d02b39e5e 100644 --- a/apps/sim/lib/mothership/request/tools/executor.ts +++ b/apps/sim/lib/mothership/request/tools/executor.ts @@ -59,6 +59,7 @@ import { requireToolCallError, setTerminalToolCallState, } from '@/lib/mothership/request/tool-call-state' +import { settleToolCallFailure } from '@/lib/mothership/request/tools/call-failure' import { desktopToolNotStarted } from '@/lib/mothership/request/tools/desktop-wait' import { type ToolExecutionLifetime, @@ -75,7 +76,6 @@ import { maybeWriteOutputToTable, maybeWriteReadCsvToTable, } from '@/lib/mothership/request/tools/tables' -import { settleToolCallFailure } from '@/lib/mothership/request/tools/tool-call-failure' import { applyCreateWorkflowOutputToContext } from '@/lib/mothership/request/tools/workflow-context' import { type ExecutionContext, diff --git a/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts b/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts index b37a660bde1..c43c2b7847f 100644 --- a/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts +++ b/apps/sim/lib/mothership/request/tools/workflow-client-settlement.ts @@ -1,5 +1,6 @@ import { ASYNC_TOOL_CONFIRMATION_STATUS, + getTerminalConfirmationStatus, isTerminalAsyncStatus, } from '@/lib/mothership/async-runs/lifecycle' import { @@ -10,7 +11,6 @@ import { publishToolConfirmation } from '@/lib/mothership/persistence/tool-confi import { createStructuralWorkflowToolCompletionData, getWorkflowToolCompletionMessage, - getWorkflowToolConfirmationStatus, } from '@/lib/mothership/tools/workflow-tools' import { getWorkflowExecutionLogStatus } from '@/lib/workflows/executor/execution-state' @@ -41,7 +41,7 @@ export async function reportSettledClientWorkflowTool({ if (logStatus !== undefined && !isTerminalAsyncStatus(logStatus)) return const executionStatus = logStatus ?? 'failed' - const status = getWorkflowToolConfirmationStatus(executionStatus) + const status = getTerminalConfirmationStatus(executionStatus) const message = getWorkflowToolCompletionMessage(status) const data = createStructuralWorkflowToolCompletionData(status, workflowId, executionId) const completed = await completeClientWorkflowToolCall( diff --git a/apps/sim/lib/mothership/tools/client/desktop-tool-authorization.integration.ts b/apps/sim/lib/mothership/tools/client/desktop-tool-authorization.integration.ts index ee9d4f1fc96..e04d8728f9f 100644 --- a/apps/sim/lib/mothership/tools/client/desktop-tool-authorization.integration.ts +++ b/apps/sim/lib/mothership/tools/client/desktop-tool-authorization.integration.ts @@ -31,7 +31,7 @@ import { NextRequest } from 'next/server' import { closeRedisConnection } from '@/lib/core/config/redis' import { SIM_TOOL_EXECUTION_VERSION } from '@/lib/mothership/async-runs/lifecycle' import { - claimToolExecution, + claimDesktopToolCall, closeStreamToolAdmission, requestRunStop, } from '@/lib/mothership/async-runs/repository' @@ -326,7 +326,7 @@ describe.runIf(Boolean(redisUrl))('desktop tool calls the server no longer admit await closeStreamToolAdmission(streamId, userId) expect( - await claimToolExecution({ toolCallId, runId, userId, claimedBy: 'desktop-terminal' }) + await claimDesktopToolCall({ toolCallId, runId, userId, claimedBy: 'desktop-terminal' }) ).toEqual({ outcome: 'closed' }) expect(await storedCall(toolCallId)).toMatchObject({ status: 'pending', claimedBy: null }) }) diff --git a/apps/sim/lib/mothership/tools/desktop-tools.ts b/apps/sim/lib/mothership/tools/desktop-tools.ts index 7d56d111dd2..f6dba0ee75c 100644 --- a/apps/sim/lib/mothership/tools/desktop-tools.ts +++ b/apps/sim/lib/mothership/tools/desktop-tools.ts @@ -1,4 +1,8 @@ -import { CURRENT_BROWSER_TOOL_NAMES, isCurrentBrowserToolName } from '@sim/browser-protocol' +import { + CURRENT_BROWSER_TOOL_NAMES, + isBrowserToolName, + isCurrentBrowserToolName, +} from '@sim/browser-protocol' import { isTerminalToolName, TERMINAL_TOOL_NAME } from '@sim/terminal-protocol' import { DESKTOP_TOOL_CLAIM_OWNER } from '@/lib/mothership/async-runs/lifecycle' import { isUserLocalVfsToolCall } from '@/lib/mothership/tools/local-filesystem' @@ -34,6 +38,14 @@ export function getDesktopToolClaimOwner(toolName: string): DesktopToolClaimOwne return undefined } +/** + * Whether a call's result is accepted only under the desktop's native claim rules: the claimed + * desktop tools, plus browser tools retired from the catalog whose calls remain in history. + */ +export function isNativeDesktopTool(toolName: string): boolean { + return getDesktopToolClaimOwner(toolName) !== undefined || isBrowserToolName(toolName) +} + /** A read of the user's machine: `read_local_file`, or a VFS read of a granted local folder. */ export function isLocalReadToolCall(toolName: string, args: Record | undefined) { return toolName === 'read_local_file' || isUserLocalVfsToolCall(toolName, args) diff --git a/apps/sim/lib/mothership/tools/workflow-tools.ts b/apps/sim/lib/mothership/tools/workflow-tools.ts index 666c111d2fd..6858c9962a1 100644 --- a/apps/sim/lib/mothership/tools/workflow-tools.ts +++ b/apps/sim/lib/mothership/tools/workflow-tools.ts @@ -171,14 +171,6 @@ export function getWorkflowToolCompletionMessage(status: AsyncConfirmationStatus return 'Workflow execution failed.' } -export function getWorkflowToolConfirmationStatus( - status: 'completed' | 'failed' | 'cancelled' -): AsyncConfirmationStatus { - if (status === 'completed') return ASYNC_TOOL_CONFIRMATION_STATUS.success - if (status === 'cancelled') return ASYNC_TOOL_CONFIRMATION_STATUS.cancelled - return ASYNC_TOOL_CONFIRMATION_STATUS.error -} - export function createStructuralWorkflowToolCompletionData( status: AsyncConfirmationStatus, workflowId?: string, diff --git a/packages/testing/src/mocks/mothership-async-runs.mock.ts b/packages/testing/src/mocks/mothership-async-runs.mock.ts index a38f16d6274..5df4837791d 100644 --- a/packages/testing/src/mocks/mothership-async-runs.mock.ts +++ b/packages/testing/src/mocks/mothership-async-runs.mock.ts @@ -33,6 +33,7 @@ export const mothershipAsyncRunsMockFns = { mockGetAsyncToolCall: vi.fn(), mockMarkAsyncToolRunning: vi.fn(), mockClaimToolExecution: vi.fn(), + mockClaimDesktopToolCall: vi.fn(), mockRenewSimToolExecutionLease: vi.fn(), mockRevokeExpiredSimToolExecutions: vi.fn(), mockSettleSimToolExecution: vi.fn(), @@ -92,6 +93,7 @@ export const mothershipAsyncRunsMock = { getAsyncToolCall: mothershipAsyncRunsMockFns.mockGetAsyncToolCall, markAsyncToolRunning: mothershipAsyncRunsMockFns.mockMarkAsyncToolRunning, claimToolExecution: mothershipAsyncRunsMockFns.mockClaimToolExecution, + claimDesktopToolCall: mothershipAsyncRunsMockFns.mockClaimDesktopToolCall, renewSimToolExecutionLease: mothershipAsyncRunsMockFns.mockRenewSimToolExecutionLease, revokeExpiredSimToolExecutions: mothershipAsyncRunsMockFns.mockRevokeExpiredSimToolExecutions, settleSimToolExecution: mothershipAsyncRunsMockFns.mockSettleSimToolExecution, From d9b5bc27e3f2c7b450fb6b550823fe7cc5d9412c Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 5 Oct 2026 20:26:21 -0700 Subject: [PATCH 3/4] test(mothership): drive the authorize unit test through the desktop claim --- apps/sim/app/api/desktop/tool/authorize/route.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/sim/app/api/desktop/tool/authorize/route.test.ts b/apps/sim/app/api/desktop/tool/authorize/route.test.ts index df6f238d96d..2e8900956ec 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.test.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.test.ts @@ -17,7 +17,7 @@ vi.mock('@/lib/mothership/async-runs/repository', () => mothershipAsyncRunsMock) import { POST } from './route' -const claimDesktopToolCall = mothershipAsyncRunsMockFns.mockClaimToolExecution +const claimDesktopToolCall = mothershipAsyncRunsMockFns.mockClaimDesktopToolCall const getAsyncToolCall = mothershipAsyncRunsMockFns.mockGetAsyncToolCall const getRunSegment = mothershipAsyncRunsMockFns.mockGetRunSegment const resolveInvocationWorkspace = mothershipWorkspaceTargetMockFns.mockResolveInvocationWorkspace From 60a4d186df7fa2108a1d5e517b7a2427d1e84581 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Mon, 5 Oct 2026 20:35:22 -0700 Subject: [PATCH 4/4] fix(mothership): settle an unclaimed desktop call whose wait ends before the grace A wait shorter than the pickup grace ended with no result and left the call claimable; it now settles it as never started like any unclaimed call. The desktop E2E fixture models claims per call, the way the server accepts a local read's repeat claim and refuses a second import claim. The admission probe follows the claim rename. --- apps/desktop/e2e/local-files.spec.ts | 11 +++++--- .../request/tools/desktop-wait.test.ts | 10 ++++++- .../mothership/request/tools/desktop-wait.ts | 2 +- .../probes/mothership-request-admission.ts | 26 +++++++++---------- 4 files changed, 31 insertions(+), 18 deletions(-) diff --git a/apps/desktop/e2e/local-files.spec.ts b/apps/desktop/e2e/local-files.spec.ts index d2645297598..0a9a08f0678 100644 --- a/apps/desktop/e2e/local-files.spec.ts +++ b/apps/desktop/e2e/local-files.spec.ts @@ -16,7 +16,8 @@ test('native file tools read and import through the installed preload without Si const png = 'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAusB9Y9Zl1sAAAAASUVORK5CYII=' writeFileSync(join(source, 'image.png'), Buffer.from(png, 'base64')) - let claimed = false + /** Calls the server saw claimed; like the server, only an import refuses a second claim. */ + const claimed = new Set() const calls: Record }> = { text: { toolName: 'read_local_file', args: { path: join(source, 'report.txt') } }, image: { toolName: 'read_local_file', args: { path: join(source, 'image.png') } }, @@ -44,11 +45,14 @@ test('native file tools read and import through the installed preload without Si for await (const chunk of request) body += chunk.toString() const input = JSON.parse(body) const call = calls[input.toolCallId] - if (!call || (input.claim && claimed)) { + if ( + !call || + (input.claim && call.toolName === 'import_local_files' && claimed.has(input.toolCallId)) + ) { response.writeHead(call ? 409 : 403, { 'Content-Type': 'application/json' }).end('{}') return } - if (input.claim) claimed = true + if (input.claim) claimed.add(input.toolCallId) response .writeHead(200, { 'Content-Type': 'application/json' }) .end(JSON.stringify({ ...call, chatId: 'org-chat' })) @@ -92,6 +96,7 @@ test('native file tools read and import through the installed preload without Si ok: true, data: { observations: [{ mediaType: 'image/png', data: png }] }, }) + expect([...claimed]).toEqual(['text', 'image']) const result = await invoke({ operation: 'manifest', toolCallId: 'import' }) if (!result.ok || result.data.kind !== 'manifest') throw new Error(JSON.stringify(result)) expect(result.data.targetWorkspaceId).toBe('target-workspace') diff --git a/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts b/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts index 5358721b65f..41c27416959 100644 --- a/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts +++ b/apps/sim/lib/mothership/request/tools/desktop-wait.test.ts @@ -45,11 +45,19 @@ describe('waitForDesktopToolCall', () => { abortSignal.addEventListener('abort', () => resolve(desktopResult), { once: true }) }) ) - completePendingAsyncToolCall.mockResolvedValueOnce(null) const answer = waitForDesktopToolCall(params) await vi.advanceTimersByTimeAsync(GRACE_MS) expect(await answer).toEqual(desktopResult) }) + + it('settles an unclaimed call as never started when its wait ends before the grace', async () => { + waitForClientToolCompletion.mockResolvedValueOnce(null) + completePendingAsyncToolCall.mockImplementationOnce(async (input) => ({ ...input })) + + const answer = await waitForDesktopToolCall({ ...params, timeoutMs: 1_000 }) + + expect(answer).toMatchObject({ status: 'error', data: { notStarted: true } }) + }) }) diff --git a/apps/sim/lib/mothership/request/tools/desktop-wait.ts b/apps/sim/lib/mothership/request/tools/desktop-wait.ts index bbb0ec4c398..39105cb2636 100644 --- a/apps/sim/lib/mothership/request/tools/desktop-wait.ts +++ b/apps/sim/lib/mothership/request/tools/desktop-wait.ts @@ -62,7 +62,7 @@ export async function waitForDesktopToolCall( })), ]) graceOver.abort() - if (first.kind === 'reported') return first.completion + if (first.kind === 'reported' && first.completion) return first.completion if (abortSignal?.aborted) return pickupWait // A result that landed as the grace ran out is already being restored: it is the answer. diff --git a/scripts/probes/mothership-request-admission.ts b/scripts/probes/mothership-request-admission.ts index a95497e5f2c..bb77819e3d7 100644 --- a/scripts/probes/mothership-request-admission.ts +++ b/scripts/probes/mothership-request-admission.ts @@ -124,7 +124,7 @@ try { await repo.upsertAsyncToolCall({ runId: run.id, toolCallId, toolName: 'run_code' }) await repo.updateRunStatus(run.id, status, { completedAt: new Date() }) assert.equal( - (await repo.claimSimToolExecution({ runId: run.id, toolCallId, userId: 'user-1' })).outcome, + (await repo.claimToolExecution({ runId: run.id, toolCallId, userId: 'user-1' })).outcome, 'closed' ) await repo.updateRunStatus(run.id, 'paused_waiting_for_tool') @@ -142,7 +142,7 @@ try { toolCallId: ownerTool.toolCallId, toolName: 'run_code', }) - assert.equal((await repo.claimSimToolExecution(ownerTool)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(ownerTool)).outcome, 'claimed') const commandIdentity = { id: generateId(), sandboxId: 'recorded-sandbox', @@ -180,7 +180,7 @@ try { toolName: 'run_code', }) const [claim] = await Promise.all([ - repo.claimSimToolExecution(tool), + repo.claimToolExecution(tool), repo.updateRunStatus(run.id, 'complete'), ]) assert.equal( @@ -207,12 +207,12 @@ try { const workbenchChatId = await createWorkbenchChat() const sessionKey = `mothership-chat:${workbenchChatId}` const prior = await createWorkbenchTool(workbenchChatId, 'workbench-prior') - assert.equal((await repo.claimSimToolExecution(prior)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(prior)).outcome, 'claimed') const priorCommand = { id: generateId(), sandboxId: 'workbench-vm', sessionKey } await repo.recordSimSandboxProcess({ ...prior, process: priorCommand }) await repo.updateRunStatus(prior.runId, 'complete') const current = await createWorkbenchTool(workbenchChatId, 'workbench-current') - assert.equal((await repo.claimSimToolExecution(current)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(current)).outcome, 'claimed') const access = { ...current, sessionKey } assert.deepEqual(await repo.prepareWorkbenchAccess(access), { handlersPending: true, @@ -229,7 +229,7 @@ try { await assert.rejects(repo.prepareWorkbenchAccess({ ...prior, sessionKey })) const sibling = { ...current, toolCallId: 'workbench-sibling' } await repo.upsertAsyncToolCall({ ...sibling, toolName: 'run_code' }) - assert.equal((await repo.claimSimToolExecution(sibling)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(sibling)).outcome, 'claimed') const siblingCommand = { ...priorCommand, id: generateId() } await repo.recordSimSandboxProcess({ ...sibling, process: siblingCommand }) assert.deepEqual(await repo.prepareWorkbenchAccess(access), ready) @@ -252,7 +252,7 @@ try { await repo.settleSimToolExecution(sibling.toolCallId) const nextTool = { runId: emptyNext.id, userId: 'user-1', toolCallId: 'workbench-next-tool' } await repo.upsertAsyncToolCall({ ...nextTool, toolName: 'run_code' }) - assert.equal((await repo.claimSimToolExecution(nextTool)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(nextTool)).outcome, 'claimed') assert.deepEqual(await repo.prepareWorkbenchAccess({ ...nextTool, sessionKey }), { handlersPending: false, processes: [{ ...siblingCommand, toolCallId: sibling.toolCallId }], @@ -266,14 +266,14 @@ try { const raceChatId = await createWorkbenchChat() const older = await createWorkbenchTool(raceChatId, `workbench-race-old-${i}`) const newer = await createWorkbenchTool(raceChatId, `workbench-race-new-${i}`) - assert.equal((await repo.claimSimToolExecution(newer)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(newer)).outcome, 'claimed') const [claim, state] = await Promise.all([ - repo.claimSimToolExecution(older), + repo.claimToolExecution(older), repo.prepareWorkbenchAccess({ ...newer, sessionKey: `mothership-chat:${raceChatId}` }), ]) assert.deepEqual(state, { handlersPending: claim.outcome === 'claimed', processes: [] }) assert.equal( - (await repo.claimSimToolExecution({ ...older, toolCallId: 'late' })).outcome, + (await repo.claimToolExecution({ ...older, toolCallId: 'late' })).outcome, 'closed' ) await assert.rejects( @@ -283,11 +283,11 @@ try { logger.info('PASS thirty predecessor-claim/workbench-takeover races without false readiness') const corruptChatId = await createWorkbenchChat() const corruptOwner = await createWorkbenchTool(corruptChatId, 'workbench-corrupt-owner') - assert.equal((await repo.claimSimToolExecution(corruptOwner)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(corruptOwner)).outcome, 'claimed') const corruptCommand = { id: generateId(), sandboxId: 'foreign', sessionKey: 'another-chat' } await repo.recordSimSandboxProcess({ ...corruptOwner, process: corruptCommand }) const successor = await createWorkbenchTool(corruptChatId, 'workbench-corrupt-successor') - assert.equal((await repo.claimSimToolExecution(successor)).outcome, 'claimed') + assert.equal((await repo.claimToolExecution(successor)).outcome, 'claimed') const successorAccess = { ...successor, sessionKey: `mothership-chat:${corruptChatId}` } await assert.rejects(repo.prepareWorkbenchAccess(successorAccess), /does not match this chat/) assert.equal( @@ -334,7 +334,7 @@ try { }) assert.equal( ( - await repo.claimSimToolExecution({ + await repo.claimToolExecution({ runId: delayed.id, toolCallId: 'late-tool', userId: 'user-1',