Skip to content

Commit d79c626

Browse files
committed
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
1 parent d7bfc58 commit d79c626

4 files changed

Lines changed: 175 additions & 53 deletions

File tree

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

Lines changed: 68 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3549,6 +3549,7 @@ describe('runCopilotLifecycle', () => {
35493549
/** Runs a turn whose import never settles, recording each leg's request body. */
35503550
function runImportTurn() {
35513551
const bodies: Record<string, unknown>[] = []
3552+
let streamContext: StreamingContext | undefined
35523553
mockForceFailHungToolCall.mockImplementation(
35533554
async (toolCallId: string, context: StreamingContext) => {
35543555
const tool = context.toolCalls.get(toolCallId)
@@ -3562,6 +3563,7 @@ describe('runCopilotLifecycle', () => {
35623563
mockRunStreamLoop.mockImplementationOnce(
35633564
async (_url: string, fetchOptions: RequestInit, context: StreamingContext) => {
35643565
bodies.push(JSON.parse(String(fetchOptions.body)))
3566+
streamContext = context
35653567
context.toolCalls.set('tool-import', {
35663568
id: 'tool-import',
35673569
name: 'import_local_files',
@@ -3601,7 +3603,11 @@ describe('runCopilotLifecycle', () => {
36013603
bodies.length === 2 &&
36023604
JSON.stringify(bodies[1].results).includes('"callId":"tool-import"') &&
36033605
JSON.stringify(bodies[1].results).includes('"success":false')
3604-
return { bodies, lifecycle, resumedWithLostImport }
3606+
const context = () => {
3607+
if (!streamContext) throw new Error('The turn did not start its stream')
3608+
return streamContext
3609+
}
3610+
return { bodies, lifecycle, resumedWithLostImport, context }
36053611
}
36063612

36073613
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', () => {
36823688
}
36833689
})
36843690

3691+
it('a lease lookup that stalls counts as failed, and gives the import up after one lease', async () => {
3692+
vi.useFakeTimers()
3693+
try {
3694+
mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockImplementation(
3695+
() => new Promise<number | null>(() => {})
3696+
)
3697+
const turn = runImportTurn()
3698+
// The lookup at the default budget (90 s) stalls; each attempt is cut off after 5 s and
3699+
// tried again, for one lease.
3700+
await vi.advanceTimersByTimeAsync(150_000)
3701+
expect(turn.bodies).toHaveLength(1)
3702+
await vi.advanceTimersByTimeAsync(20_000)
3703+
expect((await turn.lifecycle).success).toBe(true)
3704+
expect(turn.resumedWithLostImport()).toBe(true)
3705+
} finally {
3706+
mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset()
3707+
vi.useRealTimers()
3708+
}
3709+
})
3710+
3711+
it('does not fail a call replaced while its lease was being read', async () => {
3712+
vi.useFakeTimers()
3713+
try {
3714+
let finishReplacement = () => {}
3715+
const turn = runImportTurn()
3716+
mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockImplementationOnce(
3717+
async () => {
3718+
const context = turn.context()
3719+
context.pendingToolPromises.set(
3720+
'tool-import',
3721+
new Promise<{ status: 'success' }>((resolve) => {
3722+
finishReplacement = () => {
3723+
const tool = context.toolCalls.get('tool-import')
3724+
if (tool) {
3725+
tool.status = MothershipStreamV1ToolOutcome.success
3726+
tool.endTime = Date.now()
3727+
tool.result = { success: true, output: { imported: true } }
3728+
}
3729+
context.pendingToolPromises.delete('tool-import')
3730+
resolve({ status: 'success' })
3731+
}
3732+
})
3733+
)
3734+
// The old promise's lease has lapsed, but the call now belongs to the replacement.
3735+
return null
3736+
}
3737+
)
3738+
await vi.advanceTimersByTimeAsync(91_000)
3739+
expect(turn.bodies).toHaveLength(1)
3740+
// The replacement is still running: nothing has settled the call as failed.
3741+
expect(turn.context().toolCalls.get('tool-import')?.status).toBe('executing')
3742+
finishReplacement()
3743+
await vi.advanceTimersByTimeAsync(0)
3744+
expect((await turn.lifecycle).success).toBe(true)
3745+
expect(JSON.stringify(turn.bodies[1].results)).toContain('"success":true')
3746+
} finally {
3747+
mothershipAsyncRunsMockFns.mockGetChatViewDesktopLeaseRemainingMs.mockReset()
3748+
vi.useRealTimers()
3749+
}
3750+
})
3751+
36853752
it('force-fails each hung tool on its own budget while awaiting a long approval', async () => {
36863753
vi.useFakeTimers()
36873754
try {

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

Lines changed: 31 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -100,6 +100,32 @@ const logger = createLogger('CopilotLifecycle')
100100
const LEASE_RECHECK_SLACK_MS = 1_000
101101
/** How soon a failed lease lookup is tried again. */
102102
const LEASE_LOOKUP_RETRY_MS = 5_000
103+
/** How long one lease lookup may take before it counts as failed. */
104+
const LEASE_LOOKUP_TIMEOUT_MS = 5_000
105+
106+
/** A pending call's chat-view lease, or why it could not be read within `timeoutMs`. */
107+
async function readLeaseWithin(
108+
toolCallId: string,
109+
timeoutMs: number,
110+
abortSignal: AbortSignal | undefined
111+
): Promise<{ remainingMs: number | null } | { error: unknown }> {
112+
const lookup = getChatViewDesktopLeaseRemainingMs(toolCallId).then(
113+
(remainingMs) => ({ remainingMs }),
114+
(error: unknown) => ({ error })
115+
)
116+
const giveUp = new AbortController()
117+
const signal = abortSignal ? AbortSignal.any([abortSignal, giveUp.signal]) : giveUp.signal
118+
try {
119+
return await Promise.race([
120+
lookup,
121+
interruptibleSleep(timeoutMs, signal).then(() => ({
122+
error: new Error(`The lease lookup took longer than ${timeoutMs} ms`),
123+
})),
124+
])
125+
} finally {
126+
giveUp.abort()
127+
}
128+
}
103129

104130
const COPILOT_MODEL_CONTENT_PROJECTION_ERROR = 'Copilot model input could not be safely projected'
105131

@@ -1426,15 +1452,16 @@ async function runCheckpointLoop(
14261452
// A desktop import the chat view is running renews its lease while it works: its budget
14271453
// runs to the end of that lease, and only a lapsed lease fails it.
14281454
// A lookup that fails says nothing about the lease, so it is retried, for at most one lease.
1455+
// Each lookup is bounded, so a stalled read cannot hold the wait past its deadlines or Stop.
14291456
const leases = await Promise.all(
14301457
overdueTools.map(([toolCallId]) =>
1431-
getChatViewDesktopLeaseRemainingMs(toolCallId).then(
1432-
(remainingMs) => ({ remainingMs }),
1433-
(error: unknown) => ({ error })
1434-
)
1458+
readLeaseWithin(toolCallId, LEASE_LOOKUP_TIMEOUT_MS, options.abortSignal)
14351459
)
14361460
)
1461+
if (isAborted(options, context)) break
14371462
const expiredTools = overdueTools.filter(([toolCallId, watchdog], index) => {
1463+
// A call replaced while its lease was read belongs to its new watchdog.
1464+
if (context.pendingToolPromises.get(toolCallId) !== watchdog.promise) return false
14381465
const lease = leases[index]
14391466
const checkedAt = Date.now()
14401467
if (checkedAt >= watchdog.ceilingAt) return true

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

Lines changed: 52 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import {
44
apiClientRequestMockFns,
55
} from '@sim/testing/mocks/api-client-request.mock'
66
import { libDesktopMock, libDesktopMockFns } from '@sim/testing/mocks/lib-desktop.mock'
7+
import { sleep } from '@sim/utils/helpers'
78
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
89

910
const hoisted = vi.hoisted(() => ({
@@ -154,21 +155,37 @@ it('a read the user stopped while the desktop was reading reports nothing', asyn
154155
})
155156

156157
describe('an import keeps its lease while it runs', () => {
157-
/** The import's upload, which the test finishes when it chooses. */
158+
const LEASE_MS = 60_000
159+
/**
160+
* The server's side of the lease, by the rules the lease route applies: the claim (which the
161+
* desktop makes as it starts building the manifest) takes a lease, and a renewal extends it only
162+
* while it is live; a refused renewal answers 410.
163+
*/
164+
let server: {
165+
leaseUntil: number
166+
claimed: boolean
167+
stopped: boolean
168+
renewalsReceived: number
169+
transient: number[]
170+
}
171+
let scanMs: number
158172
let finishUpload: () => void
159-
/** What the server answers each lease renewal, in turn. */
160-
let renewals: Array<() => Promise<unknown>>
161-
const renewalsSent = () =>
162-
mocks.json.mock.calls.filter(([contract]) => contract === renewDesktopToolLeaseContract).length
173+
const leaseLive = () => Date.now() < server.leaseUntil
174+
const answer = (status: number) =>
175+
Promise.reject(new ApiClientError({ message: `HTTP ${status}`, status, body: {} }))
163176

164177
beforeEach(() => {
165178
vi.useFakeTimers()
166-
renewals = []
167-
mocks.invoke.mockImplementation(async (request: { operation: string }) =>
168-
request.operation === 'manifest'
169-
? { ok: true, data: manifest }
170-
: { ok: true, data: { kind: 'chunk', bytes: new Uint8Array([65, 66, 67]), eof: true } }
171-
)
179+
server = { leaseUntil: 0, claimed: false, stopped: false, renewalsReceived: 0, transient: [] }
180+
scanMs = 0
181+
mocks.invoke.mockImplementation(async (request: { operation: string }) => {
182+
if (request.operation !== 'manifest')
183+
return { ok: true, data: { kind: 'chunk', bytes: new Uint8Array([65, 66, 67]), eof: true } }
184+
server.claimed = true
185+
server.leaseUntil = Date.now() + LEASE_MS
186+
await sleep(scanMs)
187+
return { ok: true, data: manifest }
188+
})
172189
mocks.upload.mockImplementation(
173190
() =>
174191
new Promise((resolve) => {
@@ -177,44 +194,50 @@ describe('an import keeps its lease while it runs', () => {
177194
)
178195
mocks.json.mockImplementation(async (contract: unknown) => {
179196
if (contract !== renewDesktopToolLeaseContract) return { folder: { id: 'created-folder' } }
180-
const answer = renewals.shift()
181-
return answer ? answer() : { renewed: true }
197+
server.renewalsReceived += 1
198+
const transient = server.transient.shift()
199+
if (transient) return answer(transient)
200+
if (server.stopped || !server.claimed || !leaseLive()) return answer(410)
201+
server.leaseUntil = Date.now() + LEASE_MS
202+
return { renewed: true }
182203
})
183204
})
184205

185206
afterEach(() => {
186207
vi.useRealTimers()
187208
})
188209

189-
const refused = (status: number) => () =>
190-
Promise.reject(new ApiClientError({ message: `HTTP ${status}`, status, body: {} }))
191-
192-
it('renews at once, then every heartbeat, and stops when the import ends', async () => {
210+
it('keeps the lease live through a slow scan and a long upload, and lets it lapse after', async () => {
211+
scanMs = 70_000
193212
const run = executeNativeFileTool('tool', 'import_local_files')
194-
await vi.advanceTimersByTimeAsync(0)
195-
expect(renewalsSent()).toBe(1)
196-
await vi.advanceTimersByTimeAsync(40_000)
197-
expect(renewalsSent()).toBe(3)
213+
await vi.advanceTimersByTimeAsync(70_000)
214+
expect(leaseLive()).toBe(true)
215+
await vi.advanceTimersByTimeAsync(100_000)
216+
expect(leaseLive()).toBe(true)
198217
finishUpload()
199218
await run
200-
await vi.advanceTimersByTimeAsync(60_000)
201-
expect(renewalsSent()).toBe(3)
219+
await vi.advanceTimersByTimeAsync(LEASE_MS + 1_000)
220+
expect(leaseLive()).toBe(false)
202221
})
203222

204-
it('stops renewing once the server refuses the call', async () => {
205-
renewals.push(refused(410))
223+
it('stops renewing once the server refuses the claimed call', async () => {
206224
const run = executeNativeFileTool('tool', 'import_local_files')
225+
await vi.advanceTimersByTimeAsync(1_000)
226+
expect(leaseLive()).toBe(true)
227+
server.stopped = true
228+
await vi.advanceTimersByTimeAsync(20_000)
229+
const received = server.renewalsReceived
207230
await vi.advanceTimersByTimeAsync(60_000)
208-
expect(renewalsSent()).toBe(1)
231+
expect(server.renewalsReceived).toBe(received)
209232
finishUpload()
210233
await run
211234
})
212235

213236
it('keeps renewing through failures that may pass', async () => {
214-
renewals.push(refused(401), refused(429), refused(503))
237+
server.transient = [401, 429, 503]
215238
const run = executeNativeFileTool('tool', 'import_local_files')
216-
await vi.advanceTimersByTimeAsync(60_000)
217-
expect(renewalsSent()).toBe(4)
239+
await vi.advanceTimersByTimeAsync(90_000)
240+
expect(leaseLive()).toBe(true)
218241
finishUpload()
219242
await run
220243
})

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

Lines changed: 24 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -112,19 +112,20 @@ export async function importNativeFiles(
112112
}
113113

114114
/**
115-
* Renews an import's lease at once and then every heartbeat, until stopped or until the server
116-
* refuses it (410: the call was stopped, settled, or its lease already lapsed), the way a desktop
117-
* renews a bound call. The first renewal comes right after the claim, so a slow start cannot
118-
* outlast the lease the claim took.
115+
* Renews an import's lease from the moment its manifest is requested, every heartbeat, the way a
116+
* desktop renews a bound call, so a slow directory scan cannot outlast the lease the claim took.
117+
* The claim lands while the manifest is being built, so a refusal (410) counts only once the
118+
* manifest confirmed the claim: then the call was stopped, settled, or its lease lapsed, and
119+
* renewing stops. Any other failure may pass, and the next beat tries again.
119120
*/
120-
function keepImportLeased(toolCallId: string): { stop(): void } {
121+
function keepImportLeased(toolCallId: string): { claimed(): void; stop(): void } {
122+
let claimed = false
121123
let stopped = false
122124
const renew = () => {
123125
if (stopped) return
124126
requestJson(renewDesktopToolLeaseContract, { body: { toolCallId, chatView: true } }).catch(
125127
(error) => {
126-
// Any other failure may pass: keep renewing, as the lease outlasts a couple of missed beats.
127-
if (error instanceof ApiClientError && error.status === 410) {
128+
if (claimed && error instanceof ApiClientError && error.status === 410) {
128129
stop()
129130
return
130131
}
@@ -141,7 +142,13 @@ function keepImportLeased(toolCallId: string): { stop(): void } {
141142
clearInterval(timer)
142143
}
143144
renew()
144-
return { stop }
145+
return {
146+
claimed() {
147+
claimed = true
148+
renew()
149+
},
150+
stop,
151+
}
145152
}
146153

147154
/** The server claims imports before reading their manifest, preventing replayed uploads. */
@@ -166,6 +173,10 @@ export async function executeNativeFileTool(
166173
)
167174
}
168175
window.addEventListener('pagehide', onPageHide)
176+
// An import's claim takes a lease under this session: keep it renewed while the import runs, so
177+
// the turn waits for it however long it takes, and no longer than a lease once this page stops
178+
// renewing (closed, crashed, or signed out).
179+
const lease = toolName === 'import_local_files' ? keepImportLeased(toolCallId) : null
169180
try {
170181
const response = await invoke(
171182
{ operation: toolName === 'read_local_file' ? 'read' : 'manifest', toolCallId },
@@ -178,17 +189,10 @@ export async function executeNativeFileTool(
178189
if (response.data.kind === 'chunk') throw new Error('Unexpected chunk outside an import.')
179190
let completion
180191
if (response.data.kind === 'manifest') {
181-
// The import's claim took a lease under this session: keep it renewed while files transfer,
182-
// so the turn waits for the import however long it takes, and no longer than a lease once
183-
// this page stops renewing (closed, crashed, or signed out).
184-
const lease = keepImportLeased(toolCallId)
185-
try {
186-
completion = localFileImportCompletion(
187-
await importNativeFiles(toolCallId, response.data, signal)
188-
)
189-
} finally {
190-
lease.stop()
191-
}
192+
lease?.claimed()
193+
completion = localFileImportCompletion(
194+
await importNativeFiles(toolCallId, response.data, signal)
195+
)
192196
} else completion = localFileReadCompletion(response)
193197
// Cancelled by the user's Stop: Stop settles the call. Cancelled by signing out: nobody reports
194198
// 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(
220224
)
221225
settled = true
222226
} finally {
227+
lease?.stop()
223228
window.removeEventListener('pagehide', onPageHide)
224229
}
225230
}

0 commit comments

Comments
 (0)