diff --git a/packages/code/README.md b/packages/code/README.md index 636cd1a8..d348fe88 100644 --- a/packages/code/README.md +++ b/packages/code/README.md @@ -729,9 +729,10 @@ native SRT command backend. This operator switch controls availability; LibreChat tool approval hooks remain the user-facing allow/deny boundary for each invocation. -Reads reject absolute paths, traversal, escaping symlinks, non-regular files, -and files larger than 1 MiB. The opened file is checked against its canonical -in-workspace inode before it is read. Text search uses `rg` only to enumerate a +Reads reject absolute paths, traversal, escaping symlinks, and non-regular files. +The opened file is checked against its canonical in-workspace inode before it is +read. Ordinary reads stream bounded line windows even from files larger than 1 MiB; +see the [read contract](WORKSPACE-READS.md). Text search uses `rg` only to enumerate a bounded set of ignored-aware candidates with configuration and symlink following disabled. It then opens and verifies each candidate through the same confined 1 MiB read boundary before matching locally. File listing invokes `rg` without diff --git a/packages/code/WORKSPACE-READS.md b/packages/code/WORKSPACE-READS.md new file mode 100644 index 00000000..f3edbe58 --- /dev/null +++ b/packages/code/WORKSPACE-READS.md @@ -0,0 +1,32 @@ +# Workspace file read contract + +Ordinary `read_file` streams complete LF-delimited line windows through the verified +file descriptor. It defaults to 200 lines, accepts at most 500, and returns at most +1 MiB of UTF-8 content, including inter-line separators, regardless of file size. +CR in CRLF is preserved; the final LF does not add a phantom line. An empty file +has one empty line. A start beyond EOF returns an empty, non-truncated window. + +Byte-limited windows stop before the first line that cannot fit and return +`nextStartLine = endLine + 1`. If the first requested line exceeds 1 MiB, the read +fails with `READ_LIMIT_EXCEEDED`, without a repeating continuation. Scanning and +collection check cancellation and share a 10-second deadline. Far-away starts +can fail with `READ_LIMIT_EXCEEDED`; use an earlier start or `search_text`. +Cancellation and timeout settle the request without waiting for queued I/O. +Node retains an interrupted read's descriptor and chunk until that I/O drains; +file and held-root cleanup are initiated immediately and finish asynchronously. +This releases the request lane, but does not cancel kernel I/O. + +Memory is bounded by 64 KiB chunks and the returned window. Reads stop at the +opened file's initial size, so concurrent growth cannot extend the scan. Truncation +ends the scan at observed EOF; pathname replacement never switches the open +handle. In-place writes can change observed content. Pagination is not a snapshot +across requests and assumes the file stays unchanged. + +UTF-8 and BOM-marked UTF-16LE/BE are decoded incrementally. One leading BOM is +removed; malformed sequences use replacement characters, as before. Binary bytes +are decoded as text, including NUL, not classified or rejected. Neither a successful +window nor continuation validates unread content. + +Request/result shapes, protocol version, and capabilities are unchanged. Search, +preview, edit, and write limits remain unchanged. Digest-checked repository +instruction snapshots retain their separate byte-bounded, newline-preserving path. diff --git a/packages/code/src/root-access.ts b/packages/code/src/root-access.ts index 5bdddca9..2f049bcb 100644 --- a/packages/code/src/root-access.ts +++ b/packages/code/src/root-access.ts @@ -378,6 +378,12 @@ export class WorkspaceRootAccess { } const context = new AsyncLocalStorage(); +const deferredCleanup = new WeakSet(); +/** Interrupted I/O must not hold a request open while Node drains its handles. */ +export function deferWorkspaceRootCleanup(): void { + const access = context.getStore(); + if (access) deferredCleanup.add(access); +} /** Whether filesystem adapters in this call are anchored to a held root descriptor. */ export const holdsWorkspaceRoot = (): boolean => context.getStore() != null; export async function withWorkspaceRoot( @@ -390,7 +396,9 @@ export async function withWorkspaceRoot( try { return await context.run(access, action); } finally { - await access.close(); + const closing = access.close(); + if (deferredCleanup.delete(access)) void closing.catch(() => {}); + else await closing; } } diff --git a/packages/code/src/workspace-read-io.test.ts b/packages/code/src/workspace-read-io.test.ts new file mode 100644 index 00000000..2fdd2b90 --- /dev/null +++ b/packages/code/src/workspace-read-io.test.ts @@ -0,0 +1,147 @@ +import assert from 'node:assert/strict'; +import { execFile } from 'node:child_process'; +import * as fs from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import test from 'node:test'; +import { promisify } from 'node:util'; + +const execFileAsync = promisify(execFile); + +// Isolate libuv saturation from the test runner and sibling tests. +const queuedRead = ` +import assert from 'node:assert/strict'; +import * as fs from 'node:fs'; +import * as fsp from 'node:fs/promises'; +import { execFileSync } from 'node:child_process'; +import { getEventListeners } from 'node:events'; +import { join } from 'node:path'; +import { mock } from 'node:test'; +import { LocalWorkspaceTools } from ${JSON.stringify( + new URL('./workspace.js', import.meta.url).href +)}; +import { captureWorkspaceRootIdentity } from ${JSON.stringify( + new URL('./root-identity.js', import.meta.url).href +)}; +import { WorkspaceRootAccess } from ${JSON.stringify( + new URL('./root-access.js', import.meta.url).href +)}; +const [root, cause, held] = process.argv.slice(1); +const tools = await LocalWorkspaceTools.create({ workspaces: [{ + id: 'primary', root, + ...(held === 'true' ? { identity: await captureWorkspaceRootIdentity(root) } : {}), +}] }); +const fifo = join(root, 'pool'); +execFileSync('mkfifo', [fifo]); +const pipe = fs.openSync(fifo, fs.constants.O_RDWR); +const probe = await fsp.open(join(root, 'file'), 'r'); +const prototype = Object.getPrototypeOf(probe); +await probe.close(); +const originalRead = prototype.read; +const controller = new AbortController(); +let entered, blocker, pendingIo, physicalFd, readSettled = false; +const started = new Promise(resolve => { entered = resolve; }); +const closing = []; +const descriptors = []; +let elapsed = 0; +const now = performance.now(); +mock.method(performance, 'now', () => now + elapsed); +function trackClose(handle) { + const originalClose = handle.close; + mock.method(handle, 'close', function (...args) { + descriptors.push(this.fd); + const promise = originalClose.apply(this, args); + closing.push(promise); + return promise; + }); +} +const originalOpen = WorkspaceRootAccess.open; +mock.method(WorkspaceRootAccess, 'open', async (...args) => { + const access = await originalOpen(...args); + trackClose(access.handle); + return access; +}); +mock.method(prototype, 'read', function (...args) { + physicalFd = this.fd; + trackClose(this); + // This actual pipe read occupies the sole libuv worker until explicitly drained. + blocker = new Promise((resolve, reject) => fs.read(pipe, Buffer.alloc(1), 0, 1, null, + error => error ? reject(error) : resolve())); + // Invoke Node's original FileHandle.read, including its active-I/O references. + pendingIo = originalRead.apply(this, args).finally(() => { readSettled = true; }); + if (cause === 'deadline') elapsed = 10_000; + if (cause === 'early-deadline') elapsed = 9_999.75; + entered(); + return pendingIo; +}); +let released = false; +function drain() { + if (!released) { released = true; fs.writeSync(pipe, Buffer.from('x')); } +} +const wallStart = Date.now(); +const pending = tools.execute({ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'file' }, controller.signal); +const rejected = assert.rejects(pending, error => error.code === (cause === 'abort' ? 'EXECUTION_ABORTED' : 'READ_LIMIT_EXCEEDED')); +await started; +if (cause === 'abort') controller.abort(); +// Always drain, even if a regression prevents bounded settlement. +const fallback = setTimeout(drain, 2_000); +try { + await rejected; + assert.equal(released, false, 'request settled before pool was drained'); + if (cause === 'early-deadline') assert.ok(performance.now() < now + 10_000, 'timer settled before the monotonic deadline'); + assert.equal(readSettled, false, 'real descriptor read remained pending'); + assert.ok(fs.fstatSync(physicalFd).isFile(), 'Node still owns the physical descriptor'); + assert.equal(closing.length, held === 'true' ? 2 : 1, 'file and held-root closes initiated'); + assert.deepEqual(getEventListeners(controller.signal, 'abort'), []); + console.log(JSON.stringify({ cause, held: held === 'true', settledBeforeDrain: true, readSettled, elapsedMs: Date.now() - wallStart })); +} finally { + clearTimeout(fallback); + drain(); + await blocker; + await pendingIo; + await Promise.all(closing); + for (const fd of descriptors) assert.throws(() => fs.fstatSync(fd), { code: 'EBADF' }); + mock.restoreAll(); + fs.closeSync(pipe); +} +`; + +for (const cause of ['abort', 'deadline', 'early-deadline'] as const) { + for (const held of [false, true]) { + test(`${cause} settles before real queued I/O drains (${ + held ? 'held' : 'legacy' + } root)`, async t => { + if (!['linux', 'darwin'].includes(process.platform)) + return t.skip('requires POSIX FIFO and descriptor access'); + const root = await fs.realpath( + await fs.mkdtemp(join(tmpdir(), 'workspace-read-io-')) + ); + t.after(() => fs.rm(root, { recursive: true, force: true })); + await fs.writeFile(join(root, 'file'), 'line\n'.repeat(300_000)); + const { stdout } = await execFileAsync( + process.execPath, + [ + '--input-type=module', + '--eval', + queuedRead, + root, + cause, + String(held), + ], + { + // Node 20 can bypass the pool for regular files via io_uring. + env: { + ...process.env, + UV_THREADPOOL_SIZE: '1', + UV_USE_IO_URING: '0', + }, + timeout: 10_000, + } + ); + const observed = JSON.parse(stdout); + assert.equal(observed.settledBeforeDrain, true); + assert.equal(observed.readSettled, false); + t.diagnostic(stdout.trim()); + }); + } +} diff --git a/packages/code/src/workspace-read.test.ts b/packages/code/src/workspace-read.test.ts new file mode 100644 index 00000000..1c6a7f96 --- /dev/null +++ b/packages/code/src/workspace-read.test.ts @@ -0,0 +1,488 @@ +import assert from 'node:assert/strict'; +import * as fs from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import test from 'node:test'; +import type { TestContext } from 'node:test'; +import type { FileHandle } from 'node:fs/promises'; +import { LocalWorkspaceTools, WorkspaceToolError } from './workspace.js'; +import { WorkspaceRootAccess } from './root-access.js'; +import { captureWorkspaceRootIdentity } from './root-identity.js'; +import { readRepositoryInstructions } from './instructions.js'; +import { + BRIDGE_WORKSPACE_READ_MAX_BYTES as MAX_BYTES, + isWorkspaceToolResult, +} from './protocol.js'; +import type { WorkspaceReadFileRequest } from './protocol.js'; + +async function fixture(t: TestContext, held = false) { + const parent = await fs.realpath( + await fs.mkdtemp(join(tmpdir(), 'workspace-read-')) + ); + t.after(() => fs.rm(parent, { recursive: true, force: true })); + const root = join(parent, 'root'); + await fs.mkdir(root); + const tools = await LocalWorkspaceTools.create({ + workspaces: [ + { + id: 'primary', + root, + writable: true, + ...(held + ? { identity: await captureWorkspaceRootIdentity(root) } + : {}), + }, + ], + repositoryInstructions: true, + }); + const read = async ( + startLine?: number, + maxLines?: number, + signal?: AbortSignal, + path = 'file' + ) => { + const request: WorkspaceReadFileRequest = { + protocolVersion: 1, + operation: 'read_file', + workspaceId: 'primary', + path, + ...(startLine !== undefined ? { startLine } : {}), + ...(maxLines !== undefined ? { maxLines } : {}), + }; + const result = await tools.execute(request, signal); + assert.ok( + isWorkspaceToolResult(request, result), + 'current protocol accepts the window' + ); + assert.equal(result.operation, 'read_file'); + return result; + }; + const probe = await fs.open(join(root, 'file'), 'w+'); + const prototype = Object.getPrototypeOf(probe) as FileHandle; + await probe.close(); + return { parent, root, tools, read, prototype }; +} + +function hasCode(code: string) { + return (error: unknown) => + error instanceof WorkspaceToolError && error.code === code; +} + +test('small windows from large files use bounded descriptor reads, not readFile', async t => { + const { root, read, prototype } = await fixture(t, true); + await fs.writeFile( + join(root, 'file'), + 'first\nsecond\n' + 'tail\n'.repeat(300_000) + ); + const originalRead = prototype.read; + let bytesRead = 0; + t.mock.method(prototype, 'readFile', () => assert.fail('whole-file read')); + t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + ...args: Parameters + ) { + const result = await originalRead.apply(this, args); + bytesRead += result.bytesRead; + return result; + } + ); + const result = await read(1, 2); + assert.equal(result.content, 'first\nsecond'); + assert.equal(result.nextStartLine, 3); + assert.ok(bytesRead <= 64 * 1024 + 3); +}); + +test('sequential line and byte pagination reconstructs large content exactly', async t => { + const { root, read } = await fixture(t); + const lines = Array.from( + { length: 1_201 }, + (_, i) => `${i}: ${'é🦊'.repeat(300)}` + ); + await fs.writeFile(join(root, 'file'), '\ufeff' + lines.join('\n') + '\n'); + const collected: string[] = []; + let start = 1; + for (;;) { + const result = await read(start, 500); + assert.ok(Buffer.byteLength(result.content) <= MAX_BYTES); + collected.push(...result.content.split('\n')); + assert.equal(result.endLine, collected.length); + if (!result.truncated) break; + assert.ok(result.nextStartLine! > start); + start = result.nextStartLine!; + } + assert.deepEqual(collected, lines); + const defaults = await read(); + assert.equal(defaults.endLine, 200); + assert.equal(defaults.nextStartLine, 201); +}); + +test('byte ceilings include separators and never split or repeat an oversized line', async t => { + const { root, read } = await fixture(t); + const exact = 'a'.repeat(MAX_BYTES); + await fs.writeFile( + join(root, 'file'), + exact + '\nx\n' + 'z'.repeat(MAX_BYTES + 1) + ); + const first = await read(1, 500); + assert.equal(first.content, exact); + assert.equal(first.nextStartLine, 2); + const second = await read(2, 500); + assert.equal(second.content, 'x'); + assert.equal(second.nextStartLine, 3); + await assert.rejects(read(3), error => { + assert.ok(hasCode('READ_LIMIT_EXCEEDED')(error)); + assert.match( + (error as Error).message, + /line 3.*later startLine.*search_text/ + ); + assert.ok((error as Error).message.length < 300); + return true; + }); + await fs.writeFile( + join(root, 'file'), + 'a'.repeat(MAX_BYTES - 2) + '\nb\nc' + ); + const separators = await read(1, 500); + assert.equal(Buffer.byteLength(separators.content), MAX_BYTES); + assert.equal(separators.endLine, 2); + assert.equal(separators.nextStartLine, 3); +}); + +test('EOF, empty lines, CRLF, BOM-only files, and starts beyond EOF retain their contract', async t => { + const { root, read } = await fixture(t); + for (const [source, expected] of [ + ['', ['']], + ['\ufeff', ['']], + ['\n', ['']], + ['\n\n', ['', '']], + ['a\n', ['a']], + ['a\r\nb\r\n', ['a\r', 'b\r']], + ['a\nb', ['a', 'b']], + ['\ufeff\ufeffa', ['\ufeffa']], + ] as const) { + await fs.writeFile(join(root, 'file'), source); + for (let line = 1; line <= expected.length; line++) { + const result = await read(line, 1); + assert.equal(result.content, expected[line - 1]); + assert.equal(result.endLine, line); + assert.equal(result.truncated, line < expected.length); + } + const past = await read(expected.length + 1); + assert.equal(past.content, ''); + assert.equal(past.endLine, expected.length); + assert.equal(past.truncated, false); + } +}); + +test('streaming decoding handles short reads, split UTF-8, CRLF, and UTF-16 surrogates', async t => { + const { root, read, prototype } = await fixture(t); + const originalRead = prototype.read; + t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + buffer: Buffer, + offset: number, + length: number, + position: number + ) { + return originalRead.call(this, { + buffer, + offset, + length: Math.min(length, 1), + position, + }); + } + ); + const text = 'é🦊\r\n\ufeffsecond\n'; + const little = Buffer.from(text, 'utf16le'); + const big = Buffer.from(little).swap16(); + for (const source of [ + Buffer.from('\ufeff' + text), + Buffer.concat([Buffer.from([0xff, 0xfe]), little]), + Buffer.concat([Buffer.from([0xfe, 0xff]), big]), + ]) { + await fs.writeFile(join(root, 'file'), source); + const result = await read(); + assert.equal(result.content, 'é🦊\r\n\ufeffsecond'); + assert.equal(result.endLine, 2); + assert.equal(result.truncated, false); + } +}); + +test('UTF-8 replacement decoding and binary bytes are bounded text, not file validation', async t => { + const { root, read } = await fixture(t); + await fs.writeFile( + join(root, 'file'), + Buffer.from([0x61, 0x0a, 0xff, 0x00, 0xe2, 0x82]) + ); + const first = await read(1, 1); + assert.equal(first.content, 'a'); + assert.equal(first.truncated, true); + const second = await read(2); + assert.equal(second.content, '\ufffd\0\ufffd'); + await fs.writeFile(join(root, 'file'), Buffer.alloc(MAX_BYTES / 2, 0xff)); + await assert.rejects(read(), hasCode('READ_LIMIT_EXCEEDED')); +}); + +test('scanning skips huge earlier lines without accumulating them', async t => { + const { root, read } = await fixture(t); + await fs.writeFile( + join(root, 'file'), + 'x'.repeat(4 * MAX_BYTES) + '\nrequested\n' + ); + assert.equal((await read(2, 1)).content, 'requested'); +}); + +test('cancellation during scanning and collection closes every read handle', async t => { + for (const startLine of [1, 300_000]) { + const { root, read, prototype } = await fixture(t); + await fs.writeFile(join(root, 'file'), 'line\n'.repeat(300_000)); + const controller = new AbortController(); + const originalRead = prototype.read; + const handles = new Set(); + let calls = 0; + const mock = t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + ...args: Parameters + ) { + handles.add(this); + const result = await originalRead.apply(this, args); + if (++calls === 2) controller.abort(); + return result; + } + ); + await assert.rejects( + read(startLine, 500, controller.signal), + hasCode('EXECUTION_ABORTED') + ); + assert.ok(handles.size > 0); + for (const handle of handles) assert.equal(handle.fd, -1); + mock.mock.restore(); + } +}); + +test('far-away starts hit the scan deadline and close the descriptor', async t => { + const { root, read, prototype } = await fixture(t); + await fs.writeFile(join(root, 'file'), 'line\n'.repeat(300_000)); + const originalRead = prototype.read; + const now = performance.now(); + let elapsed = 0; + let handle: FileHandle | undefined; + t.mock.method(performance, 'now', () => now + elapsed); + t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + ...args: Parameters + ) { + handle = this; + const result = await originalRead.apply(this, args); + elapsed += 6_000; + return result; + } + ); + await assert.rejects(read(300_000), hasCode('READ_LIMIT_EXCEEDED')); + assert.equal(handle?.fd, -1); +}); + +test('growth is capped at the opened extent and truncation produces a valid EOF window', async t => { + for (const change of ['grow', 'truncate'] as const) { + const { root, read, prototype } = await fixture(t); + const path = join(root, 'file'); + await fs.writeFile(path, 'first\nsecond\n'); + const originalRead = prototype.read; + let calls = 0; + const mock = t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + ...args: Parameters + ) { + const result = await originalRead.apply(this, args); + if (++calls === 1) { + if (change === 'grow') + await fs.appendFile(path, 'appended\n'); + else await fs.truncate(path, 3); + } + return result; + } + ); + const result = await read(); + assert.equal( + result.content, + change === 'grow' ? 'first\nsecond' : 'fir' + ); + assert.equal(result.truncated, false); + mock.mock.restore(); + } +}); + +test('file and root replacement cannot redirect a held streaming read; descriptors close', async t => { + const { parent, root, read, prototype } = await fixture(t, true); + const path = join(root, 'file'); + await fs.writeFile(path, 'original\n' + 'inside\n'.repeat(200_000)); + await fs.writeFile(join(parent, 'secret'), 'outside secret'); + const originalRead = prototype.read; + const originalOpen = WorkspaceRootAccess.open; + let rootHandle: FileHandle | undefined; + let fileHandle: FileHandle | undefined; + t.mock.method( + WorkspaceRootAccess, + 'open', + async (...args: Parameters) => { + const access = await originalOpen(...args); + rootHandle = access.handle; + return access; + } + ); + let calls = 0; + t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + ...args: Parameters + ) { + fileHandle = this; + const result = await originalRead.apply(this, args); + if (++calls === 1) { + await fs.rename(path, join(root, 'old-file')); + await fs.symlink(join(parent, 'secret'), path); + await fs.rename(root, join(parent, 'old-root')); + await fs.mkdir(root); + await fs.writeFile(path, 'replacement'); + } + return result; + } + ); + const result = await read(2, 500); + assert.equal(result.content, Array(500).fill('inside').join('\n')); + assert.equal(result.nextStartLine, 502); + assert.equal(fileHandle?.fd, -1); + assert.equal(rootHandle?.fd, -1); + await assert.rejects(read(), hasCode('REGISTRATION_INVALID')); +}); + +test('oversized lines and I/O failures close file and held-root descriptors', async t => { + const { root, read, prototype } = await fixture(t, true); + const originalRead = prototype.read; + const originalOpen = WorkspaceRootAccess.open; + let fileHandle: FileHandle | undefined; + let rootHandle: FileHandle | undefined; + t.mock.method( + WorkspaceRootAccess, + 'open', + async (...args: Parameters) => { + const access = await originalOpen(...args); + rootHandle = access.handle; + return access; + } + ); + t.mock.method( + prototype, + 'read', + async function ( + this: FileHandle, + ...args: Parameters + ) { + fileHandle = this; + return originalRead.apply(this, args); + } + ); + await fs.writeFile(join(root, 'file'), 'x'.repeat(MAX_BYTES + 1)); + await assert.rejects(read(), hasCode('READ_LIMIT_EXCEEDED')); + assert.equal(fileHandle?.fd, -1); + assert.equal(rootHandle?.fd, -1); + t.mock.method(prototype, 'read', async function (this: FileHandle) { + fileHandle = this; + throw Object.assign(new Error('read failed'), { code: 'EIO' }); + }); + await assert.rejects(read(), hasCode('INVALID_PATH')); + assert.equal(fileHandle?.fd, -1); + assert.equal(rootHandle?.fd, -1); +}); + +test('search, preview, edit, and digest-checked instruction limits remain separate', async t => { + const { root, tools } = await fixture(t, true); + await fs.writeFile(join(root, 'file'), 'needle\n'.repeat(200_000)); + const search = await tools.execute({ + protocolVersion: 1, + operation: 'search_text', + workspaceId: 'primary', + path: 'file', + query: 'needle', + }); + assert.equal(search.operation, 'search_text'); + assert.deepEqual(search.matches, []); + for (const operation of ['preview_edit', 'edit_file'] as const) { + await assert.rejects( + tools.execute({ + protocolVersion: 1, + operation, + workspaceId: 'primary', + path: 'file', + oldText: 'needle', + newText: 'changed', + }), + hasCode( + operation === 'preview_edit' + ? 'READ_LIMIT_EXCEEDED' + : 'WRITE_LIMIT_EXCEEDED' + ) + ); + } + await fs.writeFile( + join(root, 'AGENTS.md'), + 'instructions\n'.repeat(200_000) + ); + const snapshot = await readRepositoryInstructions(root); + assert.ok(snapshot); + const request: WorkspaceReadFileRequest = { + protocolVersion: 1, + operation: 'read_file', + workspaceId: 'primary', + path: 'AGENTS.md', + instructionSha256: snapshot.descriptor.sha256, + }; + const result = await tools.execute(request); + assert.ok(isWorkspaceToolResult(request, result)); + assert.equal(result.operation, 'read_file'); + assert.equal(result.content, snapshot.content); + assert.equal(result.truncated, true); + assert.equal(result.nextStartLine, undefined); + await fs.writeFile(join(root, 'AGENTS.md'), 'changed'); + await assert.rejects(tools.execute(request), hasCode('INVALID_PATH')); +}); + +test('byte-budget pagination reconstructs every complete line', async t => { + const { root, read } = await fixture(t); + const lines = Array.from( + { length: 1_201 }, + (_, i) => `${i}: ${'é🦊'.repeat(400)}` + ); + await fs.writeFile(join(root, 'file'), lines.join('\n') + '\n'); + const collected: string[] = []; + let startLine = 1; + for (;;) { + const result = await read(startLine, 500); + const selected = result.content.split('\n'); + assert.ok(Buffer.byteLength(result.content) <= MAX_BYTES); + if (result.truncated) + assert.ok(selected.length < 500, 'bytes forced continuation'); + collected.push(...selected); + if (!result.truncated) break; + assert.ok(result.nextStartLine! > startLine); + startLine = result.nextStartLine!; + } + assert.deepEqual(collected, lines); +}); diff --git a/packages/code/src/workspace.ts b/packages/code/src/workspace.ts index e48365c8..746c7249 100644 --- a/packages/code/src/workspace.ts +++ b/packages/code/src/workspace.ts @@ -1,6 +1,6 @@ import { createHash, randomBytes } from 'node:crypto'; import { constants } from 'node:fs'; -import { holdsWorkspaceRoot, link, lstat, mkdir, open, realpath, rename, stat, unlink, spawn, withWorkspaceRoot, WorkspaceRootAccessError } from './root-access.js'; +import { deferWorkspaceRootCleanup, holdsWorkspaceRoot, link, lstat, mkdir, open, realpath, rename, stat, unlink, spawn, withWorkspaceRoot, WorkspaceRootAccessError } from './root-access.js'; import { basename, dirname, isAbsolute, relative, resolve, sep } from 'node:path'; import type { FileHandle } from 'node:fs/promises'; @@ -133,6 +133,8 @@ const MAX_SEARCH_CANDIDATE_BYTES = 1024 * 1024; const MAX_SEARCH_CANDIDATES = 20_000; const SEARCH_TIMEOUT_MS = 10_000; const LIST_TIMEOUT_MS = 10_000; +const READ_TIMEOUT_MS = 10_000; +const READ_CHUNK_BYTES = 64 * 1024; const READ_OPERATIONS = [ 'read_file', @@ -273,10 +275,18 @@ function resolveWorkspacePath(root: string, requestedPath: string): string { return candidate; } -async function readConfinedFileBuffer( +interface WorkspaceReadControl { + signal?: AbortSignal; + deadline: number; + interrupted: boolean; +} + +async function withConfinedFile( root: string, requestedPath: string, -): Promise { + read: (handle: FileHandle, size: number) => Promise, + control?: WorkspaceReadControl, +): Promise { const candidate = resolveWorkspacePath(root, requestedPath); let handle: FileHandle | undefined; try { @@ -297,7 +307,31 @@ async function readConfinedFileBuffer( ) { throw new Error('Invalid workspace path'); } - if (openedFile.size > BRIDGE_WORKSPACE_READ_MAX_BYTES) { + return await read(handle, openedFile.size); + } catch (error) { + if (error instanceof WorkspaceToolError) throw error; + if (isMissingEntry(error)) { + throw await classifyMissingWorkspacePath(root, candidate); + } + throw new WorkspaceToolError('Invalid workspace path', 'INVALID_PATH'); + } finally { + if (handle) { + const closing = handle.close(); + if (control?.interrupted) { + // Node retains the descriptor and buffer until outstanding I/O drains. + deferWorkspaceRootCleanup(); + void closing.catch(() => {}); + } else await closing; + } + } +} + +async function readConfinedFileBuffer( + root: string, + requestedPath: string, +): Promise { + return withConfinedFile(root, requestedPath, async (handle, size) => { + if (size > BRIDGE_WORKSPACE_READ_MAX_BYTES) { throw new WorkspaceToolError( 'Workspace file exceeds read limit', 'READ_LIMIT_EXCEEDED', @@ -322,31 +356,222 @@ async function readConfinedFileBuffer( ); } return buffer.subarray(0, bytesRead); - } catch (error) { - if (error instanceof WorkspaceToolError) throw error; - if (isMissingEntry(error)) { - throw await classifyMissingWorkspacePath(root, candidate); - } - throw new WorkspaceToolError('Invalid workspace path', 'INVALID_PATH'); + }); +} + +function interruptRead( + control: WorkspaceReadControl, + message: string, + code: 'EXECUTION_ABORTED' | 'READ_LIMIT_EXCEEDED', +): WorkspaceToolError { + control.interrupted = true; + return new WorkspaceToolError(message, code); +} + +function checkReadDeadline(control: WorkspaceReadControl): void { + if (control.signal?.aborted) { + throw interruptRead( + control, + 'Workspace tool execution aborted', + 'EXECUTION_ABORTED', + ); + } + if (performance.now() >= control.deadline) { + throw interruptRead( + control, + 'Workspace read exceeded its scan time limit; request an earlier startLine or use search_text to locate content', + 'READ_LIMIT_EXCEEDED', + ); + } +} + +async function readWorkspaceChunk( + handle: FileHandle, + buffer: Buffer, + length: number, + position: number, + control: WorkspaceReadControl, +): Promise { + checkReadDeadline(control); + const { signal, deadline } = control; + let timer: ReturnType | undefined; + let abort: (() => void) | undefined; + try { + return await Promise.race([ + handle + .read(buffer, 0, length, position) + .then((result) => result.bytesRead), + new Promise((_, reject) => { + abort = () => + reject( + interruptRead( + control, + 'Workspace tool execution aborted', + 'EXECUTION_ABORTED', + ), + ); + signal?.addEventListener('abort', abort, { once: true }); + if (signal?.aborted) abort(); + timer = setTimeout( + () => + reject( + interruptRead( + control, + 'Workspace read exceeded its scan time limit; retry with an earlier startLine', + 'READ_LIMIT_EXCEEDED', + ), + ), + Math.max(0, deadline - performance.now()), + ); + }), + ]); } finally { - await handle?.close(); + if (timer !== undefined) clearTimeout(timer); + if (abort) signal?.removeEventListener('abort', abort); } } async function readConfinedFile( root: string, - requestedPath: string, -): Promise { - const decoded = decodeWorkspaceText( - await readConfinedFileBuffer(root, requestedPath), - ); - if (Buffer.byteLength(decoded, 'utf8') > BRIDGE_WORKSPACE_READ_MAX_BYTES) { - throw new WorkspaceToolError( - 'Workspace file exceeds read limit', - 'READ_LIMIT_EXCEEDED', + request: WorkspaceReadFileRequest, + signal?: AbortSignal, +): Promise { + const startLine = request.startLine ?? 1; + const maxLines = request.maxLines ?? 200; + const control: WorkspaceReadControl = { + signal, + deadline: performance.now() + READ_TIMEOUT_MS, + interrupted: false, + }; + const read = async ( + handle: FileHandle, + size: number, + ): Promise => { + const buffer = Buffer.allocUnsafe(READ_CHUNK_BYTES); + // Fix the read extent at admission so an appending writer cannot extend the scan. + let position = 0; + const header = Buffer.alloc(3); + let headerBytes = 0; + while (headerBytes < header.length && position < size) { + const bytes = await readWorkspaceChunk( + handle, + buffer, + Math.min(header.length - headerBytes, size - position), + position, + control, + ); + checkReadDeadline(control); + if (bytes === 0) break; + buffer.copy(header, headerBytes, 0, bytes); + headerBytes += bytes; + position += bytes; + } + const utf16le = + headerBytes >= 2 && header[0] === 0xff && header[1] === 0xfe; + const utf16be = + headerBytes >= 2 && header[0] === 0xfe && header[1] === 0xff; + const bomBytes = + utf16le || utf16be + ? 2 + : headerBytes === 3 && + header[0] === 0xef && + header[1] === 0xbb && + header[2] === 0xbf + ? 3 + : 0; + // Preserve ordinary reads' replacement decoding and BOM-marked UTF-16 support. + const decoder = new TextDecoder( + utf16le ? 'utf-16le' : utf16be ? 'utf-16be' : 'utf-8', + { ignoreBOM: true }, ); - } - return decoded; + const lines: string[] = []; + let fragments: string[] = []; + let line = 1; + let lineBytes = 0; + let returnedBytes = 0; + let hasText = false; + let truncated = false; + const consume = (text: string): void => { + let offset = 0; + while (offset < text.length) { + checkReadDeadline(control); + const newline = text.indexOf('\n', offset); + const end = newline < 0 ? text.length : newline; + if (line >= startLine) { + if (lines.length === maxLines) { + truncated = true; + return; + } + const fragment = text.slice(offset, end); + const bytes = Buffer.byteLength(fragment, 'utf8'); + if ( + returnedBytes + (lines.length > 0 ? 1 : 0) + lineBytes + bytes > + BRIDGE_WORKSPACE_READ_MAX_BYTES + ) { + if (lines.length === 0) { + throw new WorkspaceToolError( + `Workspace file exceeds read limit: line ${line} cannot fit within ${BRIDGE_WORKSPACE_READ_MAX_BYTES} UTF-8 bytes; request a later startLine or use search_text`, + 'READ_LIMIT_EXCEEDED', + ); + } + truncated = true; + return; + } + if (fragment.length > 0) fragments.push(fragment); + lineBytes += bytes; + } + hasText ||= end > offset; + if (newline < 0) return; + if (line >= startLine) { + lines.push(fragments.join('')); + returnedBytes += (lines.length > 1 ? 1 : 0) + lineBytes; + fragments = []; + lineBytes = 0; + } + line += 1; + hasText = false; + offset = newline + 1; + } + }; + consume( + decoder.decode(header.subarray(bomBytes, headerBytes), { + stream: true, + }), + ); + while (!truncated && position < size) { + const bytes = await readWorkspaceChunk( + handle, + buffer, + Math.min(buffer.length, size - position), + position, + control, + ); + checkReadDeadline(control); + if (bytes === 0) break; + position += bytes; + consume(decoder.decode(buffer.subarray(0, bytes), { stream: true })); + } + if (!truncated) { + consume(decoder.decode()); + if (!truncated && line >= startLine && (hasText || line === 1)) { + lines.push(fragments.join('')); + } + } + checkReadDeadline(control); + const endLine = startLine + lines.length - 1; + return { + protocolVersion: BRIDGE_PROTOCOL_VERSION, + operation: 'read_file', + workspaceId: request.workspaceId, + path: request.path, + content: lines.join('\n'), + startLine, + endLine, + truncated, + ...(truncated ? { nextStartLine: endLine + 1 } : {}), + }; + }; + return withConfinedFile(root, request.path, read, control); } interface WorkspaceRoot { @@ -1831,31 +2056,7 @@ export class LocalWorkspaceTools implements WorkspaceToolExecutor { ) { throw new WorkspaceToolError('Invalid workspace read', 'INVALID_REQUEST'); } - const content = await readConfinedFile(root, request.path); - if (signal?.aborted) { - throw new WorkspaceToolError( - 'Workspace tool execution aborted', - 'EXECUTION_ABORTED', - ); - } - const lines = content.endsWith('\n') - ? content.slice(0, -1).split('\n') - : content.split('\n'); - const selected = lines.slice(startLine - 1, startLine - 1 + maxLines); - const endLine = startLine + selected.length - 1; - const truncated = endLine < lines.length; - - return { - protocolVersion: BRIDGE_PROTOCOL_VERSION, - operation: 'read_file', - workspaceId: request.workspaceId, - path: request.path, - content: selected.join('\n'), - startLine, - endLine, - truncated, - ...(truncated ? { nextStartLine: endLine + 1 } : {}), - }; + return readConfinedFile(root, request, signal); } }