diff --git a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts index e53ce45d7ae..10aa0c1a030 100644 --- a/apps/desktop/e2e/desktop-tools-live-sim.spec.ts +++ b/apps/desktop/e2e/desktop-tools-live-sim.spec.ts @@ -11,6 +11,7 @@ import { test, } from '@playwright/test' import type { SimDesktopApi } from '@sim/desktop-bridge' +import { sleep } from '@sim/utils/helpers' import { generateId } from '@sim/utils/id' import { toRecord } from '@sim/utils/object' import { @@ -34,6 +35,10 @@ import { const DESKTOP_DIR = fileURLToPath(new URL('..', import.meta.url)) const config = liveSimConfig() const PICKUP_GRACE_MS = 15_000 +/** Longer than the default tool budget (60 s) plus the resume grace (30 s). */ +const LONG_IMPORT_MS = 120_000 +/** The execution lease a running import holds and renews (`SIM_TOOL_EXECUTION_LEASE_SECONDS`). */ +const LEASE_MS = 60_000 /** How long a held request may take to arrive once the step that sends it ran. */ const ARRIVAL_MS = 60_000 /** First requests to a route compile it, which takes minutes on a cold dev app. */ @@ -472,6 +477,69 @@ test.describe('desktop tools against a live Sim', () => { expect(await db.workspaceFileNames(user.workspaceId)).toEqual(['a.txt']) }) + /** An import whose first upload the proxy holds until released, as a large file's would take. */ + async function slowImport(user: SeededUser, title: string, marker: string) { + const source = importSource() + let callId = '' + agent.script(marker, (turn) => { + callId = turn.toolCall({ + toolName: 'import_local_files', + args: { path: source, targetWorkspaceId: user.workspaceId }, + }) + turn.pause() + }) + const page = await openApp(user, title) + const firstUpload = proxy.hold(isUploadStart) + await send(page, `${marker} import my reports`) + await firstUpload.arrival(ARRIVAL_MS, 'The first upload') + return { page, firstUpload, callId: () => callId } + } + + test('an import that runs longer than the default tool budget still completes', async () => { + test.setTimeout(420_000) + const user = await db.seedUser(['Long import chat', 'Other chat']) + const chatId = user.chats['Long import chat'] + const { page, firstUpload, callId } = await slowImport( + user, + 'Long import chat', + '[long-import]' + ) + // The import is alive and working past the 60 s default budget and its 30 s grace, while the + // user is in another chat: its lease is renewed by the import, not by the chat view. + await openChat(page, user, 'Other chat') + await page.waitForTimeout(LONG_IMPORT_MS) + expect(agent.resultFor(callId())).toBeUndefined() + firstUpload.release() + + await agent.waitForResume(() => Boolean(agent.resultFor(callId())), 120_000) + const result = agent.resultFor(callId()) + expect(JSON.stringify(result?.data)).not.toContain('outcomeUnknown') + expect(result?.success).toBe(true) + await expect + .poll(() => db.workspaceFileNames(user.workspaceId), { timeout: 30_000 }) + .toEqual(['a.txt', 'b.txt']) + expect(await callState(chatId)).toMatch(/^completed/) + }) + + test('a long import whose window crashed settles as outcome unknown about one lease later', async () => { + test.setTimeout(420_000) + const user = await db.seedUser(['Crashing import chat']) + const { callId } = await slowImport(user, 'Crashing import chat', '[crash-import]') + await sleep(LONG_IMPORT_MS) + expect(agent.resultFor(callId())).toBeUndefined() + // A crash reports nothing on its way out, unlike a closed window: only the lapse of the + // lease the page was renewing tells Sim the import is gone. + const crashedAt = Date.now() + await app?.evaluate(({ webContents }) => { + for (const contents of webContents.getAllWebContents()) + if (contents.getURL().includes('/workspace/')) contents.forcefullyCrashRenderer() + }) + await agent.waitForResume(() => Boolean(agent.resultFor(callId())), 150_000) + const result = agent.resultFor(callId()) + expect(result?.data).toMatchObject({ outcomeUnknown: true }) + expect((result?.at ?? 0) - crashedAt).toBeLessThan(LEASE_MS + 20_000) + }) + test('signing out ends a desktop tool still running', async () => { const user = await db.seedUser(['Import chat']) const chatId = user.chats['Import chat'] diff --git a/apps/sim/app/api/desktop/tool/authorize/route.ts b/apps/sim/app/api/desktop/tool/authorize/route.ts index 42ee66bc92b..4326122636c 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.ts @@ -18,6 +18,7 @@ import { createUnauthorizedResponse, } from '@/lib/mothership/request/http' import { + chatViewDesktopLeaseOwnerToken, getDesktopToolClaimOwner, isDesktopToolCall, isLocalReadToolCall, @@ -52,7 +53,7 @@ function refusedClaimResponse( * that they have not allowed. */ export const POST = withRouteHandler(async (request: NextRequest) => { - const { userId, isAuthenticated } = await authenticateCopilotRequestSessionOnly() + const { userId, isAuthenticated, principal } = await authenticateCopilotRequestSessionOnly() if (!isAuthenticated || !userId) { return createUnauthorizedResponse() } @@ -114,11 +115,16 @@ export const POST = withRouteHandler(async (request: NextRequest) => { { status: 409 } ) if (toolCall.status !== 'pending') return alreadyStarted() + // An import runs as long as its files take, so its claim takes a lease this session holds: + // the chat view renews it while the import runs. Reads finish in seconds and take none. const { outcome } = await claimDesktopToolCall({ toolCallId: toolCall.toolCallId, runId: toolCall.runId, userId, claimedBy: DESKTOP_TOOL_CLAIM_OWNER.files, + ...(principal + ? { chatView: { ownerToken: chatViewDesktopLeaseOwnerToken(principal.sessionId) } } + : {}), }) if (outcome !== 'claimed') return refusedClaimResponse(outcome, alreadyStarted) } else if ( diff --git a/apps/sim/lib/api/contracts/desktop-executor.ts b/apps/sim/lib/api/contracts/desktop-executor.ts index 2cfe74f9d84..fbc10ea159f 100644 --- a/apps/sim/lib/api/contracts/desktop-executor.ts +++ b/apps/sim/lib/api/contracts/desktop-executor.ts @@ -136,11 +136,21 @@ export const claimDesktopToolContract = defineRouteContract({ error: z.object({ error: z.string() }), }) -const renewDesktopToolLeaseBodySchema = z.object({ - deviceId: desktopDeviceIdSchema, - toolCallId: desktopToolCallIdSchema, - executionToken: z.string().min(1).max(128), -}) +/** + * A device renews a call of a run bound to it under its execution token; the chat view renews an + * import it claimed (`chatView`), as the session that claimed it. + */ +const renewDesktopToolLeaseBodySchema = z.union([ + z.object({ + deviceId: desktopDeviceIdSchema, + toolCallId: desktopToolCallIdSchema, + executionToken: z.string().min(1).max(128), + }), + z.object({ + toolCallId: desktopToolCallIdSchema, + chatView: z.literal(true), + }), +]) export type RenewDesktopToolLeaseBody = z.input export const renewDesktopToolLeaseResponseSchema = z.object({ renewed: z.literal(true) }) diff --git a/apps/sim/lib/desktop/application/executor.ts b/apps/sim/lib/desktop/application/executor.ts index 0dd8fc89b39..f70c844f05e 100644 --- a/apps/sim/lib/desktop/application/executor.ts +++ b/apps/sim/lib/desktop/application/executor.ts @@ -41,6 +41,7 @@ import { import { claimDesktopToolCall, type DesktopToolCallClaim, + getAsyncToolCall, renewSimToolExecutionLease, } from '@/lib/mothership/async-runs/repository' import { sealClientToolSettlement } from '@/lib/mothership/request/tools/client-completion-seal.server' @@ -50,6 +51,7 @@ import { settleClientToolCall, } from '@/lib/mothership/request/tools/client-settlement.server' import { + chatViewDesktopLeaseOwnerToken, getDesktopExecutorClaimOwner, isDesktopToolCall, } from '@/lib/mothership/tools/desktop-tools' @@ -329,9 +331,19 @@ interface DesktopCallTokenInput extends DeviceInput { executionToken: string } -/** Keeps a running call owned; failing here always means the device must stop the action. */ +/** A desktop call the chat view is running, renewed by the session that claimed it. */ +interface ChatViewCallInput { + toolCallId: string + chatView: true +} + +/** + * Keeps a running call owned; failing here always means the device must stop the action. A + * device renews a call of a run bound to it under its execution token. The chat view renews an + * import it claimed on an unbound run: only the session that claimed it, while it runs. + */ export const renewDesktopToolLease = defineAuthorizedCredentialUserUseCase({ - // permission-group-exempt: extends only a lease this device's token already holds. + // permission-group-exempt: extends only a lease this device's token, or this session, already holds. operation: defineOperation({ id: 'desktop.executor.calls.renew', principalKinds: ['session'], @@ -342,8 +354,24 @@ export const renewDesktopToolLease = defineAuthorizedCredentialUserUseCase({ input, }: { principal: SessionPrincipal - input: DesktopCallTokenInput + input: DesktopCallTokenInput | ChatViewCallInput }) { + if ('chatView' in input) { + const call = await getAsyncToolCall(input.toolCallId) + const renewed = + call !== null && + (await renewSimToolExecutionLease( + { + toolCallId: call.toolCallId, + runId: call.runId, + userId: principal.userId, + ownerToken: chatViewDesktopLeaseOwnerToken(principal.sessionId), + }, + { chatView: true } + )) + if (!renewed) throw new DesktopCallRevokedError() + return { renewed: true as const } + } await requireBoundDevice(principal, input.deviceId) await markDesktopPresent(input.deviceId) const call = await getBoundDesktopCall( diff --git a/apps/sim/lib/mothership/async-runs/repository.ts b/apps/sim/lib/mothership/async-runs/repository.ts index 85fcd70deeb..0cd9ffd392a 100644 --- a/apps/sim/lib/mothership/async-runs/repository.ts +++ b/apps/sim/lib/mothership/async-runs/repository.ts @@ -779,6 +779,12 @@ export interface DesktopToolCallClaimant { * settlement or the next turn's workbench, whatever the device does after. */ executor?: { deviceId: string; ownerToken: string } + /** + * Set when the chat view claims a call that runs long enough to need a lease (an import): the + * claim takes the execution lease under the chat view's session token, and the chat view renews + * it while the import runs. Without renewals the lease lapses as the default tool budget would. + */ + chatView?: { ownerToken: string } } export type DesktopToolCallClaim = @@ -795,7 +801,8 @@ export type DesktopToolCallClaim = export async function claimDesktopToolCall( claimant: DesktopToolCallClaimant ): Promise { - const { executor } = claimant + const { executor, chatView } = claimant + const leaseOwnerToken = executor?.ownerToken ?? chatView?.ownerToken return await claimUnderRunAdmission( { ...claimant, desktopDeviceId: executor?.deviceId }, claimant.claimedBy, @@ -810,9 +817,9 @@ export async function claimDesktopToolCall( claimedBy: claimant.claimedBy, claimedAt, updatedAt: claimedAt, - ...(executor + ...(leaseOwnerToken ? { - executionOwnerToken: executor.ownerToken, + executionOwnerToken: leaseOwnerToken, executionLeaseExpiresAt: sql`clock_timestamp() + ${SIM_TOOL_EXECUTION_LEASE_SECONDS} * interval '1 second'`, } : {}), @@ -868,7 +875,7 @@ export async function claimDesktopToolCall( */ export async function renewSimToolExecutionLease( owner: SimToolExecutionOwner, - desktop?: { deviceId: string } + desktop?: { deviceId: string } | { chatView: true } ): Promise { const [renewed] = await db .update(copilotAsyncToolCalls) @@ -887,7 +894,12 @@ export async function renewSimToolExecutionLease( desktop ? and( eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.running), - sql`EXISTS (SELECT 1 FROM ${copilotRuns} r WHERE r.id = ${copilotAsyncToolCalls.runId} AND r.desktop_device_id = ${desktop.deviceId})` + 'deviceId' in desktop + ? sql`EXISTS (SELECT 1 FROM ${copilotRuns} r WHERE r.id = ${copilotAsyncToolCalls.runId} AND r.desktop_device_id = ${desktop.deviceId})` + : and( + eq(copilotAsyncToolCalls.claimedBy, DESKTOP_TOOL_CLAIM_OWNER.files), + sql`EXISTS (SELECT 1 FROM ${copilotRuns} r WHERE r.id = ${copilotAsyncToolCalls.runId} AND r.desktop_device_id IS NULL)` + ) ) : undefined ) @@ -896,6 +908,35 @@ export async function renewSimToolExecutionLease( return !!renewed } +/** + * How much longer the chat view's renewed lease keeps a desktop call it is running alive, in ms, + * on the database clock; null once the lease lapsed or the call is not one the chat view holds. + */ +export async function getChatViewDesktopLeaseRemainingMs( + toolCallId: string +): Promise { + const [row] = await db + .select({ + remainingMs: sql`extract(epoch from (${copilotAsyncToolCalls.executionLeaseExpiresAt} - clock_timestamp())) * 1000`, + }) + .from(copilotAsyncToolCalls) + .innerJoin(copilotRuns, eq(copilotRuns.id, copilotAsyncToolCalls.runId)) + .where( + and( + eq(copilotAsyncToolCalls.toolCallId, toolCallId), + eq(copilotAsyncToolCalls.status, ASYNC_TOOL_STATUS.running), + eq(copilotAsyncToolCalls.claimedBy, DESKTOP_TOOL_CLAIM_OWNER.files), + isNotNull(copilotAsyncToolCalls.executionOwnerToken), + isNull(copilotAsyncToolCalls.executionSettledAt), + isNull(copilotAsyncToolCalls.executionRevokedAt), + isNull(copilotRuns.desktopDeviceId), + sql`${copilotAsyncToolCalls.executionLeaseExpiresAt} > clock_timestamp()` + ) + ) + .limit(1) + return row ? Number(row.remainingMs) : null +} + /** Revocation ends local execution authority; recorded remote commands remain independently unsettled. */ export async function revokeExpiredSimToolExecutions(input: { runId: string; userId: string }) { return db.transaction((tx) => diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index 77feb41971b..56df7a5da99 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -1,4 +1,5 @@ import { resetEnvFlagsMock, resetEnvironmentUtilsMock, setEnvFlags } from '@sim/testing' +import { getMockLogger } from '@sim/testing/mocks/logger.mock' import { mothershipAgentUrlMock, mothershipAgentUrlMockFns, @@ -221,6 +222,7 @@ vi.mock('@/lib/mothership/request/enterprise-byok', () => ({ import { resetUsageGateCache } from '@/lib/billing/core/usage-gate-cache' import { buildPersistedAssistantMessage } from '@/lib/mothership/chat/persisted-message' +import { CLIENT_TOOL_RESULT_TIMEOUT_MS } from '@/lib/mothership/constants' import { MothershipStreamV1CompletionStatus, MothershipStreamV1ToolOutcome, @@ -3545,6 +3547,235 @@ describe('runCopilotLifecycle', () => { } }) + /** Runs a turn whose import never settles, recording each leg's request body. */ + function runImportTurn() { + const bodies: Record[] = [] + let streamContext: StreamingContext | undefined + mockForceFailHungToolCall.mockImplementation( + async (toolCallId: string, context: StreamingContext) => { + const tool = context.toolCalls.get(toolCallId) + if (!tool) return + tool.status = MothershipStreamV1ToolOutcome.error + tool.endTime = Date.now() + tool.result = { success: false } + tool.error = 'Tool execution hung' + } + ) + mockRunStreamLoop.mockImplementationOnce( + async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => { + bodies.push(JSON.parse(String(fetchOptions.body))) + streamContext = context + context.toolCalls.set('tool-import', { + id: 'tool-import', + name: 'import_local_files', + status: 'executing', + }) + context.pendingToolPromises.set('tool-import', new Promise(() => {})) + context.awaitingAsyncContinuation = { + checkpointId: 'ckpt-1', + pendingToolCallIds: ['tool-import'], + } + } + ) + mockRunStreamLoop.mockImplementationOnce( + async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => { + bodies.push(JSON.parse(String(fetchOptions.body))) + context.accumulatedContent = 'Done.' + } + ) + const lifecycle = runCopilotLifecycle( + { message: 'import', messageId: 'stream-1' }, + { + userId: 'user-1', + workspaceId: 'ws-1', + chatId: 'chat-1', + executionId: 'exec-1', + runId: 'run-1', + executionContext: { + userId: 'user-1', + workflowId: '', + workspaceId: 'ws-1', + chatId: 'chat-1', + }, + } + ) + /** The agent was resumed with the import given up as lost. */ + const resumedWithLostImport = () => + bodies.length === 2 && + JSON.stringify(bodies[1].results).includes('"callId":"tool-import"') && + JSON.stringify(bodies[1].results).includes('"success":false') + const context = () => { + if (!streamContext) throw new Error('The turn did not start its stream') + return streamContext + } + /** The reasons the turn logged for force-failing its calls. */ + const forceFailures = () => + getMockLogger('CopilotLifecycle') + .error.mock.calls.map(([message]) => String(message)) + .filter((message) => message.endsWith('force-failing')) + return { bodies, lifecycle, resumedWithLostImport, context, forceFailures } + } + + it('waits on a chat-view import while its lease is renewed, and fails it once the lease lapses', async () => { + vi.useFakeTimers() + try { + // Renewed once past the default budget, then the renewals stop and the lease lapses. + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs + .mockResolvedValueOnce(50_000) + .mockResolvedValue(null) + const turn = runImportTurn() + // Past the default budget (60 s + 30 s grace) the lease is still live: the agent waits. + await vi.advanceTimersByTimeAsync(91_000) + expect(turn.bodies).toHaveLength(1) + // Once that lease runs out, the agent is resumed with the import given up as lost. + await vi.advanceTimersByTimeAsync(52_000) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + expect(turn.forceFailures()).toEqual([ + 'Pending tool execution has no live lease past its budget; force-failing', + ]) + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() + vi.useRealTimers() + } + }) + + it('a lease lookup that fails is checked again instead of failing a renewed import', async () => { + vi.useFakeTimers() + try { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs + .mockRejectedValueOnce(new Error('database unavailable')) + .mockResolvedValueOnce(30_000) + .mockResolvedValue(null) + const turn = runImportTurn() + // The failed lookup at the default budget, and its retry 5 s later, leave the import running. + await vi.advanceTimersByTimeAsync(97_000) + expect(turn.bodies).toHaveLength(1) + // It is given up only once the lease the retry found has lapsed. + await vi.advanceTimersByTimeAsync(32_000) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + expect(turn.forceFailures()).toEqual([ + 'Pending tool execution has no live lease past its budget; force-failing', + ]) + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() + vi.useRealTimers() + } + }) + + it('gives up a renewed import at the client result cap, however long the renewals go on', async () => { + vi.useFakeTimers() + try { + // The page stays alive and renews, but the import itself never finishes. + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockResolvedValue(50_000) + const turn = runImportTurn() + await vi.advanceTimersByTimeAsync(CLIENT_TOOL_RESULT_TIMEOUT_MS - 1_000) + expect(turn.bodies).toHaveLength(1) + const leaseReads = + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mock.calls.length + await vi.advanceTimersByTimeAsync(2_000) + // At the cap the call is given up without reading its lease. + expect( + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mock.calls.length + ).toBe(leaseReads) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + expect(turn.forceFailures()).toEqual([ + 'Pending tool execution reached the client tool result cap; force-failing', + ]) + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() + vi.useRealTimers() + } + }) + + it('gives up an import whose lease cannot be read for a whole lease', async () => { + vi.useFakeTimers() + try { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockRejectedValue( + new Error('database unavailable') + ) + const turn = runImportTurn() + // The first failed lookup comes at the default budget (90 s); retries go on for one lease. + await vi.advanceTimersByTimeAsync(145_000) + expect(turn.bodies).toHaveLength(1) + await vi.advanceTimersByTimeAsync(10_000) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + expect(turn.forceFailures()).toEqual([ + 'Pending tool execution lease could not be read for a whole lease; force-failing', + ]) + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() + vi.useRealTimers() + } + }) + + it('a lease lookup that stalls counts as failed, and gives the import up after one lease', async () => { + vi.useFakeTimers() + try { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockImplementation( + () => new Promise(() => {}) + ) + const turn = runImportTurn() + // The lookup at the default budget (90 s) stalls; each attempt is cut off after 5 s and + // tried again, for one lease. + await vi.advanceTimersByTimeAsync(150_000) + expect(turn.bodies).toHaveLength(1) + await vi.advanceTimersByTimeAsync(20_000) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + expect(turn.forceFailures()).toEqual([ + 'Pending tool execution lease could not be read for a whole lease; force-failing', + ]) + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() + vi.useRealTimers() + } + }) + + it('does not fail a call replaced while its lease was being read', async () => { + vi.useFakeTimers() + try { + let finishReplacement = () => {} + const turn = runImportTurn() + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockImplementationOnce( + async () => { + const context = turn.context() + context.pendingToolPromises.set( + 'tool-import', + new Promise<{ status: 'success' }>((resolve) => { + finishReplacement = () => { + const tool = context.toolCalls.get('tool-import') + if (tool) { + tool.status = MothershipStreamV1ToolOutcome.success + tool.endTime = Date.now() + tool.result = { success: true, output: { imported: true } } + } + context.pendingToolPromises.delete('tool-import') + resolve({ status: 'success' }) + } + }) + ) + // The old promise's lease has lapsed, but the call now belongs to the replacement. + return null + } + ) + await vi.advanceTimersByTimeAsync(91_000) + expect(turn.bodies).toHaveLength(1) + // The replacement is still running: nothing has settled the call as failed. + expect(turn.context().toolCalls.get('tool-import')?.status).toBe('executing') + finishReplacement() + await vi.advanceTimersByTimeAsync(0) + expect((await turn.lifecycle).success).toBe(true) + expect(JSON.stringify(turn.bodies[1].results)).toContain('"success":true') + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() + vi.useRealTimers() + } + }) + it('force-fails each hung tool on its own budget while awaiting a long approval', async () => { vi.useFakeTimers() try { diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index 77eb7e078c0..5d0eee4b19c 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -15,9 +15,17 @@ import { } from '@/lib/billing/core/billing-attribution' import { env } from '@/lib/core/config/env' import { isCopilotToolPermissionsEnabled, isHosted } from '@/lib/core/config/env-flags' +import { SIM_TOOL_EXECUTION_LEASE_SECONDS } from '@/lib/mothership/async-runs/execution-lease' import type { AsyncCompletionSignal } from '@/lib/mothership/async-runs/lifecycle' -import { createRunSegment, updateRunStatus } from '@/lib/mothership/async-runs/repository' -import { TOOL_WATCHDOG_RESUME_GRACE_MS } from '@/lib/mothership/constants' +import { + createRunSegment, + getChatViewDesktopLeaseRemainingMs, + updateRunStatus, +} from '@/lib/mothership/async-runs/repository' +import { + CLIENT_TOOL_RESULT_TIMEOUT_MS, + TOOL_WATCHDOG_RESUME_GRACE_MS, +} from '@/lib/mothership/constants' import { type CopilotEnvironmentContext, prepareCopilotEnvironmentContext, @@ -88,6 +96,64 @@ import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secr const logger = createLogger('CopilotLifecycle') +/** After a renewed lease's end, the wait looks once more before it calls the call lost. */ +const LEASE_RECHECK_SLACK_MS = 1_000 +/** How soon a failed lease lookup is tried again. */ +const LEASE_LOOKUP_RETRY_MS = 5_000 +/** How long one lease lookup may take before it counts as failed. */ +const LEASE_LOOKUP_TIMEOUT_MS = 5_000 + +/** The resume wait's watch over one pending tool call. */ +interface PendingToolWatchdog { + promise: Promise + settlement: Promise<{ toolCallId: string; promise: Promise }> + deadlineAt: number + waitBudgetMs: number + /** When this wait began. */ + startedAt: number + /** Past this, the call is given up however its lease stands: any client tool's cap. */ + ceilingAt: number + /** Why the deadline was moved past the plain budget, if it was. */ + extendedBy?: 'lease' | 'lease_lookup' + /** Until when a failing lease lookup is retried before the call is given up. */ + leaseLookupRetryUntil?: number +} + +/** Why a pending call's wait ran out, in the words its force-fail is logged with. */ +function pendingToolExpiryMessage(watchdog: PendingToolWatchdog, now: number): string { + if (!watchdog.extendedBy) + return 'Pending tool execution exceeded its resume wait budget; force-failing' + if (now >= watchdog.ceilingAt) + return 'Pending tool execution reached the client tool result cap; force-failing' + if (watchdog.leaseLookupRetryUntil !== undefined) + return 'Pending tool execution lease could not be read for a whole lease; force-failing' + return 'Pending tool execution has no live lease past its budget; force-failing' +} + +/** A pending call's chat-view lease, or why it could not be read within `timeoutMs`. */ +async function readLeaseWithin( + toolCallId: string, + timeoutMs: number, + abortSignal: AbortSignal | undefined +): Promise<{ remainingMs: number | null } | { error: unknown }> { + const lookup = getChatViewDesktopLeaseRemainingMs(toolCallId).then( + (remainingMs) => ({ remainingMs }), + (error: unknown) => ({ error }) + ) + const giveUp = new AbortController() + const signal = abortSignal ? AbortSignal.any([abortSignal, giveUp.signal]) : giveUp.signal + try { + return await Promise.race([ + lookup, + interruptibleSleep(timeoutMs, signal).then(() => ({ + error: new Error(`The lease lookup took longer than ${timeoutMs} ms`), + })), + ]) + } finally { + giveUp.abort() + } +} + const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected' /** @@ -1351,18 +1417,7 @@ async function runCheckpointLoop( }) let maximumWaitBudgetMs = 0 let timedOutCount = 0 - const pendingWatchdogs = new Map< - string, - { - promise: Promise - settlement: Promise<{ - toolCallId: string - promise: Promise - }> - deadlineAt: number - waitBudgetMs: number - } - >() + const pendingWatchdogs = new Map() /** * A long-running approval must not lend its deadline to an unrelated @@ -1391,26 +1446,74 @@ async function runCheckpointLoop( ), deadlineAt: now + waitBudgetMs, waitBudgetMs, + startedAt: now, + ceilingAt: now + Math.max(waitBudgetMs, CLIENT_TOOL_RESULT_TIMEOUT_MS), }) maximumWaitBudgetMs = Math.max(maximumWaitBudgetMs, waitBudgetMs) } - const expiredTools = Array.from(pendingWatchdogs.entries()).filter( + const overdueTools = Array.from(pendingWatchdogs.entries()).filter( ([toolCallId, watchdog]) => watchdog.deadlineAt <= now && context.pendingToolPromises.get(toolCallId) === watchdog.promise ) + // A desktop import the chat view is running renews its lease while it works: its budget + // runs to the end of that lease, and only a lapsed lease fails it. At its ceiling a call is + // given up however its lease stands, so its lease is not read. + // A lookup that fails says nothing about the lease, so it is retried, for at most one lease. + // Each lookup is bounded, so a stalled read cannot hold the wait past its deadlines or Stop. + const leaseChecked = overdueTools.filter(([, watchdog]) => now < watchdog.ceilingAt) + const leases = await Promise.all( + leaseChecked.map(([toolCallId]) => + readLeaseWithin(toolCallId, LEASE_LOOKUP_TIMEOUT_MS, options.abortSignal) + ) + ) + if (isAborted(options, context)) break + const expiredTools = [ + ...overdueTools.filter( + ([toolCallId, watchdog]) => + now >= watchdog.ceilingAt && + context.pendingToolPromises.get(toolCallId) === watchdog.promise + ), + ...leaseChecked.filter(([toolCallId, watchdog], index) => { + // A call replaced while its lease was read belongs to its new watchdog. + if (context.pendingToolPromises.get(toolCallId) !== watchdog.promise) return false + const lease = leases[index] + const checkedAt = Date.now() + const extendTo = (deadlineAt: number, reason: 'lease' | 'lease_lookup') => { + watchdog.deadlineAt = Math.min(deadlineAt, watchdog.ceilingAt) + watchdog.extendedBy = reason + maximumWaitBudgetMs = Math.max( + maximumWaitBudgetMs, + watchdog.deadlineAt - watchdog.startedAt + ) + return false + } + if ('error' in lease) { + watchdog.leaseLookupRetryUntil ??= checkedAt + SIM_TOOL_EXECUTION_LEASE_SECONDS * 1000 + if (checkedAt >= watchdog.leaseLookupRetryUntil) return true + logger.warn('Could not read a pending tool call lease; checking again', { + toolCallId, + error: getErrorMessage(lease.error), + }) + return extendTo(checkedAt + LEASE_LOOKUP_RETRY_MS, 'lease_lookup') + } + watchdog.leaseLookupRetryUntil = undefined + if (typeof lease.remainingMs !== 'number' || lease.remainingMs <= 0) return true + return extendTo(checkedAt + lease.remainingMs + LEASE_RECHECK_SLACK_MS, 'lease') + }), + ] if (expiredTools.length > 0) { await Promise.all( expiredTools.map(async ([toolCallId, watchdog]) => { - logger.error( - 'Pending tool execution exceeded its resume wait budget; force-failing', - { - checkpointId: continuation.checkpointId, - toolCallId, - waitBudgetMs: watchdog.waitBudgetMs, - } - ) + const waitedMs = Date.now() - watchdog.startedAt + logger.error(pendingToolExpiryMessage(watchdog, now), { + checkpointId: continuation.checkpointId, + toolCallId, + waitBudgetMs: watchdog.waitBudgetMs, + waitedMs, + ...(watchdog.extendedBy ? { extendedBy: watchdog.extendedBy } : {}), + }) await failPendingToolCall(toolCallId, context, execContext) if (context.pendingToolPromises.get(toolCallId) === watchdog.promise) { context.pendingToolPromises.delete(toolCallId) diff --git a/apps/sim/lib/mothership/tools/client/desktop-tool-chat-view-lease.integration.ts b/apps/sim/lib/mothership/tools/client/desktop-tool-chat-view-lease.integration.ts new file mode 100644 index 00000000000..a13598bcb6c --- /dev/null +++ b/apps/sim/lib/mothership/tools/client/desktop-tool-chat-view-lease.integration.ts @@ -0,0 +1,335 @@ +/** + * An import the chat view runs keeps its turn waiting for as long as it works: its claim takes the + * execution lease under the claiming session, the chat view renews it through the same lease route + * a desktop's background executor uses, and the turn's wait budget runs to the end of that lease. + * Runs against real PostgreSQL and Redis through the production pre-persist path, the authorize + * route, the lease route and 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 (rate limits, the confirmation channel) reads 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, + desktopDevices, + permissions, + user, + workspace, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, sql } from 'drizzle-orm' +import { NextRequest } from 'next/server' +import { closeRedisConnection } from '@/lib/core/config/redis' +import { + DESKTOP_TOOL_CLAIM_OWNER, + SIM_TOOL_EXECUTION_VERSION, +} from '@/lib/mothership/async-runs/lifecycle' +import { + claimDesktopToolCall, + getChatViewDesktopLeaseRemainingMs, +} from '@/lib/mothership/async-runs/repository' +import { prePersistClientExecutableToolCall } from '@/lib/mothership/request/handlers' +import { TraceCollector } from '@/lib/mothership/request/trace' +import type { StreamingContext } from '@/lib/mothership/request/types' +import { chatViewDesktopLeaseOwnerToken } from '@/lib/mothership/tools/desktop-tools' +import { POST as confirmPOST } from '@/app/api/copilot/confirm/route' +import { POST as authorizePOST } from '@/app/api/desktop/tool/authorize/route' +import { POST as leasePOST } from '@/app/api/desktop/tool/lease/route' + +const APP_ORIGIN = 'http://localhost:3000' +const LEASE_MS = 60_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) => + post(authorizePOST, '/api/desktop/tool/authorize', { toolCallId, claim: true }) + +/** The chat view renewing an import's lease through the desktop lease route. */ +const chatViewRenews = (toolCallId: string) => + leasePOST( + new NextRequest(new URL('/api/desktop/tool/lease', APP_ORIGIN), { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ toolCallId, chatView: true }), + }) + ) + +/** Moves a call's lease end to `seconds` from now on the database clock. */ +async function leaseEndsIn(toolCallId: string, seconds: number) { + await db + .update(copilotAsyncToolCalls) + .set({ + executionLeaseExpiresAt: sql`clock_timestamp() + ${seconds} * interval '1 second'`, + }) + .where(eq(copilotAsyncToolCalls.toolCallId, toolCallId)) +} + +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 chat-view import kept alive by its lease', () => { + const userId = generateId() + const otherUserId = generateId() + const workspaceId = generateId() + const chatId = generateId() + const sessionId = generateId() + + function signedInAs(id: string, session: string) { + authMockFns.mockGetSession.mockResolvedValue({ + user: { id, email: `${id}@chat-view-lease.test`, name: 'Desktop' }, + session: { id: session, userId: id }, + }) + } + + async function startRun(desktopDeviceId?: string) { + const runId = generateId() + await db.insert(copilotRuns).values({ + id: runId, + executionId: generateId(), + chatId, + userId, + workspaceId, + streamId: generateId(), + toolExecutionVersion: SIM_TOOL_EXECUTION_VERSION, + status: 'active', + requestContext: { source: 'headless_lifecycle' }, + ...(desktopDeviceId ? { desktopDeviceId } : {}), + }) + return runId + } + + /** The production pre-persist path for a desktop call the agent issues on `runId`. */ + async function agentCalls(runId: string, toolName: string, args: Record) { + 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 }, + } + await prePersistClientExecutableToolCall( + { + type: 'tool', + payload: { + toolCallId, + toolName, + arguments: args, + executor: 'client', + mode: 'async', + phase: 'call', + }, + }, + context, + {} + ) + return toolCallId + } + + /** An import the chat view claimed, as the desktop app's manifest request does. */ + async function chatViewImports() { + const runId = await startRun() + const toolCallId = await agentCalls(runId, 'import_local_files', { + path: '/Users/fixture/Reports', + targetWorkspaceId: workspaceId, + }) + const claim = await desktopClaims(toolCallId) + expect(claim.status).toBe(200) + return { runId, toolCallId } + } + + beforeAll(async () => { + const now = new Date() + for (const id of [userId, otherUserId]) + await db.insert(user).values({ + id, + name: 'Chat-view lease fixture', + email: `${id}@chat-view-lease.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Chat-view lease fixture', + ownerId: userId, + billedAccountUserId: userId, + }) + await db.insert(permissions).values({ + id: generateId(), + userId, + entityType: 'workspace', + entityId: workspaceId, + permissionType: 'admin', + }) + await db.insert(copilotChats).values({ + id: chatId, + userId, + workspaceId, + type: 'mothership', + conversationId: generateId(), + }) + }) + + 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)) + await db.delete(user).where(eq(user.id, otherUserId)) + }) + + it('an import claim takes a lease as long as the default budget, which lapses without renewals', async () => { + signedInAs(userId, sessionId) + const { toolCallId } = await chatViewImports() + const remaining = await getChatViewDesktopLeaseRemainingMs(toolCallId) + expect(remaining).toBeGreaterThan(LEASE_MS - 5_000) + expect(remaining).toBeLessThanOrEqual(LEASE_MS) + + await leaseEndsIn(toolCallId, -1) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeNull() + // A lapsed lease is not revived: the turn has already given the import up. + expect((await chatViewRenews(toolCallId)).status).toBe(410) + }) + + it('the claiming session renews it, which keeps the turn waiting', async () => { + signedInAs(userId, sessionId) + const { toolCallId } = await chatViewImports() + await leaseEndsIn(toolCallId, 10) + const renewal = await chatViewRenews(toolCallId) + expect(renewal.status).toBe(200) + expect(await renewal.json()).toEqual({ renewed: true }) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeGreaterThan(LEASE_MS - 5_000) + }) + + it('a read claims without a lease and keeps the default budget', async () => { + signedInAs(userId, sessionId) + const runId = await startRun() + const toolCallId = await agentCalls(runId, 'read_local_file', { path: '/Users/fixture/a.txt' }) + expect((await desktopClaims(toolCallId)).status).toBe(200) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeNull() + expect((await chatViewRenews(toolCallId)).status).toBe(410) + }) + + it('another session of the same user cannot renew it', async () => { + signedInAs(userId, sessionId) + const { toolCallId } = await chatViewImports() + await leaseEndsIn(toolCallId, 10) + signedInAs(userId, generateId()) + expect((await chatViewRenews(toolCallId)).status).toBe(410) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeLessThanOrEqual(10_000) + }) + + it('another user cannot renew it', async () => { + signedInAs(userId, sessionId) + const { toolCallId } = await chatViewImports() + await leaseEndsIn(toolCallId, 10) + signedInAs(otherUserId, sessionId) + expect((await chatViewRenews(toolCallId)).status).toBe(410) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeLessThanOrEqual(10_000) + }) + + it('a settled import cannot be renewed', async () => { + signedInAs(userId, sessionId) + const { toolCallId } = await chatViewImports() + const report = await post(confirmPOST, '/api/copilot/confirm', { + toolCallId, + status: 'success', + message: 'Imported', + data: { success: true, files: [], folders: [] }, + }) + expect(report.status).toBe(200) + expect((await chatViewRenews(toolCallId)).status).toBe(410) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeNull() + }) + + it("a call a desktop device holds is not the chat view's to renew", async () => { + signedInAs(userId, sessionId) + const deviceId = generateId() + await db.insert(desktopDevices).values({ + id: deviceId, + userId, + name: 'Fixture desktop', + appVersion: '0.9.0', + platform: 'darwin-arm64', + capabilities: { executor: 1, browser: true, terminal: true, localFiles: true }, + }) + try { + const runId = await startRun(deviceId) + const toolCallId = await agentCalls(runId, 'read_local_file', { + path: '/Users/fixture/a.txt', + }) + // Offered to the device, as Sim does for a bound run's call. + await db + .update(copilotAsyncToolCalls) + .set({ pickupDeadlineAt: sql`clock_timestamp() + interval '60 seconds'` }) + .where(eq(copilotAsyncToolCalls.toolCallId, toolCallId)) + // Even under the token the chat view would use, a bound run's call stays the device's. + const claim = await claimDesktopToolCall({ + toolCallId, + runId, + userId, + claimedBy: DESKTOP_TOOL_CLAIM_OWNER.files, + executor: { deviceId, ownerToken: chatViewDesktopLeaseOwnerToken(sessionId) }, + }) + expect(claim.outcome).toBe('claimed') + expect((await chatViewRenews(toolCallId)).status).toBe(410) + expect(await getChatViewDesktopLeaseRemainingMs(toolCallId)).toBeNull() + } finally { + await db.delete(copilotRuns).where(eq(copilotRuns.desktopDeviceId, deviceId)) + await db.delete(desktopDevices).where(eq(desktopDevices.id, deviceId)) + } + }) +}) diff --git a/apps/sim/lib/mothership/tools/client/native-files.test.ts b/apps/sim/lib/mothership/tools/client/native-files.test.ts index f275c5e1b90..18d6e430c47 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.test.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.test.ts @@ -4,7 +4,8 @@ import { apiClientRequestMockFns, } from '@sim/testing/mocks/api-client-request.mock' import { libDesktopMock, libDesktopMockFns } from '@sim/testing/mocks/lib-desktop.mock' -import { beforeEach, expect, it, vi } from 'vitest' +import { sleep } from '@sim/utils/helpers' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const hoisted = vi.hoisted(() => ({ invoke: vi.fn(), @@ -23,6 +24,8 @@ vi.mock('@/lib/mothership/tools/client/completion', () => ({ })) import type { DesktopLocalFileManifest } from '@sim/desktop-bridge' +import { ApiClientError } from '@/lib/api/client/errors' +import { renewDesktopToolLeaseContract } from '@/lib/api/contracts/desktop-executor' import { executeNativeFileTool, importNativeFiles, @@ -150,3 +153,114 @@ it('a read the user stopped while the desktop was reading reports nothing', asyn await executeNativeFileTool('tool', 'read_local_file', stop.signal) expect(mocks.complete).not.toHaveBeenCalled() }) + +describe('an import keeps its lease while it runs', () => { + const LEASE_MS = 60_000 + /** + * The server's side of the lease, by the rules the lease route applies: the claim (which the + * desktop makes before it scans, within its 8 s authorization timeout) takes a lease, and a + * renewal extends it only while it is live; a refused renewal answers 410. + */ + let server: { + leaseUntil: number + claimed: boolean + stopped: boolean + renewalsReceived: number + /** How the server answers the renewals it receives next: a failure that may pass, or as usual. */ + answers: (number | 'usual')[] + } + let claimMs: number + let scanMs: number + let finishUpload: () => void + const leaseLive = () => Date.now() < server.leaseUntil + const answer = (status: number) => + Promise.reject(new ApiClientError({ message: `HTTP ${status}`, status, body: {} })) + + beforeEach(() => { + vi.useFakeTimers() + server = { + leaseUntil: 0, + claimed: false, + stopped: false, + renewalsReceived: 0, + answers: [], + } + claimMs = 0 + scanMs = 0 + mocks.invoke.mockImplementation(async (request: { operation: string }) => { + if (request.operation !== 'manifest') + return { ok: true, data: { kind: 'chunk', bytes: new Uint8Array([65, 66, 67]), eof: true } } + await sleep(claimMs) + server.claimed = true + server.leaseUntil = Date.now() + LEASE_MS + await sleep(scanMs) + return { ok: true, data: manifest } + }) + mocks.upload.mockImplementation( + () => + new Promise((resolve) => { + finishUpload = () => resolve({ id: 'saved-file', name: 'report.txt' }) + }) + ) + mocks.json.mockImplementation(async (contract: unknown) => { + if (contract !== renewDesktopToolLeaseContract) return { folder: { id: 'created-folder' } } + server.renewalsReceived += 1 + const transient = server.answers.shift() + if (typeof transient === 'number') return answer(transient) + if (server.stopped || !server.claimed || !leaseLive()) return answer(410) + server.leaseUntil = Date.now() + LEASE_MS + return { renewed: true } + }) + }) + + afterEach(() => { + vi.useRealTimers() + }) + + it('keeps the lease live through a slow scan and a long upload, and lets it lapse after', async () => { + scanMs = 70_000 + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(70_000) + expect(leaseLive()).toBe(true) + await vi.advanceTimersByTimeAsync(100_000) + expect(leaseLive()).toBe(true) + finishUpload() + await run + await vi.advanceTimersByTimeAsync(LEASE_MS + 1_000) + expect(leaseLive()).toBe(false) + }) + + it('keeps renewing when the desktop takes most of its authorization timeout to claim', async () => { + claimMs = 7_000 + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(90_000) + expect(leaseLive()).toBe(true) + finishUpload() + await run + }) + + it('stops renewing once the server refuses the claimed call', async () => { + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(1_000) + expect(leaseLive()).toBe(true) + server.stopped = true + await vi.advanceTimersByTimeAsync(20_000) + const received = server.renewalsReceived + await vi.advanceTimersByTimeAsync(60_000) + expect(server.renewalsReceived).toBe(received) + finishUpload() + await run + }) + + it('keeps renewing through failures that may pass', async () => { + // Beats at 20 s, 60 s and 100 s fail; the beats between them renew. Stopping at any of the + // failures would let the lease lapse by 140 s. + server.answers = [401, 'usual', 429, 'usual', 503] + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(150_000) + expect(server.renewalsReceived).toBe(7) + expect(leaseLive()).toBe(true) + finishUpload() + await run + }) +}) diff --git a/apps/sim/lib/mothership/tools/client/native-files.ts b/apps/sim/lib/mothership/tools/client/native-files.ts index 884aeca5a6d..e26daa5e3f4 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -15,11 +15,13 @@ import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' import { ApiClientError } from '@/lib/api/client/errors' import { requestJson } from '@/lib/api/client/request' +import { renewDesktopToolLeaseContract } from '@/lib/api/contracts/desktop-executor' import { createWorkspaceFileFolderContract, listWorkspaceFileFoldersContract, } from '@/lib/api/contracts/workspace-file-folders' import { getDesktopBridge } from '@/lib/desktop' +import { SIM_TOOL_EXECUTION_HEARTBEAT_MS } from '@/lib/mothership/async-runs/execution-lease' import { ASYNC_TOOL_CONFIRMATION_STATUS } from '@/lib/mothership/async-runs/lifecycle' import { reportClientToolCompletion, @@ -109,6 +111,32 @@ export async function importNativeFiles( } } +/** + * Renews an import's lease every heartbeat from the moment its manifest is requested, the way a + * desktop renews a bound call, so a slow directory scan cannot outlast the lease the claim took. + * The desktop claims the call before it scans, within its authorization timeout, which is well + * inside one heartbeat: every renewal follows the claim, and the claim's own lease covers the + * first beat. A refusal (410) means the call was stopped, settled, or its lease lapsed, and + * renewing stops. Any other failure may pass, and the next beat tries again. + */ +function keepImportLeased(toolCallId: string): { stop(): void } { + const timer = setInterval(() => { + requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( + (error) => { + if (error instanceof ApiClientError && error.status === 410) { + clearInterval(timer) + return + } + logger.warn('Could not renew the import lease; trying again next beat', { + toolCallId, + error: getErrorMessage(error), + }) + } + ) + }, SIM_TOOL_EXECUTION_HEARTBEAT_MS) + return { stop: () => clearInterval(timer) } +} + /** The server claims imports before reading their manifest, preventing replayed uploads. */ export async function executeNativeFileTool( toolCallId: string, @@ -131,6 +159,10 @@ export async function executeNativeFileTool( ) } window.addEventListener('pagehide', onPageHide) + // An import's claim takes a lease under this session: keep it renewed while the import runs, so + // the turn waits for it however long it takes, and no longer than a lease once this page stops + // renewing (closed, crashed, or signed out). + const lease = toolName === 'import_local_files' ? keepImportLeased(toolCallId) : null try { const response = await invoke( { operation: toolName === 'read_local_file' ? 'read' : 'manifest', toolCallId }, @@ -141,12 +173,15 @@ export async function executeNativeFileTool( throw new Error(response.error) } if (response.data.kind === 'chunk') throw new Error('Unexpected chunk outside an import.') - const completion = - response.data.kind === 'manifest' - ? localFileImportCompletion(await importNativeFiles(toolCallId, response.data, signal)) - : localFileReadCompletion(response) - // Cancelled by the user's Stop or by signing out: whoever cancelled it settles the call, as for - // browser actions and granted-folder reads. A failure reported here would race Stop's record. + let completion + if (response.data.kind === 'manifest') { + completion = localFileImportCompletion( + await importNativeFiles(toolCallId, response.data, signal) + ) + } else completion = localFileReadCompletion(response) + // Cancelled by the user's Stop: Stop settles the call. Cancelled by signing out: nobody reports + // it, and the server's resume watchdog settles it once its budget or lease runs out. Either way + // a failure reported here would race the settlement that decides the call. if (signal?.aborted) return await reportClientToolCompletion( toolCallId, @@ -174,6 +209,7 @@ export async function executeNativeFileTool( ) settled = true } finally { + lease?.stop() window.removeEventListener('pagehide', onPageHide) } } diff --git a/apps/sim/lib/mothership/tools/desktop-tools.ts b/apps/sim/lib/mothership/tools/desktop-tools.ts index 32f86b2f769..b250fb993bf 100644 --- a/apps/sim/lib/mothership/tools/desktop-tools.ts +++ b/apps/sim/lib/mothership/tools/desktop-tools.ts @@ -87,3 +87,11 @@ export const STOPPED_BEFORE_START_MESSAGE = /** What the model learns about a desktop call that Stop cancelled after the desktop picked it up. */ export const STOPPED_WHILE_RUNNING_MESSAGE = 'Stopped by the user while the Sim desktop app was running this action. It may already have taken effect; inspect the current state before repeating it.' + +/** + * The lease owner token of a desktop call the chat view claims under a session: only that session + * renews the lease. + */ +export function chatViewDesktopLeaseOwnerToken(sessionId: string): string { + return `chat-view:${sessionId}` +} diff --git a/packages/testing/src/mocks/mothership-async-runs.mock.ts b/packages/testing/src/mocks/mothership-async-runs.mock.ts index 50ce923b138..94e62736d88 100644 --- a/packages/testing/src/mocks/mothership-async-runs.mock.ts +++ b/packages/testing/src/mocks/mothership-async-runs.mock.ts @@ -35,6 +35,10 @@ export const mothershipAsyncRunsMockFns = { mockClaimToolExecution: vi.fn(), mockClaimDesktopToolCall: vi.fn(), mockRenewSimToolExecutionLease: vi.fn(), + /** No chat-view lease by default: a pending call keeps its plain wait budget. */ + mockGetChatViewDesktopLeaseRemainingMs: vi.fn( + async (_toolCallId: string): Promise => null + ), mockRevokeExpiredSimToolExecutions: vi.fn(), mockSettleSimToolExecution: vi.fn(), mockSettleClientWorkflowToolExecution: vi.fn(), @@ -96,6 +100,8 @@ export const mothershipAsyncRunsMock = { claimToolExecution: mothershipAsyncRunsMockFns.mockClaimToolExecution, claimDesktopToolCall: mothershipAsyncRunsMockFns.mockClaimDesktopToolCall, renewSimToolExecutionLease: mothershipAsyncRunsMockFns.mockRenewSimToolExecutionLease, + getChatViewDesktopLeaseRemainingMs: + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs, revokeExpiredSimToolExecutions: mothershipAsyncRunsMockFns.mockRevokeExpiredSimToolExecutions, settleSimToolExecution: mothershipAsyncRunsMockFns.mockSettleSimToolExecution, settleClientWorkflowToolExecution: