Skip to content

Commit e310c4c

Browse files
committed
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
1 parent 198a8b7 commit e310c4c

3 files changed

Lines changed: 118 additions & 71 deletions

File tree

‎apps/sim/lib/mothership/request/lifecycle/run.test.ts‎

Lines changed: 84 additions & 63 deletions
Original file line numberDiff line numberDiff line change
@@ -3545,78 +3545,99 @@ describe('runCopilotLifecycle', () => {
35453545
}
35463546
})
35473547

3548+
/** Runs a turn whose import never settles, recording each leg's request body. */
3549+
function runImportTurn() {
3550+
const bodies: Record<string, unknown>[] = []
3551+
mockForceFailHungToolCall.mockImplementation(
3552+
async (toolCallId: string, context: StreamingContext) => {
3553+
const tool = context.toolCalls.get(toolCallId)
3554+
if (!tool) return
3555+
tool.status = MothershipStreamV1ToolOutcome.error
3556+
tool.endTime = Date.now()
3557+
tool.result = { success: false }
3558+
tool.error = 'Tool execution hung'
3559+
}
3560+
)
3561+
mockRunStreamLoop.mockImplementationOnce(
3562+
async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => {
3563+
bodies.push(JSON.parse(String(fetchOptions.body)))
3564+
context.toolCalls.set('tool-import', {
3565+
id: 'tool-import',
3566+
name: 'import_local_files',
3567+
status: 'executing',
3568+
})
3569+
context.pendingToolPromises.set('tool-import', new Promise(() => {}))
3570+
context.awaitingAsyncContinuation = {
3571+
checkpointId: 'ckpt-1',
3572+
pendingToolCallIds: ['tool-import'],
3573+
}
3574+
}
3575+
)
3576+
mockRunStreamLoop.mockImplementationOnce(
3577+
async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => {
3578+
bodies.push(JSON.parse(String(fetchOptions.body)))
3579+
context.accumulatedContent = 'Done.'
3580+
}
3581+
)
3582+
const lifecycle = runCopilotLifecycle(
3583+
{ message: 'import', messageId: 'stream-1' },
3584+
{
3585+
userId: 'user-1',
3586+
workspaceId: 'ws-1',
3587+
chatId: 'chat-1',
3588+
executionId: 'exec-1',
3589+
runId: 'run-1',
3590+
executionContext: {
3591+
userId: 'user-1',
3592+
workflowId: '',
3593+
workspaceId: 'ws-1',
3594+
chatId: 'chat-1',
3595+
},
3596+
}
3597+
)
3598+
/** The agent was resumed with the import given up as lost. */
3599+
const resumedWithLostImport = () =>
3600+
bodies.length === 2 &&
3601+
JSON.stringify(bodies[1].results).includes('"callId":"tool-import"') &&
3602+
JSON.stringify(bodies[1].results).includes('"success":false')
3603+
return { bodies, lifecycle, resumedWithLostImport }
3604+
}
3605+
35483606
it('waits on a chat-view import while its lease is renewed, and fails it once the lease lapses', async () => {
35493607
vi.useFakeTimers()
35503608
try {
3551-
const bodies: Record<string, unknown>[] = []
3552-
mockForceFailHungToolCall.mockImplementation(
3553-
async (toolCallId: string, context: StreamingContext) => {
3554-
const tool = context.toolCalls.get(toolCallId)
3555-
if (!tool) return
3556-
tool.status = MothershipStreamV1ToolOutcome.error
3557-
tool.endTime = Date.now()
3558-
tool.result = { success: false }
3559-
tool.error = 'Tool execution hung'
3560-
}
3561-
)
35623609
// Renewed once past the default budget, then the renewals stop and the lease lapses.
35633610
mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs
35643611
.mockResolvedValueOnce(50_000)
35653612
.mockResolvedValue(null)
3566-
mockRunStreamLoop.mockImplementationOnce(
3567-
async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => {
3568-
bodies.push(JSON.parse(String(fetchOptions.body)))
3569-
context.toolCalls.set('tool-import', {
3570-
id: 'tool-import',
3571-
name: 'import_local_files',
3572-
status: 'executing',
3573-
})
3574-
context.pendingToolPromises.set('tool-import', new Promise(() => {}))
3575-
context.awaitingAsyncContinuation = {
3576-
checkpointId: 'ckpt-1',
3577-
pendingToolCallIds: ['tool-import'],
3578-
}
3579-
}
3580-
)
3581-
mockRunStreamLoop.mockImplementationOnce(
3582-
async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => {
3583-
bodies.push(JSON.parse(String(fetchOptions.body)))
3584-
context.accumulatedContent = 'Done.'
3585-
}
3586-
)
3587-
3588-
const lifecycle = runCopilotLifecycle(
3589-
{ message: 'import', messageId: 'stream-1' },
3590-
{
3591-
userId: 'user-1',
3592-
workspaceId: 'ws-1',
3593-
chatId: 'chat-1',
3594-
executionId: 'exec-1',
3595-
runId: 'run-1',
3596-
executionContext: {
3597-
userId: 'user-1',
3598-
workflowId: '',
3599-
workspaceId: 'ws-1',
3600-
chatId: 'chat-1',
3601-
},
3602-
}
3603-
)
3604-
3605-
// Past the default budget (60 s + 30 s grace) the lease is still live: no force-fail.
3613+
const turn = runImportTurn()
3614+
// Past the default budget (60 s + 30 s grace) the lease is still live: the agent waits.
36063615
await vi.advanceTimersByTimeAsync(91_000)
3607-
expect(mockForceFailHungToolCall).not.toHaveBeenCalled()
3608-
// Once the lease the renewals kept alive runs out, the import is failed as lost.
3616+
expect(turn.bodies).toHaveLength(1)
3617+
// Once that lease runs out, the agent is resumed with the import given up as lost.
36093618
await vi.advanceTimersByTimeAsync(52_000)
3610-
const result = await lifecycle
3611-
expect(mockForceFailHungToolCall).toHaveBeenCalledWith(
3612-
'tool-import',
3613-
expect.anything(),
3614-
expect.objectContaining({ userId: 'user-1' })
3615-
)
3616-
expect(bodies[1].results).toEqual([
3617-
expect.objectContaining({ callId: 'tool-import', success: false }),
3618-
])
3619-
expect(result.success).toBe(true)
3619+
expect((await turn.lifecycle).success).toBe(true)
3620+
expect(turn.resumedWithLostImport()).toBe(true)
3621+
} finally {
3622+
vi.useRealTimers()
3623+
}
3624+
})
3625+
3626+
it('a lease lookup that fails is checked again instead of failing a renewed import', async () => {
3627+
vi.useFakeTimers()
3628+
try {
3629+
mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs
3630+
.mockRejectedValueOnce(new Error('database unavailable'))
3631+
.mockResolvedValueOnce(30_000)
3632+
.mockResolvedValue(null)
3633+
const turn = runImportTurn()
3634+
// The failed lookup at the default budget, and its retry 5 s later, leave the import running.
3635+
await vi.advanceTimersByTimeAsync(97_000)
3636+
expect(turn.bodies).toHaveLength(1)
3637+
// It is given up only once the lease the retry found has lapsed.
3638+
await vi.advanceTimersByTimeAsync(32_000)
3639+
expect((await turn.lifecycle).success).toBe(true)
3640+
expect(turn.resumedWithLostImport()).toBe(true)
36203641
} finally {
36213642
vi.useRealTimers()
36223643
}

‎apps/sim/lib/mothership/request/lifecycle/run.ts‎

Lines changed: 26 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
} from '@/lib/billing/core/billing-attribution'
1616
import { env } from '@/lib/core/config/env'
1717
import { isCopilotToolPermissionsEnabled, isHosted } from '@/lib/core/config/env-flags'
18+
import { SIM_TOOL_EXECUTION_LEASE_SECONDS } from '@/lib/mothership/async-runs/execution-lease'
1819
import type { AsyncCompletionSignal } from '@/lib/mothership/async-runs/lifecycle'
1920
import {
2021
createRunSegment,
@@ -94,6 +95,8 @@ const logger = createLogger('CopilotLifecycle')
9495

9596
/** After a renewed lease's end, the wait looks once more before it calls the call lost. */
9697
const LEASE_RECHECK_SLACK_MS = 1_000
98+
/** How soon a failed lease lookup is tried again. */
99+
const LEASE_LOOKUP_RETRY_MS = 5_000
97100

98101
const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected'
99102

@@ -1368,6 +1371,8 @@ async function runCheckpointLoop(
13681371
}>
13691372
deadlineAt: number
13701373
waitBudgetMs: number
1374+
/** Until when a failing lease lookup is retried before the call is given up. */
1375+
leaseLookupRetryUntil?: number
13711376
}
13721377
>()
13731378

@@ -1409,15 +1414,31 @@ async function runCheckpointLoop(
14091414
)
14101415
// A desktop import the chat view is running renews its lease while it works: its budget
14111416
// runs to the end of that lease, and only a lapsed lease fails it.
1417+
// A lookup that fails says nothing about the lease, so it is retried, for at most one lease.
14121418
const leases = await Promise.all(
14131419
overdueTools.map(([toolCallId]) =>
1414-
getChatViewDesktopLeaseRemainingMs(toolCallId).catch(() => null)
1420+
getChatViewDesktopLeaseRemainingMs(toolCallId).then(
1421+
(remainingMs) => ({ remainingMs }),
1422+
(error: unknown) => ({ error })
1423+
)
14151424
)
14161425
)
1417-
const expiredTools = overdueTools.filter(([, watchdog], index) => {
1418-
const remainingMs = leases[index]
1419-
if (typeof remainingMs !== 'number' || remainingMs <= 0) return true
1420-
watchdog.deadlineAt = Date.now() + remainingMs + LEASE_RECHECK_SLACK_MS
1426+
const expiredTools = overdueTools.filter(([toolCallId, watchdog], index) => {
1427+
const lease = leases[index]
1428+
const checkedAt = Date.now()
1429+
if ('error' in lease) {
1430+
watchdog.leaseLookupRetryUntil ??= checkedAt + SIM_TOOL_EXECUTION_LEASE_SECONDS * 1000
1431+
if (checkedAt >= watchdog.leaseLookupRetryUntil) return true
1432+
logger.warn('Could not read a pending tool call lease; checking again', {
1433+
toolCallId,
1434+
error: getErrorMessage(lease.error),
1435+
})
1436+
watchdog.deadlineAt = checkedAt + LEASE_LOOKUP_RETRY_MS
1437+
return false
1438+
}
1439+
watchdog.leaseLookupRetryUntil = undefined
1440+
if (typeof lease.remainingMs !== 'number' || lease.remainingMs <= 0) return true
1441+
watchdog.deadlineAt = checkedAt + lease.remainingMs + LEASE_RECHECK_SLACK_MS
14211442
return false
14221443
})
14231444
if (expiredTools.length > 0) {

‎apps/sim/lib/mothership/tools/client/native-files.ts‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -112,18 +112,23 @@ export async function importNativeFiles(
112112
}
113113

114114
/**
115-
* Renews an import's lease every heartbeat until stopped, or until the server refuses (the call
115+
* Renews an import's lease every heartbeat until stopped, or until the server refuses it (the call
116116
* was stopped, settled, or its lease already lapsed), the way a desktop renews a bound call.
117117
*/
118118
function keepImportLeased(toolCallId: string): { stop(): void } {
119119
const timer = setInterval(() => {
120120
requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch(
121121
(error) => {
122-
logger.warn('Could not renew the import lease; it will lapse', {
122+
// 410: the server refuses the call (stopped, settled, or its lease lapsed), so stop. Any
123+
// other failure may pass: keep renewing, as the lease outlasts a couple of missed beats.
124+
if (error instanceof ApiClientError && error.status === 410) {
125+
clearInterval(timer)
126+
return
127+
}
128+
logger.warn('Could not renew the import lease; trying again next beat', {
123129
toolCallId,
124130
error: getErrorMessage(error),
125131
})
126-
clearInterval(timer)
127132
}
128133
)
129134
}, SIM_TOOL_EXECUTION_HEARTBEAT_MS)

0 commit comments

Comments
 (0)