From 3c802184ba6e7bb108d7a98be70d5989605c7d4f Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 03:48:55 -0700 Subject: [PATCH 1/7] fix(desktop): keep a chat-view import alive while it works, by the lease its session renews An import the chat view runs was failed as lost once it ran past the default tool budget (60 s plus the 30 s resume grace), though it was still uploading. Its claim now takes the execution lease under the claiming session, the chat view renews it through the existing lease route while the import runs, and the turn's wait budget runs to the end of that lease. Without renewals the lease lapses with the default budget, so a closed or crashed window still settles within about a lease. --- .../e2e/desktop-tools-live-sim.spec.ts | 68 ++++ .../app/api/desktop/tool/authorize/route.ts | 8 +- .../sim/lib/api/contracts/desktop-executor.ts | 20 +- apps/sim/lib/desktop/application/executor.ts | 34 +- .../lib/mothership/async-runs/repository.ts | 51 ++- .../mothership/request/lifecycle/run.test.ts | 77 ++++ .../lib/mothership/request/lifecycle/run.ts | 24 +- ...esktop-tool-chat-view-lease.integration.ts | 335 ++++++++++++++++++ .../mothership/tools/client/native-files.ts | 44 ++- .../sim/lib/mothership/tools/desktop-tools.ts | 17 + .../src/mocks/mothership-async-runs.mock.ts | 6 + 11 files changed, 662 insertions(+), 22 deletions(-) create mode 100644 apps/sim/lib/mothership/tools/client/desktop-tool-chat-view-lease.integration.ts 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..5265fc7b666 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.ts @@ -18,8 +18,10 @@ import { createUnauthorizedResponse, } from '@/lib/mothership/request/http' import { + chatViewDesktopLeaseOwnerToken, getDesktopToolClaimOwner, isDesktopToolCall, + isLeasedChatViewDesktopTool, isLocalReadToolCall, } from '@/lib/mothership/tools/desktop-tools' @@ -52,7 +54,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 +116,15 @@ export const POST = withRouteHandler(async (request: NextRequest) => { { status: 409 } ) if (toolCall.status !== 'pending') return alreadyStarted() + // The import's lease is held by this session; the chat view renews it while the import runs. const { outcome } = await claimDesktopToolCall({ toolCallId: toolCall.toolCallId, runId: toolCall.runId, userId, claimedBy: DESKTOP_TOOL_CLAIM_OWNER.files, + ...(principal && isLeasedChatViewDesktopTool(toolCall.toolName) + ? { 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..a47b2a39d1f 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -3545,6 +3545,83 @@ describe('runCopilotLifecycle', () => { } }) + it('waits on a chat-view import while its lease is renewed, and fails it once the lease lapses', async () => { + vi.useFakeTimers() + try { + const bodies: Record[] = [] + 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' + } + ) + // Renewed once past the default budget, then the renewals stop and the lease lapses. + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs + .mockResolvedValueOnce(50_000) + .mockResolvedValue(null) + mockRunStreamLoop.mockImplementationOnce( + async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => { + bodies.push(JSON.parse(String(fetchOptions.body))) + 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', + }, + } + ) + + // Past the default budget (60 s + 30 s grace) the lease is still live: no force-fail. + await vi.advanceTimersByTimeAsync(91_000) + expect(mockForceFailHungToolCall).not.toHaveBeenCalled() + // Once the lease the renewals kept alive runs out, the import is failed as lost. + await vi.advanceTimersByTimeAsync(52_000) + const result = await lifecycle + expect(mockForceFailHungToolCall).toHaveBeenCalledWith( + 'tool-import', + expect.anything(), + expect.objectContaining({ userId: 'user-1' }) + ) + expect(bodies[1].results).toEqual([ + expect.objectContaining({ callId: 'tool-import', success: false }), + ]) + expect(result.success).toBe(true) + } finally { + 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..213f11e7c1a 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -16,7 +16,11 @@ import { import { env } from '@/lib/core/config/env' import { isCopilotToolPermissionsEnabled, isHosted } from '@/lib/core/config/env-flags' import type { AsyncCompletionSignal } from '@/lib/mothership/async-runs/lifecycle' -import { createRunSegment, updateRunStatus } from '@/lib/mothership/async-runs/repository' +import { + createRunSegment, + getChatViewDesktopLeaseRemainingMs, + updateRunStatus, +} from '@/lib/mothership/async-runs/repository' import { TOOL_WATCHDOG_RESUME_GRACE_MS } from '@/lib/mothership/constants' import { type CopilotEnvironmentContext, @@ -88,6 +92,9 @@ 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 + const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected' /** @@ -1395,11 +1402,24 @@ async function runCheckpointLoop( 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. + const leases = await Promise.all( + overdueTools.map(([toolCallId]) => + getChatViewDesktopLeaseRemainingMs(toolCallId).catch(() => null) + ) + ) + const expiredTools = overdueTools.filter(([, watchdog], index) => { + const remainingMs = leases[index] + if (typeof remainingMs !== 'number' || remainingMs <= 0) return true + watchdog.deadlineAt = Date.now() + remainingMs + LEASE_RECHECK_SLACK_MS + return false + }) if (expiredTools.length > 0) { await Promise.all( expiredTools.map(async ([toolCallId, watchdog]) => { 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..3462d8ac745 --- /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, 2) + 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.ts b/apps/sim/lib/mothership/tools/client/native-files.ts index 884aeca5a6d..492436b727f 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,25 @@ export async function importNativeFiles( } } +/** + * Renews an import's lease every heartbeat until stopped, or until the server refuses (the call + * was stopped, settled, or its lease already lapsed), the way a desktop renews a bound call. + */ +function keepImportLeased(toolCallId: string): { stop(): void } { + const timer = setInterval(() => { + requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( + (error) => { + logger.warn('Could not renew the import lease; it will lapse', { + toolCallId, + error: getErrorMessage(error), + }) + clearInterval(timer) + } + ) + }, 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, @@ -141,12 +162,23 @@ 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') { + // The import's claim took a lease under this session: keep it renewed while files transfer, + // so the turn waits for the import however long it takes, and no longer than a lease once + // this page stops renewing (closed, crashed, or signed out). + const lease = keepImportLeased(toolCallId) + try { + completion = localFileImportCompletion( + await importNativeFiles(toolCallId, response.data, signal) + ) + } finally { + lease.stop() + } + } 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, diff --git a/apps/sim/lib/mothership/tools/desktop-tools.ts b/apps/sim/lib/mothership/tools/desktop-tools.ts index 32f86b2f769..245b6ecab9f 100644 --- a/apps/sim/lib/mothership/tools/desktop-tools.ts +++ b/apps/sim/lib/mothership/tools/desktop-tools.ts @@ -87,3 +87,20 @@ 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}` +} + +/** + * Desktop calls the chat view keeps alive with a renewed lease while they run: imports, which + * take as long as their files do (up to 1,000 entries of up to 64 MB each). Reads finish in + * seconds and keep the default budget. + */ +export function isLeasedChatViewDesktopTool(toolName: string): boolean { + return toolName === 'import_local_files' +} 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: From 313fe15197edecaac1dd0a6a5179f50d9dd6d1f6 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 04:31:03 -0700 Subject: [PATCH 2/7] fix(desktop): keep renewing an import's lease through transient failures, and retry a failed lease lookup - The chat view stops renewing only when the server refuses the call (410) - The resume watchdog retries a failed lease lookup for up to one lease instead of treating it as a lapse - The lifecycle tests assert what the agent is resumed with, and when --- .../mothership/request/lifecycle/run.test.ts | 147 ++++++++++-------- .../lib/mothership/request/lifecycle/run.ts | 31 +++- .../mothership/tools/client/native-files.ts | 11 +- 3 files changed, 118 insertions(+), 71 deletions(-) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index a47b2a39d1f..9e479e14d47 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -3545,78 +3545,99 @@ describe('runCopilotLifecycle', () => { } }) + /** Runs a turn whose import never settles, recording each leg's request body. */ + function runImportTurn() { + const bodies: Record[] = [] + 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))) + 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') + return { bodies, lifecycle, resumedWithLostImport } + } + it('waits on a chat-view import while its lease is renewed, and fails it once the lease lapses', async () => { vi.useFakeTimers() try { - const bodies: Record[] = [] - 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' - } - ) // Renewed once past the default budget, then the renewals stop and the lease lapses. mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs .mockResolvedValueOnce(50_000) .mockResolvedValue(null) - mockRunStreamLoop.mockImplementationOnce( - async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => { - bodies.push(JSON.parse(String(fetchOptions.body))) - 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', - }, - } - ) - - // Past the default budget (60 s + 30 s grace) the lease is still live: no force-fail. + 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(mockForceFailHungToolCall).not.toHaveBeenCalled() - // Once the lease the renewals kept alive runs out, the import is failed as lost. + 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) - const result = await lifecycle - expect(mockForceFailHungToolCall).toHaveBeenCalledWith( - 'tool-import', - expect.anything(), - expect.objectContaining({ userId: 'user-1' }) - ) - expect(bodies[1].results).toEqual([ - expect.objectContaining({ callId: 'tool-import', success: false }), - ]) - expect(result.success).toBe(true) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + } finally { + 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) } finally { vi.useRealTimers() } diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index 213f11e7c1a..d9c7d4d069b 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -15,6 +15,7 @@ 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, @@ -94,6 +95,8 @@ 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 const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected' @@ -1368,6 +1371,8 @@ async function runCheckpointLoop( }> deadlineAt: number waitBudgetMs: number + /** Until when a failing lease lookup is retried before the call is given up. */ + leaseLookupRetryUntil?: number } >() @@ -1409,15 +1414,31 @@ async function runCheckpointLoop( ) // 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. + // A lookup that fails says nothing about the lease, so it is retried, for at most one lease. const leases = await Promise.all( overdueTools.map(([toolCallId]) => - getChatViewDesktopLeaseRemainingMs(toolCallId).catch(() => null) + getChatViewDesktopLeaseRemainingMs(toolCallId).then( + (remainingMs) => ({ remainingMs }), + (error: unknown) => ({ error }) + ) ) ) - const expiredTools = overdueTools.filter(([, watchdog], index) => { - const remainingMs = leases[index] - if (typeof remainingMs !== 'number' || remainingMs <= 0) return true - watchdog.deadlineAt = Date.now() + remainingMs + LEASE_RECHECK_SLACK_MS + const expiredTools = overdueTools.filter(([toolCallId, watchdog], index) => { + const lease = leases[index] + const checkedAt = Date.now() + 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), + }) + watchdog.deadlineAt = checkedAt + LEASE_LOOKUP_RETRY_MS + return false + } + watchdog.leaseLookupRetryUntil = undefined + if (typeof lease.remainingMs !== 'number' || lease.remainingMs <= 0) return true + watchdog.deadlineAt = checkedAt + lease.remainingMs + LEASE_RECHECK_SLACK_MS return false }) if (expiredTools.length > 0) { diff --git a/apps/sim/lib/mothership/tools/client/native-files.ts b/apps/sim/lib/mothership/tools/client/native-files.ts index 492436b727f..492d22e7b2b 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -112,18 +112,23 @@ export async function importNativeFiles( } /** - * Renews an import's lease every heartbeat until stopped, or until the server refuses (the call + * Renews an import's lease every heartbeat until stopped, or until the server refuses it (the call * was stopped, settled, or its lease already lapsed), the way a desktop renews a bound call. */ function keepImportLeased(toolCallId: string): { stop(): void } { const timer = setInterval(() => { requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( (error) => { - logger.warn('Could not renew the import lease; it will lapse', { + // 410: the server refuses the call (stopped, settled, or its lease lapsed), so stop. Any + // other failure may pass: keep renewing, as the lease outlasts a couple of missed beats. + 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), }) - clearInterval(timer) } ) }, SIM_TOOL_EXECUTION_HEARTBEAT_MS) From d7bfc58bd6c3d634bd5a62336375f31f0909f4b3 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 05:23:03 -0700 Subject: [PATCH 3/7] fix(desktop): cap a renewed import's wait at the client tool limit, renew at once, and report the extended wait - A chat-view import's lease extends its wait only up to the cap every client tool has (CLIENT_TOOL_RESULT_TIMEOUT_MS), so an import that hangs with its page alive still settles - The page renews the lease as soon as the import starts, then every heartbeat - The force-fail log names an extended wait and how long it lasted; the wait span's budget includes the extension - Tests for the cap, the bound on failed lease lookups, and the client heartbeat --- .../app/api/desktop/tool/authorize/route.ts | 6 +- .../mothership/request/lifecycle/run.test.ts | 39 ++++++++++ .../lib/mothership/request/lifecycle/run.ts | 36 ++++++++-- .../tools/client/native-files.test.ts | 71 ++++++++++++++++++- .../mothership/tools/client/native-files.ts | 25 ++++--- .../sim/lib/mothership/tools/desktop-tools.ts | 9 --- 6 files changed, 159 insertions(+), 27 deletions(-) diff --git a/apps/sim/app/api/desktop/tool/authorize/route.ts b/apps/sim/app/api/desktop/tool/authorize/route.ts index 5265fc7b666..4326122636c 100644 --- a/apps/sim/app/api/desktop/tool/authorize/route.ts +++ b/apps/sim/app/api/desktop/tool/authorize/route.ts @@ -21,7 +21,6 @@ import { chatViewDesktopLeaseOwnerToken, getDesktopToolClaimOwner, isDesktopToolCall, - isLeasedChatViewDesktopTool, isLocalReadToolCall, } from '@/lib/mothership/tools/desktop-tools' @@ -116,13 +115,14 @@ export const POST = withRouteHandler(async (request: NextRequest) => { { status: 409 } ) if (toolCall.status !== 'pending') return alreadyStarted() - // The import's lease is held by this session; the chat view renews it while the import runs. + // 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 && isLeasedChatViewDesktopTool(toolCall.toolName) + ...(principal ? { chatView: { ownerToken: chatViewDesktopLeaseOwnerToken(principal.sessionId) } } : {}), }) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index 9e479e14d47..62b51316e51 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -221,6 +221,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, @@ -3619,6 +3620,7 @@ describe('runCopilotLifecycle', () => { expect((await turn.lifecycle).success).toBe(true) expect(turn.resumedWithLostImport()).toBe(true) } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() vi.useRealTimers() } }) @@ -3639,6 +3641,43 @@ describe('runCopilotLifecycle', () => { expect((await turn.lifecycle).success).toBe(true) expect(turn.resumedWithLostImport()).toBe(true) } 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) + await vi.advanceTimersByTimeAsync(2_000) + expect((await turn.lifecycle).success).toBe(true) + expect(turn.resumedWithLostImport()).toBe(true) + } 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) + } finally { + mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset() vi.useRealTimers() } }) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index d9c7d4d069b..ac750885686 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -22,7 +22,10 @@ import { getChatViewDesktopLeaseRemainingMs, updateRunStatus, } from '@/lib/mothership/async-runs/repository' -import { TOOL_WATCHDOG_RESUME_GRACE_MS } from '@/lib/mothership/constants' +import { + CLIENT_TOOL_RESULT_TIMEOUT_MS, + TOOL_WATCHDOG_RESUME_GRACE_MS, +} from '@/lib/mothership/constants' import { type CopilotEnvironmentContext, prepareCopilotEnvironmentContext, @@ -1371,6 +1374,12 @@ async function runCheckpointLoop( }> 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 } @@ -1403,6 +1412,8 @@ async function runCheckpointLoop( ), deadlineAt: now + waitBudgetMs, waitBudgetMs, + startedAt: now, + ceilingAt: now + Math.max(waitBudgetMs, CLIENT_TOOL_RESULT_TIMEOUT_MS), }) maximumWaitBudgetMs = Math.max(maximumWaitBudgetMs, waitBudgetMs) } @@ -1426,6 +1437,16 @@ async function runCheckpointLoop( const expiredTools = overdueTools.filter(([toolCallId, watchdog], index) => { const lease = leases[index] const checkedAt = Date.now() + if (checkedAt >= watchdog.ceilingAt) return true + 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 @@ -1433,23 +1454,26 @@ async function runCheckpointLoop( toolCallId, error: getErrorMessage(lease.error), }) - watchdog.deadlineAt = checkedAt + LEASE_LOOKUP_RETRY_MS - return false + return extendTo(checkedAt + LEASE_LOOKUP_RETRY_MS, 'lease_lookup') } watchdog.leaseLookupRetryUntil = undefined if (typeof lease.remainingMs !== 'number' || lease.remainingMs <= 0) return true - watchdog.deadlineAt = checkedAt + lease.remainingMs + LEASE_RECHECK_SLACK_MS - return false + return extendTo(checkedAt + lease.remainingMs + LEASE_RECHECK_SLACK_MS, 'lease') }) if (expiredTools.length > 0) { await Promise.all( expiredTools.map(async ([toolCallId, watchdog]) => { + const waitedMs = Date.now() - watchdog.startedAt logger.error( - 'Pending tool execution exceeded its resume wait budget; force-failing', + watchdog.extendedBy + ? 'Pending tool execution outlived its renewed lease or its cap; force-failing' + : 'Pending tool execution exceeded its resume wait budget; force-failing', { checkpointId: continuation.checkpointId, toolCallId, waitBudgetMs: watchdog.waitBudgetMs, + waitedMs, + ...(watchdog.extendedBy ? { extendedBy: watchdog.extendedBy } : {}), } ) await failPendingToolCall(toolCallId, context, execContext) 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..97a66eda300 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,7 @@ 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 { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const hoisted = vi.hoisted(() => ({ invoke: vi.fn(), @@ -23,6 +23,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 +152,70 @@ 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', () => { + /** The import's upload, which the test finishes when it chooses. */ + let finishUpload: () => void + /** What the server answers each lease renewal, in turn. */ + let renewals: Array<() => Promise> + const renewalsSent = () => + mocks.json.mock.calls.filter(([contract]) => contract === renewDesktopToolLeaseContract).length + + beforeEach(() => { + vi.useFakeTimers() + renewals = [] + mocks.invoke.mockImplementation(async (request: { operation: string }) => + request.operation === 'manifest' + ? { ok: true, data: manifest } + : { ok: true, data: { kind: 'chunk', bytes: new Uint8Array([65, 66, 67]), eof: true } } + ) + 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' } } + const answer = renewals.shift() + return answer ? answer() : { renewed: true } + }) + }) + + afterEach(() => { + vi.useRealTimers() + }) + + const refused = (status: number) => () => + Promise.reject(new ApiClientError({ message: `HTTP ${status}`, status, body: {} })) + + it('renews at once, then every heartbeat, and stops when the import ends', async () => { + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(0) + expect(renewalsSent()).toBe(1) + await vi.advanceTimersByTimeAsync(40_000) + expect(renewalsSent()).toBe(3) + finishUpload() + await run + await vi.advanceTimersByTimeAsync(60_000) + expect(renewalsSent()).toBe(3) + }) + + it('stops renewing once the server refuses the call', async () => { + renewals.push(refused(410)) + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(60_000) + expect(renewalsSent()).toBe(1) + finishUpload() + await run + }) + + it('keeps renewing through failures that may pass', async () => { + renewals.push(refused(401), refused(429), refused(503)) + const run = executeNativeFileTool('tool', 'import_local_files') + await vi.advanceTimersByTimeAsync(60_000) + expect(renewalsSent()).toBe(4) + 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 492d22e7b2b..786fecbaffe 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -112,17 +112,20 @@ export async function importNativeFiles( } /** - * Renews an import's lease every heartbeat until stopped, or until the server refuses it (the call - * was stopped, settled, or its lease already lapsed), the way a desktop renews a bound call. + * Renews an import's lease at once and then every heartbeat, until stopped or until the server + * refuses it (410: the call was stopped, settled, or its lease already lapsed), the way a desktop + * renews a bound call. The first renewal comes right after the claim, so a slow start cannot + * outlast the lease the claim took. */ function keepImportLeased(toolCallId: string): { stop(): void } { - const timer = setInterval(() => { + let stopped = false + const renew = () => { + if (stopped) return requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( (error) => { - // 410: the server refuses the call (stopped, settled, or its lease lapsed), so stop. Any - // other failure may pass: keep renewing, as the lease outlasts a couple of missed beats. + // Any other failure may pass: keep renewing, as the lease outlasts a couple of missed beats. if (error instanceof ApiClientError && error.status === 410) { - clearInterval(timer) + stop() return } logger.warn('Could not renew the import lease; trying again next beat', { @@ -131,8 +134,14 @@ function keepImportLeased(toolCallId: string): { stop(): void } { }) } ) - }, SIM_TOOL_EXECUTION_HEARTBEAT_MS) - return { stop: () => clearInterval(timer) } + } + const timer = setInterval(renew, SIM_TOOL_EXECUTION_HEARTBEAT_MS) + const stop = () => { + stopped = true + clearInterval(timer) + } + renew() + return { stop } } /** The server claims imports before reading their manifest, preventing replayed uploads. */ diff --git a/apps/sim/lib/mothership/tools/desktop-tools.ts b/apps/sim/lib/mothership/tools/desktop-tools.ts index 245b6ecab9f..b250fb993bf 100644 --- a/apps/sim/lib/mothership/tools/desktop-tools.ts +++ b/apps/sim/lib/mothership/tools/desktop-tools.ts @@ -95,12 +95,3 @@ export const STOPPED_WHILE_RUNNING_MESSAGE = export function chatViewDesktopLeaseOwnerToken(sessionId: string): string { return `chat-view:${sessionId}` } - -/** - * Desktop calls the chat view keeps alive with a renewed lease while they run: imports, which - * take as long as their files do (up to 1,000 entries of up to 64 MB each). Reads finish in - * seconds and keep the default budget. - */ -export function isLeasedChatViewDesktopTool(toolName: string): boolean { - return toolName === 'import_local_files' -} From d79c6264baf7def6a534092f6944d25d66be67de Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 05:58:22 -0700 Subject: [PATCH 4/7] fix(desktop): bound lease lookups, never fail a replaced call, and renew from the start of an import - Each lease lookup gets 5 s (and Stop) before it counts as failed, so a stalled read cannot hold the wait past its deadlines - A call replaced while its lease was read is left to its new watchdog - The page renews from the moment it asks for the manifest; a refusal counts only once the claim is confirmed, and renewing stops on every exit - Heartbeat tests check the lease a fake server holds, not request counts --- .../mothership/request/lifecycle/run.test.ts | 69 +++++++++++++++- .../lib/mothership/request/lifecycle/run.ts | 35 +++++++- .../tools/client/native-files.test.ts | 81 ++++++++++++------- .../mothership/tools/client/native-files.ts | 43 +++++----- 4 files changed, 175 insertions(+), 53 deletions(-) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index 62b51316e51..72ce406c2bb 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.test.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.test.ts @@ -3549,6 +3549,7 @@ 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) @@ -3562,6 +3563,7 @@ describe('runCopilotLifecycle', () => { 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', @@ -3601,7 +3603,11 @@ describe('runCopilotLifecycle', () => { bodies.length === 2 && JSON.stringify(bodies[1].results).includes('"callId":"tool-import"') && JSON.stringify(bodies[1].results).includes('"success":false') - return { bodies, lifecycle, resumedWithLostImport } + const context = () => { + if (!streamContext) throw new Error('The turn did not start its stream') + return streamContext + } + return { bodies, lifecycle, resumedWithLostImport, context } } it('waits on a chat-view import while its lease is renewed, and fails it once the lease lapses', async () => { @@ -3682,6 +3688,67 @@ describe('runCopilotLifecycle', () => { } }) + 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) + } 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 ac750885686..71f7968737e 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -100,6 +100,32 @@ const logger = createLogger('CopilotLifecycle') 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 + +/** 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' @@ -1426,15 +1452,16 @@ async function runCheckpointLoop( // 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. // 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 leases = await Promise.all( overdueTools.map(([toolCallId]) => - getChatViewDesktopLeaseRemainingMs(toolCallId).then( - (remainingMs) => ({ remainingMs }), - (error: unknown) => ({ error }) - ) + readLeaseWithin(toolCallId, LEASE_LOOKUP_TIMEOUT_MS, options.abortSignal) ) ) + if (isAborted(options, context)) break const expiredTools = overdueTools.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() if (checkedAt >= watchdog.ceilingAt) return true 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 97a66eda300..4d05ce612f2 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.test.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.test.ts @@ -4,6 +4,7 @@ import { apiClientRequestMockFns, } from '@sim/testing/mocks/api-client-request.mock' import { libDesktopMock, libDesktopMockFns } from '@sim/testing/mocks/lib-desktop.mock' +import { sleep } from '@sim/utils/helpers' import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' const hoisted = vi.hoisted(() => ({ @@ -154,21 +155,37 @@ it('a read the user stopped while the desktop was reading reports nothing', asyn }) describe('an import keeps its lease while it runs', () => { - /** The import's upload, which the test finishes when it chooses. */ + 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 as it starts building the manifest) 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 + transient: number[] + } + let scanMs: number let finishUpload: () => void - /** What the server answers each lease renewal, in turn. */ - let renewals: Array<() => Promise> - const renewalsSent = () => - mocks.json.mock.calls.filter(([contract]) => contract === renewDesktopToolLeaseContract).length + const leaseLive = () => Date.now() < server.leaseUntil + const answer = (status: number) => + Promise.reject(new ApiClientError({ message: `HTTP ${status}`, status, body: {} })) beforeEach(() => { vi.useFakeTimers() - renewals = [] - mocks.invoke.mockImplementation(async (request: { operation: string }) => - request.operation === 'manifest' - ? { ok: true, data: manifest } - : { ok: true, data: { kind: 'chunk', bytes: new Uint8Array([65, 66, 67]), eof: true } } - ) + server = { leaseUntil: 0, claimed: false, stopped: false, renewalsReceived: 0, transient: [] } + 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 } } + server.claimed = true + server.leaseUntil = Date.now() + LEASE_MS + await sleep(scanMs) + return { ok: true, data: manifest } + }) mocks.upload.mockImplementation( () => new Promise((resolve) => { @@ -177,8 +194,12 @@ describe('an import keeps its lease while it runs', () => { ) mocks.json.mockImplementation(async (contract: unknown) => { if (contract !== renewDesktopToolLeaseContract) return { folder: { id: 'created-folder' } } - const answer = renewals.shift() - return answer ? answer() : { renewed: true } + server.renewalsReceived += 1 + const transient = server.transient.shift() + if (transient) return answer(transient) + if (server.stopped || !server.claimed || !leaseLive()) return answer(410) + server.leaseUntil = Date.now() + LEASE_MS + return { renewed: true } }) }) @@ -186,35 +207,37 @@ describe('an import keeps its lease while it runs', () => { vi.useRealTimers() }) - const refused = (status: number) => () => - Promise.reject(new ApiClientError({ message: `HTTP ${status}`, status, body: {} })) - - it('renews at once, then every heartbeat, and stops when the import ends', async () => { + 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(0) - expect(renewalsSent()).toBe(1) - await vi.advanceTimersByTimeAsync(40_000) - expect(renewalsSent()).toBe(3) + await vi.advanceTimersByTimeAsync(70_000) + expect(leaseLive()).toBe(true) + await vi.advanceTimersByTimeAsync(100_000) + expect(leaseLive()).toBe(true) finishUpload() await run - await vi.advanceTimersByTimeAsync(60_000) - expect(renewalsSent()).toBe(3) + await vi.advanceTimersByTimeAsync(LEASE_MS + 1_000) + expect(leaseLive()).toBe(false) }) - it('stops renewing once the server refuses the call', async () => { - renewals.push(refused(410)) + 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(renewalsSent()).toBe(1) + expect(server.renewalsReceived).toBe(received) finishUpload() await run }) it('keeps renewing through failures that may pass', async () => { - renewals.push(refused(401), refused(429), refused(503)) + server.transient = [401, 429, 503] const run = executeNativeFileTool('tool', 'import_local_files') - await vi.advanceTimersByTimeAsync(60_000) - expect(renewalsSent()).toBe(4) + await vi.advanceTimersByTimeAsync(90_000) + 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 786fecbaffe..b7d18626e5e 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -112,19 +112,20 @@ export async function importNativeFiles( } /** - * Renews an import's lease at once and then every heartbeat, until stopped or until the server - * refuses it (410: the call was stopped, settled, or its lease already lapsed), the way a desktop - * renews a bound call. The first renewal comes right after the claim, so a slow start cannot - * outlast the lease the claim took. + * Renews an import's lease from the moment its manifest is requested, every heartbeat, the way a + * desktop renews a bound call, so a slow directory scan cannot outlast the lease the claim took. + * The claim lands while the manifest is being built, so a refusal (410) counts only once the + * manifest confirmed the claim: then 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 } { +function keepImportLeased(toolCallId: string): { claimed(): void; stop(): void } { + let claimed = false let stopped = false const renew = () => { if (stopped) return requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( (error) => { - // Any other failure may pass: keep renewing, as the lease outlasts a couple of missed beats. - if (error instanceof ApiClientError && error.status === 410) { + if (claimed && error instanceof ApiClientError && error.status === 410) { stop() return } @@ -141,7 +142,13 @@ function keepImportLeased(toolCallId: string): { stop(): void } { clearInterval(timer) } renew() - return { stop } + return { + claimed() { + claimed = true + renew() + }, + stop, + } } /** The server claims imports before reading their manifest, preventing replayed uploads. */ @@ -166,6 +173,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 }, @@ -178,17 +189,10 @@ export async function executeNativeFileTool( if (response.data.kind === 'chunk') throw new Error('Unexpected chunk outside an import.') let completion if (response.data.kind === 'manifest') { - // The import's claim took a lease under this session: keep it renewed while files transfer, - // so the turn waits for the import however long it takes, and no longer than a lease once - // this page stops renewing (closed, crashed, or signed out). - const lease = keepImportLeased(toolCallId) - try { - completion = localFileImportCompletion( - await importNativeFiles(toolCallId, response.data, signal) - ) - } finally { - lease.stop() - } + lease?.claimed() + 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 @@ -220,6 +224,7 @@ export async function executeNativeFileTool( ) settled = true } finally { + lease?.stop() window.removeEventListener('pagehide', onPageHide) } } From 0b20373a3bc413539d108201c0b9675c4893300b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 06:15:08 -0700 Subject: [PATCH 5/7] fix(desktop): judge a lease refusal by whether the claim was confirmed when the renewal was sent --- .../tools/client/native-files.test.ts | 31 +++++++++++++++++-- .../mothership/tools/client/native-files.ts | 5 ++- 2 files changed, 32 insertions(+), 4 deletions(-) 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 4d05ce612f2..cee52ed4ecb 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.test.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.test.ts @@ -167,6 +167,8 @@ describe('an import keeps its lease while it runs', () => { stopped: boolean renewalsReceived: number transient: number[] + /** How long the server takes to answer the next renewal it receives. */ + nextAnswerDelayMs: number } let scanMs: number let finishUpload: () => void @@ -176,7 +178,14 @@ describe('an import keeps its lease while it runs', () => { beforeEach(() => { vi.useFakeTimers() - server = { leaseUntil: 0, claimed: false, stopped: false, renewalsReceived: 0, transient: [] } + server = { + leaseUntil: 0, + claimed: false, + stopped: false, + renewalsReceived: 0, + transient: [], + nextAnswerDelayMs: 0, + } scanMs = 0 mocks.invoke.mockImplementation(async (request: { operation: string }) => { if (request.operation !== 'manifest') @@ -195,10 +204,14 @@ describe('an import keeps its lease while it runs', () => { mocks.json.mockImplementation(async (contract: unknown) => { if (contract !== renewDesktopToolLeaseContract) return { folder: { id: 'created-folder' } } server.renewalsReceived += 1 + const delayMs = server.nextAnswerDelayMs + server.nextAnswerDelayMs = 0 const transient = server.transient.shift() + const refused = server.stopped || !server.claimed || !leaseLive() + if (!transient && !refused) server.leaseUntil = Date.now() + LEASE_MS + await sleep(delayMs) if (transient) return answer(transient) - if (server.stopped || !server.claimed || !leaseLive()) return answer(410) - server.leaseUntil = Date.now() + LEASE_MS + if (refused) return answer(410) return { renewed: true } }) }) @@ -220,6 +233,18 @@ describe('an import keeps its lease while it runs', () => { expect(leaseLive()).toBe(false) }) + it('a refusal of a renewal sent before the claim does not stop it, however late it arrives', async () => { + // The first renewal goes out before the desktop claims the call; its 410 arrives only after + // the manifest confirmed the claim. + server.nextAnswerDelayMs = 5_000 + scanMs = 1_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) diff --git a/apps/sim/lib/mothership/tools/client/native-files.ts b/apps/sim/lib/mothership/tools/client/native-files.ts index b7d18626e5e..0da0c99e6bf 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -123,9 +123,12 @@ function keepImportLeased(toolCallId: string): { claimed(): void; stop(): void } let stopped = false const renew = () => { if (stopped) return + // Whether the claim was confirmed when this renewal was sent: a refusal of one sent before it + // may arrive after, and says nothing about the running import. + const sentAfterClaim = claimed requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( (error) => { - if (claimed && error instanceof ApiClientError && error.status === 410) { + if (sentAfterClaim && error instanceof ApiClientError && error.status === 410) { stop() return } From a1985a7ba7362b8040fbe3385eba7e6c698130a9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 07:36:03 -0700 Subject: [PATCH 6/7] fix(desktop): renew an import's lease only after its claim, and give a capped call up without reading its lease The first heartbeat now comes one beat in, after the desktop's bounded claim, so every renewal follows the claim and a refusal always means the call was stopped, settled, or lapsed. A call at its ceiling is given up before its lease is read, and the force-fail log names whether the budget, the cap, a lapsed lease, or failed lease lookups ended the wait. --- .../mothership/request/lifecycle/run.test.ts | 29 +++- .../lib/mothership/request/lifecycle/run.ts | 135 ++++++++++-------- .../tools/client/native-files.test.ts | 41 +++--- .../mothership/tools/client/native-files.ts | 38 ++--- 4 files changed, 130 insertions(+), 113 deletions(-) diff --git a/apps/sim/lib/mothership/request/lifecycle/run.test.ts b/apps/sim/lib/mothership/request/lifecycle/run.test.ts index 72ce406c2bb..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, @@ -3607,7 +3608,12 @@ describe('runCopilotLifecycle', () => { if (!streamContext) throw new Error('The turn did not start its stream') return streamContext } - return { bodies, lifecycle, resumedWithLostImport, context } + /** 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 () => { @@ -3625,6 +3631,9 @@ describe('runCopilotLifecycle', () => { 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() @@ -3646,6 +3655,9 @@ describe('runCopilotLifecycle', () => { 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() @@ -3660,9 +3672,18 @@ describe('runCopilotLifecycle', () => { 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() @@ -3682,6 +3703,9 @@ describe('runCopilotLifecycle', () => { 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() @@ -3702,6 +3726,9 @@ describe('runCopilotLifecycle', () => { 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() diff --git a/apps/sim/lib/mothership/request/lifecycle/run.ts b/apps/sim/lib/mothership/request/lifecycle/run.ts index 71f7968737e..5d0eee4b19c 100644 --- a/apps/sim/lib/mothership/request/lifecycle/run.ts +++ b/apps/sim/lib/mothership/request/lifecycle/run.ts @@ -103,6 +103,33 @@ 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, @@ -1390,26 +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 - /** 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 - } - >() + const pendingWatchdogs = new Map() /** * A long-running approval must not lend its deadline to an unrelated @@ -1450,59 +1458,62 @@ async function runCheckpointLoop( 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. + // 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( - overdueTools.map(([toolCallId]) => + leaseChecked.map(([toolCallId]) => readLeaseWithin(toolCallId, LEASE_LOOKUP_TIMEOUT_MS, options.abortSignal) ) ) if (isAborted(options, context)) break - const expiredTools = overdueTools.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() - if (checkedAt >= watchdog.ceilingAt) return true - 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') - }) + 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]) => { const waitedMs = Date.now() - watchdog.startedAt - logger.error( - watchdog.extendedBy - ? 'Pending tool execution outlived its renewed lease or its cap; force-failing' - : 'Pending tool execution exceeded its resume wait budget; force-failing', - { - checkpointId: continuation.checkpointId, - toolCallId, - waitBudgetMs: watchdog.waitBudgetMs, - waitedMs, - ...(watchdog.extendedBy ? { extendedBy: watchdog.extendedBy } : {}), - } - ) + 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/native-files.test.ts b/apps/sim/lib/mothership/tools/client/native-files.test.ts index cee52ed4ecb..18d6e430c47 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.test.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.test.ts @@ -158,18 +158,18 @@ 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 as it starts building the manifest) takes a lease, and a renewal extends it only - * while it is live; a refused renewal answers 410. + * 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 - transient: number[] - /** How long the server takes to answer the next renewal it receives. */ - nextAnswerDelayMs: 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 @@ -183,13 +183,14 @@ describe('an import keeps its lease while it runs', () => { claimed: false, stopped: false, renewalsReceived: 0, - transient: [], - nextAnswerDelayMs: 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) @@ -204,14 +205,10 @@ describe('an import keeps its lease while it runs', () => { mocks.json.mockImplementation(async (contract: unknown) => { if (contract !== renewDesktopToolLeaseContract) return { folder: { id: 'created-folder' } } server.renewalsReceived += 1 - const delayMs = server.nextAnswerDelayMs - server.nextAnswerDelayMs = 0 - const transient = server.transient.shift() - const refused = server.stopped || !server.claimed || !leaseLive() - if (!transient && !refused) server.leaseUntil = Date.now() + LEASE_MS - await sleep(delayMs) - if (transient) return answer(transient) - if (refused) return answer(410) + 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 } }) }) @@ -233,11 +230,8 @@ describe('an import keeps its lease while it runs', () => { expect(leaseLive()).toBe(false) }) - it('a refusal of a renewal sent before the claim does not stop it, however late it arrives', async () => { - // The first renewal goes out before the desktop claims the call; its 410 arrives only after - // the manifest confirmed the claim. - server.nextAnswerDelayMs = 5_000 - scanMs = 1_000 + 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) @@ -259,9 +253,12 @@ describe('an import keeps its lease while it runs', () => { }) it('keeps renewing through failures that may pass', async () => { - server.transient = [401, 429, 503] + // 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(90_000) + 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 0da0c99e6bf..e26daa5e3f4 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -112,24 +112,19 @@ export async function importNativeFiles( } /** - * Renews an import's lease from the moment its manifest is requested, every heartbeat, the way a + * 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 claim lands while the manifest is being built, so a refusal (410) counts only once the - * manifest confirmed the claim: then the call was stopped, settled, or its lease lapsed, and + * 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): { claimed(): void; stop(): void } { - let claimed = false - let stopped = false - const renew = () => { - if (stopped) return - // Whether the claim was confirmed when this renewal was sent: a refusal of one sent before it - // may arrive after, and says nothing about the running import. - const sentAfterClaim = claimed +function keepImportLeased(toolCallId: string): { stop(): void } { + const timer = setInterval(() => { requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch( (error) => { - if (sentAfterClaim && error instanceof ApiClientError && error.status === 410) { - stop() + if (error instanceof ApiClientError && error.status === 410) { + clearInterval(timer) return } logger.warn('Could not renew the import lease; trying again next beat', { @@ -138,20 +133,8 @@ function keepImportLeased(toolCallId: string): { claimed(): void; stop(): void } }) } ) - } - const timer = setInterval(renew, SIM_TOOL_EXECUTION_HEARTBEAT_MS) - const stop = () => { - stopped = true - clearInterval(timer) - } - renew() - return { - claimed() { - claimed = true - renew() - }, - stop, - } + }, SIM_TOOL_EXECUTION_HEARTBEAT_MS) + return { stop: () => clearInterval(timer) } } /** The server claims imports before reading their manifest, preventing replayed uploads. */ @@ -192,7 +175,6 @@ export async function executeNativeFileTool( if (response.data.kind === 'chunk') throw new Error('Unexpected chunk outside an import.') let completion if (response.data.kind === 'manifest') { - lease?.claimed() completion = localFileImportCompletion( await importNativeFiles(toolCallId, response.data, signal) ) From c259f70f53db3f696c24d0fbe248e8c5511a5287 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 07:55:49 -0700 Subject: [PATCH 7/7] test(desktop): give the renewal test's lease room for a slow round trip --- .../tools/client/desktop-tool-chat-view-lease.integration.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 index 3462d8ac745..a13598bcb6c 100644 --- 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 @@ -247,7 +247,7 @@ describe.runIf(Boolean(redisUrl))('a chat-view import kept alive by its lease', it('the claiming session renews it, which keeps the turn waiting', async () => { signedInAs(userId, sessionId) const { toolCallId } = await chatViewImports() - await leaseEndsIn(toolCallId, 2) + await leaseEndsIn(toolCallId, 10) const renewal = await chatViewRenews(toolCallId) expect(renewal.status).toBe(200) expect(await renewal.json()).toEqual({ renewed: true })