diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b8868fcc..ac397c6f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -88,7 +88,7 @@ jobs: # namespace and every sandboxed command dies with # "loopback: Failed RTM_NEWADDR: Operation not permitted". sudo sysctl -w kernel.apparmor_restrict_unprivileged_userns=0 || true - bwrap --ro-bind / / --unshare-net --dev /dev --proc /proc -- /bin/sh -lc 'echo sandbox-ok' + bwrap --ro-bind / / --unshare-net --unshare-pid --dev /dev --proc /proc -- /bin/sh -lc 'echo sandbox-ok' - run: npm ci --prefer-offline --no-audit --no-fund # Coverage thresholds live in vitest.config.ts (ratcheted floor below the # measured baseline — see that file's comment). Previously `npm test` had @@ -146,7 +146,7 @@ jobs: - name: Bubblewrap alone must not be assumed sufficient run: | set -uo pipefail - if bwrap --ro-bind / / --unshare-net --dev /dev --proc /proc -- /usr/bin/true 2>/dev/null; then + if bwrap --ro-bind / / --unshare-net --unshare-pid --dev /dev --proc /proc -- /usr/bin/true 2>/dev/null; then echo "::warning::The sandbox now works without lifting the AppArmor restriction — the README's sysctl step may be obsolete on this image" else echo "Confirmed: installing bubblewrap alone does not give a working sandbox on this runner" @@ -157,5 +157,5 @@ jobs: sudo sysctl -w kernel.apparmor_restrict_unprivileged_userns=0 # The same namespaces src/verify/runner.ts sets up. Probing a subset # would pass on a host that denies one of the others. - bwrap --ro-bind / / --unshare-net --dev /dev --proc /proc -- /usr/bin/true + bwrap --ro-bind / / --unshare-net --unshare-pid --dev /dev --proc /proc -- /usr/bin/true echo "OS verification sandbox is available on ubuntu-latest with bubblewrap + the AppArmor sysctl" diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index b4e2c83f..49d97652 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -76,7 +76,7 @@ jobs: if: steps.ver.outputs.exists == 'false' run: | sudo sysctl -w kernel.apparmor_restrict_unprivileged_userns=0 || true - bwrap --ro-bind / / --unshare-net --dev /dev --proc /proc -- /bin/sh -lc 'echo sandbox-ok' + bwrap --ro-bind / / --unshare-net --unshare-pid --dev /dev --proc /proc -- /bin/sh -lc 'echo sandbox-ok' - if: steps.ver.outputs.exists == 'false' run: npm ci diff --git a/CHANGELOG.md b/CHANGELOG.md index 567f6b09..4096ec83 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,11 @@ ## [Unreleased] +### Added + +- **`advisor` role: a second, independent review of the same diff.** Ported from the harness agent patterns. A separate, independently-prompted model is asked one narrow question — *what concrete defects did the reviewer miss?* — and its only permitted effect is to make the gate **more** cautious: it may append concrete findings the reviewer did not report and raise severity, but the merged decision is the max rank of the two (`approve < revise < reject`), so an advisor `approve` can never soften a reviewer `revise`/`reject`, and a raised severity with no concrete finding is discarded. Any error, timeout, empty, or unparseable output fails **open** — `ran: false`, the reviewer's result untouched, no exit code changed. `review`, `review --max` (per area), and the `--fix` re-review loop all run it before dedupe, so its findings are deduped against history like any other. Disabled by default (a second paid call per review), configured under `autonomous.defaultRoles.advisor`; its model must come from a different family than the reviewer's or it is a second identical opinion. +- **Declarative per-role subagent settings: `tools` and `effort`.** Each role (`worker`/`reviewer`/`advisor`/…) may declare `tools.allow` / `tools.deny` and `effort`. An allow-list can only **narrow** the role's default tool set — it can never grant a tool the role would not otherwise have, so a misconfiguration cannot hand a read-only reviewer `bash`. `deny` is applied after `allow` (deny wins) and supports a trailing `*` (`scratch_*`). `effort` selects the native-loop reasoning level for the stage. + ## 0.24.3 — 2026-09-28 ### Added diff --git a/config.example.yaml b/config.example.yaml index 9605cd64..ddd39532 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -183,6 +183,18 @@ autonomous: # run ends on completion, the repeated-tool-call guard, or timeoutMs — a turn # count is not a property of the task. Set a number only to cap cost hard. # maxTurns: 0 + # + # Declarative tool scoping per role. `allow` can only NARROW the role's + # default tools — it can never grant one the role would not otherwise have, + # so an allow-list naming `bash` on a read-only role stays withheld. `deny` + # is applied after `allow` (deny wins), and supports a trailing `*` + # (`scratch_*`). Use this to keep a stage's blast radius explicit instead + # of relying on the stage's defaults. + # tools: + # allow: [read_file, write_file, edit_file, bash] + # deny: [web_fetch, web_search] + # reasoning effort for this role's native-loop adapter (low|medium|high) + # effort: medium reviewer: enabled: true model: gpt-5.6-sol # Correctness gate; light profile lowers this to Terra @@ -192,6 +204,23 @@ autonomous: # diff-scaled default (300s base); a slow model needs more — measured # 2026-09-17: half the reviews died at 300s. maxTurns 0 = no ceiling. timeoutMs: 600000 + # Second, independently-prompted review of the SAME diff. Its only permitted + # effect is to add findings the reviewer missed and to raise severity — it + # can never soften the reviewer's verdict, and any raised severity without a + # concrete finding is discarded. DISABLED by default: it is a second paid + # call on every review. + # + # Its model must come from a DIFFERENT family than the reviewer's. On the + # reviewer's own model the advisor is a second identical opinion — the same + # weights re-deriving the same blind spots (the trap modelCompat.ts + # documents for `escalate`). The shipped reviewer default is + # deepseek/deepseek-v4-flash; z-ai/glm-5.2 is the family-independent + # alternative measured at 100% detect / 6% false-reject (6s avg) on the + # planted-defect fixtures where the reviewer scored 0% false-reject at 36s. + advisor: + enabled: false + model: z-ai/glm-5.2 + timeoutMs: 45000 # 45s single-turn ceiling, as the guard arbiter uses tester: enabled: false model: gpt-5.6-terra # Used only if deterministic verify cannot run diff --git a/src/adapters/agenticLoop.test.ts b/src/adapters/agenticLoop.test.ts index 0fbcd15e..549649fa 100644 --- a/src/adapters/agenticLoop.test.ts +++ b/src/adapters/agenticLoop.test.ts @@ -534,6 +534,135 @@ describe('runAgenticLoop tool exposure options', () => { expect(toolNames).not.toContain('search_memory'); }); + it('exposes only the allow-listed tools, and a deny entry removes what allow kept', async () => { + let toolNames: string[] = []; + + await runAgenticLoop({ + prompt: 'x', + cwd: process.cwd(), + model: 'test', + webTools: false, + memoryTools: false, + // `search_files` is allow-listed and then denied: allow can never add a + // tool back that deny removed, not even a member of `allow` itself. + toolAllow: ['read_file', 'search_files', 'write_file'], + toolDeny: ['search_files'], + maxTurns: 1, + callApi: async (_messages, tools) => { + toolNames = tools.map((tool) => tool.function.name); + return finalResp('done'); + }, + }); + + expect(toolNames).toEqual(['read_file', 'write_file']); + }); + + it('a deny wildcard withholds the whole scratch family', async () => { + let toolNames: string[] = []; + + await runAgenticLoop({ + prompt: 'x', + cwd: process.cwd(), + model: 'test', + webTools: false, + memoryTools: false, + scratchpadRunId: 'AGT-0000', + toolDeny: ['scratch_*'], + maxTurns: 1, + callApi: async (_messages, tools) => { + toolNames = tools.map((tool) => tool.function.name); + return finalResp('done'); + }, + }); + + expect(toolNames).not.toContain('scratch_write'); + expect(toolNames).not.toContain('scratch_read'); + expect(toolNames).toContain('read_file'); + }); + + it('an allow-list cannot resurrect bash on a read-only run', async () => { + // The narrowing-only invariant: `allow` intersects with the composition, so + // naming a withheld tool does not expose it. A read-only run (the reviewer's + // shape) keeps bash hidden no matter what the role declared. + let toolNames: string[] = []; + + await runAgenticLoop({ + prompt: 'x', + cwd: process.cwd(), + model: 'test', + readOnly: true, + webTools: false, + memoryTools: false, + toolAllow: ['bash', 'read_file', 'write_file'], + maxTurns: 1, + callApi: async (_messages, tools) => { + toolNames = tools.map((tool) => tool.function.name); + return finalResp('done'); + }, + }); + + expect(toolNames).toEqual(['read_file']); + }); + + it('refuses a tool the role scope narrowed away, at dispatch as well as in the schema', async () => { + // The schema above and the dispatch set below are the same narrowed array, + // so a provider that emits a withheld name anyway is answered, not obeyed. + let turn = 0; + let deniedResult = ''; + + await runAgenticLoop({ + prompt: 'x', + cwd: process.cwd(), + model: 'test', + webTools: false, + memoryTools: false, + toolAllow: ['read_file'], + maxTurns: 2, + callApi: async (messages, tools) => { + if (turn++ === 0) { + expect(tools.map((tool) => tool.function.name)).toEqual(['read_file']); + return toolCallResp('hidden-write', 'write_file', { path: 'out.txt', content: 'x' }); + } + deniedResult = messages.at(-1)?.content ?? ''; + return finalResp('done'); + }, + }); + + expect(deniedResult).toContain('TOOL_NOT_ALLOWED'); + expect(deniedResult).toContain('write_file'); + }); + + it('leaves MCP and coordination tools to their own flags, not the role list', async () => { + // RoleConfig.tools names built-ins (see its doc): a role narrowing to the + // file tools must not silently lose its board access, which the MCP and + // coordination flags govern. + let toolNames: string[] = []; + + await runAgenticLoop({ + prompt: 'x', + cwd: process.cwd(), + model: 'test', + webTools: false, + memoryTools: false, + toolAllow: ['read_file'], + maxTurns: 1, + mcpTools: [{ + type: 'function', + function: { name: 'linear__get_issue', description: '', parameters: { type: 'object' } }, + }], + coordinationContext: { repository: '/repo', taskId: 'supervisor', actor: 'orchestrator' }, + callApi: async (_messages, tools) => { + toolNames = tools.map((tool) => tool.function.name); + return finalResp('done'); + }, + }); + + expect(toolNames).toContain('read_file'); + expect(toolNames).toContain('linear__get_issue'); + expect(toolNames).toContain('coordination_read'); + expect(toolNames).not.toContain('write_file'); + }); + it('withholds the scratch tools when the run has no scratchpad', async () => { let toolNames: string[] = []; await runAgenticLoop({ diff --git a/src/adapters/agenticLoop.ts b/src/adapters/agenticLoop.ts index 91226a41..d38a6a96 100644 --- a/src/adapters/agenticLoop.ts +++ b/src/adapters/agenticLoop.ts @@ -202,6 +202,16 @@ export interface AgenticLoopOptions { webTools?: boolean; /** Expose search_memory (default true). Disabled for isolated/temp repo benchmarks. */ memoryTools?: boolean; + /** + * Declarative per-role tool scope (RoleConfig.tools), naming BUILT-IN tools. + * `toolAllow` keeps only the names it lists; `toolDeny` then removes names from + * what remains, a trailing `*` standing for a prefix (`scratch_*`). Both run over + * the built-ins left by the readOnly/shell/web rules, so they only ever NARROW — + * `allow` cannot resurrect `bash` on a readOnly run. MCP and coordination tools + * keep their own flags; the narrowed set is also the dispatch allow-list. + */ + toolAllow?: string[]; + toolDeny?: string[]; /** * Run whose scratchpad `scratch_write`/`scratch_read` address (AGT-4459). * Absent means no scratchpad: the two tools are withheld from the model and @@ -287,6 +297,33 @@ export interface AgenticLoopResult { // ============ 에이전틱 루프 ============ +/** + * Apply a role's declarative `tools.allow` / `tools.deny` (RoleConfig) to the + * built-in tools the readOnly / scratch / shell / web filters have already shaped. + * + * It runs LAST and only removes entries — an allow-list cannot resurrect a tool an + * earlier rule withheld (a `bash` in `allow` stays hidden on a readOnly run), and + * `deny` follows `allow`, so the two cannot contradict each other. A `deny` entry + * ending in `*` matches by prefix (`scratch_*`), how the scratch tools are + * addressed as a family. + */ +function applyRoleToolScope( + tools: ToolDefinition[], + toolAllow?: string[], + toolDeny?: string[], +): ToolDefinition[] { + let scoped = tools; + if (toolAllow && toolAllow.length > 0) { + const allowed = new Set(toolAllow); + scoped = scoped.filter((tool) => allowed.has(tool.function.name)); + } + return toolDeny && toolDeny.length > 0 + ? scoped.filter((tool) => !toolDeny.some((denied) => denied.endsWith('*') + ? tool.function.name.startsWith(denied.slice(0, -1)) + : tool.function.name === denied)) + : scoped; +} + /** * 에이전틱 도구 루프 실행 * @@ -359,6 +396,8 @@ async function runAgenticLoopInner( bashTimeoutMs, webTools = true, memoryTools = true, + toolAllow, + toolDeny, scratchpadRunId, shellTools: requestedShellTools = true, sandboxExecutorSessionFactory, @@ -444,23 +483,31 @@ async function runAgenticLoopInner( const visibleBaseTools = readOnly ? shellFilteredTools.filter((t) => !['write_file', 'edit_file', 'bash', 'remember'].includes(t.function.name)) : shellFilteredTools; - const tools = enableTools + const builtinTools = enableTools ? [ ...visibleBaseTools, ...(filesystemTools && applyPatch && editFormat === 'json' && !readOnly ? [APPLY_PATCH_TOOL] : []), // Not in readOnly: it spawns compiler subprocesses, matching bash's exclusion. ...(filesystemTools && diagnosticsTool && !readOnly && shellTools ? [DIAGNOSTICS_TOOL] : []), - // Both are withheld in readOnly. A read-only run exists because the - // material under inspection is untrusted, and a fetch is an outbound - // channel for anything the agent can read — the provider credential - // included. MCP servers are withheld for the mirror reason: OpenSwarm's - // own memory server exposes writes, so injected content could leave - // something behind for a later run. (INT-3189) + // Also withheld in readOnly: the material under inspection is untrusted, + // and a fetch is an outbound channel for anything the agent can read — + // the provider credential included. (INT-3189) ...(webTools && !readOnly ? WEB_TOOL_DEFINITIONS : []), + ] + : []; + // MCP and coordination tools keep their own flags; the role list names built-ins. + // MCP is withheld in readOnly for the mirror reason: our memory server exposes + // writes, so injected content could leave something behind. (INT-3189) + const externalTools = enableTools + ? [ ...(readOnly ? [] : humanSurfaceFilteredMcp.tools), ...(readOnly || !coordinationContext ? [] : COORDINATION_TOOL_DEFINITIONS), ] : []; + // The role's declared scope goes LAST, over the built-ins every rule above has + // already shaped, so it intersects with them instead of overriding one — see + // applyRoleToolScope. `allowedToolNames` below is built from the same result. + const tools = [...applyRoleToolScope(builtinTools, toolAllow, toolDeny), ...externalTools]; // The provider-visible schema is not an enforcement boundary. Carry the // exact same set into dispatch so a hidden tool call cannot reach a globally // registered MCP route (or another built-in withheld for this run). diff --git a/src/adapters/atlascloud.ts b/src/adapters/atlascloud.ts index 49ec8623..f968bcf8 100644 --- a/src/adapters/atlascloud.ts +++ b/src/adapters/atlascloud.ts @@ -164,6 +164,8 @@ export class AtlasCloudCliAdapter implements CliAdapter { webTools: options.webTools, memoryTools: options.memoryTools, shellTools: options.shellTools, + toolAllow: options.toolAllow, + toolDeny: options.toolDeny, filesystemTools: options.filesystemTools, diagnosticsTool: options.diagnosticsTool, readOnly: options.readOnly, diff --git a/src/adapters/base.spawn.test.ts b/src/adapters/base.spawn.test.ts index dc69230a..3227661e 100644 --- a/src/adapters/base.spawn.test.ts +++ b/src/adapters/base.spawn.test.ts @@ -13,7 +13,7 @@ vi.mock('node:child_process', async (importOriginal) => ({ spawn: spawnMock, })); -import { spawnCli, terminateCliProcessTree } from './base.js'; +import { spawnCli, terminateCliProcessTree, CLI_OUTPUT_MAX_BYTES } from './base.js'; import { prepareCliProcessTreeSpawn, trackCliProcessTree, @@ -533,6 +533,71 @@ describe('delegated-CLI capability guards', () => { ).rejects.toThrow(/cannot withhold shell access/); }); + it('warns that a role tool allow/deny list is inert on a delegated CLI, without failing the run', async () => { + // The delegated CLI owns its own tools, so the list cannot be enforced here. + // Warning (not throwing) keeps existing claude/codex configs running while + // making sure a fence nobody applies is never silent. (AGT-4444 class) + const warns: string[] = []; + const warn = vi.spyOn(console, 'warn').mockImplementation((line: unknown) => { warns.push(String(line)); }); + const proc = Object.assign(new EventEmitter(), { + pid: 311, + stdout: new PassThrough(), + stderr: new PassThrough(), + stdin: Object.assign(new EventEmitter(), { end: vi.fn() }), + kill: vi.fn(), + }); + spawnMock.mockImplementationOnce(() => { + queueMicrotask(() => { + proc.stdout.end('ok'); + proc.emit('close', 0); + }); + return proc; + }); + try { + await expect(spawnCli(delegated(), { + prompt: 'p', cwd: process.cwd(), toolAllow: ['read_file'], toolDeny: ['scratch_*'], + })).resolves.toMatchObject({ stdout: 'ok' }); + } finally { + warn.mockRestore(); + } + expect(warns.join('\n')).toContain('the role tool allow/deny list'); + }); + + it('warns that protectedFiles and forbidPublication are inert on a delegated CLI (AGT-4444)', async () => { + // Same failure class as the tool list: these fences live in the in-process + // tool executor, so a role routed to a delegated CLI silently loses them. + // The warning is the defect's whole point — silence is what made a dead + // fence look like a working one. + const warns: string[] = []; + const warn = vi.spyOn(console, 'warn').mockImplementation((line: unknown) => { warns.push(String(line)); }); + const proc = Object.assign(new EventEmitter(), { + pid: 312, + stdout: new PassThrough(), + stderr: new PassThrough(), + stdin: Object.assign(new EventEmitter(), { end: vi.fn() }), + kill: vi.fn(), + }); + spawnMock.mockImplementationOnce(() => { + queueMicrotask(() => { + proc.stdout.end('ok'); + proc.emit('close', 0); + }); + return proc; + }); + try { + await expect(spawnCli(delegated(), { + prompt: 'p', cwd: process.cwd(), + protectedFiles: ['secrets.env'], + forbidPublication: true, + })).resolves.toMatchObject({ stdout: 'ok' }); + } finally { + warn.mockRestore(); + } + const text = warns.join('\n'); + expect(text).toContain('1 protected path(s)'); + expect(text).toContain('the publication fence'); + }); + it('does not construct or spawn a delegated fake CLI in strict mode, even with HOME credentials', async () => { const buildCommand = vi.fn(() => ({ command: 'fake-codex', args: [] })); const adapter = { ...delegated(), name: 'fake-codex', buildCommand } satisfies CliAdapter; @@ -599,3 +664,143 @@ describe('delegated-CLI capability guards', () => { warn.mockRestore(); }); }); + +describe('bounded CLI output retention', () => { + /** An adapter whose CLI shells out, with optional incremental stream parsing. */ + const fixture = (parseStreamingChunk?: CliAdapter['parseStreamingChunk']): CliAdapter => ({ + name: 'fixture-cli', + capabilities: { + supportsStreaming: !!parseStreamingChunk, + supportsJsonOutput: false, + supportsModelSelection: false, + managedGit: false, + supportedSkills: [], + }, + isAvailable: async () => true, + getDefaultModel: async () => 'fixture', + buildCommand: () => ({ command: 'fixture-cli', args: [] }), + parseStreamingChunk, + parseWorkerOutput: () => ({ success: true, summary: '', filesChanged: [], commands: [], output: '' }), + parseReviewerOutput: () => ({ decision: 'approve', feedback: '', issues: [], suggestions: [] }), + }); + + const mockProc = (pid: number, emit: (proc: { stdout: PassThrough; stderr: PassThrough }) => void) => { + const proc = Object.assign(new EventEmitter(), { + pid, + stdout: new PassThrough(), + stderr: new PassThrough(), + stdin: Object.assign(new EventEmitter(), { end: vi.fn() }), + kill: vi.fn(), + }); + spawnMock.mockImplementationOnce(() => { + queueMicrotask(() => { + emit(proc); + proc.emit('close', 0); + }); + return proc; + }); + return proc; + }; + + it('resolves a flooding CLI and marks the retained output as truncated', async () => { + // A wedged or adversarial CLI can write for the whole timeout window. The + // retained copy is capped, so the daemon holds a bounded tail instead of + // everything the child ever printed. (The marker must be in the data: a + // silently short stdout reads as "the CLI said nothing".) + const floodBytes = 3 * 1024 * 1024; + mockProc(911, ({ stdout }) => { + stdout.write('a'.repeat(floodBytes)); + stdout.end('FINAL_RESULT_LINE'); + }); + + const result = await spawnCli(fixture(), { prompt: 'p', cwd: process.cwd() }); + + expect(result.exitCode).toBe(0); + // The tail is what every parser reads (`messages.at(-1)`, the stream-json + // result event), so the terminal line must survive the clip. + expect(result.stdout.endsWith('FINAL_RESULT_LINE')).toBe(true); + expect(result.stdout).toContain('[openswarm-cli-output:'); + expect(result.stdout).toMatch(/\[openswarm-cli-output: \d+ bytes omitted from the head\]/); + expect(Buffer.byteLength(result.stdout, 'utf8')).toBeLessThanOrEqual(CLI_OUTPUT_MAX_BYTES); + + // The count is the real one: it plus what was kept accounts for every byte. + const [, dropped] = result.stdout.match(/\[openswarm-cli-output: (\d+) bytes omitted/)!; + const keptBytes = Buffer.byteLength(result.stdout.slice(result.stdout.indexOf('\n') + 1), 'utf8'); + expect(Number(dropped) + keptBytes).toBe(floodBytes + 'FINAL_RESULT_LINE'.length); + // The marker line itself is what the operator sees instead of the head. + expect(result.stdout.split('\n')[0]).toMatch(/^\[openswarm-cli-output:/); + }); + + it('bounds stderr independently and keeps the diagnostic tail', async () => { + // stderr is where a failing CLI explains itself; the failure path below + // reports a snippet from it, so the tail must be the part retained. + const floodBytes = 2.5 * 1024 * 1024; + mockProc(912, ({ stderr }) => { + stderr.write('x'.repeat(floodBytes)); + stderr.end('Error: ENOENT: no such file or directory'); + }); + + const result = await spawnCli(fixture(), { prompt: 'p', cwd: process.cwd() }); + + expect(result.stderr.endsWith('Error: ENOENT: no such file or directory')).toBe(true); + expect(result.stderr).toContain('[openswarm-cli-output:'); + expect(Buffer.byteLength(result.stderr, 'utf8')).toBeLessThanOrEqual(CLI_OUTPUT_MAX_BYTES); + // Each stream carries its own ceiling — a flooding stdout must not eat the + // stderr budget, and vice versa. + expect(result.stdout).toBe(''); + }); + + it('leaves a normal run byte-for-byte untouched', async () => { + mockProc(913, ({ stdout, stderr }) => { + stdout.end('{"type":"result","result":"all good"}'); + stderr.end('warning: deprecated flag'); + }); + + const result = await spawnCli(fixture(), { prompt: 'p', cwd: process.cwd() }); + + expect(result.stdout).toBe('{"type":"result","result":"all good"}'); + expect(result.stderr).toBe('warning: deprecated flag'); + expect(result.stdout).not.toContain('[openswarm-cli-output:'); + expect(result.stderr).not.toContain('[openswarm-cli-output:'); + }); + + it('still feeds every byte to an incremental stream parser', async () => { + // The bound applies to the RETAINED copy only. The live log is built by the + // streaming parser from each chunk as it arrives; clipping its input would + // silently drop the middle of a long assistant message from the dashboard. + let seen = 0; + const parseStreamingChunk = vi.fn((chunk: string, _onLog: (line: string) => void, buffer = '') => { + seen += chunk.length; + return buffer; + }); + const floodBytes = 3 * 1024 * 1024; + mockProc(914, ({ stdout }) => stdout.end('b'.repeat(floodBytes))); + + const logged: string[] = []; + const result = await spawnCli(fixture(parseStreamingChunk), { + prompt: 'p', cwd: process.cwd(), onLog: (line) => logged.push(line), + }); + + expect(seen).toBe(floodBytes); + expect(result.stdout).toContain('[openswarm-cli-output:'); + }); + + it('bounds multibyte output by bytes, not characters, and keeps the tail decodable', async () => { + // 3 bytes per char: a char-wise cap would retain 3x the ceiling in bytes, + // and clipping the string rather than the buffer would also leave a broken + // half-character at the cut. Both are invisible with ASCII fixtures. + const floodBytes = 3 * 1024 * 1024; + mockProc(915, ({ stdout }) => { + stdout.write('한'.repeat(floodBytes / 3)); + stdout.end('\n결과: 완료'); + }); + + const result = await spawnCli(fixture(), { prompt: 'p', cwd: process.cwd() }); + + expect(Buffer.byteLength(result.stdout, 'utf8')).toBeLessThanOrEqual(CLI_OUTPUT_MAX_BYTES); + expect(result.stdout.endsWith('\n결과: 완료')).toBe(true); + expect(result.stdout).toContain('[openswarm-cli-output:'); + // No replacement character: the cut did not split a character in place. + expect(result.stdout).not.toContain('\uFFFD'); + }); +}); diff --git a/src/adapters/base.ts b/src/adapters/base.ts index c33a73cf..8530f277 100644 --- a/src/adapters/base.ts +++ b/src/adapters/base.ts @@ -28,6 +28,96 @@ import { createSessionRecorder, type SessionRecorder } from '../support/sessionL export { terminateCliProcessTree } from './processTree.js'; +/** + * Byte ceiling on the output kept for ONE stream. + * + * The data handlers below used to concatenate every chunk for the whole worker + * lifetime, so a CLI that is verbose, wedged, or adversarial could hold hundreds + * of MB in the daemon until its timeout fired — and streaming parsing does not + * reduce that, it only decides what is logged. 2 MiB is the ceiling this repo + * already puts on other bytes read from a source we do not control (web fetch + * bodies, codex MCP enumeration), it is 16x the 128 KiB CLI buffer in + * automation/scheduler.ts — which was sized for a stderr snippet, not for + * parseable stdout — and it sits below the 8 MiB session log cap, so two + * retained streams can never dominate a session record. + */ +export const CLI_OUTPUT_MAX_BYTES = 2 * 1024 * 1024; +/** Head room for the truncation marker, so a clipped stream stays under the ceiling. */ +const CLI_OUTPUT_MARKER_RESERVE = 128; +const CLI_OUTPUT_KEEP_BYTES = CLI_OUTPUT_MAX_BYTES - CLI_OUTPUT_MARKER_RESERVE; +/** + * Greppable opener of the marker a clipped stream carries in place of its head. + * + * The cut is stated in the data itself, so an operator (or a parser reading the + * retained text) sees that bytes are missing and how many, rather than inferring + * it from output that silently never arrived. Kept free of failure phrasings: + * `isExplicitFailure` scans raw stdout for real failure declarations. + */ +export const CLI_OUTPUT_TRUNCATION_MARKER = '[openswarm-cli-output:'; + +/** How much of the tail a cut keeps, so the next cut is a megabyte away (below). */ +const CLI_OUTPUT_TRIM_BYTES = Math.ceil(CLI_OUTPUT_KEEP_BYTES / 2); + +/** + * A stream's retained tail, its byte length, and the head bytes already dropped. + * The retained text carries no marker of its own: the drop count lives here so a + * second cut reports the total, not just the last one. + */ +interface RetainedCliOutput { + text: string; + bytes: number; + droppedBytes: number; +} + +/** + * Append a chunk, keeping only the TAIL once the ceiling is passed. + * + * The tail, not the head, because every consumer of a delegated CLI's output + * reads the terminal event: `extractResultFromStreamJson` (claude stream-json), + * `extractCodexMessageText` (`messages.at(-1)`), `extractCursorFinalText`, + * `detectRateLimit`, and `extractStreamJsonError` all need the LAST lines — a + * head clip would throw away the result event of exactly the verbose runs this + * bound exists for. Same direction as `tailWithinBytes` + * (agents/verificationEvidence.ts) and the tail-keeping `appendBounded` in + * automation/scheduler.ts. + * + * A cut keeps half the budget rather than exactly filling it: re-slicing on + * every chunk after the ceiling would copy 2 MiB per chunk, which a chatty CLI + * turns into gigabytes of memcpy while it floods the pipe. Each cut therefore + * buys a full megabyte of appends, and the retained text still never exceeds the + * ceiling — an append that would cross it is trimmed in the same call. + */ +function appendCliOutput(output: RetainedCliOutput, chunk: string): void { + const chunkBytes = Buffer.byteLength(chunk, 'utf8'); + const total = output.bytes + chunkBytes; + if (total <= CLI_OUTPUT_KEEP_BYTES) { + output.text += chunk; + output.bytes = total; + return; + } + // A single chunk can exceed what we keep on its own; then the retained text is + // going to be discarded entirely, and concatenating it first would be a copy + // of a payload already known to be thrown away. + const source = chunkBytes >= CLI_OUTPUT_TRIM_BYTES ? chunk : output.text + chunk; + // The cut can land mid-character, so the retained length is re-measured rather + // than assumed. Decoding turns at most a handful of stray UTF-8 bytes into + // replacement characters (3 bytes each), so the retained text can exceed + // CLI_OUTPUT_TRIM_BYTES by a few bytes — never by more than + // CLI_OUTPUT_MARKER_RESERVE, which is why the ceiling itself still holds. + output.text = Buffer.from(source, 'utf8') + .subarray(-CLI_OUTPUT_TRIM_BYTES) + .toString('utf8'); + output.bytes = Buffer.byteLength(output.text, 'utf8'); + output.droppedBytes += total - output.bytes; +} + +/** The retained output, preceded by the marker when the head was clipped. */ +function renderCliOutput(output: RetainedCliOutput): string { + return output.droppedBytes > 0 + ? `${CLI_OUTPUT_TRUNCATION_MARKER} ${output.droppedBytes} bytes omitted from the head]\n${output.text}` + : output.text; +} + /** * Spawn a CLI process using the given adapter and options. * Handles: temp file write, argv-safe spawn, timeout/SIGKILL, @@ -98,11 +188,24 @@ export async function spawnCli( // Below this line the adapter runs its own tool loop inside its own CLI, so // anything OpenSwarm assembles for *our* loop is dropped. Silence there is // how a configured MCP grant or an `ask_human` escape hatch turns into an - // agent that quietly never had it — say it out loud instead. - if (options.mcpTools?.length || options.coordinationContext) { + // agent that quietly never had it — say it out loud instead. The same goes for + // a role's `tools.allow`/`tools.deny`, and for `protectedFiles` / + // `forbidPublication`: this path cannot honor any of them (the CLI owns its + // tools), and a silently inert fence is worse than none. (AGT-4444) + if ( + options.mcpTools?.length + || options.coordinationContext + || options.toolAllow?.length + || options.toolDeny?.length + || options.protectedFiles?.length + || options.forbidPublication + ) { const dropped = [ options.mcpTools?.length ? `${options.mcpTools.length} MCP tool(s)` : '', options.coordinationContext ? 'coordination tools' : '', + options.toolAllow?.length || options.toolDeny?.length ? 'the role tool allow/deny list' : '', + options.protectedFiles?.length ? `${options.protectedFiles.length} protected path(s)` : '', + options.forbidPublication ? 'the publication fence' : '', ].filter(Boolean).join(' and '); console.warn( `[Adapter] '${adapter.name}' delegates to its own CLI tool loop; ${dropped} will not be available to this run. ` @@ -221,13 +324,15 @@ export async function spawnCli( }, proc); } - let stdout = ''; - let stderr = ''; + // Retained only for the final parse and the transcript; bounded per stream + // so a flooding CLI cannot hold the daemon's memory for its whole run. + const stdoutOutput: RetainedCliOutput = { text: '', bytes: 0, droppedBytes: 0 }; + const stderrOutput: RetainedCliOutput = { text: '', bytes: 0, droppedBytes: 0 }; let streamBuffer = ''; proc.stdout?.on('data', (data: Buffer) => { const text = data.toString(); - stdout += text; + appendCliOutput(stdoutOutput, text); if (options.onLog && adapter.capabilities.supportsStreaming) { streamBuffer = adapter.parseStreamingChunk ? adapter.parseStreamingChunk(text, options.onLog, streamBuffer) @@ -236,7 +341,7 @@ export async function spawnCli( }); proc.stderr?.on('data', (data: Buffer) => { - stderr += data.toString(); + appendCliOutput(stderrOutput, data.toString()); }); let exitDrainTimer: NodeJS.Timeout | null = null; @@ -252,7 +357,7 @@ export async function spawnCli( cleanupLifecycle(); terminateCliProcessTree(proc); const reason = lifecycleController.signal.reason; - session?.record({ type: 'assistant', rawStdout: stdout, rawStderr: stderr }); + session?.record({ type: 'assistant', rawStdout: renderCliOutput(stdoutOutput), rawStderr: renderCliOutput(stderrOutput) }); session?.close({ outcome: 'aborted', durationMs: Date.now() - startTime, error: reason instanceof Error ? `${reason.name}: ${reason.message}` : String(reason), @@ -272,6 +377,11 @@ export async function spawnCli( : parseCliStreamChunk('\n', options.onLog, streamBuffer); } + // Rendered once per settling path: the marker belongs in the transcript + // and in the returned result, but the retained text itself carries none. + const stdout = renderCliOutput(stdoutOutput); + const stderr = renderCliOutput(stderrOutput); + session?.record({ type: 'assistant', rawStdout: stdout, rawStderr: stderr }); session?.close({ outcome: code === 0 || code === null ? 'returned' : 'exit_nonzero', diff --git a/src/adapters/chatStream.test.ts b/src/adapters/chatStream.test.ts index c1913e53..3520de3c 100644 --- a/src/adapters/chatStream.test.ts +++ b/src/adapters/chatStream.test.ts @@ -1,5 +1,46 @@ import { describe, it, expect, vi } from 'vitest'; -import { reduceChatChunks } from './chatStream.js'; +import { consumeChatCompletionsStream, reduceChatChunks } from './chatStream.js'; + +const encoder = new TextEncoder(); + +/** An SSE body delivered as exactly the given chunks. */ +const sseBody = (chunks: string[]) => + new Response( + new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(encoder.encode(chunk)); + controller.close(); + }, + }), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ); + +/** + * A lazy byte flood: `chunks` copies of one buffer, pulled on demand. Lazy so a + * test can describe an over-cap stream without allocating it up front, and so a + * bounded reader that cancels stops the source instead of draining it. The + * returned `pulls` reads -1 once the reader cancelled the body. + */ +const floodResponse = (chunk: Uint8Array, chunks: number) => { + let pulls = 0; + const res = new Response( + new ReadableStream({ + pull(controller) { + if (pulls >= chunks) { + controller.close(); + return; + } + pulls += 1; + controller.enqueue(chunk); + }, + cancel() { + pulls = -1; + }, + }), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ); + return { res, pulls: () => pulls }; +}; describe('reduceChatChunks', () => { it('keeps the model the server reports serving', () => { @@ -94,3 +135,81 @@ describe('reduceChatChunks usage accounting (AGT-4178)', () => { expect(out.usage && 'cost' in out.usage).toBe(false); }); }); + +describe('consumeChatCompletionsStream bounds (AGT-3429)', () => { + it('aborts a non-newline flood instead of buffering it forever', async () => { + // Frames are split on '\n', so newline-free bytes are carried until one + // arrives. Without a ceiling this stream just accumulates — 2 MiB of 'x' in + // 64 KiB chunks — and the call only ends when the server gives up. Real SSE + // frames for this adapter are a few hundred bytes. + const { res } = floodResponse(encoder.encode('x'.repeat(64 * 1024)), 32); + + await expect(consumeChatCompletionsStream(res)).rejects.toThrow( + /partial frame exceeded the 1 MiB limit/, + ); + }); + + it('caps total bytes read even when every frame is well formed', async () => { + // Every line parses, so the carry never trips — only the raw ceiling can stop + // an endpoint that streams valid-but-endless frames (a runaway server loop). + const frame = `data: {"choices":[{"delta":{"content":"${'x'.repeat(4_096)}"}}]}\n`; + const burst = encoder.encode(frame.repeat(8)); + const { res } = floodResponse(burst, Math.ceil((17 * 1024 * 1024) / burst.byteLength)); + + await expect(consumeChatCompletionsStream(res)).rejects.toThrow(/exceeded the 16 MiB limit/); + }); + + it('cancels the body when it gives up, so the source stops producing', async () => { + // The bound must stop the endpoint, not just stop storing what it sends — + // otherwise a flooding server keeps a connection busy draining bytes. + const { res, pulls } = floodResponse(encoder.encode('y'.repeat(64 * 1024)), 32); + + await expect(consumeChatCompletionsStream(res)).rejects.toThrow(/1 MiB limit/); + expect(pulls()).toBe(-1); + }); + + it('does not mistake one chunk carrying many complete frames for a partial frame', async () => { + // The carry ceiling must be applied AFTER the split: a burst of complete + // frames arriving in one chunk can legitimately exceed 1 MiB in total (the + // whole point of the raw ceiling being 16 MiB), while the carry is empty. + const frame = `data: {"choices":[{"delta":{"content":"${'q'.repeat(1_024)}"}}]}\n`; + const frames = 2_048; // 2 MiB of complete frames in a single chunk + const res = sseBody([`${frame.repeat(frames)}data: [DONE]\n`]); + + const out = await consumeChatCompletionsStream(res); + + expect(out.choices[0].message.content).toBe('q'.repeat(1_024 * frames)); + }); + + it('parses a normal multi-frame stream end to end', async () => { + // The bound must not disturb the ordinary case: content deltas streamed in + // several frames, a tool call split across frames, usage last — and a single + // frame as large as a real tool-call fragment (256 KiB) staying well inside + // the carry ceiling. + const onToken = vi.fn(); + const bigFrame = `data: ${JSON.stringify({ + choices: [{ delta: { content: 'z'.repeat(256 * 1024) } }], + })}\n`; + const res = sseBody([ + 'data: {"model":"deepseek-v4.1-flash","choices":[{"delta":{"content":"Hel"}}]}\n', + 'data: {"choices":[{"delta":{"content":"lo"}}]}\n', + bigFrame, + 'data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"edit_file","arguments":"{\\"p"}}]}}]}\n', + 'data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"ath\\":\\"x\\"}"}}]}}]}\n', + 'data: {"choices":[{"delta":{},"finish_reason":"tool_calls"}]}\n', + 'data: {"choices":[],"usage":{"prompt_tokens":7,"completion_tokens":3,"total_tokens":10}}\n', + 'data: [DONE]\n', + ]); + + const out = await consumeChatCompletionsStream(res, onToken); + + expect(out.choices[0].message.content).toBe(`Hello${'z'.repeat(256 * 1024)}`); + expect(out.choices[0].finish_reason).toBe('tool_calls'); + expect(out.choices[0].message.tool_calls).toEqual([ + { id: 'call_1', type: 'function', function: { name: 'edit_file', arguments: '{"path":"x"}' } }, + ]); + expect(out.usage).toEqual({ prompt_tokens: 7, completion_tokens: 3, total_tokens: 10 }); + expect(out.model).toBe('deepseek-v4.1-flash'); + expect(onToken.mock.calls.map((c) => c[0])).toEqual(['Hel', 'lo', 'z'.repeat(256 * 1024)]); + }); +}); diff --git a/src/adapters/chatStream.ts b/src/adapters/chatStream.ts index 5f9ab688..8b2526ae 100644 --- a/src/adapters/chatStream.ts +++ b/src/adapters/chatStream.ts @@ -152,6 +152,25 @@ function parseChunkLine(line: string): StreamChunk | null { } } +/** + * Byte ceilings on ONE streamed chat/completions call. + * + * Frames are split on '\n', so bytes from an endpoint that never sends one pile + * up in the carry buffer for as long as it keeps writing — the request deadline + * is the only other thing that would stop it, and that is minutes away. 1 MiB + * matches the partial-line bound this repo already keeps for the same reason + * (streamBuffer.ts MAX_LINE_BUFFER_BYTES); a real chunk here is a few hundred + * bytes, so nothing legitimate approaches it. + * + * The raw ceiling bounds bytes read from a source we do not control, the same + * class as the 2 MiB web-fetch body cap (webTools.ts) and the 16 MiB + * GIT_OUTPUT_LIMIT in applyPatch.ts. It is deliberately the same 16 MiB the + * Responses adapter applies, so the two streaming parsers this repo runs against + * model endpoints share one ceiling: below the model's own output window, and + * far above any real answer. Both limits are stated in the error. + */ +const MAX_CHAT_STREAM_BYTES = 16 * 1024 * 1024; +const MAX_CHAT_PARTIAL_FRAME_BYTES = 1024 * 1024; /** Read a chat/completions SSE body and reduce it, emitting content deltas live. */ export async function consumeChatCompletionsStream( res: Response, @@ -165,6 +184,7 @@ export async function consumeChatCompletionsStream( const chunks: StreamChunk[] = []; const decoder = new TextDecoder(); let buffer = ''; + let readBytes = 0; const handle = (c: StreamChunk | null) => { if (!c) return; const delta = c.choices?.[0]?.delta?.content; @@ -175,9 +195,21 @@ export async function consumeChatCompletionsStream( const { done, value } = await reader.read(); if (done) break; onBytes?.(); + readBytes += value.byteLength; + if (readBytes > MAX_CHAT_STREAM_BYTES) { + await reader.cancel('chat stream exceeded retention limit').catch(() => {}); + throw new Error(`chat stream exceeded the 16 MiB limit (${readBytes} bytes read)`); + } buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); + // The carry is what a non-newline flood grows, so the ceiling is applied to + // the carry rather than to the decoded chunk: a burst of many complete + // frames in one chunk is not a partial frame and must not trip this bound. buffer = lines.pop() ?? ''; + if (Buffer.byteLength(buffer, 'utf8') > MAX_CHAT_PARTIAL_FRAME_BYTES) { + await reader.cancel('chat stream partial frame exceeded limit').catch(() => {}); + throw new Error('chat stream partial frame exceeded the 1 MiB limit (no newline in the stream)'); + } for (const line of lines) handle(parseChunkLine(line)); } handle(parseChunkLine(buffer)); diff --git a/src/adapters/codexResponses.test.ts b/src/adapters/codexResponses.test.ts index e0aea2ff..cd55f888 100644 --- a/src/adapters/codexResponses.test.ts +++ b/src/adapters/codexResponses.test.ts @@ -694,3 +694,153 @@ describe('401 refresh scope', () => { expect(responsesCalls).toBe(2); }); }); + +describe('Responses stream byte bounds (AGT-3429)', () => { + type CreateApiCaller = ( + initialToken: string, + accountId: string, + store: unknown, + model: string, + onToken?: (delta: string) => void, + signal?: AbortSignal, + onReasoning?: (line: string) => void, + ) => (messages: ChatMessage[], tools: ToolDefinition[]) => Promise; + + interface ChatLikeResponseShape { + choices: Array<{ message: { content: string | null; tool_calls?: unknown[] }; finish_reason: string }>; + usage?: unknown; + } + + const encoder = new TextEncoder(); + const caller = (onToken?: (delta: string) => void, onReasoning?: (line: string) => void) => { + const adapter = new CodexResponsesAdapter() as unknown as { createApiCaller: CreateApiCaller }; + return adapter.createApiCaller('token', 'account', {}, 'gpt-5.6-terra', onToken, undefined, onReasoning); + }; + + /** + * A lazy byte flood: `chunks` copies of one buffer, pulled on demand. Lazy so + * a test can describe an over-cap stream without allocating it up front, and + * so a reader that gives up stops the source instead of draining it. + */ + const floodResponse = (chunk: Uint8Array, chunks: number) => { + let pulls = 0; + const res = new Response( + new ReadableStream({ + pull(controller) { + if (pulls >= chunks) { + controller.close(); + return; + } + pulls += 1; + controller.enqueue(chunk); + }, + cancel() { + pulls = -1; + }, + }), + { status: 200, headers: { 'Content-Type': 'text/event-stream' } }, + ); + return { res, pulls: () => pulls }; + }; + + it('aborts a non-newline flood instead of carrying it forever', async () => { + // The carry buffer grows with every byte until a '\n' arrives. An endpoint + // that never sends one (malformed, or adversarial) used to accumulate the + // whole body — the request deadline is minutes away. 2 MiB of 'x' in 64 KiB + // chunks; real frames here are a few dozen bytes. + const { res } = floodResponse(encoder.encode('x'.repeat(64 * 1024)), 32); + vi.stubGlobal('fetch', vi.fn(async () => res)); + + await expect(caller()([{ role: 'user', content: 'hi' }], [])).rejects.toThrow( + /partial frame exceeded the 1 MiB limit/, + ); + }); + + it('cancels the body when it gives up, so the endpoint stops producing', async () => { + const { res, pulls } = floodResponse(encoder.encode('z'.repeat(64 * 1024)), 32); + vi.stubGlobal('fetch', vi.fn(async () => res)); + + await expect(caller()([{ role: 'user', content: 'hi' }], [])).rejects.toThrow(/1 MiB limit/); + expect(pulls()).toBe(-1); + }); + + it('caps total bytes read even when every frame is well formed', async () => { + // Valid frames never trip the carry bound, so only the raw ceiling stops an + // endpoint that loops on well-formed events instead of ending the stream. + const burst = encoder.encode( + `${'data: {"type":"response.output_text.delta","delta":"ok"}\n'.repeat(32_768)}`, + ); + const { res } = floodResponse(burst, Math.ceil((17 * 1024 * 1024) / burst.byteLength)); + vi.stubGlobal('fetch', vi.fn(async () => res)); + + await expect(caller()([{ role: 'user', content: 'hi' }], [])).rejects.toThrow( + /exceeded the 16 MiB limit/, + ); + }); + + it('parses a normal multi-frame stream through the streaming path', async () => { + // The bound and the incremental reducer must leave ordinary streams alone: + // content deltas, a tool call assembled from fragments, reasoning summary + // lines, and usage last — plus one frame as large as a real long content + // delta (256 KiB) staying well inside the carry ceiling. + const bigFrame = `data: ${JSON.stringify({ + type: 'response.output_text.delta', + delta: 'z'.repeat(256 * 1024), + })}`; + const body = [ + 'data: {"type":"response.created"}', + 'data: {"type":"response.output_text.delta","delta":"Hel"}', + 'data: {"type":"response.output_text.delta","delta":"lo"}', + bigFrame, + 'data: {"type":"response.reasoning_summary_text.delta","delta":"weighing options\\n"}', + 'data: {"type":"response.output_item.added","item":{"type":"function_call","id":"fc_1","call_id":"call_9","name":"edit_file"}}', + 'data: {"type":"response.function_call_arguments.delta","item_id":"fc_1","delta":"{\\"path\\":"}', + 'data: {"type":"response.function_call_arguments.delta","item_id":"fc_1","delta":"\\"x.ts\\"}"}', + 'data: {"type":"response.completed","response":{"usage":{"input_tokens":11,"output_tokens":4}}}', + 'data: [DONE]', + '', + ].join('\n'); + vi.stubGlobal('fetch', vi.fn(async () => new Response(body, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }))); + + const tokens: string[] = []; + const thoughts: string[] = []; + const out = await caller((d) => tokens.push(d), (l) => thoughts.push(l))( + [{ role: 'user', content: 'hi' }], + [], + ); + + expect(out.choices[0].message.content).toBe(`Hello${'z'.repeat(256 * 1024)}`); + expect(out.choices[0].finish_reason).toBe('tool_calls'); + expect(out.choices[0].message.tool_calls).toEqual([ + { id: 'call_9', type: 'function', function: { name: 'edit_file', arguments: '{"path":"x.ts"}' } }, + ]); + expect(out.usage).toEqual({ prompt_tokens: 11, completion_tokens: 4, total_tokens: 15, cached_tokens: 0 }); + expect(tokens).toEqual(['Hel', 'lo', 'z'.repeat(256 * 1024)]); + expect(thoughts).toEqual(['weighing options']); + }); + + it('folds a long stream incrementally, matching the batch reduce exactly', async () => { + // Thousands of frames through the split/carry loop, which the handful-of- + // frames test above does not exercise: the carry must survive every chunk + // boundary, and the incremental reducer must produce exactly what the batch + // form does for the same sequence. (~8.5 MiB of body, under the 16 MiB cap.) + const events = Array.from({ length: 8_000 }, (_, i) => ({ + type: 'response.output_text.delta', + delta: `chunk-${i}-${'y'.repeat(1_000)}`, + })); + const body = `${events.map((e) => `data: ${JSON.stringify(e)}`).join('\n')}\ndata: [DONE]\n`; + vi.stubGlobal('fetch', vi.fn(async () => new Response(body, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }))); + + const streamed = await caller()([{ role: 'user', content: 'hi' }], []); + const batched = reduceResponsesEvents(events); + + expect(streamed).toEqual(batched); + expect(streamed.choices[0].message.content).toBe(events.map((e) => e.delta).join('')); + }); +}); diff --git a/src/adapters/codexResponses.ts b/src/adapters/codexResponses.ts index e9584819..d2796fcf 100644 --- a/src/adapters/codexResponses.ts +++ b/src/adapters/codexResponses.ts @@ -36,6 +36,28 @@ export const DEFAULT_MODEL = 'gpt-5.6-terra'; const PROFILE_KEY = 'openai-gpt:default'; const SPARK_MODEL = 'gpt-5.3-codex-spark'; +/** + * Byte ceilings on ONE streamed Responses call. + * + * The partial-frame ceiling is the one a hostile endpoint drives: frames are + * split on '\n', so bytes that never contain a newline pile up in the carry + * buffer for as long as the server keeps sending them, and the request deadline + * only fires when it fires. 1 MiB matches the partial-line bound this repo + * already keeps for the same reason (streamBuffer.ts MAX_LINE_BUFFER_BYTES); a + * real frame — `response.output_text.delta` or a tool-argument fragment — is a + * few dozen bytes, so nothing legitimate comes close. + * + * The raw ceiling bounds bytes read from a source we do not control, the same + * class as the 2 MiB web-fetch body cap and the 16 MiB GIT_OUTPUT_LIMIT in + * applyPatch.ts. It is per API call, and it stays a ceiling on the STREAM rather + * than on retained state: the reducer below keeps only accumulated text, + * tool-call state and final usage, so >16 MiB in one call means an endpoint that + * will not stop, not a long answer (the model's window caps output far below + * that). Both limits are stated in the error, so an operator sees which one hit. + */ +const MAX_STREAM_BYTES = 16 * 1024 * 1024; +const MAX_PARTIAL_FRAME_BYTES = 1024 * 1024; + // ---- Responses API wire types (the subset we send/receive) ---- interface ResponsesTool { @@ -148,8 +170,28 @@ interface SseEvent { /** * Reduce parsed Responses SSE events → a chat-completions-shaped response. * Exported so the SSE→chat mapping is unit-testable without a live stream. + * + * A batch is just the incremental reducer fed up front: the live path + * (`consumeResponsesStream`) must NOT retain the event array, because a long + * stream duplicates the whole raw body again as a parsed object graph plus one + * repeated delta string per token. Only accumulated text, tool-call state and + * the final usage survive a stream here. */ export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse { + const reducer = createResponsesEventReducer(); + for (const ev of events) reducer.accept(ev); + return reducer.finish(); +} + +/** + * The streaming form of `reduceResponsesEvents`: one event in, no history kept. + * `accept` folds an event into the retained state; `finish` renders the + * chat-completions shape the agentic loop consumes. + */ +function createResponsesEventReducer(): { + accept: (event: SseEvent) => void; + finish: () => ChatLikeResponse; +} { let text = ''; // Keyed by the streaming item id; the emitted tool-call id is the call_id so it // round-trips back as `function_call_output.call_id` on the next turn. @@ -157,7 +199,7 @@ export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse { let usage: ChatLikeResponse['usage']; const getOnlyCall = () => calls.size === 1 ? calls.values().next().value : undefined; - for (const ev of events) { + const accept = (ev: SseEvent) => { switch (ev.type) { case 'response.output_text.delta': if (ev.delta) text += ev.delta; @@ -197,26 +239,31 @@ export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse { break; } } - } - - const toolCalls: ApiToolCallShape[] = [...calls.values()].map((c) => ({ - id: c.callId, - type: 'function', - function: { name: c.name, arguments: c.args }, - })); + }; return { - choices: [ - { - message: { - role: 'assistant', - content: text || null, - tool_calls: toolCalls.length > 0 ? toolCalls : undefined, - }, - finish_reason: toolCalls.length > 0 ? 'tool_calls' : 'stop', - }, - ], - usage, + accept, + finish: () => { + const toolCalls: ApiToolCallShape[] = [...calls.values()].map((c) => ({ + id: c.callId, + type: 'function', + function: { name: c.name, arguments: c.args }, + })); + + return { + choices: [ + { + message: { + role: 'assistant', + content: text || null, + tool_calls: toolCalls.length > 0 ? toolCalls : undefined, + }, + finish_reason: toolCalls.length > 0 ? 'tool_calls' : 'stop', + }, + ], + usage, + }; + }, }; } @@ -243,12 +290,13 @@ async function consumeResponsesStream( onToken?: (delta: string) => void, onReasoning?: (line: string) => void, ): Promise { - const events: SseEvent[] = []; + const reducer = createResponsesEventReducer(); const reader = res.body?.getReader(); if (!reader) throw new Error('Codex responses: empty stream body'); const decoder = new TextDecoder(); let buffer = ''; + let readBytes = 0; let terminalError: string | undefined; // Reasoning summary streams token-by-token; buffer and emit whole lines so the // live log shows readable thoughts instead of one-word-per-line spam. @@ -265,7 +313,7 @@ async function consumeResponsesStream( }; const handle = (ev: SseEvent | null) => { if (!ev) return; - events.push(ev); + reducer.accept(ev); if (ev.type === 'response.incomplete') { terminalError = `Responses stream incomplete${ev.response?.incomplete_details?.reason ? `: ${ev.response.incomplete_details.reason}` : ''}`; } else if (ev.type === 'response.failed') { @@ -284,9 +332,21 @@ async function consumeResponsesStream( for (;;) { const { done, value } = await reader.read(); if (done) break; + readBytes += value.byteLength; + if (readBytes > MAX_STREAM_BYTES) { + await reader.cancel('responses stream exceeded retention limit').catch(() => {}); + throw new Error(`Codex response stream exceeded the 16 MiB limit (${readBytes} bytes read)`); + } buffer += decoder.decode(value, { stream: true }); const lines = buffer.split('\n'); + // The carry is what a non-newline flood grows, so the ceiling is applied to + // the carry rather than to the whole decoded chunk — a legitimate burst of + // many complete frames in one chunk must not trip a partial-frame bound. buffer = lines.pop() ?? ''; + if (Buffer.byteLength(buffer, 'utf8') > MAX_PARTIAL_FRAME_BYTES) { + await reader.cancel('responses partial frame exceeded limit').catch(() => {}); + throw new Error('Codex response partial frame exceeded the 1 MiB limit (no newline in the stream)'); + } for (const line of lines) handle(parseSseLine(line)); } handle(parseSseLine(buffer)); @@ -294,7 +354,7 @@ async function consumeResponsesStream( if (terminalError) throw new Error(terminalError); - return reduceResponsesEvents(events); + return reducer.finish(); } // ---- Adapter ---- @@ -417,6 +477,8 @@ export class CodexResponsesAdapter implements CliAdapter { webTools: options.webTools, memoryTools: options.memoryTools, shellTools: options.shellTools, + toolAllow: options.toolAllow, + toolDeny: options.toolDeny, filesystemTools: options.filesystemTools, diagnosticsTool: options.diagnosticsTool, mcpTools: options.mcpTools, diff --git a/src/adapters/gpt.ts b/src/adapters/gpt.ts index 165928f5..c8c891ce 100644 --- a/src/adapters/gpt.ts +++ b/src/adapters/gpt.ts @@ -164,6 +164,8 @@ export class GptCliAdapter implements CliAdapter { webTools: options.webTools, memoryTools: options.memoryTools, shellTools: options.shellTools, + toolAllow: options.toolAllow, + toolDeny: options.toolDeny, filesystemTools: options.filesystemTools, diagnosticsTool: options.diagnosticsTool, readOnly: options.readOnly, diff --git a/src/adapters/local.ts b/src/adapters/local.ts index 83d7388e..160c8215 100644 --- a/src/adapters/local.ts +++ b/src/adapters/local.ts @@ -202,6 +202,8 @@ export class LocalModelAdapter implements CliAdapter { webTools: options.webTools, memoryTools: options.memoryTools, shellTools: options.shellTools, + toolAllow: options.toolAllow, + toolDeny: options.toolDeny, filesystemTools: options.filesystemTools, diagnosticsTool: options.diagnosticsTool, readOnly: options.readOnly, diff --git a/src/adapters/modelCompat.ts b/src/adapters/modelCompat.ts index 377e99e8..bcffc317 100644 --- a/src/adapters/modelCompat.ts +++ b/src/adapters/modelCompat.ts @@ -33,7 +33,7 @@ const CLAUDE_ALIASES = new Set(['sonnet', 'opus', 'haiku']); * without a home here is a compile error, not a silent fall to the bulk model. */ export type ModelRole = - | 'worker' | 'reviewer' | 'tester' | 'documenter' | 'auditor' | 'skill-documenter' + | 'worker' | 'reviewer' | 'advisor' | 'tester' | 'documenter' | 'auditor' | 'skill-documenter' | 'planner' | 'orchestrator' | 'escalate'; // Verified against `cursor-agent --list-models` (2026.09.08, vela, 2026-09-10): @@ -45,6 +45,13 @@ const CURSOR_JUDGE_MODEL = 'cursor-grok-4.6-high'; // `reviewer` because a tier that resolves to the same model as the tier it // escalates FROM is a log line claiming work that did not happen. const CURSOR_ESCALATE_MODEL = 'cursor-grok-4.6-xhigh'; +// The advisor must not resolve to the reviewer's model either — same reason as +// the escalation above, one tier over: an advisor on `cursor-grok-4.6-high` +// would be the reviewer's own weights asked the same question twice. cursor's +// catalogue carries exactly two judge-grade ids, so the advisor takes the other +// one (`-xhigh`); the collision with escalate is the catalogue's limit, not a +// routing decision, and the two never run as the same role. +const CURSOR_ADVISOR_MODEL = CURSOR_ESCALATE_MODEL; const ADAPTER_DEFAULT_MODEL: Partial> = { 'codex-responses': 'gpt-5.6-terra', @@ -98,6 +105,7 @@ const CURSOR_ROLE_MODEL: Readonly> = { documenter: CURSOR_BULK_MODEL, 'skill-documenter': CURSOR_BULK_MODEL, reviewer: CURSOR_JUDGE_MODEL, + advisor: CURSOR_ADVISOR_MODEL, auditor: CURSOR_JUDGE_MODEL, planner: CURSOR_JUDGE_MODEL, orchestrator: CURSOR_JUDGE_MODEL, diff --git a/src/adapters/ollamaCloud.ts b/src/adapters/ollamaCloud.ts index b82f7877..111c522b 100644 --- a/src/adapters/ollamaCloud.ts +++ b/src/adapters/ollamaCloud.ts @@ -406,6 +406,8 @@ export class OllamaCloudAdapter extends LocalModelAdapter implements CliAdapter webTools: options.webTools, memoryTools: options.memoryTools, shellTools: options.shellTools, + toolAllow: options.toolAllow, + toolDeny: options.toolDeny, filesystemTools: options.filesystemTools, diagnosticsTool: options.diagnosticsTool, readOnly: options.readOnly, diff --git a/src/adapters/openrouter.ts b/src/adapters/openrouter.ts index 3afb415d..c915f08f 100644 --- a/src/adapters/openrouter.ts +++ b/src/adapters/openrouter.ts @@ -198,6 +198,8 @@ export class OpenRouterCliAdapter implements CliAdapter { webTools: options.webTools, memoryTools: options.memoryTools, shellTools: options.shellTools, + toolAllow: options.toolAllow, + toolDeny: options.toolDeny, filesystemTools: options.filesystemTools, diagnosticsTool: options.diagnosticsTool, readOnly: options.readOnly, diff --git a/src/adapters/rateLimitError.test.ts b/src/adapters/rateLimitError.test.ts index afd9563a..b7f19e69 100644 --- a/src/adapters/rateLimitError.test.ts +++ b/src/adapters/rateLimitError.test.ts @@ -347,6 +347,25 @@ describe('classifyLimitResponse — spent quota vs short-window throttle (INT-29 expect(classifyLimitResponse(new Headers({ 'retry-after': '7' }), '').retryAfterSeconds).toBe(7); expect(classifyLimitResponse(new Headers(), '').retryAfterSeconds).toBeUndefined(); }); + + // RFC 9110 allows delta-seconds AND an HTTP-date. parseInt alone turned the + // dated form into NaN, so the throttle path discarded the provider's wait and + // retried on the short local backoff before the window had cleared. (AGT-3442) + it('parses an HTTP-date Retry-After into seconds-from-now, and that wait drives the backoff', () => { + const c = classifyLimitResponse(new Headers({ 'retry-after': new Date(Date.now() + 30_000).toUTCString() }), ''); + expect(Number.isFinite(c.retryAfterSeconds)).toBe(true); + expect(c.retryAfterSeconds).toBeGreaterThanOrEqual(28); + expect(c.retryAfterSeconds).toBeLessThanOrEqual(30); + // The honored wait, not the 5-6s local fallback throttleWaitMs would pick. + expect(throttleWaitMs(0, c.retryAfterSeconds)).toBeGreaterThanOrEqual(28_000); + }); + + it('429 with an HTTP-date Retry-After yields the parsed reset time, not NaN', () => { + const err = rateLimitFromHttpResponse(429, new Headers({ 'retry-after': new Date(Date.now() + 45_000).toUTCString() }), ''); + const nowSec = Math.floor(Date.now() / 1000); + expect(err?.resetsAt).toBeGreaterThanOrEqual(nowSec + 43); + expect(err?.resetsAt).toBeLessThanOrEqual(nowSec + 45); + }); }); describe('resolveLimitResponse gating (INT-2907)', () => { diff --git a/src/adapters/rateLimitError.ts b/src/adapters/rateLimitError.ts index eb1b27d2..684acf4c 100644 --- a/src/adapters/rateLimitError.ts +++ b/src/adapters/rateLimitError.ts @@ -110,24 +110,40 @@ export function parseResetsAtFromBody(text: string): number | undefined { return m ? parseInt(m[1], 10) : undefined; } +/** + * Parse an RFC 9110 `Retry-After` value into seconds-from-now. The header has TWO + * legal forms — delta-seconds ("120") and an HTTP-date ("Wed, 21 Oct 2015 + * 07:28:00 GMT") — and parseInt alone turns the dated form into NaN, which + * discarded the wait the provider asked for and fell through to the short local + * backoff (a retry before the window had cleared). Every caller shares this + * parse so both forms are honored identically. (AGT-3442) + */ +function parseRetryAfterSeconds(value: string | null | undefined): number | undefined { + if (value == null) return undefined; + const trimmed = value.trim(); + if (!trimmed) return undefined; + // delta-seconds — the common form, and the only one parseInt may see. + if (/^\d+$/.test(trimmed)) return parseInt(trimmed, 10); + // HTTP-date (RFC 1123). A past date means "retry now": 0, which the caller's + // backoff turns into a bounded wait rather than an instant retry storm. + const at = Date.parse(trimmed); + return Number.isFinite(at) ? Math.max(0, Math.round((at - Date.now()) / 1000)) : undefined; +} + /** Pull a unix reset timestamp (seconds) out of headers or a JSON body, if present. */ function extractResetsAt(headers: Headers | undefined, body: string): number | undefined { - const fromHeader = (k: string): number | undefined => { - const v = headers?.get(k); - const n = v == null ? NaN : parseInt(v, 10); - return Number.isFinite(n) ? n : undefined; - }; // Only headers/fields that are genuinely UNIX-epoch seconds or seconds-from-now: // - x-codex-primary-reset-at: epoch seconds - // - Retry-After: seconds-from-now (→ convert to epoch) + // - Retry-After: seconds-from-now (→ convert to epoch; delta-seconds OR HTTP-date) // - body "resets_at": epoch seconds // Deliberately NOT x-ratelimit-reset-requests/-tokens: OpenAI returns those as // DURATION strings ("1s", "6ms", "2m59s"), not epoch — parseInt would yield a // 1970 timestamp and defeat the pause. Omitting them falls back to the safe // 60s default, which is correct rather than wrong. (INT-2520 review) - const codexReset = fromHeader('x-codex-primary-reset-at'); - if (codexReset != null) return codexReset; - const retryAfter = fromHeader('retry-after'); + const codexResetRaw = headers?.get('x-codex-primary-reset-at'); + const codexReset = codexResetRaw == null ? NaN : parseInt(codexResetRaw, 10); + if (Number.isFinite(codexReset)) return codexReset; + const retryAfter = parseRetryAfterSeconds(headers?.get('retry-after')); if (retryAfter != null) return Math.floor(Date.now() / 1000) + retryAfter; return parseResetsAtFromBody(body); } @@ -201,7 +217,10 @@ export function classifyLimitResponse(headers: Headers | undefined, body: string return Number.isFinite(n) ? n : undefined; }; const usedPercent = num('x-codex-primary-used-percent'); - const retryAfterSeconds = num('retry-after'); + // Delta-seconds AND HTTP-date: parseInt alone left a dated Retry-After as NaN, + // so the throttle path discarded the provider's wait and retried on the short + // local backoff before the window cleared. (AGT-3442) + const retryAfterSeconds = parseRetryAfterSeconds(headers?.get('retry-after')); const lower = body.toLowerCase(); const quota = QUOTA_EXHAUSTED_SUBSTRINGS.some((s) => lower.includes(s)) || diff --git a/src/adapters/resultParsing.test.ts b/src/adapters/resultParsing.test.ts index 3bfe09b0..5ce304d0 100644 --- a/src/adapters/resultParsing.test.ts +++ b/src/adapters/resultParsing.test.ts @@ -1,5 +1,5 @@ import { describe, it, expect } from 'vitest'; -import { parseReviewerResult, parseWorkerResult } from './resultParsing.js'; +import { findStringAwareJsonObject, parseReviewerResult, parseWorkerResult } from './resultParsing.js'; import { t } from '../locale/index.js'; const wrap = (obj: unknown) => '```json\n' + JSON.stringify(obj) + '\n```'; @@ -321,3 +321,54 @@ describe('a limitation report is not a failure declaration (AGT-4534)', () => { expect(isExplicitFailure('Failed to apply the patch.')).toBe(true); }); }); + +describe('findStringAwareJsonObject (AGT-3466)', () => { + const find = (text: string) => findStringAwareJsonObject(text, '"success"'); + + it('does not end the object at a brace inside a quoted string', () => { + // The unscanned variant sliced to the `}` inside the summary, so + // JSON.parse failed on an unterminated string and every field was lost. + expect(find('{"success": true, "summary": "Use `{}` here"}')) + .toBe('{"success": true, "summary": "Use `{}` here"}'); + expect(find('{"success": true, "summary": "trailing } here"}')) + .toBe('{"success": true, "summary": "trailing } here"}'); + expect(find('{"success": true, "summary": "the { never closes"}')) + .toBe('{"success": true, "summary": "the { never closes"}'); + }); + + it('honors backslash escapes when deciding string boundaries', () => { + // Without escape state the quote inside \"}\" closes the string early, and + // the brace that follows is then read as structure. + expect(find('{"success": true, "summary": "say \\"}\\" here"}')) + .toBe('{"success": true, "summary": "say \\"}\\" here"}'); + // An escaped backslash is a literal, so the next quote really does close. + expect(find('{"success": true, "summary": "path C:\\\\ then } here"}')) + .toBe('{"success": true, "summary": "path C:\\\\ then } here"}'); + }); + + it('returns the object enclosing the marker, not a nested object before it', () => { + // lastIndexOf('{') would settle on the inner object and hand back a + // fragment with no `success` field in it. + expect(find('{"a": {"b": 1}, "success": true}')).toBe('{"a": {"b": 1}, "success": true}'); + }); + + it('skips prose that names the marker before the object appears', () => { + expect(find('The "success" flag was set. Result: {"success":true,"summary":"done"}')) + .toBe('{"success":true,"summary":"done"}'); + }); + + it('skips stray braces in the prose before the object', () => { + // A `{` in prose opens a candidate that never balances; a `}` closes one + // that never opened. Neither may abort the scan or become the result. + expect(find('noise } then {"success":true}')).toBe('{"success":true}'); + expect(find('noise { then {"success":true}')).toBe('{"success":true}'); + expect(find('oops { {"success":true}')).toBe('{"success":true}'); + expect(find('sibling {"x":1} then {"success":true}')).toBe('{"success":true}'); + }); + + it('returns null when there is no enclosing object', () => { + expect(find('no json at all')).toBeNull(); + expect(find('{"success": true')).toBeNull(); + expect(find('"success" with no object')).toBeNull(); + }); +}); diff --git a/src/adapters/resultParsing.ts b/src/adapters/resultParsing.ts index cd6344f5..630e3b2f 100644 --- a/src/adapters/resultParsing.ts +++ b/src/adapters/resultParsing.ts @@ -166,7 +166,13 @@ function extractBulletsAfter(text: string, heading: RegExp): string[] { return items; } -/** Brace-balanced scan for the JSON object containing `marker`. */ +/** + * Brace-balanced scan for the JSON object containing `marker`, counting braces + * blindly. Kept as the worker/reviewer adapters' path (their output is a + * structured completion, not prose); `findStringAwareJsonObject` below is the + * variant for result JSON that carries free text. Fixing brace handling in one + * is not fixing the other — check both call sites. (AGT-3466) + */ function findJsonObject(text: string, marker: string): string | null { const idx = text.indexOf(marker); if (idx < 0) return null; @@ -187,6 +193,66 @@ function findJsonObject(text: string, marker: string): string | null { return null; } +/** + * String-aware brace-balanced scan for the top-level JSON object containing + * `marker` — the variant for callers whose JSON carries prose. Counting braces + * without tracking quoted strings reads a brace in a value as structure: + * `"summary": "Use `{}` here"` sliced to the brace inside the string, so + * JSON.parse threw on an unterminated string and the caller fell back to its + * lossy text heuristic, losing every structured field. Escape state is tracked + * for the same reason — the quote in `\"}\"` would otherwise close the string + * early and expose the brace as structure. + * + * Returns the object that ENCLOSES the marker, not the nearest `{` before it, so + * a nested object earlier in the text (`{"a": {"b": 1}, "success": true}`) is + * not mistaken for the result. Every occurrence of the marker is tried, because + * prose can name the field before the object appears. (AGT-3466) + */ +export function findStringAwareJsonObject(text: string, marker: string): string | null { + for (let idx = text.indexOf(marker); idx >= 0; idx = text.indexOf(marker, idx + 1)) { + const found = enclosingObject(text, idx); + if (found) return found; + } + return null; +} + +/** The outermost brace-balanced object spanning `idx`, or null. */ +function enclosingObject(text: string, idx: number): string | null { + // Each `{` at or before the marker is a candidate, tried in position order so + // the outermost one wins. A candidate that closes before the marker is a + // sibling object, and one that never closes is an unbalanced `{` in prose; + // both are skipped rather than aborting the scan — prose before the result + // routinely carries stray braces of either kind. + for (let start = text.indexOf('{'); start >= 0 && start <= idx; start = text.indexOf('{', start + 1)) { + const end = endOfBalancedObject(text, start); + if (end !== null && end > idx) return text.slice(start, end); + } + return null; +} + +/** Index just past the `}` matching the `{` at `start`, tracking quoted strings. */ +function endOfBalancedObject(text: string, start: number): number | null { + let depth = 0; + let inString = false; + let escaped = false; + + for (let i = start; i < text.length; i++) { + const ch = text[i]; + + if (escaped) { escaped = false; continue; } + if (ch === '\\') { escaped = true; continue; } + if (ch === '"') { inString = !inString; continue; } + if (inString) continue; + + if (ch === '{') depth++; + if (ch === '}') { + depth--; + if (depth === 0) return i + 1; + } + } + return null; +} + /** * A statement of what the agent could not check, which the worker prompt * requires as a "Could not verify" section (locale/prompts/*.ts). It names a diff --git a/src/adapters/shellCommandGuard.test.ts b/src/adapters/shellCommandGuard.test.ts new file mode 100644 index 00000000..7e04b11a --- /dev/null +++ b/src/adapters/shellCommandGuard.test.ts @@ -0,0 +1,160 @@ +import { describe, it, expect, beforeAll, afterAll } from 'vitest'; +import fs from 'node:fs/promises'; +import { executeTool, ToolCall } from './tools.js'; +import { isCommandBlocked } from './shellCommandGuard.js'; + +/** Helper to build a ToolCall object */ +function makeCall(name: string, args: Record): ToolCall { + return { id: 'tc-1', function: { name, arguments: JSON.stringify(args) } }; +} + +const TMP_DIR = await fs.mkdtemp('/tmp/openswarm-guard-test-'); + +beforeAll(async () => { + await fs.mkdir(TMP_DIR, { recursive: true }); +}); + +afterAll(async () => { + await fs.rm(TMP_DIR, { recursive: true, force: true }); +}); + +// ────────────────────────────────────────────── +// AGT-3436 — destructive commands that never contain the literal +// ────────────────────────────────────────────── + +/** + * Each of these executes a destructive command while containing no literal any + * of the old patterns looked for. Bash rewrites the line before running it: + * quote removal (`r"m"`), backslash escapes (`\rm`), empty-quote splicing + * (`g''it`), and brace expansion (`r{m,}`) all reconstruct the verb. + */ +const rewriteBypasses = [ + 'r"m" -rf /foo', + '\\rm -rf /foo', + "g''it clean -fdx", + 'r{m,} -rf /foo', + "$'\\x72\\x6d' -rf /foo", + 'rm -r -f /foo', + 'rm -fr /foo', + 'git clean -fd', + 'git clean -f -d', +]; + +describe('destructive-command guard sees what the shell will run (AGT-3436)', () => { + it.each(rewriteBypasses)('blocks the shell-rewritten form: %s', (command) => { + expect(isCommandBlocked(command)).toBe(true); + }); + + it.each(rewriteBypasses)('refuses it through the bash tool too: %s', async (command) => { + const result = await executeTool(makeCall('bash', { command }), TMP_DIR); + expect(result.is_error).toBe(true); + expect(result.content).toContain('BLOCKED'); + // A refused command ran nothing, so it is no evidence of anything. + expect(result.executed).toBeUndefined(); + }); + + // Destructive verbs reached through a launcher or a nested shell are still + // that verb; the guard follows both. + it.each([ + "sh -c 'rm -rf /foo'", + 'bash -c "git reset --hard"', + 'sudo -u root rm -rf /foo', + 'env FOO=1 rm -rf /foo', + 'echo "x" | xargs rm -rf', + 'FOO=bar rm -rf /foo', + 'cd /tmp && rm -rf /foo', + 'true; rm -rf /foo', + 'rm -rf /foo > /dev/sda', + ])('blocks a destructive command reached indirectly: %s', (command) => { + expect(isCommandBlocked(command)).toBe(true); + }); + + /** + * The other half of AGT-3436: text that merely MENTIONS a destructive command + * is data, not a command. Refusing it teaches the model to route around the + * guard instead of respecting it. + */ + it.each([ + '# rm -rf /tmp/x', + 'echo "rm -rf is blocked"', + 'echo "run rm -rf only when you mean it"', + 'git status', + 'git log --oneline -5', + 'grep -rn "rm -rf" docs', + 'chmod 755 script.sh', + "python -c 'print(1)'", + 'VERSION=$(cat package.json)', + 'for f in $(ls); do echo "$f"; done', + 'echo `date`', + 'git commit -m "fix: rename variable"', + // `rm` without the recursive flag, and `git clean` without force, are + // ordinary parts of a build loop — flagging them would make the guard noise. + 'rm -f ./dist/bundle.js', + 'rm build/output.txt', + 'git clean -n', + 'git clean -nd', + 'chown user:group file.txt', + 'kill -0 1234', + 'pkill -f local-server', + // Ordinary build/verification commands, the guard's main traffic. + 'npm test', + 'npx vitest run src/adapters/tools.test.ts', + 'npx tsc --noEmit', + 'git status --porcelain', + 'git diff HEAD~1', + 'git log --oneline -5 | head -20', + 'rg -n "pattern" src | head -30', + 'sed -n \'1,50p\' src/adapters/tools.ts', + 'python3 -m pytest tests/ -q', + 'cat package.json | jq .version', + 'mkdir -p a/b && touch a/b/c', + 'echo "hello" > out.txt', + 'ls nonexistent 2>&1 | head -3', + 'node -e "console.log(1+1)"', + 'for f in src/*.ts; do echo "$f"; done', + 'git add -A && git commit -m "fix: thing"', + 'git stash', + "curl -sS https://example.com -o /tmp/out.html", + 'find src -name "*.ts" -type f | wc -l', + 'timeout 30 npm test', + 'env NODE_ENV=test npm test', + "bash -c 'echo hello'", + 'sudo -n true 2>/dev/null || echo nope', + "awk '{print $1}' file.txt", + "printf 'a\\nb\\n' > f.txt", + 'npx oxlint src/adapters/tools.ts', + ])('allows a command that only mentions one: %s', (command) => { + expect(isCommandBlocked(command)).toBe(false); + }); + + it('does not refuse a mention that actually runs', async () => { + const result = await executeTool(makeCall('bash', { command: 'echo "rm -rf is blocked"' }), TMP_DIR); + expect(result.is_error).toBe(false); + expect(result.content).toContain('rm -rf is blocked'); + expect(result.content).not.toContain('BLOCKED'); + }); + + // Text the guard cannot resolve is refused rather than guessed at: a false + // positive costs a retry, a false negative costs the working tree. + it.each([ + 'echo "unterminated', + 'echo "r$(true)m -rf /foo"', + 'rm -rf /{a,b,c,d,e,f,g,h,i,j,k,l,m,n,o,p,q,r,s,t,u,v}', + ])('refuses what it cannot resolve: %s', (command) => { + expect(isCommandBlocked(command)).toBe(true); + }); + + // A brace that is not an expansion group must not send the scan looking for + // the previous one forever: `awk '{print $1}'` is an ordinary command, and a + // guard that never returns is a denial of service on every bash call. + it.each([ + "awk '{print $1}' file.txt", + "awk '{print}' f.txt", + 'echo "{a,b}"', + 'grep -E "{2,3}" file', + 'echo "}"', + 'echo "{unclosed"', + ])('returns promptly for a brace that is not a group: %s', (command) => { + expect(isCommandBlocked(command)).toBe(false); + }); +}); diff --git a/src/adapters/shellCommandGuard.ts b/src/adapters/shellCommandGuard.ts new file mode 100644 index 00000000..9c148619 --- /dev/null +++ b/src/adapters/shellCommandGuard.ts @@ -0,0 +1,407 @@ +// ============================================ +// OpenSwarm - Destructive shell-command guard +// Split out of tools.ts, which sits near the 1500-line pre-commit cap. +// Purpose: decide whether a bash tool command would run something destructive, +// after the shell's own rewriting (quotes, escapes, braces) is applied. +// ============================================ + +import path from 'node:path'; + +/** + * Destructive-command guard (AGT-3436). + * + * This was a regex sweep over the raw command text, and that shape was wrong in + * both directions: + * + * - It missed what bash does before running anything. Quote removal, backslash + * escapes, `$'...'` decoding and brace expansion all rewrite the command + * first, so `r"m" -rf /`, `\rm -rf /`, `$'\x72\x6d' -rf /`, `r{m,} -rf /` + * and `git clean -fdx` each execute a destructive command while containing + * no literal those patterns looked for. + * - It fired on text that is only data. `echo "rm -rf stays blocked"` and + * `# rm -rf /tmp/x` were refused, which is how a model learns to route + * around a guard rather than respect it. + * + * So the command is now resolved the way bash resolves it — quotes and escapes + * removed, `$'...'` decoded, comments dropped, braces expanded — and matched by + * WORD: the first word of a simple command is the program that runs, so a + * destructive verb is one only where a program name sits. Anything that cannot + * be resolved (an unclosed quote or substitution, a substitution spliced into a + * word, a brace expansion past its cap) is refused rather than guessed at: a + * false positive costs a retry, a false negative costs the working tree. + */ + +/** How far the guard follows `$(...)`, backticks and `sh -c` scripts before refusing. */ +const GUARD_MAX_DEPTH = 4; +/** Candidates `{a,b}` expansion may produce before the command is refused. */ +const GUARD_MAX_EXPANSIONS = 32; +/** Words scanned for a launcher's real command (`sudo -u root rm -rf /`). */ +const GUARD_MAX_WORDS = 64; + +/** One simple command, split out of a `;`/`&&`/`||`/`|`/newline chain. */ +interface ResolvedCommand { + /** Words after quote removal, backslash escapes and brace expansion. */ + words: string[]; + /** Targets of `>`/`<` redirections, kept apart from arguments. */ + redirects: string[]; + /** Bodies of `$(...)`/backtick substitutions — each runs a command of its own. */ + nested: string[]; +} + +/** Programs that only launch another command: the real one is in the arguments. */ +const COMMAND_LAUNCHERS: Record = { + sudo: true, doas: true, su: true, command: true, builtin: true, env: true, + nohup: true, nice: true, ionice: true, stdbuf: true, setsid: true, time: true, + timeout: true, watch: true, flock: true, chroot: true, exec: true, xargs: true, + find: true, +}; + +/** Programs whose `-c` argument is a script the shell runs. */ +const SCRIPT_HOSTS: Record = { sh: true, bash: true, zsh: true, dash: true, ksh: true, su: true }; + +function isWordChar(char: string | undefined): boolean { + return char !== undefined && /[A-Za-z0-9_]/.test(char); +} + +/** Innermost `{a,b}` group bash would expand, or null when the word has none. */ +function innermostBraceGroup(word: string): { start: number; end: number; alternatives: string[] } | null { + let start = word.lastIndexOf('{'); + while (start >= 0) { + const close = word.indexOf('}', start + 1); + if (close >= 0) { + const body = word.slice(start + 1, close); + // Not innermost (the inner group expands first) and no alternative list + // (`{x}` is literal to bash) both mean this brace is not a group. + if (!body.includes('{') && body.includes(',')) { + return { start, end: close + 1, alternatives: body.split(',') }; + } + } + // NB: `lastIndexOf('{', -1)` clamps to 0 and would rescan index 0 forever, + // so the walk stops explicitly rather than relying on a negative fromIndex. + if (start === 0) break; + start = word.lastIndexOf('{', start - 1); + } + return null; +} + +/** + * Every word `{a,b}` expansion can produce, or null when the count explodes. + * Each candidate is a word bash may run, so an unresolvable expansion is + * refused instead of being matched as its own literal text. + */ +function expandBraces(word: string): string[] | null { + let candidates = [word]; + for (;;) { + const next: string[] = []; + let expanded = false; + for (const candidate of candidates) { + const group = innermostBraceGroup(candidate); + if (!group) { + next.push(candidate); + continue; + } + expanded = true; + for (const alternative of group.alternatives) { + next.push(candidate.slice(0, group.start) + alternative + candidate.slice(group.end)); + } + } + if (!expanded) return next; + if (next.length > GUARD_MAX_EXPANSIONS) return null; + candidates = next; + } +} + +/** Index of the `'` closing the `$'...'` quote that starts at `start`, or -1. */ +function closingAnsiCQuote(command: string, start: number): number { + for (let i = start + 1; i < command.length; i++) { + if (command[i] === '\\') { i++; continue; } + if (command[i] === "'") return i; + } + return -1; +} + +/** + * What `$'...'` resolves to: bash decodes backslash escapes there, so + * `$'\x72\x6d' -rf /` runs `rm -rf /`. Decoding keeps the guard looking at the + * characters the process will actually see. + */ +function decodeAnsiCQuote(body: string): string { + return body.replace( + /\\(x[0-9a-fA-F]{1,2}|[0-7]{1,3}|u[0-9a-fA-F]{4}|U[0-9a-fA-F]{8}|[\s\S])/g, + (_all, escape: string) => { + const kind = escape[0]; + const code = kind === 'x' || kind === 'u' || kind === 'U' + ? parseInt(escape.slice(1), 16) + : kind >= '0' && kind <= '7' ? parseInt(escape, 8) : -1; + if (code >= 0 && code <= 0x10ffff) return String.fromCodePoint(code); + if (escape === 'n') return '\n'; + if (escape === 't') return '\t'; + if (escape === 'r') return '\r'; + return escape; + }, + ); +} + +/** Index just past the `$(...)`, `${...}` or `` `...` `` span at `start`, or -1 when it never closes. */ +function endOfSubstitution(command: string, start: number): number { + if (command[start] === '`') { + for (let i = start + 1; i < command.length; i++) { + if (command[i] === '\\') { i++; continue; } + if (command[i] === '`') return i + 1; + } + return -1; + } + const opens = command[start + 1]; + const closes = opens === '(' ? ')' : '}'; + let depth = 0; + let quote: '"' | "'" | null = null; + for (let i = start + 1; i < command.length; i++) { + const char = command[i]; + if (quote) { + if (char === '\\' && quote === '"') { i++; continue; } + if (char === quote) quote = null; + continue; + } + if (char === '\\') { i++; continue; } + if (char === '"' || char === "'") { quote = char; continue; } + if (char === opens) depth++; + else if (char === closes && --depth === 0) return i + 1; + } + return -1; +} + +interface SubstitutionSpan { + /** The span as written, which is all the guard can know about its output. */ + text: string; + /** The command inside `$(...)`/backticks, for recursive inspection. */ + body: string; + /** Index just past the span. */ + end: number; + runsCommand: boolean; +} + +/** + * The `$(...)`/`${...}`/`` `...` `` span at `i`: `'none'` when there is none, + * `'unclosed'` when it never terminates (the caller refuses). + */ +function substitutionSpanAt(command: string, i: number): SubstitutionSpan | 'none' | 'unclosed' { + const char = command[i]; + const dollar = char === '$' && (command[i + 1] === '(' || command[i + 1] === '{'); + if (char !== '`' && !dollar) return 'none'; + const end = endOfSubstitution(command, i); + if (end < 0) return 'unclosed'; + const runsCommand = char === '`' || command[i + 1] === '('; + const text = command.slice(i, end); + return { text, body: runsCommand ? text.slice(char === '`' ? 1 : 2, -1) : text, end, runsCommand }; +} + +/** + * The simple commands bash will run for `command`, each word resolved the way + * bash resolves it. Returns null when the text cannot be resolved — see the + * block comment above for why that is refused rather than matched. + */ +function resolveCommands(command: string): ResolvedCommand[] | null { + const commands: ResolvedCommand[] = []; + let words: string[] = []; + let redirects: string[] = []; + let nested: string[] = []; + let word = ''; + let wordStarted = false; + let atWordStart = true; + let redirectTarget = false; + let unresolvable = false; + let quote: '"' | "'" | null = null; + + const endWord = (): void => { + if (!wordStarted) return; + if (redirectTarget) redirects.push(word); + else { + const expanded = expandBraces(word); + if (expanded === null) unresolvable = true; + else words.push(...expanded); + } + word = ''; + wordStarted = false; + redirectTarget = false; + }; + const endCommand = (): void => { + endWord(); + if (words.length || redirects.length || nested.length) commands.push({ words, redirects, nested }); + words = []; + redirects = []; + nested = []; + atWordStart = true; + }; + + for (let i = 0; i < command.length; i++) { + const char = command[i]; + + if (quote === "'") { + if (char === "'") quote = null; + else word += char; + continue; + } + + if (quote === '"') { + if (char === '"') { quote = null; continue; } + if (char === '\\') { + const next = command[i + 1]; + if (next === undefined) return null; + if (next === '"' || next === '\\' || next === '$' || next === '`') { word += next; i++; } + else if (next === '\n') i++; + else word += char; + continue; + } + const span = substitutionSpanAt(command, i); + if (span === 'unclosed') return null; + if (span !== 'none') { + // `r"$(true)"m` is `rm` once the quotes come off: a splice is refused + // because the guard cannot know what the substitution yields. + if (isWordChar(word[word.length - 1]) || isWordChar(command[span.end])) return null; + word += span.text; + if (span.runsCommand) nested.push(span.body); + i = span.end - 1; + continue; + } + word += char; + continue; + } + + const ifs = /^\$\{IFS\}|\$IFS(?![A-Za-z0-9_])/.exec(command.slice(i, i + 6)); + if (ifs) { + // `rm$IFS-rf` is `rm -rf`: bash splits the word where IFS expands. + endWord(); + atWordStart = true; + i += ifs[0].length - 1; + continue; + } + if (char === '\\') { + const next = command[i + 1]; + if (next === undefined) return null; + if (next !== '\n') { word += next; wordStarted = true; atWordStart = false; } + i++; + continue; + } + if (char === "'" || char === '"') { quote = char; wordStarted = true; atWordStart = false; continue; } + if (char === '$' && command[i + 1] === "'") { + const close = closingAnsiCQuote(command, i + 1); + if (close < 0) return null; + word += decodeAnsiCQuote(command.slice(i + 2, close)); + wordStarted = true; + atWordStart = false; + i = close; + continue; + } + if (char === ' ' || char === '\t') { endWord(); atWordStart = true; continue; } + if (char === '\n' || char === ';' || char === '&' || char === '|') { + endCommand(); + if (char !== '\n' && command[i + 1] === char) i++; + continue; + } + if (char === '>' || char === '<') { + endWord(); + atWordStart = true; + redirectTarget = false; + if (command[i + 1] === char) { i++; redirectTarget = char === '>'; } // `>> device` + else if (command[i + 1] === '&') i++; // `>&2` duplicates a descriptor + else redirectTarget = char === '>'; + continue; + } + if (char === '#' && atWordStart) { + const newline = command.indexOf('\n', i); + endCommand(); + if (newline < 0) break; + i = newline; + continue; + } + const span = substitutionSpanAt(command, i); + if (span === 'unclosed') return null; + if (span !== 'none') { + if (isWordChar(word[word.length - 1]) || isWordChar(command[span.end])) return null; + word += span.text; + wordStarted = true; + atWordStart = false; + if (span.runsCommand) nested.push(span.body); + i = span.end - 1; + continue; + } + word += char; + wordStarted = true; + atWordStart = false; + } + + if (unresolvable) return null; + if (quote !== null) return null; // an unclosed quote swallows the rest of the line + endCommand(); + return commands; +} + +/** + * The destructive shapes, as predicates over (program, arguments). Kept apart + * from the resolution above so the two questions stay separable: what will run, + * and is that thing destructive. + */ +const DESTRUCTIVE_RULES: ReadonlyArray<(name: string, args: string[]) => boolean> = [ + // `rm -rf`, `rm -fr`, `rm -R`, `rm --recursive` — the recursive flag is the destructive part. + (name, args) => name === 'rm' && args.some((arg) => arg === '--recursive' || /^-[A-Za-z]*[rR][A-Za-z]*$/.test(arg)), + (name, args) => name === 'git' && args.includes('reset') && args.includes('--hard'), + // `git clean` deletes untracked files once force meets directories: `-fd`, `-fdx`, `-d -f`. + (name, args) => name === 'git' && args.includes('clean') + && args.some((arg) => arg === '--force' || /^-[A-Za-z]*f[A-Za-z]*$/.test(arg)) + && args.some((arg) => /^-[A-Za-z]*d[A-Za-z]*$/.test(arg)), + (name, args) => name === 'chmod' && args.some((arg) => /^[0-7]*777[0-7]*$/.test(arg)), + (name, args) => name === 'chown' + && args.some((arg) => arg === '--recursive' || /^-[A-Za-z]*R[A-Za-z]*$/.test(arg)), + (name, args) => name === 'dd' && args.some((arg) => arg.startsWith('if=')), + (name, args) => (name === 'kill' || name === 'pkill') + && args.some((arg) => /^-[A-Za-z]*9$/.test(arg) || /^-(SIG)?KILL$/i.test(arg)), +]; + +/** Is one resolved simple command destructive? */ +function resolvedCommandIsBlocked(command: ResolvedCommand, depth: number): boolean { + if (command.redirects.some((target) => target.startsWith('/dev/sd'))) return true; + // SQL verbs are destructive wherever they sit in the line: a client reads + // them as a statement, so position cannot separate `psql -c "drop database + // app"` from a word that merely mentions one. + const joined = command.words.join(' '); + if (/\bdrop\s+database\b/i.test(joined) || /\btruncate\s+table\b/i.test(joined)) return true; + + if (depth < GUARD_MAX_DEPTH) { + for (const script of command.nested) { + if (isCommandBlocked(script, depth + 1)) return true; + } + } + + let start = 0; + while (command.words[start] !== undefined && /^[A-Za-z_][A-Za-z0-9_]*=/.test(command.words[start])) start++; // `FOO=bar rm -rf /` + const name = path.basename(command.words[start] ?? ''); + const args = command.words.slice(start + 1); + + if (depth < GUARD_MAX_DEPTH) { + // `sh -c '