diff --git a/apps/desktop/e2e/background-executor.spec.ts b/apps/desktop/e2e/background-executor.spec.ts index 4dd67c50458..3f452fc8c74 100644 --- a/apps/desktop/e2e/background-executor.spec.ts +++ b/apps/desktop/e2e/background-executor.spec.ts @@ -1,5 +1,6 @@ import { execFileSync } from 'node:child_process' -import { mkdtempSync, readFileSync, writeFileSync } from 'node:fs' +import { createHash } from 'node:crypto' +import { mkdirSync, mkdtempSync, readFileSync, writeFileSync } from 'node:fs' import { createServer, type IncomingMessage, type Server, type ServerResponse } from 'node:http' import { tmpdir } from 'node:os' import { join } from 'node:path' @@ -17,7 +18,8 @@ import { sleep } from '@sim/utils/helpers' /** * The Sim desktop app's background executor against a fixture Sim that speaks the executor's - * device protocol (register, inbox, doorbell, claim, lease, complete) the way Sim's routes do. + * device protocol (register, inbox, doorbell, claim, lease, complete, import) the way Sim's + * routes do. * The window navigates, reloads and leaves the chats while their calls run: nothing in this * suite depends on a chat view, which is the point. Each scenario's checks land in a JSON * report at BACKGROUND_EXECUTOR_REPORT_PATH. @@ -66,12 +68,26 @@ interface ReportCheck { const report: ReportCheck[] = [] +interface ImportedEntry { + toolCallId: string + kind: string + sourceName: string + relativePath: string + sha256?: string + bytes?: number +} + +function sha256(bytes: Buffer): string { + return createHash('sha256').update(bytes).digest('hex') +} + class FixtureSim { readonly calls = new Map() readonly devices = new Map() readonly requests: string[] = [] readonly hits = new Map() readonly streams = new Map>() + readonly imported: ImportedEntry[] = [] enabled = true offline = false droppedWhileOffline = 0 @@ -100,6 +116,7 @@ class FixtureSim { this.devices.clear() this.requests.length = 0 this.hits.clear() + this.imported.length = 0 this.enabled = true this.offline = false this.droppedWhileOffline = 0 @@ -156,12 +173,60 @@ class FixtureSim { } } + private async raw(request: IncomingMessage): Promise { + const chunks: Buffer[] = [] + for await (const chunk of request) chunks.push(Buffer.from(chunk)) + return Buffer.concat(chunks) + } + private async body(request: IncomingMessage): Promise> { - let text = '' - for await (const chunk of request) text += chunk.toString() + const text = (await this.raw(request)).toString() return text ? (JSON.parse(text) as Record) : {} } + /** Stores one entry of a claimed, running import, as Sim's import route does. */ + private async importEntry( + url: URL, + request: IncomingMessage, + response: ServerResponse + ): Promise { + const query = Object.fromEntries(url.searchParams) + const call = this.calls.get(query.toolCallId ?? '') + if ( + request.method !== 'PUT' || + !call || + call.deviceId !== query.deviceId || + call.toolName !== 'import_local_files' || + call.token !== request.headers['x-sim-execution-token'] || + call.status !== 'running' + ) { + this.json(response, 404, { error: 'Desktop import not found' }) + return + } + // As Sim's route does: a file must declare its length, and arrive whole. + if (query.kind === 'file' && request.headers['content-length'] === undefined) { + this.json(response, 411, { error: 'A file import must declare its length' }) + return + } + const content = await this.raw(request) + if (query.kind === 'file' && content.length !== Number(request.headers['content-length'])) { + this.json(response, 400, { error: 'The file did not arrive whole' }) + return + } + const entry: ImportedEntry = { + toolCallId: call.toolCallId, + kind: query.kind ?? '', + sourceName: query.sourceName ?? '', + relativePath: query.relativePath ?? '', + ...(query.kind === 'file' ? { sha256: sha256(content), bytes: content.length } : {}), + } + this.imported.push(entry) + this.json(response, 200, { + id: `entry-${this.imported.length}`, + name: entry.relativePath.split('/').at(-1) || entry.sourceName, + }) + } + private json(response: ServerResponse, status: number, body: unknown): void { response.writeHead(status, { 'Content-Type': 'application/json' }) response.end(JSON.stringify(body)) @@ -261,6 +326,10 @@ class FixtureSim { this.json(response, 200, { items }) return } + if (path === '/api/desktop/tool/import') { + await this.importEntry(url, request, response) + return + } if (path.startsWith('/api/desktop/tool/')) { const body = await this.body(request) const call = this.calls.get(String(body.toolCallId)) @@ -585,6 +654,67 @@ test.describe('background executor', () => { }) }) + test('D: a folder import lands in Sim while the user is in another chat', async () => { + const userData = mkdtempSync(join(tmpdir(), 'sim-executor-d-')) + const launched = await launch(userData) + app = launched.app + const deviceId = await registeredDevice() + const source = join(userData, 'Reports') + mkdirSync(join(source, 'q3'), { recursive: true }) + writeFileSync(join(source, 'notes.txt'), 'remember the numbers') + // Larger than one 8 MB read, so the file crosses Electron in several chunks. + const large = Buffer.alloc(9 * 1024 * 1024, 7) + writeFileSync(join(source, 'q3', 'export.bin'), large) + await launched.window.goto(`${sim.origin}/workspace/${WORKSPACE}/chat/${CHAT_C}`) + + const call = sim.issue(deviceId, CHAT_B, 'import_local_files', { + path: source, + targetWorkspaceId: WORKSPACE, + folderId: 'folder-e2e', + }) + + await check('D: the import completes with every entry it stored', async () => { + const completion = await settled(call, 60_000) + expect(completion.status).toBe('success') + expect(completion.data).toMatchObject({ + success: true, + workspaceId: WORKSPACE, + folders: [ + { id: 'entry-1', relativePath: '' }, + { id: 'entry-3', relativePath: 'q3' }, + ], + files: [ + { id: 'entry-2', relativePath: 'notes.txt' }, + { id: 'entry-4', relativePath: 'q3/export.bin' }, + ], + }) + }) + + await check('D: Sim received the tree with the bytes on disk, once', async () => { + expect(sim.imported).toEqual([ + { toolCallId: call, kind: 'directory', sourceName: 'Reports', relativePath: '' }, + { + toolCallId: call, + kind: 'file', + sourceName: 'Reports', + relativePath: 'notes.txt', + sha256: sha256(Buffer.from('remember the numbers')), + bytes: 20, + }, + { toolCallId: call, kind: 'directory', sourceName: 'Reports', relativePath: 'q3' }, + { + toolCallId: call, + kind: 'file', + sourceName: 'Reports', + relativePath: 'q3/export.bin', + sha256: sha256(large), + bytes: large.length, + }, + ]) + expect(sim.requireCall(call).claims).toBe(1) + }) + }) + test('E: Stop from another chat stops a running browser wait and terminal command', async () => { const userData = mkdtempSync(join(tmpdir(), 'sim-executor-e-')) app = (await launch(userData)).app diff --git a/apps/desktop/src/main/desktop-executor/client.test.ts b/apps/desktop/src/main/desktop-executor/client.test.ts index 3dcf4e08f23..f91c3b4e4a9 100644 --- a/apps/desktop/src/main/desktop-executor/client.test.ts +++ b/apps/desktop/src/main/desktop-executor/client.test.ts @@ -1,3 +1,4 @@ +import { DESKTOP_IMPORT_TOKEN_HEADER } from '@sim/desktop-bridge' import { describe, expect, it } from 'vitest' import { createDesktopExecutorClient, UnsendableRequestError } from '@/main/desktop-executor/client' @@ -48,4 +49,38 @@ describe('desktop executor client', () => { const loneSurrogate = /[\uD800-\uDBFF](?![\uDC00-\uDFFF])|(? { + const received: Array<{ url: string; token: string | null }> = [] + const client = createDesktopExecutorClient({ + origin: () => 'https://sim.test', + fetch: async (url, init) => { + received.push({ url, token: new Headers(init.headers).get(DESKTOP_IMPORT_TOKEN_HEADER) }) + return Response.json({ id: 'file-1', name: 'notes.txt' }) + }, + deviceId: '00000000-0000-4000-8000-000000000000', + }) + + await client.importEntry( + { + call: { + toolCallId: 'call-1', + toolName: 'import_local_files', + args: {}, + chatId: 'chat-1', + workspaceId: 'ws-1', + executionToken: 'secret-token-1', + }, + kind: 'file', + sourceName: 'notes.txt', + relativePath: '', + content: new Blob(['hello']), + }, + new AbortController().signal + ) + + expect(received).toHaveLength(1) + expect(received[0]?.token).toBe('secret-token-1') + expect(decodeURIComponent(received[0]?.url ?? '')).not.toContain('secret-token-1') + }) }) diff --git a/apps/desktop/src/main/desktop-executor/client.ts b/apps/desktop/src/main/desktop-executor/client.ts index 7431ddd8f3c..0a17990af6f 100644 --- a/apps/desktop/src/main/desktop-executor/client.ts +++ b/apps/desktop/src/main/desktop-executor/client.ts @@ -3,6 +3,7 @@ * own session cookie, which is the session the device registered under; Sim refuses any other. */ +import { DESKTOP_IMPORT_TOKEN_HEADER } from '@sim/desktop-bridge' import { getErrorMessage } from '@sim/utils/errors' import { parseRetryAfter } from '@sim/utils/retry' import { truncateAtCodePoint } from '@sim/utils/string' @@ -13,14 +14,19 @@ import { type DesktopCompletionRequest, type DesktopDeviceRegistration, type DesktopExecutorTiming, + type DesktopImportEntryRequest, + type DesktopImportedEntry, type DesktopInboxItem, parseClaim, parseCompletionOutcome, + parseImportedEntry, parseInbox, parseRegistration, } from '@/main/desktop-executor/protocol' const REQUEST_TIMEOUT_MS = 15_000 +/** An import entry carries up to a 64 MB file, so it gets longer than a control request. */ +const IMPORT_TIMEOUT_MS = 300_000 /** A request Sim answered with a failure, or that never got an answer (`status` 0). */ export class DeviceRequestError extends Error { @@ -73,6 +79,10 @@ export interface DesktopExecutorClient { claim(toolCallId: string): Promise renewLease(toolCallId: string, executionToken: string): Promise complete(request: DesktopCompletionRequest): Promise + importEntry( + request: DesktopImportEntryRequest, + signal: AbortSignal + ): Promise } interface DesktopExecutorClientOptions { @@ -97,13 +107,16 @@ export function createDesktopExecutorClient( const { deviceId } = options async function send( - method: 'GET' | 'POST', + method: 'GET' | 'POST' | 'PUT', path: string, - body?: Record, - signal?: AbortSignal + body?: Record | Blob, + signal?: AbortSignal, + timeoutMs = REQUEST_TIMEOUT_MS, + extraHeaders: Record = {} ): Promise { - const encoded = body ? encode(body) : undefined - const timeout = AbortSignal.timeout(REQUEST_TIMEOUT_MS) + const timeout = AbortSignal.timeout(timeoutMs) + const raw = body instanceof Blob + const encoded = body === undefined ? undefined : raw ? body : encode(body) let response: Response try { response = await options.fetch(`${options.origin()}${path}`, { @@ -111,7 +124,10 @@ export function createDesktopExecutorClient( credentials: 'include', headers: { Accept: 'application/json', - ...(body ? { 'Content-Type': 'application/json' } : {}), + ...(body + ? { 'Content-Type': raw ? 'application/octet-stream' : 'application/json' } + : {}), + ...extraHeaders, }, ...(encoded !== undefined ? { body: encoded } : {}), signal: signal ? AbortSignal.any([signal, timeout]) : timeout, @@ -204,5 +220,25 @@ export function createDesktopExecutorClient( if (!outcome) throw malformed('completion') return outcome }, + async importEntry({ call, kind, sourceName, relativePath, content }, signal) { + const query = new URLSearchParams({ + deviceId, + toolCallId: call.toolCallId, + kind, + sourceName, + relativePath, + }) + const response = await send( + 'PUT', + `/api/desktop/tool/import?${query}`, + content, + signal, + IMPORT_TIMEOUT_MS, + { [DESKTOP_IMPORT_TOKEN_HEADER]: call.executionToken } + ) + const entry = parseImportedEntry(await response.json().catch(() => null)) + if (!entry) throw malformed('import') + return entry + }, } } diff --git a/apps/desktop/src/main/desktop-executor/executor.test.ts b/apps/desktop/src/main/desktop-executor/executor.test.ts index e512119a266..f1ed7801d55 100644 --- a/apps/desktop/src/main/desktop-executor/executor.test.ts +++ b/apps/desktop/src/main/desktop-executor/executor.test.ts @@ -67,6 +67,9 @@ class FakeSim { openInboxStream: async () => { throw new Error('not used') }, + importEntry: async () => { + throw new Error('not used') + }, claim: async (toolCallId) => { this.claims.push(toolCallId) const error = this.claimError?.(toolCallId) diff --git a/apps/desktop/src/main/desktop-executor/protocol.ts b/apps/desktop/src/main/desktop-executor/protocol.ts index 34b162dcaf0..23e127097ca 100644 --- a/apps/desktop/src/main/desktop-executor/protocol.ts +++ b/apps/desktop/src/main/desktop-executor/protocol.ts @@ -1,6 +1,6 @@ /** * The device side of Sim's background executor protocol (`/api/desktop/devices`, `/inbox`, - * `/inbox/stream`, `/tool/claim`, `/tool/lease`, `/tool/complete`). Sim's contracts are the + * `/inbox/stream`, `/tool/claim`, `/tool/lease`, `/tool/complete`, `/tool/import`). Sim's contracts are the * source of truth; responses are parsed defensively here because a malformed one must never * reach a tool. */ @@ -62,6 +62,25 @@ export interface ClaimedDesktopCall { export type DesktopCompletionOutcome = 'recorded' | 'duplicate' | 'superseded' +/** + * One entry of a claimed import, stored under the call's target folder: `sourceName` is the + * import source's own name and `relativePath` the entry's place inside it (`''` for the source). + */ +export interface DesktopImportEntryRequest { + call: ClaimedDesktopCall + kind: 'file' | 'directory' + sourceName: string + relativePath: string + /** A file's bytes; a directory has none. */ + content?: Blob +} + +/** What Sim stored an import entry as: the file, or the folder it reused or created. */ +export interface DesktopImportedEntry { + id: string + name: string +} + export interface DesktopCompletionRequest { toolCallId: string executionToken: string @@ -175,3 +194,13 @@ export function parseCompletionOutcome(body: unknown): DesktopCompletionOutcome ? body.outcome : null } + +/** + * Reads Sim's answer to an import entry: the id and name it stored the entry as. Null when either + * is missing or empty, which the client reports as a malformed response. + */ +export function parseImportedEntry(body: unknown): DesktopImportedEntry | null { + if (!isRecordLike(body) || typeof body.id !== 'string' || typeof body.name !== 'string') + return null + return body.id && body.name ? { id: body.id, name: body.name } : null +} diff --git a/apps/desktop/src/main/desktop-executor/runner.test.ts b/apps/desktop/src/main/desktop-executor/runner.test.ts index 90fd0685890..f10ac5a3748 100644 --- a/apps/desktop/src/main/desktop-executor/runner.test.ts +++ b/apps/desktop/src/main/desktop-executor/runner.test.ts @@ -1,7 +1,15 @@ +import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import type { TerminalToolResponse } from '@sim/terminal-protocol' import { afterEach, describe, expect, it, vi } from 'vitest' -import type { ClaimedDesktopCall } from '@/main/desktop-executor/protocol' +import { DeviceRequestError } from '@/main/desktop-executor/client' +import type { + ClaimedDesktopCall, + DesktopImportEntryRequest, +} from '@/main/desktop-executor/protocol' import { createDesktopToolRunner, type DesktopToolRunnerDeps } from '@/main/desktop-executor/runner' +import { executeLocalFileRequest } from '@/main/local-files' function terminalCall(toolCallId: string, operation: string): ClaimedDesktopCall { return { @@ -25,7 +33,11 @@ function runner(overrides: Partial = {}) { restoreScope: vi.fn(), }, terminal: { executeTool: vi.fn(), cancelTool: vi.fn(async () => true) }, - localFiles: { read: vi.fn() }, + localFiles: { + request: (call, request) => + executeLocalFileRequest(request, { toolName: call.toolName, args: call.args }), + }, + imports: { importEntry: vi.fn() }, localFilesystem: { handle: vi.fn(), vfsRoot: () => 'user-local/x--1' }, ...overrides, }) @@ -121,3 +133,166 @@ describe('local file calls', () => { expect(completion.message).not.toContain('settings') }) }) + +describe('background imports', () => { + const roots: string[] = [] + + afterEach(async () => { + await Promise.all(roots.splice(0).map((root) => rm(root, { recursive: true, force: true }))) + }) + + /** A `Reports` folder holding `q3/summary.txt` and `notes.txt`. */ + async function reportsFolder(): Promise { + const root = await mkdtemp(join(tmpdir(), 'sim-runner-import-')) + roots.push(root) + const reports = join(root, 'Reports') + await mkdir(join(reports, 'q3'), { recursive: true }) + await writeFile(join(reports, 'q3', 'summary.txt'), 'quarterly numbers') + await writeFile(join(reports, 'notes.txt'), 'remember') + return reports + } + + function importCall(path: string): ClaimedDesktopCall { + return { + toolCallId: 'import-1', + toolName: 'import_local_files', + args: { path, targetWorkspaceId: 'ws-1', folderId: 'folder-1' }, + chatId: 'chat-b', + workspaceId: 'ws-1', + executionToken: 'token-import-1', + } + } + + /** + * A fake Sim import route: records every entry it stores, with each file's bytes as text, and + * the token each request presented. `answer` can refuse a request instead. + */ + function recordingSim(answer?: (request: DesktopImportEntryRequest) => Error | null) { + const stored: Array<{ kind: string; relativePath: string; text?: string }> = [] + const tokens: string[] = [] + let requests = 0 + const importEntry = async (request: DesktopImportEntryRequest) => { + requests += 1 + tokens.push(request.call.executionToken) + const refusal = answer?.(request) + if (refusal) throw refusal + stored.push({ + kind: request.kind, + relativePath: request.relativePath, + ...(request.content ? { text: await request.content.text() } : {}), + }) + return { id: `id-${stored.length}`, name: request.relativePath || request.sourceName } + } + return { stored, tokens, importEntry, requests: () => requests } + } + + it("stores a folder's tree in Sim, each file with the bytes on disk", async () => { + const sim = recordingSim() + const completion = await runner({ imports: { importEntry: sim.importEntry } }).run( + importCall(await reportsFolder()), + new AbortController().signal + ) + + expect(completion.status).toBe('success') + expect(sim.stored).toEqual([ + { kind: 'directory', relativePath: '' }, + { kind: 'file', relativePath: 'notes.txt', text: 'remember' }, + { kind: 'directory', relativePath: 'q3' }, + { kind: 'file', relativePath: 'q3/summary.txt', text: 'quarterly numbers' }, + ]) + expect(new Set(sim.tokens)).toEqual(new Set(['token-import-1'])) + expect(completion.data).toMatchObject({ + success: true, + workspaceId: 'ws-1', + folders: [ + { id: 'id-1', relativePath: '' }, + { id: 'id-3', relativePath: 'q3' }, + ], + files: [ + { id: 'id-2', relativePath: 'notes.txt' }, + { id: 'id-4', relativePath: 'q3/summary.txt' }, + ], + }) + }) + + it('reports what landed when an import stops part way, and not to retry it', async () => { + const sim = recordingSim((request) => + request.relativePath === 'q3/summary.txt' ? new Error('Sim refused the entry') : null + ) + const completion = await runner({ imports: { importEntry: sim.importEntry } }).run( + importCall(await reportsFolder()), + new AbortController().signal + ) + + expect(completion.status).toBe('error') + expect(completion.data).toMatchObject({ + success: false, + partial: true, + doNotRetry: true, + outcomeUnknown: true, + error: 'Sim refused the entry', + files: [{ id: 'id-2', relativePath: 'notes.txt' }], + }) + }) + + it('stores nothing more once the call is stopped', async () => { + const controller = new AbortController() + const sim = recordingSim(() => { + controller.abort() + return null + }) + const completion = await runner({ imports: { importEntry: sim.importEntry } }).run( + importCall(await reportsFolder()), + controller.signal + ) + + expect(completion.status).toBe('error') + expect(sim.stored).toEqual([{ kind: 'directory', relativePath: '' }]) + }) + + it("waits out Sim's rate limit instead of failing the import part way", async () => { + let limited = false + const sim = recordingSim((request) => { + if (request.relativePath !== 'notes.txt' || limited) return null + limited = true + return new DeviceRequestError(429, 'Too many requests', 10) + }) + const completion = await runner({ imports: { importEntry: sim.importEntry } }).run( + importCall(await reportsFolder()), + new AbortController().signal + ) + + expect(completion.status).toBe('success') + expect(sim.stored.map((entry) => entry.relativePath)).toEqual([ + '', + 'notes.txt', + 'q3', + 'q3/summary.txt', + ]) + }) + + it('reports an import the rate limit never let start as safe to ask for again', async () => { + const sim = recordingSim(() => new DeviceRequestError(429, 'Too many requests', 1)) + const completion = await runner({ imports: { importEntry: sim.importEntry } }).run( + importCall(await reportsFolder()), + new AbortController().signal + ) + + expect(completion.status).toBe('error') + expect(completion.data).toMatchObject({ partial: false, files: [], folders: [] }) + expect(completion.data).not.toHaveProperty('doNotRetry') + expect(completion.data).not.toHaveProperty('outcomeUnknown') + }) + + it('fails without storing anything when the source cannot be read', async () => { + const sim = recordingSim() + const completion = await runner({ imports: { importEntry: sim.importEntry } }).run( + importCall(join(tmpdir(), 'sim-runner-import-missing', 'Reports')), + new AbortController().signal + ) + + expect(completion.status).toBe('error') + expect(sim.requests()).toBe(0) + expect(completion.data).toMatchObject({ workspaceId: 'ws-1', partial: false }) + }) +}) diff --git a/apps/desktop/src/main/desktop-executor/runner.ts b/apps/desktop/src/main/desktop-executor/runner.ts index fc700251a3b..1bb3663d765 100644 --- a/apps/desktop/src/main/desktop-executor/runner.ts +++ b/apps/desktop/src/main/desktop-executor/runner.ts @@ -10,20 +10,26 @@ import { isCurrentBrowserToolName, } from '@sim/browser-protocol' import type { + DesktopLocalFileRequest, DesktopLocalFileResponse, LocalFilesystemRequest, LocalFilesystemResponse, } from '@sim/desktop-bridge' import { runUserLocalFilesystemTool } from '@sim/desktop-bridge/local-filesystem-tools' import { + assertImportableManifest, browserSessionClosedCompletion, browserToolCompletion, browserToolFailure, browserToolNeedsLivePage, browserToolTimeoutMessage, + type DesktopLocalFileImportResult, type DesktopToolCompletion, + localFileImportCompletion, + localFileImportFailure, localFileReadCompletion, localFilesystemToolCompletion, + readImportEntry, terminalOperationTimeoutMs, terminalToolCompletion, terminalToolFailure, @@ -38,12 +44,21 @@ import { import { getErrorMessage } from '@sim/utils/errors' import { interruptibleSleep } from '@sim/utils/helpers' import { isRecordLike } from '@sim/utils/object' +import { backoffWithJitter } from '@sim/utils/retry' +import { DeviceRequestError } from '@/main/desktop-executor/client' import type { DesktopToolRunner } from '@/main/desktop-executor/executor' -import type { ClaimedDesktopCall } from '@/main/desktop-executor/protocol' +import type { + ClaimedDesktopCall, + DesktopImportEntryRequest, + DesktopImportedEntry, +} from '@/main/desktop-executor/protocol' const logger = createLogger('DesktopExecutorRunner') const USER_LOCAL_TOOLS: ReadonlySet = new Set(['read', 'grep', 'glob']) +/** A rate-limited import entry is retried this many times; Sim never stored a refused one. */ +const IMPORT_RATE_LIMIT_ATTEMPTS = 8 +const IMPORT_RATE_LIMIT_MAX_WAIT_MS = 30_000 /** The model learns a call never ran because a surface is switched off on this machine. */ function surfaceOff(surface: string): DesktopToolCompletion { @@ -87,8 +102,18 @@ export interface DesktopToolRunnerDeps { ): Promise cancelTool(scope: string, toolCallId: string): Promise } + /** Reads a `read_local_file` or `import_local_files` source, authorized by the call itself. */ localFiles: { - read(call: ClaimedDesktopCall): Promise + request( + call: ClaimedDesktopCall, + request: DesktopLocalFileRequest + ): Promise + } + imports: { + importEntry( + request: DesktopImportEntryRequest, + signal: AbortSignal + ): Promise } localFilesystem: { handle(request: LocalFilesystemRequest): Promise @@ -224,6 +249,86 @@ export function createDesktopToolRunner(deps: DesktopToolRunnerDeps): DesktopToo } } + /** + * Imports the call's source into its workspace one entry at a time: each directory as a folder, + * each file read in chunks and checked against the manifest that listed it. + */ + /** + * Sends one entry, waiting out Sim's rate limit: a large tree can outrun it, and a 429 means + * nothing was stored, so sending the entry again cannot duplicate it. + */ + async function importEntry( + request: DesktopImportEntryRequest, + signal: AbortSignal + ): Promise { + for (let attempt = 1; ; attempt++) { + try { + return await deps.imports.importEntry(request, signal) + } catch (error) { + if ( + !(error instanceof DeviceRequestError) || + error.status !== 429 || + attempt >= IMPORT_RATE_LIMIT_ATTEMPTS + ) { + throw error + } + await interruptibleSleep( + backoffWithJitter(attempt, error.retryAfterMs, { maxMs: IMPORT_RATE_LIMIT_MAX_WAIT_MS }), + signal + ) + signal.throwIfAborted() + } + } + } + + async function runImport( + call: ClaimedDesktopCall, + signal: AbortSignal + ): Promise { + const files: DesktopLocalFileImportResult['files'] = [] + const folders: DesktopLocalFileImportResult['folders'] = [] + const targetWorkspaceId = + typeof call.args.targetWorkspaceId === 'string' ? call.args.targetWorkspaceId : '' + const read = (request: DesktopLocalFileRequest) => deps.localFiles.request(call, request) + try { + const response = await read({ operation: 'manifest', toolCallId: call.toolCallId }) + if (!response.ok) throw new Error(response.error) + if (response.data.kind !== 'manifest') throw new Error('Unexpected file manifest response.') + const manifest = response.data + assertImportableManifest(manifest) + for (const entry of manifest.entries) { + signal.throwIfAborted() + const target = { call, sourceName: manifest.name, relativePath: entry.relativePath } + if (entry.kind === 'directory') { + const folder = await importEntry({ ...target, kind: 'directory' }, signal) + folders.push({ id: folder.id, relativePath: entry.relativePath }) + continue + } + const parts = await readImportEntry(call.toolCallId, entry, read, signal) + const file = await importEntry( + { ...target, kind: 'file', content: new Blob(parts) }, + signal + ) + files.push({ id: file.id, name: file.name, relativePath: entry.relativePath }) + } + return localFileImportCompletion({ + success: true, + workspaceId: manifest.targetWorkspaceId, + files, + folders, + }) + } catch (error) { + return localFileImportCompletion( + localFileImportFailure( + { workspaceId: targetWorkspaceId, files, folders }, + getErrorMessage(error), + // A rate limit that outlasted every retry refused the entry outright: nothing of it landed. + { outcomeKnown: error instanceof DeviceRequestError && error.status === 429 } + ) + ) + } + } + return { async run(call, signal) { try { @@ -231,8 +336,14 @@ export function createDesktopToolRunner(deps: DesktopToolRunnerDeps): DesktopToo if (call.toolName === 'terminal') return await runTerminal(call, signal) if (!deps.accountDataAvailable()) return localAccessUnavailable() if (call.toolName === 'read_local_file') { - return localFileReadCompletion(await deps.localFiles.read(call)) + return localFileReadCompletion( + await deps.localFiles.request(call, { + operation: 'read', + toolCallId: call.toolCallId, + }) + ) } + if (call.toolName === 'import_local_files') return await runImport(call, signal) if (USER_LOCAL_TOOLS.has(call.toolName)) return await runUserLocal(call, signal) return unsupported(call.toolName) } catch (error) { diff --git a/apps/desktop/src/main/desktop-executor/service.ts b/apps/desktop/src/main/desktop-executor/service.ts index c0dd079d36c..b1390bfbc9b 100644 --- a/apps/desktop/src/main/desktop-executor/service.ts +++ b/apps/desktop/src/main/desktop-executor/service.ts @@ -36,6 +36,8 @@ import { createExecutorJournal } from '@/main/desktop-executor/journal' import { DESKTOP_EXECUTOR_PROTOCOL_VERSION, type DesktopExecutorTiming, + type DesktopImportEntryRequest, + type DesktopImportedEntry, } from '@/main/desktop-executor/protocol' const logger = createLogger('DesktopExecutorService') @@ -64,6 +66,11 @@ export interface DesktopExecutorService { /** Re-registers after a sign-in, a session change, or a change to what this device can run. */ refreshRegistration(): void getDevice(): DesktopExecutorDevice | null + /** Stores one entry of a claimed import, as this device's registered session. */ + importEntry( + request: DesktopImportEntryRequest, + signal: AbortSignal + ): Promise /** Sign-out: stops every action, forgets every call, and retires this install id. */ signOut(): Promise } @@ -433,6 +440,10 @@ export function createDesktopExecutorService( getDevice() { return device }, + importEntry(request, signal) { + if (!client) throw new Error('The Sim desktop app is not signed in to Sim.') + return client.importEntry(request, signal) + }, async signOut() { // A registration in flight now answers for a session that is gone; it must not restart. generation += 1 diff --git a/apps/desktop/src/main/index.ts b/apps/desktop/src/main/index.ts index 990fa6724aa..4243e97f722 100644 --- a/apps/desktop/src/main/index.ts +++ b/apps/desktop/src/main/index.ts @@ -603,11 +603,11 @@ function main(): void { }, terminal, localFiles: { - read: (call) => - executeLocalFileRequest( - { operation: 'read', toolCallId: call.toolCallId }, - { toolName: call.toolName, args: call.args } - ), + request: (call, request) => + executeLocalFileRequest(request, { toolName: call.toolName, args: call.args }), + }, + imports: { + importEntry: (request, signal) => desktopExecutor.importEntry(request, signal), }, localFilesystem: { handle: (request) => localFilesystem.handle(request), diff --git a/apps/sim/app/api/desktop/tool/file/route.test.ts b/apps/sim/app/api/desktop/tool/file/route.test.ts index e4d33dfa061..b2c054b55c3 100644 --- a/apps/sim/app/api/desktop/tool/file/route.test.ts +++ b/apps/sim/app/api/desktop/tool/file/route.test.ts @@ -1,6 +1,7 @@ import { authMockFns } from '@sim/testing' +import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock' import { NextRequest } from 'next/server' -import { beforeEach, describe, expect, it, vi } from 'vitest' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' import { OrchestrationError } from '@/lib/core/orchestration/types' const { mockAdmit, mockRead, mockSave } = vi.hoisted(() => ({ @@ -40,6 +41,7 @@ const put = (query: string, body: Uint8Array, headers: Record = headers: { 'content-length': String(body.byteLength), ...headers }, body, }) +const MCP_HOST = 'mcp.sim.test' describe('/api/desktop/tool/file', () => { beforeEach(() => { @@ -51,6 +53,26 @@ describe('/api/desktop/tool/file', () => { }) }) + afterEach(() => { + resetEnvMock() + }) + + it('serves neither a read nor a save on the dedicated MCP host', async () => { + setEnv({ SIM_MCP_URL: `https://${MCP_HOST}/mcp` }) + const onMcpHost = (request: NextRequest) => + new NextRequest(request.url.replace('localhost', MCP_HOST), { + method: request.method, + headers: { ...Object.fromEntries(request.headers), host: MCP_HOST }, + body: request.body, + duplex: 'half', + }) + + expect((await POST(onMcpHost(post({ toolCallId: 'call-1', index: 0 })), {})).status).toBe(404) + const save = onMcpHost(put('toolCallId=call-2&name=report.csv', new Uint8Array([1, 2, 3]))) + expect((await PUT(save)).status).toBe(404) + expect(save.bodyUsed).toBe(false) + }) + it('authenticates before parsing and conceals an unknown transfer', async () => { mockGetSession.mockResolvedValueOnce(null) expect((await POST(post({ toolCallId: 'call-1', index: 0 }), {})).status).toBe(401) @@ -87,4 +109,15 @@ describe('/api/desktop/tool/file', () => { expect(request.bodyUsed).toBe(false) expect(mockSave).not.toHaveBeenCalled() }) + + it('saves nothing from a download that does not match its declared length', async () => { + const res = await PUT( + put('toolCallId=call-2&name=report.csv', new Uint8Array([1, 2, 3]), { + 'content-length': '10', + }) + ) + + expect(res.status).toBe(400) + expect(mockSave).not.toHaveBeenCalled() + }) }) diff --git a/apps/sim/app/api/desktop/tool/file/route.ts b/apps/sim/app/api/desktop/tool/file/route.ts index b21cdeea748..9829a64e681 100644 --- a/apps/sim/app/api/desktop/tool/file/route.ts +++ b/apps/sim/app/api/desktop/tool/file/route.ts @@ -4,6 +4,7 @@ import { readBrowserUploadFileContract, saveBrowserDownloadContract, } from '@/lib/api/contracts/desktop-browser-files' +import { isOffAppHost } from '@/lib/api/mcp/host-routing' import { parseRequest } from '@/lib/api/server' import { defineInternalBinaryRoute, @@ -31,7 +32,7 @@ export const dynamic = 'force-dynamic' * The desktop main process fetches one file a claimed `browser_upload_file` call attaches to a * page. The call's persisted arguments name the file; the request only picks which of them. */ -export const POST = defineInternalBinaryRoute({ +const readUploadFile = defineInternalBinaryRoute({ contract: readBrowserUploadFileContract, auth: internalSessionAuth, operation: readBrowserUploadFile.operation, @@ -49,6 +50,14 @@ export const POST = defineInternalBinaryRoute({ }), }) +/** Off the proxy (see `proxy.ts`), so the dedicated MCP host is refused here, as the proxy would. */ +function notFoundOffAppHost(): NextResponse { + return NextResponse.json(withRequestId({ error: 'Not found' }), { status: 404 }) +} + +export const POST: typeof readUploadFile = (request, context) => + isOffAppHost(request) ? Promise.resolve(notFoundOffAppHost()) : readUploadFile(request, context) + /** * PUT /api/desktop/tool/file?toolCallId=…&name=… * @@ -58,6 +67,7 @@ export const POST = defineInternalBinaryRoute({ * re-validates and claims that call. */ export const PUT = withRouteHandler(async (request: NextRequest) => { + if (isOffAppHost(request)) return notFoundOffAppHost() let principal try { principal = await internalSessionAuth.authenticate() @@ -67,8 +77,13 @@ export const PUT = withRouteHandler(async (request: NextRequest) => { } throw error } - const declaredLength = Number(request.headers.get('content-length')) - if (!Number.isFinite(declaredLength) || declaredLength > BROWSER_FILE_TRANSFER_MAX_BYTES) { + // Shells already in use may send a download without a declared length; that stays accepted. + const lengthHeader = request.headers.get('content-length') + const declaredLength = lengthHeader === null ? null : Number(lengthHeader) + if ( + declaredLength !== null && + (!Number.isFinite(declaredLength) || declaredLength > BROWSER_FILE_TRANSFER_MAX_BYTES) + ) { return NextResponse.json(withRequestId({ error: 'Download is too large to save' }), { status: 413, }) @@ -98,6 +113,13 @@ export const PUT = withRouteHandler(async (request: NextRequest) => { } throw error } + // Anything between the device and here that cut the body short must not become a saved file. + if (declaredLength !== null && content.length !== declaredLength) { + return NextResponse.json( + withRequestId({ error: 'The download did not arrive whole; nothing was saved' }), + { status: 400 } + ) + } try { const { file } = await saveBrowserDownload.execute({ diff --git a/apps/sim/app/api/desktop/tool/import/route.test.ts b/apps/sim/app/api/desktop/tool/import/route.test.ts new file mode 100644 index 00000000000..83229c8d79d --- /dev/null +++ b/apps/sim/app/api/desktop/tool/import/route.test.ts @@ -0,0 +1,162 @@ +import { authMockFns } from '@sim/testing' +import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock' +import { NextRequest } from 'next/server' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors' + +const { admitted, stored, admitError } = vi.hoisted(() => ({ + /** Entries the route admitted, in order. */ + admitted: [] as Array>, + /** Entries the route stored, with how many bytes each file carried. */ + stored: [] as Array<{ executionToken: unknown; bytes: number | null }>, + admitError: { next: null as Error | null }, +})) + +vi.mock('@/lib/desktop/application/import', () => ({ + admitDesktopImportEntry: async (_principal: unknown, input: Record) => { + const error = admitError.next + admitError.next = null + if (error) throw error + admitted.push(input) + }, + importDesktopEntry: { + execute: async ({ input }: { input: { executionToken: unknown; content?: Buffer } }) => { + stored.push({ + executionToken: input.executionToken, + bytes: input.content ? input.content.length : null, + }) + return { id: 'file-1', name: 'notes.txt' } + }, + }, +})) + +import { PUT } from '@/app/api/desktop/tool/import/route' + +const DEVICE = '00000000-0000-4000-8000-000000000001' +const QUERY = `deviceId=${DEVICE}&toolCallId=call-1&kind=file&sourceName=Reports&relativePath=notes.txt` + +/** A request whose body counts how much of it the route read. */ +function put(options: { + body?: Uint8Array + length?: string | null + token?: string | null + host?: string + query?: string +}) { + const body = options.body ?? new TextEncoder().encode('hello') + let pulled = 0 + const stream = new ReadableStream( + { + pull(controller) { + pulled += 1 + controller.enqueue(body) + controller.close() + }, + }, + // Pulled only when the route reads it, never ahead of time. + { highWaterMark: 0 } + ) + const headers: Record = {} + const length = options.length === undefined ? String(body.byteLength) : options.length + if (length !== null) headers['content-length'] = length + if (options.token !== null) headers['x-sim-execution-token'] = options.token ?? 'token-1' + if (options.host) headers.host = options.host + const request = new NextRequest( + `http://${options.host ?? 'localhost'}/api/desktop/tool/import?${options.query ?? QUERY}`, + { method: 'PUT', headers, body: stream, duplex: 'half' } + ) + return { request, pulled: () => pulled } +} + +describe('PUT /api/desktop/tool/import', () => { + beforeEach(() => { + authMockFns.mockGetSession.mockResolvedValue({ + user: { id: 'u1' }, + session: { id: 's1' }, + }) + admitted.length = 0 + stored.length = 0 + admitError.next = null + }) + + afterEach(() => { + resetEnvMock() + }) + + it('serves nothing on the dedicated MCP host, reading no body', async () => { + setEnv({ SIM_MCP_URL: 'https://mcp.sim.test/mcp' }) + const { request, pulled } = put({ host: 'mcp.sim.test' }) + + expect((await PUT(request, {})).status).toBe(404) + expect(pulled()).toBe(0) + expect(stored).toEqual([]) + }) + + it('stores a whole file under the claim the token header names', async () => { + const { request } = put({}) + + const response = await PUT(request, {}) + + expect(response.status).toBe(200) + expect(stored).toEqual([{ executionToken: 'token-1', bytes: 5 }]) + }) + + it('refuses a body that does not match its declared length, storing nothing', async () => { + const { request } = put({ length: '10' }) + + const response = await PUT(request, {}) + + expect(response.status).toBe(400) + expect(stored).toEqual([]) + }) + + it('refuses a file that does not declare its length', async () => { + const { request } = put({ length: null }) + + expect((await PUT(request, {})).status).toBe(411) + expect(stored).toEqual([]) + }) + + it('refuses a file over the import limit before admitting it', async () => { + const { request, pulled } = put({ length: String(65 * 1024 * 1024) }) + + expect((await PUT(request, {})).status).toBe(413) + expect(admitted).toEqual([]) + expect(pulled()).toBe(0) + }) + + it('reads no body for an entry it does not admit', async () => { + admitError.next = new DesktopDeviceUnrecognizedError() + const { request, pulled } = put({}) + + const response = await PUT(request, {}) + + expect(response.status).toBe(401) + expect(pulled()).toBe(0) + expect(stored).toEqual([]) + }) + + it('answers an unauthenticated request with 401', async () => { + authMockFns.mockGetSession.mockResolvedValueOnce(null) + const { request } = put({}) + + expect((await PUT(request, {})).status).toBe(401) + expect(admitted).toEqual([]) + }) + + it('refuses a request without the execution token header', async () => { + const { request } = put({ token: null }) + + expect((await PUT(request, {})).status).toBe(400) + expect(admitted).toEqual([]) + }) + + it('refuses a name a workspace cannot store before admitting it', async () => { + const { request } = put({ + query: `deviceId=${DEVICE}&toolCallId=call-1&kind=file&sourceName=Reports&relativePath=${encodeURIComponent('.. /notes.txt')}`, + }) + + expect((await PUT(request, {})).status).toBe(400) + expect(admitted).toEqual([]) + }) +}) diff --git a/apps/sim/app/api/desktop/tool/import/route.ts b/apps/sim/app/api/desktop/tool/import/route.ts new file mode 100644 index 00000000000..5db6dd8aec9 --- /dev/null +++ b/apps/sim/app/api/desktop/tool/import/route.ts @@ -0,0 +1,112 @@ +import { DESKTOP_IMPORT_TOKEN_HEADER, MAX_DESKTOP_IMPORT_FILE_BYTES } from '@sim/desktop-bridge' +import { type NextRequest, NextResponse } from 'next/server' +import { importDesktopEntryContract } from '@/lib/api/contracts/desktop-executor' +import { isOffAppHost } from '@/lib/api/mcp/host-routing' +import { parseRequest } from '@/lib/api/server' +import { desktopExecutorRateLimit } from '@/lib/api/server/routes/desktop-executor' +import { + InternalUnauthenticatedError, + internalSessionAuth, +} from '@/lib/api/server/routes/internal-json-route' +import { withRequestId } from '@/lib/api/server/routes/request-id' +import { PayloadSizeLimitError, readStreamToBufferWithLimit } from '@/lib/core/utils/stream-limits' +import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { admitDesktopImportEntry, importDesktopEntry } from '@/lib/desktop/application/import' +import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors' +import { internalFileErrorPolicies } from '@/lib/workspace-files/api' + +export const dynamic = 'force-dynamic' + +const TOO_LARGE = `Desktop imports support files up to ${MAX_DESKTOP_IMPORT_FILE_BYTES / 1024 / 1024} MB` + +/** + * PUT /api/desktop/tool/import?… + * + * Raw `withRouteHandler`: a file's bytes are the request body, so the entry is admitted by its + * declared length and its claimed import before the body is read under a byte ceiling, then + * handed to the use case that re-validates the claim and stores it. + */ +export const PUT = withRouteHandler(async (request: NextRequest) => { + if (isOffAppHost(request)) + return NextResponse.json(withRequestId({ error: 'Not found' }), { status: 404 }) + let principal + try { + principal = await internalSessionAuth.authenticate() + } catch (error) { + if (error instanceof InternalUnauthenticatedError) { + return NextResponse.json(withRequestId({ error: error.message }), { status: 401 }) + } + throw error + } + const limited = await desktopExecutorRateLimit.enforce(request, principal) + if (limited) return limited + const parsed = await parseRequest(importDesktopEntryContract, request, {}) + if (!parsed.success) return parsed.response + const query = { + ...parsed.data.query, + executionToken: parsed.data.headers[DESKTOP_IMPORT_TOKEN_HEADER], + } + const lengthHeader = request.headers.get('content-length') + const declaredLength = Number(lengthHeader ?? 0) + if (!Number.isFinite(declaredLength) || declaredLength > MAX_DESKTOP_IMPORT_FILE_BYTES) { + return NextResponse.json(withRequestId({ error: TOO_LARGE }), { status: 413 }) + } + if (query.kind === 'file' && lengthHeader === null) { + return NextResponse.json(withRequestId({ error: 'A file import must declare its length' }), { + status: 411, + }) + } + + try { + await admitDesktopImportEntry(principal, query) + } catch (error) { + return importErrorResponse(error) + } + + let content: Buffer | undefined + if (query.kind === 'file') { + try { + content = await readStreamToBufferWithLimit(request.body, { + maxBytes: MAX_DESKTOP_IMPORT_FILE_BYTES, + label: 'Desktop import', + signal: request.signal, + }) + } catch (error) { + if (error instanceof PayloadSizeLimitError) { + return NextResponse.json(withRequestId({ error: TOO_LARGE }), { status: 413 }) + } + throw error + } + // Anything between the device and here that cut the body short must not become a stored file. + if (content.length !== declaredLength) { + return NextResponse.json( + withRequestId({ error: 'The file did not arrive whole; nothing was stored' }), + { status: 400 } + ) + } + } + + try { + const entry = await importDesktopEntry.execute({ + principal, + input: { ...query, ...(content ? { content } : {}) }, + request, + }) + return NextResponse.json({ id: entry.id, name: entry.name }) + } catch (error) { + return importErrorResponse(error) + } +}) + +/** A device Sim no longer recognizes registers again; everything else hides what it guards. */ +function importErrorResponse(error: unknown): NextResponse { + if (error instanceof DesktopDeviceUnrecognizedError) { + return NextResponse.json(withRequestId({ error: error.message }), { status: 401 }) + } + const response = internalFileErrorPolicies.concealContentAuthorization.project(error) + if (!response) throw error + return NextResponse.json(withRequestId(response.body), { + status: response.status, + headers: response.headers, + }) +} diff --git a/apps/sim/lib/api/contracts/desktop-executor.ts b/apps/sim/lib/api/contracts/desktop-executor.ts index 626bf2486b2..2cfe74f9d84 100644 --- a/apps/sim/lib/api/contracts/desktop-executor.ts +++ b/apps/sim/lib/api/contracts/desktop-executor.ts @@ -1,3 +1,4 @@ +import { DESKTOP_IMPORT_TOKEN_HEADER, isStorableImportName } from '@sim/desktop-bridge' import { z } from 'zod' import { desktopToolCallIdSchema } from '@/lib/api/contracts/desktop-tool-authorization' import { workspaceIdSchema } from '@/lib/api/contracts/primitives' @@ -211,3 +212,52 @@ export const listDesktopActivityContract = defineRouteContract({ response: { mode: 'json', schema: desktopActivityResponseSchema }, error: z.object({ error: z.string() }), }) + +/** Names a workspace cannot hold, though a macOS or Linux file name can be any of them. */ +const UNSTORABLE_NAME = + 'Sim cannot store a file or folder whose name is blank, "." or "..", or contains a backslash' + +/** One relative path inside an import source, as the device's manifest lists it. */ +const desktopImportRelativePathSchema = z + .string() + .max(4096, 'Relative path is too long') + .refine( + (path) => + path === '' || + path.split('/').every((segment) => segment !== '' && segment !== '.' && segment !== '..'), + 'Relative path must stay inside the import source' + ) + .refine((path) => path === '' || path.split('/').every(isStorableImportName), UNSTORABLE_NAME) + +const importDesktopEntryQuerySchema = z.object({ + deviceId: desktopDeviceIdSchema, + toolCallId: desktopToolCallIdSchema, + kind: z.enum(['file', 'directory']), + /** The import source's own name: the folder a directory import lands in, or the file. */ + sourceName: z + .string() + .trim() + .min(1, 'Source name is required') + .max(255) + .refine(isStorableImportName, UNSTORABLE_NAME), + relativePath: desktopImportRelativePathSchema, +}) + +const importDesktopEntryResponseSchema = z.object({ + id: z.string().min(1), + name: z.string().min(1), +}) + +/** + * Stores one entry of a claimed `import_local_files` call in the call's target workspace and + * folder, for the device's background executor. A file's bytes are the raw request body; a + * directory has none. Directories that already exist are reused, so a tree merges into them. + */ +export const importDesktopEntryContract = defineRouteContract({ + method: 'PUT', + path: '/api/desktop/tool/import', + query: importDesktopEntryQuerySchema, + headers: z.object({ [DESKTOP_IMPORT_TOKEN_HEADER]: z.string().min(1).max(128) }), + response: { mode: 'json', schema: importDesktopEntryResponseSchema }, + error: z.object({ error: z.string() }), +}) diff --git a/apps/sim/lib/api/mcp/host-routing.test.ts b/apps/sim/lib/api/mcp/host-routing.test.ts index 4f82a5296d2..df8d0faa04d 100644 --- a/apps/sim/lib/api/mcp/host-routing.test.ts +++ b/apps/sim/lib/api/mcp/host-routing.test.ts @@ -1,7 +1,7 @@ import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock' import { resetUrlsMock, urlsMockFns } from '@sim/testing/mocks/urls.mock' import { afterAll, beforeEach, describe, expect, it } from 'vitest' -import { resolveSimMcpHostPath } from '@/lib/api/mcp/host-routing' +import { isOffAppHost, resolveSimMcpHostPath } from '@/lib/api/mcp/host-routing' import { getSimMcpUrl } from '@/lib/api/mcp/urls' urlsMockFns.mockGetBaseUrl.mockReturnValue('https://sim.ai') @@ -73,3 +73,28 @@ describe('Sim MCP host routing', () => { }) }) }) + +describe('routes the proxy does not run on', () => { + const request = (host: string, path: string) => ({ + headers: new Headers({ host }), + url: `https://${host}${path}`, + }) + + afterAll(() => { + setEnv({ SIM_MCP_URL: undefined }) + }) + + it('refuses the desktop upload routes on a dedicated MCP host, and serves them on the app host', () => { + setEnv({ SIM_MCP_URL: 'https://mcp.sim.ai/mcp/' }) + + expect(isOffAppHost(request('mcp.sim.ai', '/api/desktop/tool/import'))).toBe(true) + expect(isOffAppHost(request('mcp.sim.ai', '/api/desktop/tool/file'))).toBe(true) + expect(isOffAppHost(request('sim.ai', '/api/desktop/tool/import'))).toBe(false) + }) + + it('serves them everywhere while the MCP server shares the app host', () => { + setEnv({ SIM_MCP_URL: undefined }) + + expect(isOffAppHost(request('sim.ai', '/api/desktop/tool/import'))).toBe(false) + }) +}) diff --git a/apps/sim/lib/api/mcp/host-routing.ts b/apps/sim/lib/api/mcp/host-routing.ts index 0a0311900c4..06579a60269 100644 --- a/apps/sim/lib/api/mcp/host-routing.ts +++ b/apps/sim/lib/api/mcp/host-routing.ts @@ -47,3 +47,15 @@ export function resolveSimMcpHostPath( if (!dedicated) return null return pathname === AUTHORIZATION_SERVER_METADATA ? pathname : 'not_found' } + +/** + * For a route the proxy does not run on: whether the request came in through the dedicated MCP + * host, which serves nothing but the MCP server. Such a route answers it as not found itself, as + * the proxy would have. + */ +export function isOffAppHost(request: { headers: Headers; url: string }): boolean { + return ( + resolveSimMcpHostPath(request.headers.get('host'), new URL(request.url).pathname) === + 'not_found' + ) +} diff --git a/apps/sim/lib/desktop/application/import.integration.ts b/apps/sim/lib/desktop/application/import.integration.ts new file mode 100644 index 00000000000..b726197731d --- /dev/null +++ b/apps/sim/lib/desktop/application/import.integration.ts @@ -0,0 +1,275 @@ +/** + * A desktop background executor importing a claimed `import_local_files` call's files into the + * workspace, against real PostgreSQL and file storage on disk: where entries land, how a tree + * merges into folders already there, and every way a request that does not own the import is + * refused. + */ +import { mkdtempSync } from 'node:fs' +import { rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import path from 'node:path' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' + +const fixtureStorage = vi.hoisted(() => ({ root: '' })) +vi.mock('@/lib/uploads/core/setup.server', () => ({ + get UPLOAD_DIR_SERVER() { + return fixtureStorage.root + }, +})) + +import type { SessionPrincipal } from '@sim/auth/principal' +import { db } from '@sim/db' +import { + auditLog, + copilotAsyncToolCalls, + copilotChats, + copilotRuns, + desktopDevices, + permissions, + session, + user, + userStats, + workspace, + workspaceFiles, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { and, eq, inArray, isNull } from 'drizzle-orm' +import { importDesktopEntry } from '@/lib/desktop/application/import' +import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors' +import { + createWorkspaceFileFolder, + listWorkspaceFileFolders, +} from '@/lib/uploads/contexts/workspace/workspace-file-folder-manager' + +describe('desktop imports', () => { + const userIds: string[] = [] + + beforeAll(() => { + fixtureStorage.root = mkdtempSync(path.join(tmpdir(), 'sim-desktop-import-')) + }) + + afterAll(async () => { + if (userIds.length) { + // Audit rows outlive their actor (the foreign key sets null), so they go first. + await db.delete(auditLog).where(inArray(auditLog.actorId, userIds)) + await db.delete(workspace).where(inArray(workspace.ownerId, userIds)) + await db.delete(user).where(inArray(user.id, userIds)) + } + await rm(fixtureStorage.root, { recursive: true, force: true }) + }) + + /** A signed-in desktop whose turn claimed an import into a folder of its workspace. */ + async function claimedImport() { + const userId = generateId() + const workspaceId = generateId() + const sessionId = generateId() + const deviceId = generateId() + const now = new Date() + userIds.push(userId) + await db.insert(user).values({ + id: userId, + name: 'Desktop import fixture', + email: `${userId}@desktop-import.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(userStats).values({ id: generateId(), userId }) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Desktop import fixture', + ownerId: userId, + billedAccountUserId: userId, + }) + await db.insert(permissions).values({ + id: generateId(), + userId, + entityType: 'workspace', + entityId: workspaceId, + permissionType: 'admin', + }) + await db.insert(session).values({ + id: sessionId, + userId, + token: generateId(), + expiresAt: new Date(now.getTime() + 3_600_000), + createdAt: now, + updatedAt: now, + }) + await db.insert(desktopDevices).values({ + id: deviceId, + userId, + sessionId, + name: 'Studio Mac', + appVersion: '0.9.0', + platform: 'darwin-arm64', + capabilities: { executor: 1, browser: true, terminal: true, localFiles: true }, + }) + const target = await createWorkspaceFileFolder({ + workspaceId, + userId, + name: 'Imports', + }) + const chatId = generateId() + const runId = generateId() + await db.insert(copilotChats).values({ id: chatId, userId, workspaceId, type: 'mothership' }) + await db.insert(copilotRuns).values({ + id: runId, + executionId: generateId(), + chatId, + userId, + workspaceId, + streamId: generateId(), + status: 'paused_waiting_for_tool', + desktopDeviceId: deviceId, + }) + const toolCallId = generateId() + const executionToken = generateId() + await db.insert(copilotAsyncToolCalls).values({ + runId, + toolCallId, + toolName: 'import_local_files', + args: { path: '~/Reports', targetWorkspaceId: workspaceId, folderId: target.id }, + status: 'running', + claimedBy: 'files', + executionOwnerToken: executionToken, + }) + const principal: SessionPrincipal = { kind: 'session', userId, sessionId } + return { principal, workspaceId, deviceId, toolCallId, executionToken, targetId: target.id } + } + + function entry( + claimed: Awaited>, + kind: 'file' | 'directory', + relativePath: string, + content?: string + ) { + return importDesktopEntry.execute({ + principal: claimed.principal, + input: { + deviceId: claimed.deviceId, + toolCallId: claimed.toolCallId, + executionToken: claimed.executionToken, + kind, + sourceName: 'Reports', + relativePath, + ...(content !== undefined ? { content: Buffer.from(content) } : {}), + }, + }) + } + + async function activeFile(workspaceId: string, fileId: string) { + const [file] = await db + .select() + .from(workspaceFiles) + .where( + and( + eq(workspaceFiles.workspaceId, workspaceId), + eq(workspaceFiles.id, fileId), + isNull(workspaceFiles.deletedAt) + ) + ) + return file + } + + it("lands a directory import's tree under the call's target folder", async () => { + const claimed = await claimedImport() + + const root = await entry(claimed, 'directory', '') + const q3 = await entry(claimed, 'directory', 'q3') + const report = await entry(claimed, 'file', 'q3/summary.txt', 'quarterly numbers') + + const folders = await listWorkspaceFileFolders(claimed.workspaceId) + const byId = new Map(folders.map((folder) => [folder.id, folder])) + expect(byId.get(root.id)?.parentId).toBe(claimed.targetId) + expect(byId.get(q3.id)?.parentId).toBe(root.id) + const stored = await activeFile(claimed.workspaceId, report.id) + expect(stored?.folderId).toBe(q3.id) + expect(stored).toBeDefined() + }) + + it('records the folders an import creates, and only those', async () => { + const claimed = await claimedImport() + + const root = await entry(claimed, 'directory', '') + await entry(claimed, 'directory', '') + + const auditedRoot = async () => + ( + await db + .select({ resourceId: auditLog.resourceId, action: auditLog.action }) + .from(auditLog) + .where(eq(auditLog.workspaceId, claimed.workspaceId)) + ).filter((row) => row.resourceId === root.id) + await expect.poll(auditedRoot).toEqual([{ resourceId: root.id, action: 'folder.created' }]) + }) + + it('refuses a file entry with no bytes instead of storing an empty file', async () => { + const claimed = await claimedImport() + + await expect(entry(claimed, 'file', 'empty.txt')).rejects.toMatchObject({ + code: 'validation', + }) + }) + + it('merges into folders that already exist and never overwrites a file', async () => { + const claimed = await claimedImport() + const first = await entry(claimed, 'directory', '') + const again = await entry(claimed, 'directory', '') + const original = await entry(claimed, 'file', 'notes.txt', 'v1') + const second = await entry(claimed, 'file', 'notes.txt', 'v2') + + expect(again.id).toBe(first.id) + expect(second.id).not.toBe(original.id) + expect(second.name).not.toBe(original.name) + }) + + it('refuses an import presented with a token that did not claim it', async () => { + const claimed = await claimedImport() + + await expect( + importDesktopEntry.execute({ + principal: claimed.principal, + input: { + deviceId: claimed.deviceId, + toolCallId: claimed.toolCallId, + executionToken: 'someone-else', + kind: 'file', + sourceName: 'Reports', + relativePath: '', + content: Buffer.from('x'), + }, + }) + ).rejects.toMatchObject({ code: 'not_found' }) + }) + + it('refuses an import that is no longer running', async () => { + const claimed = await claimedImport() + await db + .update(copilotAsyncToolCalls) + .set({ status: 'cancelled' }) + .where(eq(copilotAsyncToolCalls.toolCallId, claimed.toolCallId)) + + await expect(entry(claimed, 'file', '', 'late')).rejects.toMatchObject({ code: 'not_found' }) + }) + + it("refuses another user's import, even with its device id and token", async () => { + const claimed = await claimedImport() + const other = await claimedImport() + + await expect( + importDesktopEntry.execute({ + principal: other.principal, + input: { + deviceId: claimed.deviceId, + toolCallId: claimed.toolCallId, + executionToken: claimed.executionToken, + kind: 'file', + sourceName: 'Reports', + relativePath: '', + content: Buffer.from('x'), + }, + }) + ).rejects.toBeInstanceOf(DesktopDeviceUnrecognizedError) + }) +}) diff --git a/apps/sim/lib/desktop/application/import.ts b/apps/sim/lib/desktop/application/import.ts new file mode 100644 index 00000000000..7c1cdab9cb2 --- /dev/null +++ b/apps/sim/lib/desktop/application/import.ts @@ -0,0 +1,168 @@ +import { AuditAction, AuditResourceType } from '@sim/audit' +import type { Principal } from '@sim/auth/principal' +import { toRecord } from '@sim/utils/object' +import { OrchestrationError } from '@/lib/core/orchestration/types' +import { DesktopDeviceUnrecognizedError } from '@/lib/desktop/executor/errors' +import { getBoundDesktopCall, getBoundDesktopDevice } from '@/lib/desktop/executor/repository' +import { ASYNC_TOOL_STATUS } from '@/lib/mothership/async-runs/lifecycle' +import { getAsyncToolCall } from '@/lib/mothership/async-runs/repository' +import { notifyWorkspaceFilesChanged } from '@/lib/realtime/notify' +import { + ensureWorkspaceFileChildFolder, + loadActiveWorkspaceContext, +} from '@/lib/uploads/contexts/workspace' +import { resolveEffectiveMimeType } from '@/lib/uploads/utils/file-utils' +import { defineAuthorizedWorkspaceFileUseCase } from '@/lib/workspace-files/application/authorized-workspace-file-use-case' +import { + createAuthorizedWorkspaceFile, + projectCreateWorkspaceFileAudit, +} from '@/lib/workspace-files/application/create-workspace-file' +import { fileOperations } from '@/lib/workspace-files/application/operations' + +/** + * One entry of a claimed `import_local_files` call, as the device's background executor sends + * it: the claim it belongs to, and where it lands inside the import source. + */ +export interface ImportDesktopEntryInput { + deviceId: string + toolCallId: string + executionToken: string + kind: 'file' | 'directory' + sourceName: string + relativePath: string + /** A file's bytes; a directory has none. */ + content?: Buffer +} + +const IMPORT_NOT_FOUND = 'Desktop import not found' + +/** + * The claimed import an entry belongs to: a running `import_local_files` call bound to this + * device, presented with the token that claimed it, on the caller's own turn. The target + * workspace and folder come from Sim's record of the call, never from the request; writing to + * that workspace is then authorized as any file creation is. + */ +async function resolveImportContext({ + principal, + input, +}: { + principal: Principal + input: ImportDesktopEntryInput +}) { + if (principal.kind !== 'session') throw new OrchestrationError('not_found', IMPORT_NOT_FOUND) + const device = await getBoundDesktopDevice({ + deviceId: input.deviceId, + userId: principal.userId, + sessionId: principal.sessionId, + }) + if (!device) throw new DesktopDeviceUnrecognizedError() + const [call, row] = await Promise.all([ + getBoundDesktopCall({ deviceId: input.deviceId, userId: principal.userId }, input.toolCallId), + getAsyncToolCall(input.toolCallId), + ]) + if ( + !call || + !row || + call.toolName !== 'import_local_files' || + call.ownerToken !== input.executionToken || + row.status !== ASYNC_TOOL_STATUS.running + ) { + throw new OrchestrationError('not_found', IMPORT_NOT_FOUND) + } + const args = toRecord(call.args) + if (typeof args.targetWorkspaceId !== 'string') { + throw new OrchestrationError('not_found', IMPORT_NOT_FOUND) + } + const workspace = await loadActiveWorkspaceContext(args.targetWorkspaceId) + if (!workspace) throw new OrchestrationError('not_found', 'Workspace not found') + return { + ...workspace, + rootFolderId: typeof args.folderId === 'string' ? args.folderId : null, + importingUserId: principal.userId, + } +} + +const admitDesktopImportEntryUseCase = defineAuthorizedWorkspaceFileUseCase({ + operation: fileOperations.create, + resolveContext: resolveImportContext, + async execute() {}, +}) + +/** + * Admits an entry before its body is read, so a caller without a claimed, running import of its + * own is refused without buffering the file. {@link importDesktopEntry} re-validates it. + */ +export async function admitDesktopImportEntry( + principal: Principal, + input: ImportDesktopEntryInput +): Promise { + await admitDesktopImportEntryUseCase.execute({ principal, input }) +} + +/** + * Stores one entry of a desktop import. Each path segment is a folder under the call's target + * folder: the source's own folder for a directory import, then the entry's parents. Existing + * folders are reused; a file never overwrites one already there. + */ +export const importDesktopEntry = defineAuthorizedWorkspaceFileUseCase({ + operation: fileOperations.create, + resolveContext: resolveImportContext, + async execute({ principal, input, context }) { + const segments = input.relativePath + ? [input.sourceName, ...input.relativePath.split('/')] + : [input.sourceName] + const folderSegments = input.kind === 'directory' ? segments : segments.slice(0, -1) + let folderId = context.rootFolderId + const createdFolders: Array<{ id: string; name: string }> = [] + for (const name of folderSegments) { + const folder = await ensureWorkspaceFileChildFolder({ + workspaceId: context.workspaceId, + userId: context.importingUserId, + parentId: folderId, + name, + }) + if (folder.created) createdFolders.push({ id: folder.id, name: folder.name }) + folderId = folder.id + } + const name = segments[segments.length - 1] ?? input.sourceName + if (input.kind === 'directory') { + if (!folderId) throw new OrchestrationError('validation', 'A directory needs a name') + return { kind: 'directory' as const, id: folderId, name, createdFolders } + } + if (!input.content) throw new OrchestrationError('validation', 'A file import needs its bytes') + const created = await createAuthorizedWorkspaceFile({ + principal, + input: { + workspaceId: context.workspaceId, + name, + contentType: resolveEffectiveMimeType(undefined, name), + folderId, + exactName: false, + }, + content: input.content, + workspace: context, + }) + return { + kind: 'file' as const, + id: created.file.id, + name: created.file.name, + createdFolders, + created, + } + }, + // Folders an import creates are recorded like folders made by hand; the file records itself. + projectAudit: ({ result }) => [ + ...result.createdFolders.map((folder) => ({ + action: AuditAction.FOLDER_CREATED, + resourceType: AuditResourceType.FOLDER, + resourceId: folder.id, + resourceName: folder.name, + description: `Created file folder "${folder.name}"`, + })), + ...(result.kind === 'file' ? [projectCreateWorkspaceFileAudit(result.created)] : []), + ], + async afterSuccess({ context, result }) { + // A new file announces itself; a new folder is announced here, as folder creation does. + if (result.createdFolders.length > 0) await notifyWorkspaceFilesChanged(context.workspaceId) + }, +}) diff --git a/apps/sim/lib/mothership/tools/client/native-files.ts b/apps/sim/lib/mothership/tools/client/native-files.ts index 73ed8e8227e..b1d4eff830b 100644 --- a/apps/sim/lib/mothership/tools/client/native-files.ts +++ b/apps/sim/lib/mothership/tools/client/native-files.ts @@ -3,7 +3,14 @@ import type { DesktopLocalFileRequest, DesktopLocalFileResponse, } from '@sim/desktop-bridge' -import { MAX_DESKTOP_IMPORT_FILE_BYTES } from '@sim/desktop-bridge' +import { + assertImportableManifest, + type DesktopLocalFileImportResult, + localFileImportCompletion, + localFileImportFailure, + localFileReadCompletion, + readImportEntry, +} from '@sim/desktop-bridge/tool-results' import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' import { ApiClientError } from '@/lib/api/client/errors' @@ -34,34 +41,17 @@ async function invoke( return response } -interface ImportedFile { - id: string - name: string - relativePath: string -} -interface ImportedFolder { - id: string - relativePath: string -} - /** The manifest is produced from canonical pending-tool arguments in Electron, never renderer paths. */ export async function importNativeFiles( toolCallId: string, manifest: DesktopLocalFileManifest, signal?: AbortSignal -) { - const files: ImportedFile[] = [] - const folders: ImportedFolder[] = [] +): Promise { + const files: DesktopLocalFileImportResult['files'] = [] + const folders: DesktopLocalFileImportResult['folders'] = [] const parents = new Map([['', manifest.folderId]]) try { - if ( - manifest.entries.some( - (entry) => entry.kind === 'file' && entry.size > MAX_DESKTOP_IMPORT_FILE_BYTES - ) - ) - throw new Error( - 'Desktop imports support files up to 64 MB. Use the file uploader for larger files.' - ) + assertImportableManifest(manifest) for (const entry of manifest.entries) { signal?.throwIfAborted() const segments = entry.relativePath.split('/').filter(Boolean) @@ -96,32 +86,12 @@ export async function importNativeFiles( folders.push({ id: folderId, relativePath: entry.relativePath }) continue } - const parts: Uint8Array[] = [] - let offset = 0 - do { - const response = await invoke( - { - operation: 'chunk', - toolCallId, - relativePath: entry.relativePath, - offset, - revision: entry.revision, - }, - signal - ) - if (!response.ok) throw new Error(response.error) - if (response.data.kind !== 'chunk') throw new Error('Unexpected file chunk response.') - const bytes = new Uint8Array(response.data.bytes) - parts.push(bytes) - offset += bytes.length - if ( - offset > entry.size || - (response.data.eof && offset !== entry.size) || - (!response.data.eof && bytes.length === 0) - ) - throw new Error('The local file changed or its transfer was incomplete.') - if (response.data.eof) break - } while (offset < entry.size) + const parts = await readImportEntry( + toolCallId, + entry, + (request) => invoke(request, signal), + signal + ) const saved = await uploadWorkspaceFileSession({ workspaceId: manifest.targetWorkspaceId, folderId: parentId, @@ -132,16 +102,10 @@ export async function importNativeFiles( } return { success: true, workspaceId: manifest.targetWorkspaceId, files, folders } } catch (error) { - return { - success: false, - workspaceId: manifest.targetWorkspaceId, - files, - folders, - error: getErrorMessage(error), - partial: files.length > 0 || folders.length > 0, - doNotRetry: true, - outcomeUnknown: true, - } + return localFileImportFailure( + { workspaceId: manifest.targetWorkspaceId, files, folders }, + getErrorMessage(error) + ) } } @@ -176,20 +140,16 @@ export async function executeNativeFileTool( if (response.code === 'ALREADY_STARTED') return throw new Error(response.error) } - const result = + if (response.data.kind === 'chunk') throw new Error('Unexpected chunk outside an import.') + const completion = response.data.kind === 'manifest' - ? await importNativeFiles(toolCallId, response.data, signal) - : response.data - if ('kind' in result && result.kind === 'chunk') - throw new Error('Unexpected chunk outside an import.') - const failed = 'success' in result && result.success === false + ? localFileImportCompletion(await importNativeFiles(toolCallId, response.data, signal)) + : localFileReadCompletion(response) await reportClientToolCompletion( toolCallId, - failed ? ASYNC_TOOL_CONFIRMATION_STATUS.error : ASYNC_TOOL_CONFIRMATION_STATUS.success, - failed - ? 'Some files could not be imported; inspect the partial result.' - : 'Local file operation completed.', - result + completion.status, + completion.message, + completion.data ) settled = true } catch (error) { diff --git a/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts b/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts index 62dd595c92b..3ae639bcf1c 100644 --- a/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts +++ b/apps/sim/lib/uploads/contexts/workspace/workspace-file-folder-manager.ts @@ -620,6 +620,35 @@ export async function createWorkspaceFileFolder(params: { return mapFolderWithPath(params.workspaceId, folder) } +/** + * The active folder `name` directly under `parentId` (the workspace root when null), created + * when it does not exist yet, and whether this call created it. Writers that copy a tree in, one + * entry at a time, merge into a folder that is already there instead of making a numbered sibling. + */ +export async function ensureWorkspaceFileChildFolder(params: { + workspaceId: string + userId: string + parentId: string | null + name: string +}): Promise<{ id: string; name: string; created: boolean }> { + const name = normalizeWorkspaceFileItemName(params.name, 'Folder') + const existing = await findRawWorkspaceFileFolderByName(params.workspaceId, name, params.parentId) + if (existing) return { id: existing.id, name: existing.name, created: false } + try { + const created = await createWorkspaceFileFolder({ ...params, name, exactName: true }) + return { id: created.id, name: created.name, created: true } + } catch (error) { + if (!(error instanceof WorkspaceFileFolderConflictError)) throw error + const concurrent = await findRawWorkspaceFileFolderByName( + params.workspaceId, + name, + params.parentId + ) + if (!concurrent) throw error + return { id: concurrent.id, name: concurrent.name, created: false } + } +} + /** * Outcome of {@link ensureWorkspaceFileFolderPath}. `createdFolderIds` lists only the * folders this call actually inserted, outermost-first, so a caller that has to unwind diff --git a/apps/sim/proxy.test.ts b/apps/sim/proxy.test.ts index 9bd5bec5c1e..b36466e25e4 100644 --- a/apps/sim/proxy.test.ts +++ b/apps/sim/proxy.test.ts @@ -1,7 +1,8 @@ import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock' +import { unstable_doesMiddlewareMatch } from 'next/experimental/testing/server' import { NextRequest } from 'next/server' import { afterAll, describe, expect, it } from 'vitest' -import { proxy, resolveApiCorsPolicy } from '@/proxy' +import { config, proxy, resolveApiCorsPolicy } from '@/proxy' setEnv({ NEXT_PUBLIC_APP_URL: 'https://app.sim.test', @@ -148,3 +149,20 @@ describe('proxy on the dedicated MCP host', () => { expect(proxy(mcpRequest('/login', 'GET')).status).toBe(404) }) }) + +describe('proxy matcher', () => { + const matches = (path: string) => + unstable_doesMiddlewareMatch({ config, url: `https://sim.test${path}` }) + + it('leaves the desktop raw upload routes unbuffered, so a whole file reaches the route', () => { + expect(matches('/api/desktop/tool/import')).toBe(false) + expect(matches('/api/desktop/tool/file')).toBe(false) + }) + + it('still runs for every other API route', () => { + expect(matches('/api/desktop/tool/claim')).toBe(true) + expect(matches('/api/desktop/tool/import/other')).toBe(true) + expect(matches('/api/workflows/wf-1/execute')).toBe(true) + expect(matches('/api/files/upload')).toBe(true) + }) +}) diff --git a/apps/sim/proxy.ts b/apps/sim/proxy.ts index d7cea16f614..2879c6804dc 100644 --- a/apps/sim/proxy.ts +++ b/apps/sim/proxy.ts @@ -457,7 +457,9 @@ export const config = { '/login', '/signup', '/invite/:path*', // Match invitation routes - '/api/:path*', // Runtime CORS + // Runtime CORS. The desktop's raw upload routes are left out: running the proxy makes Next + // buffer the body, cutting it off at its 10 MB proxy limit, and those bodies are whole files. + '/api/((?!desktop/tool/(?:import|file)$).*)', // Catch-all for other pages, excluding static assets and public directories '/((?!api/|api$|_next/static|_next/image|ingest|favicon.ico|logo/|landing/|static/|footer/|social/|enterprise/|favicon/|twitter/|robots.txt|sitemap.xml).*)', ], diff --git a/packages/desktop-bridge/src/index.ts b/packages/desktop-bridge/src/index.ts index 1ffed49b294..aeeb6578db1 100644 --- a/packages/desktop-bridge/src/index.ts +++ b/packages/desktop-bridge/src/index.ts @@ -1151,7 +1151,11 @@ export interface SimDesktopApi { /** Reads and selects Terminal.app or iTerm2 color profiles on macOS. */ terminalThemes?: SimDesktopTerminalThemesApi } -export { MAX_DESKTOP_IMPORT_FILE_BYTES } from './local-files' +export { + DESKTOP_IMPORT_TOKEN_HEADER, + isStorableImportName, + MAX_DESKTOP_IMPORT_FILE_BYTES, +} from './local-files' export { applyDesktopTitleBarMode, DESKTOP_TITLE_BAR_ATTRIBUTE, diff --git a/packages/desktop-bridge/src/local-files.ts b/packages/desktop-bridge/src/local-files.ts index c5820a49721..6a76b2b1aa0 100644 --- a/packages/desktop-bridge/src/local-files.ts +++ b/packages/desktop-bridge/src/local-files.ts @@ -1,5 +1,26 @@ export const MAX_DESKTOP_IMPORT_FILE_BYTES = 64 * 1024 * 1024 +/** + * The header a background import presents its claim's execution token in. A header, not the + * query, so the token stays out of load balancer, CDN and trace URLs. + */ +export const DESKTOP_IMPORT_TOKEN_HEADER = 'x-sim-execution-token' + +/** + * Whether a single file or folder name can be stored as a workspace file name: something left + * after trimming, not a dot segment, and no path separator. + */ +export function isStorableImportName(name: string): boolean { + const trimmed = name.trim() + return ( + trimmed !== '' && + trimmed !== '.' && + trimmed !== '..' && + !trimmed.includes('/') && + !trimmed.includes('\\') + ) +} + /** Native desktop file operations use OS permissions and canonical pending chat calls. */ export type DesktopLocalFileRequest = | { operation: 'read' | 'manifest'; toolCallId: string } diff --git a/packages/desktop-bridge/src/tool-results.test.ts b/packages/desktop-bridge/src/tool-results.test.ts index ba45614060f..c294f385df0 100644 --- a/packages/desktop-bridge/src/tool-results.test.ts +++ b/packages/desktop-bridge/src/tool-results.test.ts @@ -1,5 +1,10 @@ import { describe, expect, it } from 'vitest' -import { sanitizeBrowserToolResultForModel } from './tool-results' +import type { DesktopLocalFileEntry, DesktopLocalFileResponse } from './local-files' +import { + assertImportableManifest, + readImportEntry, + sanitizeBrowserToolResultForModel, +} from './tool-results' describe('browser screenshot model projection', () => { it('keeps an image usable when an older desktop omits coordinate metadata', () => { @@ -89,3 +94,64 @@ describe('browser screenshot model projection', () => { expect(sanitizeBrowserToolResultForModel('browser_snapshot', undefined)).toBeUndefined() }) }) + +describe('import file reads', () => { + const entry: DesktopLocalFileEntry = { + relativePath: 'notes.txt', + kind: 'file', + size: 6, + revision: 'rev-1', + } + + function chunks(...parts: Array<{ bytes: string; eof: boolean }>) { + const queue = [...parts] + return async (): Promise => { + const next = queue.shift() + if (!next) throw new Error('read past the end') + return { + ok: true, + data: { kind: 'chunk', bytes: new TextEncoder().encode(next.bytes), eof: next.eof }, + } + } + } + + it('reads a file across chunks', async () => { + const parts = await readImportEntry( + 'call-1', + entry, + chunks({ bytes: 'abc', eof: false }, { bytes: 'def', eof: true }) + ) + expect(new TextDecoder().decode(Buffer.concat(parts))).toBe('abcdef') + }) + + it.each([ + ['ends early', [{ bytes: 'abc', eof: true }]], + ['grows past its listed size', [{ bytes: 'abcdefg', eof: true }]], + ['stalls', [{ bytes: '', eof: false }]], + ['keeps going past its listed size', [{ bytes: 'abcdef', eof: false }]], + ])('refuses a file that %s since the manifest listed it', async (_case, parts) => { + await expect(readImportEntry('call-1', entry, chunks(...parts))).rejects.toThrow( + 'The local file changed or its transfer was incomplete.' + ) + }) +}) + +describe('importable manifests', () => { + it.each([ + ['a backslash', 'q3\\draft.txt'], + ['a blank name', 'q3/ '], + ['a dot segment after trimming', '.. /notes.txt'], + ])('refuses %s before anything is imported', (_case, relativePath) => { + expect(() => + assertImportableManifest({ + kind: 'manifest', + name: 'Reports', + targetWorkspaceId: 'ws-1', + entries: [ + { relativePath: '', kind: 'directory', size: 0, revision: 'r0' }, + { relativePath, kind: 'file', size: 1, revision: 'r1' }, + ], + }) + ).toThrow('Sim cannot store a file or folder named') + }) +}) diff --git a/packages/desktop-bridge/src/tool-results.ts b/packages/desktop-bridge/src/tool-results.ts index ff6f864995b..7abfcbf0ada 100644 --- a/packages/desktop-bridge/src/tool-results.ts +++ b/packages/desktop-bridge/src/tool-results.ts @@ -7,7 +7,14 @@ import type { BrowserToolName } from '@sim/browser-protocol' import type { TerminalOperation, TerminalToolResponse } from '@sim/terminal-protocol' import { isRecordLike, toRecordOrNull } from '@sim/utils/object' -import type { DesktopLocalFileResponse } from './local-files' +import { + type DesktopLocalFileEntry, + type DesktopLocalFileManifest, + type DesktopLocalFileRequest, + type DesktopLocalFileResponse, + isStorableImportName, + MAX_DESKTOP_IMPORT_FILE_BYTES, +} from './local-files' type DesktopToolCompletionStatus = 'success' | 'error' | 'cancelled' @@ -224,6 +231,118 @@ export function localFileReadCompletion(response: DesktopLocalFileResponse): Des } } +/** + * Refuses, before anything lands, an import Sim could not store whole: a file larger than desktop + * imports carry, or a name with a backslash, which Sim's file names cannot hold. + */ +export function assertImportableManifest(manifest: DesktopLocalFileManifest): void { + if ( + manifest.entries.some( + (entry) => entry.kind === 'file' && entry.size > MAX_DESKTOP_IMPORT_FILE_BYTES + ) + ) { + throw new Error( + 'Desktop imports support files up to 64 MB. Use the file uploader for larger files.' + ) + } + const unstorable = [ + manifest.name, + ...manifest.entries.flatMap((entry) => + entry.relativePath === '' ? [] : entry.relativePath.split('/') + ), + ].find((name) => !isStorableImportName(name)) + if (unstorable !== undefined) { + throw new Error( + `Sim cannot store a file or folder named "${unstorable}": a name needs visible characters, and cannot be "." or ".." or contain a backslash. Rename it, or import the rest separately.` + ) + } +} + +/** + * Reads one manifest file in chunks, refusing a file that changed since the manifest listed it: + * a different size, an early end, or a chunk past its length. + */ +export async function readImportEntry( + toolCallId: string, + entry: DesktopLocalFileEntry, + readChunk: (request: DesktopLocalFileRequest) => Promise, + signal?: AbortSignal +): Promise[]> { + const parts: Uint8Array[] = [] + let offset = 0 + do { + signal?.throwIfAborted() + const response = await readChunk({ + operation: 'chunk', + toolCallId, + relativePath: entry.relativePath, + offset, + revision: entry.revision, + }) + if (!response.ok) throw new Error(response.error) + if (response.data.kind !== 'chunk') throw new Error('Unexpected file chunk response.') + const bytes = new Uint8Array(response.data.bytes) + parts.push(bytes) + offset += bytes.length + if ( + offset > entry.size || + (response.data.eof && offset !== entry.size) || + (!response.data.eof && (bytes.length === 0 || offset >= entry.size)) + ) { + throw new Error('The local file changed or its transfer was incomplete.') + } + if (response.data.eof) break + } while (offset < entry.size) + return parts +} + +/** What an `import_local_files` call created, and how far it got when it stopped short. */ +export interface DesktopLocalFileImportResult { + success: boolean + workspaceId: string + files: Array<{ id: string; name: string; relativePath: string }> + folders: Array<{ id: string; relativePath: string }> + error?: string + partial?: boolean + doNotRetry?: true + outcomeUnknown?: true +} + +/** + * An import that stopped part way. Files it already created stay. Usually whether more landed than + * it reports is unknown, so the model inspects the workspace instead of importing again. With + * `outcomeKnown` (Sim refused the next entry outright), the list is exact, and an import where + * nothing landed can simply be asked for again. + */ +export function localFileImportFailure( + partial: Pick, + error: string, + options: { outcomeKnown?: boolean } = {} +): DesktopLocalFileImportResult { + const landed = partial.files.length > 0 || partial.folders.length > 0 + return { + success: false, + ...partial, + error, + partial: landed, + ...(landed || !options.outcomeKnown ? { doNotRetry: true as const } : {}), + ...(options.outcomeKnown ? {} : { outcomeUnknown: true as const }), + } +} + +/** An `import_local_files` call's outcome. */ +export function localFileImportCompletion( + result: DesktopLocalFileImportResult +): DesktopToolCompletion { + return result.success + ? { status: 'success', message: 'Local file operation completed.', data: { ...result } } + : { + status: 'error', + message: 'Some files could not be imported; inspect the partial result.', + data: { ...result }, + } +} + /** A user-local folder read (`read`, `grep`, `glob`) the desktop finished or failed. */ export function localFilesystemToolCompletion( outcome: { ok: true; data: Record } | { ok: false; error: string }