Skip to content
68 changes: 68 additions & 0 deletions apps/desktop/e2e/desktop-tools-live-sim.spec.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { dirname, join } from 'node:path'
Expand All @@ -11,6 +11,7 @@
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 {
Expand All @@ -34,6 +35,10 @@
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. */
Expand Down Expand Up @@ -472,6 +477,69 @@
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']
Expand Down
8 changes: 7 additions & 1 deletion apps/sim/app/api/desktop/tool/authorize/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import {
createUnauthorizedResponse,
} from '@/lib/mothership/request/http'
import {
chatViewDesktopLeaseOwnerToken,
getDesktopToolClaimOwner,
isDesktopToolCall,
isLocalReadToolCall,
Expand Down Expand Up @@ -52,7 +53,7 @@ function refusedClaimResponse(
* that they have not allowed.
*/
export const POST = withRouteHandler(async (request: NextRequest) => {
const { userId, isAuthenticated } = await authenticateCopilotRequestSessionOnly()
const { userId, isAuthenticated, principal } = await authenticateCopilotRequestSessionOnly()
if (!isAuthenticated || !userId) {
return createUnauthorizedResponse()
}
Expand Down Expand Up @@ -114,11 +115,16 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
{ status: 409 }
)
if (toolCall.status !== 'pending') return alreadyStarted()
// An import runs as long as its files take, so its claim takes a lease this session holds:
// the chat view renews it while the import runs. Reads finish in seconds and take none.
const { outcome } = await claimDesktopToolCall({
toolCallId: toolCall.toolCallId,
runId: toolCall.runId,
userId,
claimedBy: DESKTOP_TOOL_CLAIM_OWNER.files,
...(principal
? { chatView: { ownerToken: chatViewDesktopLeaseOwnerToken(principal.sessionId) } }
: {}),
})
if (outcome !== 'claimed') return refusedClaimResponse(outcome, alreadyStarted)
} else if (
Expand Down
20 changes: 15 additions & 5 deletions apps/sim/lib/api/contracts/desktop-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof renewDesktopToolLeaseBodySchema>

export const renewDesktopToolLeaseResponseSchema = z.object({ renewed: z.literal(true) })
Expand Down
34 changes: 31 additions & 3 deletions apps/sim/lib/desktop/application/executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -50,6 +51,7 @@ import {
settleClientToolCall,
} from '@/lib/mothership/request/tools/client-settlement.server'
import {
chatViewDesktopLeaseOwnerToken,
getDesktopExecutorClaimOwner,
isDesktopToolCall,
} from '@/lib/mothership/tools/desktop-tools'
Expand Down Expand Up @@ -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'],
Expand All @@ -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(
Expand Down
51 changes: 46 additions & 5 deletions apps/sim/lib/mothership/async-runs/repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand All @@ -795,7 +801,8 @@ export type DesktopToolCallClaim =
export async function claimDesktopToolCall(
claimant: DesktopToolCallClaimant
): Promise<DesktopToolCallClaim> {
const { executor } = claimant
const { executor, chatView } = claimant
const leaseOwnerToken = executor?.ownerToken ?? chatView?.ownerToken
return await claimUnderRunAdmission(
{ ...claimant, desktopDeviceId: executor?.deviceId },
claimant.claimedBy,
Expand All @@ -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'`,
}
: {}),
Expand Down Expand Up @@ -868,7 +875,7 @@ export async function claimDesktopToolCall(
*/
export async function renewSimToolExecutionLease(
owner: SimToolExecutionOwner,
desktop?: { deviceId: string }
desktop?: { deviceId: string } | { chatView: true }
): Promise<boolean> {
const [renewed] = await db
.update(copilotAsyncToolCalls)
Expand All @@ -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
)
Expand All @@ -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<number | null> {
const [row] = await db
.select({
remainingMs: sql<number>`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) =>
Expand Down
Loading
Loading