diff --git a/.env.example b/.env.example index ee0084f..364e9cb 100644 --- a/.env.example +++ b/.env.example @@ -43,6 +43,8 @@ FOREMAN_VALIDATION_COMMANDS=[] FOREMAN_FORMAT_COMMAND= FOREMAN_VALIDATION_TIMEOUT_MS=120000 FOREMAN_VALIDATION_MAX_OUTPUT_BYTES=1048576 +# Times a Worker is resumed in its own session with failing checks before its turn ends (0 turns this off). +FOREMAN_WORKER_CHECK_ROUNDS=2 # Every validation command runs inside a bubblewrap (bwrap) sandbox that hides your home # directory, credentials, Foreman's data dir, the source checkout, /tmp and /run. Validation # fails if bwrap is missing. Set none only on hosts without bubblewrap (for example macOS); diff --git a/README.md b/README.md index 6c4c25c..ae154a7 100644 --- a/README.md +++ b/README.md @@ -65,6 +65,8 @@ Validation commands run in a disposable copy of the Worker workspace. Every conf **Checks that also fail on the base commit.** When validation fails on a Worker snapshot that has changes, Foreman runs the configured commands up to the last failed check (so installs and builds a check depends on run too) once against the unchanged pinned base commit, in the same sandbox, and caches the result on the run (pinned base + command digest). Each failed check is marked `failsOnBase` and shown with a "Fails on base" badge. If every failed check also fails there (for example `pnpm run smoke:install` with no network), a Worker retry cannot fix it, so Foreman does not spend another Worker attempt: the run stops with "Checks also fail on the base commit, so a Worker retry cannot fix them: ...". Fix the check or its network setting, then use **Retry validation** (which also discards the cached baseline). If some failures are Worker-caused, the normal automatic follow-up runs, but its correction note lists only those and names the base-failing checks as not the Worker's to fix. This never relaxes promotion: every check must still pass. If the baseline cannot run, Foreman falls back to the ordinary follow-up. +**Worker check gate.** Workers are interchangeable, so none of them is trusted to run the repository's checks: Antigravity and Claude Code have file tools only, and Codex's sandbox has no network and no installed dependencies. Instead, when a Worker's turn completes the bridge keeps the task in progress and pauses it, and Foreman runs the configured checks on the paused workspace (the same snapshot verification, format step and sandboxed validation as the final pass). If a check the Worker could have caused fails, Foreman sends the command and a head-and-tail output excerpt back, and the bridge resumes the same CLI session (`--resume` for Claude Code, `exec resume` for Codex, `--conversation` for Antigravity) so the Worker fixes it with its full context. A change outside the allowed scope counts as a failure; checks that fail on the base commit and any infrastructure problem do not. This repeats up to `FOREMAN_WORKER_CHECK_ROUNDS` times (default 2) within one Worker attempt, so it spends no Worker attempt and no Orchestrator turn. The gate's verdicts are feedback only: Foreman's validation after the turn is still the result the Reviewer and you see. It is off when automatic validation correction is off for the run. Each round is recorded as a `worker.check_round` event and in the bridge response's `metadata.worker_gate`. + **Format step (optional).** Antigravity and Claude Code Workers can only edit files, so they cannot run the repository's formatter, and a `prettier --check` validation command would fail on their output. When a `formatCommand` is configured (`{"name","command","args","cwd?","network?"}`; the open-repository dialog suggests ` run format`, `format:write` or `prettier:write` from `package.json` with an on/off toggle, `FOREMAN_FORMAT_COMMAND` sets it for the single-repository configuration), Foreman runs it after it has verified the Worker snapshot and before validation. It materializes the verified snapshot in the same bubblewrap sandbox as validation, runs the dependency installs from the validation list (`pnpm`, `npm` or `yarn` install, `ci` or `add`), then the formatter (offline unless it sets `network: true`). Only the files the Worker added, modified or renamed are read back, keeping the Worker's file mode; anything else the formatter touched or created is ignored. The new snapshot is re-verified against the pinned base and allowed scope, and validation, the Reviewer, approval and promotion all use the formatted bytes. The run's evidence records `formatting` (status `applied`, `unchanged` or `failed`, the formatted paths and the formatter's bounded output) and the UI shows "Foreman formatted N files". A formatter failure never fails the run: the Worker's snapshot is kept and validation reports the real problem. The step applies to live Worker snapshots only, not to recorded replays. The suggested allowed scope covers the top-level tracked paths except CI configuration (`.github/`, `.gitlab-ci.yml`, `.circleci/`, `.buildkite/`, `azure-pipelines.yml`, `Jenkinsfile`, `.travis.yml`), which runs with repository secrets once pushed. Add those paths by hand if a task really has to change them. @@ -155,6 +157,7 @@ The standard local flow requires no environment variables. The following setting | `FOREMAN_FORMAT_COMMAND` | — | Optional JSON object `{"name","command","args","cwd?","network?"}`: the formatter Foreman runs on the Worker's changed files before validation. Offline unless `network` is `true`. See the format step above. | | `FOREMAN_VALIDATION_TIMEOUT_MS` | `120000` | Per-command time limit (max 600,000 ms). | | `FOREMAN_VALIDATION_MAX_OUTPUT_BYTES` | `1048576` | Per-command output capture bound (max 16 MiB). | +| `FOREMAN_WORKER_CHECK_ROUNDS` | `2` | Worker check gate: how many times a Worker is resumed in its own session with failing checks before its turn ends (max 3). `0` turns the gate off. See the check gate above. | | `FOREMAN_VALIDATION_SANDBOX` | `bwrap` | `bwrap` runs every validation command in a bubblewrap sandbox and fails if it is unavailable. `none` runs them directly on the host with your credentials, and is unsafe. See [Validation sandbox](#validation-sandbox). | | `FOREMAN_VALIDATION_SANDBOX_RO_PATHS` | — | Comma-separated absolute paths mounted read-only in the sandbox, for toolchains under a hidden directory such as `$HOME/.volta`. Paths that do not exist are ignored. | diff --git a/config.schema.json b/config.schema.json index f48ba34..372cf3d 100644 --- a/config.schema.json +++ b/config.schema.json @@ -23,6 +23,7 @@ "FOREMAN_FORMAT_COMMAND": { "type": "string", "description": "Optional JSON object {name, command, args, cwd?, network?} for a formatter Foreman runs in the sandboxed validation workspace (after the install commands) on the Worker's added or modified files, before validation. The formatted bytes become the evidence that validation, the Reviewer and promotion use. Runs without network access unless network is true" }, "FOREMAN_VALIDATION_TIMEOUT_MS": { "type": "integer", "minimum": 1, "maximum": 600000, "default": 120000 }, "FOREMAN_VALIDATION_MAX_OUTPUT_BYTES": { "type": "integer", "minimum": 1, "maximum": 16777216, "default": 1048576 }, + "FOREMAN_WORKER_CHECK_ROUNDS": { "type": "integer", "minimum": 0, "maximum": 3, "default": 2, "description": "Worker check gate: after a Worker turn, Foreman runs the configured checks on the paused workspace and resumes the same Worker session with any failures, up to this many times. 0 turns the gate off" }, "FOREMAN_VALIDATION_SANDBOX": { "type": "string", "enum": ["bwrap", "none"], "default": "bwrap", "description": "bwrap runs every validation command inside a bubblewrap sandbox and fails closed when it is unavailable; none runs them unsandboxed on the host (unsafe opt-out)" }, "FOREMAN_VALIDATION_SANDBOX_RO_PATHS": { "type": "string", "description": "Comma-separated absolute paths (for example a toolchain under your home directory) mounted read-only into the validation sandbox" } } diff --git a/investigations/local-cli-uhp/README.md b/investigations/local-cli-uhp/README.md index 0928509..bd9e05e 100644 --- a/investigations/local-cli-uhp/README.md +++ b/investigations/local-cli-uhp/README.md @@ -234,6 +234,12 @@ idempotency key, and a persisted workspace ID for replay. `LOCAL_CLI_UHP_SMOKE_KEY_FILE` and `LOCAL_CLI_UHP_EVIDENCE_FILE` select the key and evidence paths. +## Worker check gate + +A Worker request may carry `metadata.foreman_worker_gate: {"max_rounds": 1-3}` (Reviewer and Planner/Orchestrator requests are refused with `worker_gate_role_unsupported`). Discovery advertises it as `extensions.foreman_worker_gate_v1`. + +After each completed Worker turn the task stays `in_progress` and the bridge records a `response.activity` event with `kind: "worker_gate"` and `gate_round: N`. While paused, no CLI is running, so the workspace snapshot endpoint serves a complete snapshot for Foreman's checks (overlay stays refused). Foreman answers with `POST /extensions/foreman-workspace/v1/responses/{id}/worker-gate` and `{"round": N, "status": "passed" | "failed" | "skipped", "feedback"?, "failed_checks"?}`; a verdict for any other round is refused with 409. On `failed`, the bridge resumes the same native session with the feedback: `--resume` for Claude Code (the transcript directories are bound per workspace, like a role session), `exec resume` without `--ephemeral` for Codex (its per-workspace `CODEX_HOME` is kept), and `--conversation` for Antigravity (its per-workspace state directory is kept). Any other verdict, no verdict within `LOCAL_CLI_UHP_WORKER_GATE_TIMEOUT_MS` (default 15 minutes), a failed turn or a cancellation ends the task as the last turn left it. `metadata.worker_gate` records each round; usage is summed across Claude Code and Codex turns and taken from the last Antigravity turn, whose usage is cumulative per conversation. `cli_invocation` stays the first turn's invocation. The kept session state is deleted when the task ends. + ## Recorded live proof The first Claude workspace attempt failed closed at `boundary_probe` before CLI diff --git a/investigations/local-cli-uhp/cli-args.mjs b/investigations/local-cli-uhp/cli-args.mjs index 879941e..fc7402d 100644 --- a/investigations/local-cli-uhp/cli-args.mjs +++ b/investigations/local-cli-uhp/cli-args.mjs @@ -4,8 +4,9 @@ export function claudeCodeCliArgs(model, { reviewer = false, sessionId, persiste return ['-p', '--output-format', 'stream-json', '--verbose', '--model', model, '--max-turns', String(Math.min(maxStep, 10)), '--restricted', '--strict-mcp-config', '--permission-mode', persistentContext ? 'plan' : 'acceptEdits', '--tools', persistentContext ? 'Read,Grep,Glob' : 'Read,Edit,Write', ...(sessionId ? ['--resume', sessionId] : [])]; } -export function codexCliArgs(model, { reviewer = false, sessionId, persistentContext = false } = {}) { +/** `keepSession` keeps a Worker's session on disk (no --ephemeral) so a Worker check round can resume it. */ +export function codexCliArgs(model, { reviewer = false, sessionId, persistentContext = false, keepSession = false } = {}) { // --sandbox is a top-level Codex option. `codex exec resume` rejects it when // placed after `resume`, before any session can be reported. - return ['--ask-for-approval', 'never', '--sandbox', reviewer || persistentContext ? 'read-only' : 'workspace-write', 'exec', ...(sessionId ? ['resume', sessionId] : []), '--json', ...(!persistentContext ? ['--ephemeral'] : []), '--ignore-user-config', ...(reviewer ? ['--ignore-rules'] : []), '--skip-git-repo-check', '--model', model, '-']; + return ['--ask-for-approval', 'never', '--sandbox', reviewer || persistentContext ? 'read-only' : 'workspace-write', 'exec', ...(sessionId ? ['resume', sessionId] : []), '--json', ...(!persistentContext && !keepSession ? ['--ephemeral'] : []), '--ignore-user-config', ...(reviewer ? ['--ignore-rules'] : []), '--skip-git-repo-check', '--model', model, '-']; } diff --git a/investigations/local-cli-uhp/server.mjs b/investigations/local-cli-uhp/server.mjs index 542034b..dcf331f 100644 --- a/investigations/local-cli-uhp/server.mjs +++ b/investigations/local-cli-uhp/server.mjs @@ -17,6 +17,7 @@ import { roleSessionBinding, roleStatePath, resolvePreviousRoleSession, withRole import { claudeCodeCliArgs, codexCliArgs } from './cli-args.mjs'; import { bwrapBaseArgs } from './bwrap-args.mjs'; import { readClaudeControlUsage } from './claude-quota.mjs'; +import { parseWorkerGateRequest, parseWorkerGateVerdict, workerGatePrompt, addUsage, MAX_WORKER_GATE_ROUNDS } from './worker-gate.mjs'; const VERSION = '2026-09-12'; const utf8 = new TextDecoder('utf-8', { fatal: true }); @@ -35,6 +36,8 @@ const BWRAP = process.env.LOCAL_CLI_UHP_BWRAP ?? 'bwrap'; const MAX_PROMPT = 16_000; const MAX_OUTPUT = 64_000; const MAX_TIMEOUT = 900; +// How long a gated Worker waits for Foreman's check verdict before ending its turn as it is. +const WORKER_GATE_TIMEOUT_MS = (() => { const value = Number(process.env.LOCAL_CLI_UHP_WORKER_GATE_TIMEOUT_MS); return Number.isFinite(value) && value > 0 ? value : 15 * 60_000; })(); const MAX_REVIEW_DIFF = 48_000; const SSE_KEEPALIVE_MS = (() => { const value = Number(process.env.LOCAL_CLI_UHP_KEEPALIVE_MS); return Number.isFinite(value) && value > 0 ? Math.min(value, 60_000) : 10_000; })(); const AGY_WORKER_AGENT = 'foreman-worker'; @@ -199,11 +202,12 @@ function safeToolName(value) { const name = value.replace(/[^A-Za-z0-9_.-]/g, '').slice(0, 40); return name ? `the ${name} tool` : 'a tool'; } -function recordActivity(task, summary) { - if (task.activityTotal >= ACTIVITY_TOTAL_LIMIT) return; +function recordActivity(task, summary, gateRound) { + if (task.activityTotal >= ACTIVITY_TOTAL_LIMIT && !gateRound) return; const safe = safeActivityText(summary); if (!safe || task.activity.at(-1)?.summary === safe) return; - task.activity.push({ index: ++task.activityTotal, kind: /^Using /.test(safe) ? 'tool' : 'commentary', summary: safe }); + // A Worker check-gate pause is always recorded: Foreman acts on it, so the activity cap must not drop it. + task.activity.push({ index: ++task.activityTotal, kind: gateRound ? 'worker_gate' : /^Using /.test(safe) ? 'tool' : 'commentary', summary: safe, ...(gateRound ? { gate_round: gateRound } : {}) }); if (task.activity.length > ACTIVITY_LIMIT) task.activity.splice(0, task.activity.length - ACTIVITY_LIMIT); for (const notify of task.activityListeners) notify(); } @@ -479,9 +483,9 @@ function parseAgy(text, stderr = '') { const distinctToolSteps = new Set(toolEvents.map((update, ordinal) => Number.isSafeInteger(update.step_index) ? String(update.step_index) : `unindexed:${ordinal}`)); return { text: responseText, model, session: conversation, turnCompleted: !!result, isError: result?.status !== 'SUCCESS', mutationAttempted, malformedOutput, unrecognizedOutput, usage: result?.usage && typeof result.usage === 'object' ? result.usage : undefined, diagnostic:{ permission_mode:permissionMode, observed_agent:agent ?? 'unreported', cwd, available_tools:initTools, available_tools_semantics:'headless_init_tools_available_to_cli_not_profile_allowlist', tool_events:toolEvents, tool_lifecycle_update_count:toolUpdates.length, distinct_tool_step_count:distinctToolSteps.size, reported_cli_turns:reportedCliTurns, soft_denial_observed:softDenialObserved, result_status:typeof result?.status === 'string' && /^(SUCCESS|ERROR|CANCELED|INTERRUPTED|INVALID|WAITING|RUNNING)$/.test(result.status) ? result.status : 'unreported', response_empty:responseText.length === 0, response_characters:responseText.length, streamed_agent_text_characters:chunks.reduce((n,s)=>n+s.length,0) } }; } -function cliArgs(kind, model, timeout, maxStep, reviewer = false, sessionId, persistentContext = false) { +function cliArgs(kind, model, timeout, maxStep, reviewer = false, sessionId, persistentContext = false, keepSession = false) { if (kind === 'claude') return claudeCodeCliArgs(model, { reviewer, sessionId, persistentContext, maxStep }); - return codexCliArgs(model, { reviewer, sessionId, persistentContext }); + return codexCliArgs(model, { reviewer, sessionId, persistentContext, keepSession }); } function agyArgs(model, timeout, conversationId, reviewer = false, worker = false) { return ['--output-format', 'stream-json', '--model', model, '--print-timeout', `${timeout}s`, `--mode=${reviewer ? 'plan' : 'accept-edits'}`, ...(worker ? ['--add-dir','/workspace','--agent', AGY_WORKER_AGENT, ...(AGY_WORKER_EFFORT ? ['--effort', AGY_WORKER_EFFORT] : [])] : []), ...(conversationId ? ['--conversation', conversationId] : [])]; @@ -783,9 +787,11 @@ async function claudeAuthMountArgs(ws) { } // Claude may create these directories even on its first turn. Keep its role // transcript/state writable while the host sign-in/config view stays RO. - if (ws.roleStatePath) { + // A gated Worker keeps its transcript the same way, per workspace, so a check round can resume the session. + const statePath = ws.roleStatePath ?? ws.workerSessionPath; + if (statePath) { for (const name of CLAUDE_ROLE_DIRS) { - const source = join(ws.roleStatePath, `claude-${name}`); + const source = join(statePath, `claude-${name}`); await mkdir(source, { recursive:true, mode:0o700 }); args.push('--dir',`/auth/${name}`,'--bind',source,`/auth/${name}`); } @@ -1041,6 +1047,63 @@ async function resolveBinary(name) { } async function runTask(record, prompt) { + const active = tasks.get(record.id), gate = active.gate; + if (!gate) return (await runCliTurn(record, prompt)).finalize(); + return runGatedWorkerTask(record, prompt, active, gate); +} + +/** + * A Worker turn with Foreman's check gate: after each completed CLI turn the task stays in progress while Foreman + * checks the workspace. A failed verdict resumes the same CLI session with the failures, up to the requested number + * of rounds. Any other verdict, a timeout, a failed turn or a cancellation ends the task as the last turn left it. + */ +async function runGatedWorkerTask(record, prompt, active, gate) { + const ws = workspaces.get(record.metadata.workspace_id); + if (ws) { ws.workerSessionPath = join(ROOT, `${ws.id}.worker-session`); await mkdir(ws.workerSessionPath, { recursive: true, mode: 0o700 }); } + let turn; + try { + turn = await runCliTurn(record, prompt, { hold: true, keepSession: true }); + const firstInvocation = record.metadata.cli_invocation, cumulativeUsage = record.metadata.harness_id === 'antigravity-cli'; + let usage = record.usage; + const rounds = []; + const publish = state => { record.metadata.worker_gate = { max_rounds: gate.maxRounds, state, rounds: rounds.map(r => ({ ...r })) }; }; + publish('running'); + while (turn.status === 'completed' && rounds.length < gate.maxRounds && !active.cancelRequested) { + const round = rounds.length + 1; + publish('awaiting_check'); + const verdict = await awaitWorkerGateVerdict(active, round, gate.maxRounds); + rounds.push({ round, status: verdict.status, ...(verdict.failedChecks ? { failed_checks: verdict.failedChecks } : {}), ...(verdict.reason ? { reason: verdict.reason } : {}) }); + if (verdict.status !== 'failed' || !record.session_id || active.cancelRequested) break; + publish('fixing'); + recordActivity(active, `Foreman checks failed; resuming the Worker to fix them (round ${round} of ${gate.maxRounds})`); + turn = await runCliTurn(record, workerGatePrompt(verdict.feedback, round, gate.maxRounds), { hold: true, keepSession: true, resumeSession: record.session_id }); + usage = cumulativeUsage ? record.usage : addUsage(usage, record.usage); + } + publish('finished'); + if (usage) record.usage = usage; + if (firstInvocation) record.metadata.cli_invocation = firstInvocation; + record.status = active.cancelRequested ? 'cancelled' : turn.status; + if (record.status === 'cancelled' && !record.error) record.error = { message: 'CLI task was cancelled' }; + } finally { + gate.pending = undefined; + if (ws?.workerSessionPath) { await rm(ws.workerSessionPath, { recursive: true, force: true }).catch(() => undefined); ws.workerSessionPath = undefined; } + } + await turn.finalize(); +} + +/** Pause until Foreman posts the verdict for `round`, the task is cancelled, or the wait times out (then the turn simply ends). */ +function awaitWorkerGateVerdict(active, round, maxRounds) { + return new Promise(resolveVerdict => { + const settle = verdict => { if (active.gate.pending?.round !== round) return; clearTimeout(timer); active.gate.pending = undefined; resolveVerdict(verdict); }; + const timer = setTimeout(() => settle({ round, status: 'skipped', reason: 'verdict_timeout' }), WORKER_GATE_TIMEOUT_MS); + timer.unref?.(); + active.gate.pending = { round, settle }; + recordActivity(active, `Waiting for Foreman checks (round ${round} of ${maxRounds})`, round); + if (active.cancelRequested) settle({ round, status: 'skipped', reason: 'cancelled' }); + }); +} + +async function runCliTurn(record, prompt, turn = {}) { const h = cliFor(record.metadata.harness_id), kind = h.id === 'claude-code' ? 'claude' : h.id === 'codex-cli' ? 'codex-cli' : 'agy'; const reviewer = record.metadata.foreman_review_mode === 'read_only'; const persistentRole = record.metadata.role_session_binding !== undefined; @@ -1064,16 +1127,16 @@ async function runTask(record, prompt) { const inheritedNames = new Set(['PATH','HOME','USER','LOGNAME','LANG','LC_ALL','TERM','TMPDIR','TMP','TEMP','XDG_RUNTIME_DIR','SSL_CERT_FILE','SSL_CERT_DIR','NODE_EXTRA_CA_CERTS','HTTP_PROXY','HTTPS_PROXY','ALL_PROXY','NO_PROXY','DBUS_SESSION_BUS_ADDRESS']); const env = Object.fromEntries(Object.entries(process.env).filter(([name]) => inheritedNames.has(name))); Object.assign(env, kind === 'claude' ? { CLAUDE_CONFIG_DIR: '/auth' } : kind === 'codex-cli' ? { CODEX_HOME: reviewer ? h.authDir : '/auth' } : {}); - const nativeSessionId = record.metadata.conversation_id; + const nativeSessionId = turn.resumeSession ?? record.metadata.conversation_id; const agyWorker = kind === 'agy' && !reviewer && !persistentRole; - const args = kind === 'agy' ? agyArgs(record.requested_model, record.timeout_seconds, nativeSessionId, reviewer || persistentRole, agyWorker) : cliArgs(kind, record.requested_model, record.timeout_seconds, record.max_step, reviewer, nativeSessionId, persistentRole); + const args = kind === 'agy' ? agyArgs(record.requested_model, record.timeout_seconds, nativeSessionId, reviewer || persistentRole, agyWorker) : cliArgs(kind, record.requested_model, record.timeout_seconds, record.max_step, reviewer, nativeSessionId, persistentRole, turn.keepSession === true); if (kind === 'agy') { record.metadata.cli_invocation = { executable: '/opt/agy', host_executable: h.bin, args: ['-p', '', ...args] }; record.metadata.actual_model_status = 'unavailable'; } if (kind === 'claude' && persistentRole) record.metadata.cli_invocation = { executable: '/opt/claude', host_executable: h.bin, args: [...args] }; const input = reviewer ? reviewerPrompt(prompt, record.metadata.review_evidence) : kind === 'agy' && !persistentRole ? agyWorkerPrompt(prompt) : prompt; - if (kind !== 'claude' && !reviewer) record.metadata.submitted_prompt_sha256 = createHash('sha256').update(input).digest('hex'); + if (kind !== 'claude' && !reviewer && !turn.resumeSession) record.metadata.submitted_prompt_sha256 = createHash('sha256').update(input).digest('hex'); const reviewerState = reviewer ? { dir: work, authDir: kind === 'claude' ? await realpath(h.authDir) : undefined } : undefined; let codexRuntime; if (kind === 'codex-cli') { @@ -1134,7 +1197,9 @@ async function runTask(record, prompt) { const invalidAgyWorkerPolicy = !!agyWorker && (!agyWorkerPolicy?.selected_agent_matches || !agyWorkerPolicy.executed_tools_within_profile); if (agyWorkerPolicy) record.metadata.agy_worker_tool_policy = { ...agyWorkerPolicy, configured_agent: AGY_WORKER_AGENT, command_execution_policy: 'off', available_tools_are_diagnostic_only:true, execution_observations_passed: !invalidAgyWorkerPolicy }; if (agyWorkerPolicy?.tolerated_tool_events?.length) record.metadata.tolerated_tool_warning = `Worker used tolerated AGY built-in tool(s) (assumed no filesystem/command effect; snapshot diff and commandExecutionPolicy "off" are the safety boundary): ${agyWorkerPolicy.tolerated_tool_events.join(', ')}`; - record.status = active.cancelRequested ? 'cancelled' : spawnError ? 'failed' : parsed.isError || invalidJsonStream || invalidAgyWorkerPolicy || missingReportedIdentity || reviewerFailed ? 'failed' : exit.code === 0 ? 'completed' : exit.signal === 'SIGTERM' ? 'incomplete' : 'failed'; + const status = active.cancelRequested ? 'cancelled' : spawnError ? 'failed' : parsed.isError || invalidJsonStream || invalidAgyWorkerPolicy || missingReportedIdentity || reviewerFailed ? 'failed' : exit.code === 0 ? 'completed' : exit.signal === 'SIGTERM' ? 'incomplete' : 'failed'; + // A held (gated) turn leaves the task in progress; the gate loop sets the final status. + if (!turn.hold) record.status = status; record.output_text = parsed.text; if (parsed.model) record.model = parsed.model; if (kind === 'codex-cli') record.metadata.actual_model_status = parsed.model ? 'observed' : 'unavailable'; @@ -1167,7 +1232,7 @@ async function runTask(record, prompt) { } if (reviewer) { record.metadata.foreman_review_mode = 'read_only'; record.metadata.reviewer_mutation_attempted = parsed.mutationAttempted === true; record.metadata.reviewer_output_overflow = cliOutputOverflow; record.metadata.reviewer_validation = record.metadata.review_evidence.controllerValidation; if (kind === 'agy') record.metadata.reviewer_boundary = { ...record.metadata.reviewer_boundary, tool_attempts_blocked: parsed.mutationAttempted === true, proven: true }; } if (kind !== 'claude' && !reviewer) record.metadata.cli_output_overflow = cliOutputOverflow; - if (record.status !== 'completed') record.error = { message: spawnError ? 'CLI could not be started' : invalidAgyWorkerPolicy ? agyWorkerPolicyFailureMessage({ ...parsed.diagnostic, ...agyWorkerPolicy }) : reviewer && parsed.mutationAttempted ? 'Reviewer attempted to use a tool or mutate state' : cliOutputOverflow ? 'CLI output exceeded the bounded stream limit' : (reviewer || kind !== 'claude') && (parsed.malformedOutput || parsed.unrecognizedOutput) ? 'CLI output contained malformed or unrecognized stream records' : kind === 'codex-cli' && !reviewer && !parsed.session ? exit.code !== 0 || exit.signal ? 'Codex CLI exited before reporting a session id' : 'Codex CLI did not report a session id' : (kind === 'codex-cli' || kind === 'agy') && !reviewer && !parsed.turnCompleted ? `${kind === 'agy' ? 'Antigravity CLI' : 'Codex CLI'} exited without completing a turn` : parsed.isError ? kind === 'claude' ? 'Claude Code reported an unsuccessful task' : kind === 'agy' ? 'Antigravity CLI reported an unsuccessful task' : 'Codex CLI reported an unsuccessful task' : missingReportedIdentity ? 'CLI did not report an actual model and session id' : active.cancelRequested ? 'CLI task was cancelled' : exit.signal ? `CLI terminated by ${exit.signal}` : 'CLI exited unsuccessfully' }; + if (status !== 'completed') record.error = { message: spawnError ? 'CLI could not be started' : invalidAgyWorkerPolicy ? agyWorkerPolicyFailureMessage({ ...parsed.diagnostic, ...agyWorkerPolicy }) : reviewer && parsed.mutationAttempted ? 'Reviewer attempted to use a tool or mutate state' : cliOutputOverflow ? 'CLI output exceeded the bounded stream limit' : (reviewer || kind !== 'claude') && (parsed.malformedOutput || parsed.unrecognizedOutput) ? 'CLI output contained malformed or unrecognized stream records' : kind === 'codex-cli' && !reviewer && !parsed.session ? exit.code !== 0 || exit.signal ? 'Codex CLI exited before reporting a session id' : 'Codex CLI did not report a session id' : (kind === 'codex-cli' || kind === 'agy') && !reviewer && !parsed.turnCompleted ? `${kind === 'agy' ? 'Antigravity CLI' : 'Codex CLI'} exited without completing a turn` : parsed.isError ? kind === 'claude' ? 'Claude Code reported an unsuccessful task' : kind === 'agy' ? 'Antigravity CLI reported an unsuccessful task' : 'Codex CLI reported an unsuccessful task' : missingReportedIdentity ? 'CLI did not report an actual model and session id' : active.cancelRequested ? 'CLI task was cancelled' : exit.signal ? `CLI terminated by ${exit.signal}` : 'CLI exited unsuccessfully' }; if (record.error && kind === 'agy') { const denied = (parsed.diagnostic?.tool_events ?? []).filter(e => e.error_category === 'permission_denied'); if (denied.length >= 3) { @@ -1179,6 +1244,7 @@ async function runTask(record, prompt) { if (persistentRole && parsed.session) { record.metadata.role_session = { binding: record.metadata.role_session_binding, cli_session_id: parsed.session, state_path: record.metadata.role_session_state_path }; } + const finalize = async () => { await persist(); touchWorkspace(ws); touchWorkspace(roWs); for (const finish of active.finishListeners ?? []) finish(); @@ -1188,6 +1254,8 @@ async function runTask(record, prompt) { } if (kind === 'agy' && taskWorkspace.agyStateDir) { await rm(taskWorkspace.agyStateDir, { recursive: true, force: true }).catch(() => undefined); taskWorkspace.agyStateDir = undefined; } if (reviewer) { await rm(work, { recursive: true, force: true }); tasks.get(record.id).reviewerWorkDir = undefined; } + }; + return { status, finalize }; } function classifyCodexFailure(stderr, exit, spawnError) { @@ -1291,7 +1359,7 @@ async function runAgySandboxed(ws, reviewer, args, env, prompt, preflight = fals const configPath = `${isolatedHome}/.gemini/antigravity-cli`; const stateDir = ws.roleStatePath ? join(ws.roleStatePath, 'agy-state') : join(ROOT, `${ws.id}.agy-state`); const writableConfigDirs = ['log','crashes','brain','conversations','cache','updater','presence','annotations','implicit','scratch']; - await mkdir(stateDir, { recursive: persistentContext, mode: 0o700 }); + await mkdir(stateDir, { recursive: persistentContext || !!ws.workerSessionPath, mode: 0o700 }); const stateMounts = []; for (const name of writableConfigDirs) { const source = join(stateDir, name); await mkdir(source, { recursive: true, mode: 0o700 }); @@ -1421,7 +1489,7 @@ const server = createServer(async (req, res) => { let url; try { url = new URL(req.url, 'http://localhost'); } catch { return send(res, 400, { error: { code: 'invalid_url' } }); } try { - if (req.method === 'GET' && url.pathname === '/v1/uhp') return send(res, 200, { protocol: 'uhp', versions: [VERSION], default_version: VERSION, implementation: { name: 'local-cli-uhp', experimental: true }, capabilities: { streaming: true, idempotency: true, sessions: true, cancellation: true, readOnlyReviewer: true, extensions: { foreman_workspace_bridge_v1: { version: 1, seed: !!SOURCE_REPO, complete_snapshot: true, execution_boundary: 'bubblewrap' } } } }); + if (req.method === 'GET' && url.pathname === '/v1/uhp') return send(res, 200, { protocol: 'uhp', versions: [VERSION], default_version: VERSION, implementation: { name: 'local-cli-uhp', experimental: true }, capabilities: { streaming: true, idempotency: true, sessions: true, cancellation: true, readOnlyReviewer: true, extensions: { foreman_workspace_bridge_v1: { version: 1, seed: !!SOURCE_REPO, complete_snapshot: true, execution_boundary: 'bubblewrap' }, foreman_worker_gate_v1: { version: 1, max_rounds: MAX_WORKER_GATE_ROUNDS } } } }); if (req.method === 'POST' && url.pathname === '/extensions/foreman-workspace/v1/workspaces') { const b = await body(req); const ws = await seedWorkspace(b.base_commit); return send(res, 201, { workspace_id: ws.id, base_commit: ws.baseCommit }); @@ -1448,9 +1516,11 @@ const server = createServer(async (req, res) => { const ws = workspaces.get(decodeURIComponent(snapshotPath[1])); if (!ws) return send(res, 404, { error: { code: 'workspace_not_found' } }); const response = ws.responseId ? state.responses[ws.responseId] : undefined; - if (response?.status === 'in_progress') return send(res, 409, { error: { code: 'workspace_task_in_progress' } }); + // A Worker paused for Foreman's checks has no CLI running, so its workspace can be read for those checks. + const pausedForChecks = response?.status === 'in_progress' && !!tasks.get(ws.responseId)?.gate?.pending; + if (response?.status === 'in_progress' && !pausedForChecks) return send(res, 409, { error: { code: 'workspace_task_in_progress' } }); const snapshot = await withWorkspaceOp(ws, () => snapshotWorkspace(ws)); - if (ws.responseId && response?.status !== 'completed') { + if (ws.responseId && response?.status !== 'completed' && !pausedForChecks) { snapshot.complete = false; snapshot.errors.push({ path: '.', error: `task_status_${response?.status ?? 'unknown'}` }); } @@ -1474,8 +1544,18 @@ const server = createServer(async (req, res) => { } const responsePath = url.pathname.match(/^\/v1\/responses\/([^/]+)$/); if (req.method === 'GET' && responsePath) { const r = state.responses[decodeURIComponent(responsePath[1])]; return r ? send(res, 200, r) : send(res, 404, { error: { code: 'not_found' } }); } + const gatePath = url.pathname.match(/^\/extensions\/foreman-workspace\/v1\/responses\/([^/]+)\/worker-gate$/); + if (req.method === 'POST' && gatePath) { + const t = tasks.get(decodeURIComponent(gatePath[1])); + if (!t?.gate) return send(res, 404, { error: { code: 'worker_gate_not_found' } }); + let verdict; + try { verdict = parseWorkerGateVerdict(await body(req)); } catch (error) { return send(res, 400, { error: { code: 'worker_gate_verdict_invalid', message: error.message } }); } + if (t.gate.pending?.round !== verdict.round) return send(res, 409, { error: { code: 'worker_gate_not_awaiting', message: `the task is not waiting for check round ${verdict.round}` } }); + t.gate.pending.settle(verdict); + return send(res, 202, { id: t.response.id, round: verdict.round, accepted: true }); + } const cancelPath = url.pathname.match(/^\/v1\/responses\/([^/]+)\/cancel$/); - if (req.method === 'POST' && cancelPath) { const t = tasks.get(decodeURIComponent(cancelPath[1])); if (!t) return send(res, 404, { error: { code: 'not_found' } }); if (t.child && t.response.status === 'in_progress') { t.cancelRequested = true; t.child.kill('SIGTERM'); } return send(res, 200, { status: t.response.status === 'in_progress' ? 'cancelling' : t.response.status }); } + if (req.method === 'POST' && cancelPath) { const t = tasks.get(decodeURIComponent(cancelPath[1])); if (!t) return send(res, 404, { error: { code: 'not_found' } }); if (t.child && t.response.status === 'in_progress') { t.cancelRequested = true; t.child.kill('SIGTERM'); } if (t.gate?.pending) { t.cancelRequested = true; t.gate.pending.settle({ round: t.gate.pending.round, status: 'skipped', reason: 'cancelled' }); } return send(res, 200, { status: t.response.status === 'in_progress' ? 'cancelling' : t.response.status }); } if (req.method === 'POST' && url.pathname === '/v1/responses') { const key = req.headers['idempotency-key']; if (typeof key !== 'string' || !key) return send(res, 400, { error: { code: 'idempotency_key_required' } }); const prior = state.keys[key]; @@ -1497,6 +1577,9 @@ const server = createServer(async (req, res) => { const roleId = b.metadata?.foreman_role_id ?? b.metadata?.role_id; const persistentRole = roleId === 'planner' || roleId === 'orchestrator'; const roWorkspaceId = typeof b.metadata?.foreman_read_only_workspace_id === 'string' ? b.metadata.foreman_read_only_workspace_id : undefined; + let workerGate; + try { workerGate = parseWorkerGateRequest(b.metadata.foreman_worker_gate); } catch (error) { return send(res, 400, { error: { code: 'worker_gate_invalid', message: error.message } }); } + if (workerGate && (reviewer || persistentRole)) return send(res, 400, { error: { code: 'worker_gate_role_unsupported' } }); let roleSession; if (b.previous_response_id !== undefined && !persistentRole) return send(res, 400, { error: { code: 'session_role_unsupported' } }); if (persistentRole) { @@ -1536,7 +1619,7 @@ const server = createServer(async (req, res) => { const record = { id, object: 'response', status: 'in_progress', requested_model: b.model, metadata: { ...b.metadata, harness_id: h.id, ...(workspaceId ? { workspace_id: workspaceId } : {}) }, timeout_seconds: timeout, max_step: steps }; if (ws) { ws.responseId = id; state.workspaces[ws.id].responseId = id; touchWorkspace(ws); } touchWorkspace(workspaces.get(roWorkspaceId)); - state.keys[key] = id; state.responses[id] = record; const task = { response: record, finishListeners: new Set(), activity: [], activityTotal: 0, activityListeners: new Set(), activityRemainder: '' }; recordActivity(task, 'Preparing isolated runtime'); tasks.set(id, task); + state.keys[key] = id; state.responses[id] = record; const task = { response: record, finishListeners: new Set(), activity: [], activityTotal: 0, activityListeners: new Set(), activityRemainder: '', ...(workerGate ? { gate: workerGate } : {}) }; recordActivity(task, 'Preparing isolated runtime'); tasks.set(id, task); await persist(); // durable idempotency intent before spawning any CLI process runTask(record, b.input).catch(async error => { record.status = 'failed'; record.error = { message: 'CLI task failed internally' }; if(task.reviewerWorkDir) { await rm(task.reviewerWorkDir,{recursive:true,force:true}).catch(()=>undefined); task.reviewerWorkDir=undefined; } const ws = workspaces.get(record.metadata.workspace_id); if (ws) { record.metadata.execution_boundary = ws.boundary ?? { proven: false, error: 'execution_boundary_unavailable' }; record.metadata.execution_stage = ws.executionStage ?? 'task_setup'; if (ws.boundaryDiagnostic) record.metadata.execution_boundary_diagnostic = ws.boundaryDiagnostic; else if (/^boundary_probe_failed:/.test(String(error?.message))) record.metadata.execution_boundary_diagnostic = { category: String(error.message).split(':')[1] ?? 'probe_failed', exit_code: Number(String(error.message).split(':')[2]) || null, signal: null }; if (record.metadata.harness_id === 'codex-cli' && record.metadata.foreman_review_mode !== 'read_only') { record.metadata.cli_exit = { exit_code: null, signal: null }; record.metadata.cli_failure_category = record.metadata.execution_stage === 'resolve_auth' ? 'host_auth_unavailable' : record.metadata.execution_stage === 'boundary_probe' ? 'boundary_setup' : record.metadata.execution_stage === 'runtime_mount' ? 'runtime_setup' : 'cli_setup'; await rm(ws.codexHome, { recursive: true, force: true }).catch(() => undefined); } } await persist(); for (const finish of task.finishListeners) finish(); }); return streamResponse(res, record, task); @@ -1559,7 +1642,7 @@ function streamResponse(res, record, task = tasks.get(record.id)) { const activity = items.find(item => item.index === activityCursor); activityCursor++; if (!activity) continue; - event(res, 'response.activity', sequence++, { id: record.id, object: 'response', status: 'in_progress', activity: { kind: activity.kind, summary: activity.summary } }); + event(res, 'response.activity', sequence++, { id: record.id, object: 'response', status: 'in_progress', activity: { kind: activity.kind, summary: activity.summary, ...(activity.gate_round ? { gate_round: activity.gate_round } : {}) } }); } }; const stopHeartbeat = () => { if (heartbeat) { clearInterval(heartbeat); heartbeat = undefined; } }; diff --git a/investigations/local-cli-uhp/test.mjs b/investigations/local-cli-uhp/test.mjs index c2c7582..147c10d 100644 --- a/investigations/local-cli-uhp/test.mjs +++ b/investigations/local-cli-uhp/test.mjs @@ -1221,3 +1221,130 @@ test('the TTL sweep removes idle workspaces but never one with a running respons assert.ok(await until(async () => !(await exists(busyDir))), 'the finished workspace was never swept once idle'); const persisted = JSON.parse(await readFile(env.LOCAL_CLI_UHP_STATE, 'utf8')).workspaces; assert.equal(persisted[idle], undefined); assert.equal(persisted[busy], undefined); }); + +// Worker check gate: the bridge pauses a completed Worker turn, Foreman checks the paused workspace, and a failed verdict +// resumes the same CLI session with the failures. +async function streamGatedWorker(base, { harness, model, key, workspaceId, rounds = 2, onGate }) { + const r = await fetch(`${base}/v1/responses`, { method:'POST', headers:{ 'Content-Type':'application/json', Accept:'text/event-stream', 'UHP-Version':'2026-09-12', 'Idempotency-Key':key }, body:JSON.stringify({ input:'Change README.md', model, metadata:{ harness_id:harness, workspace_id:workspaceId, foreman_worker_gate:{ max_rounds:rounds } }, stream:true, timeout_seconds:5, max_step:1 }) }); + if (r.status !== 200) assert.fail(await r.text()); + const reader = r.body.getReader(), decoder = new TextDecoder(), events = []; let buffer = ''; + for (;;) { + const { value, done } = await reader.read(); if (done) break; + buffer += decoder.decode(value, { stream:true }); + let cut; while ((cut = buffer.indexOf('\n\n')) >= 0) { + const chunk = buffer.slice(0, cut); buffer = buffer.slice(cut + 2); + const line = chunk.split('\n').find(x => x.startsWith('data: ')); if (!line) continue; + const item = JSON.parse(line.slice(6)); events.push(item); + if (item.type === 'response.activity' && item.response.activity.kind === 'worker_gate') await onGate(item.response.activity.gate_round, item.response.id); + } + } + return events; +} +const postVerdict = (base, responseId, verdict) => fetch(`${base}/extensions/foreman-workspace/v1/responses/${responseId}/worker-gate`, { method:'POST', headers:{ 'Content-Type':'application/json' }, body:JSON.stringify(verdict) }); +const seedWorkspace = async (base, baseCommit) => (await (await fetch(`${base}/extensions/foreman-workspace/v1/workspaces`, { method:'POST', headers:{ 'Content-Type':'application/json' }, body:JSON.stringify({ base_commit:baseCommit }) })).json()).workspace_id; + +test('check gate resumes a Claude Worker in its own kept session with the failure feedback', async t => { + // First turn writes a wrong README and records a transcript; the resumed turn must see that transcript and the feedback. + const claudeBody = `import {writeFileSync,readFileSync,existsSync,appendFileSync,mkdirSync} from 'node:fs'; let prompt='';process.stdin.setEncoding('utf8');process.stdin.on('data',c=>prompt+=c);process.stdin.on('end',()=>{const resume=process.argv.includes('--resume')?process.argv[process.argv.indexOf('--resume')+1]:''; const transcript=process.env.CLAUDE_CONFIG_DIR+'/projects/gate-transcript'; const kept=existsSync(transcript); mkdirSync(process.env.CLAUDE_CONFIG_DIR+'/projects',{recursive:true}); writeFileSync(transcript,'turn'); appendFileSync('.gate-log',JSON.stringify({resume,kept,feedback:prompt.includes('lint failed: README must say fixed')})+'\\n'); writeFileSync('README.md',resume?'fixed\\n':'first\\n'); const model='claude-actual'; console.log(JSON.stringify({type:'system',subtype:'init',model,session_id:'gate-session'})); console.log(JSON.stringify({type:'result',subtype:'success',is_error:false,result:resume?'Fixed the check.':'Changed README.',model,session_id:'gate-session',usage:{input_tokens:7,output_tokens:3,cache_read_input_tokens:2}}));});`; + const { base, baseCommit, env } = await setup(t, { claudeBody }); + const workspaceId = await seedWorkspace(base, baseCommit); + const seen = []; + const events = await streamGatedWorker(base, { harness:'claude-code', model:'claude-requested', key:'gate-claude', workspaceId, onGate: async (round, responseId) => { + const snapshot = await (await fetch(`${base}/extensions/foreman-workspace/v1/workspaces/${workspaceId}/snapshot`)).json(); + const readme = Buffer.from(snapshot.entries.find(e => e.path === 'README.md').contentBase64, 'base64').toString('utf8'); + seen.push({ round, complete:snapshot.complete, errors:snapshot.errors.length, readme }); + assert.equal((await postVerdict(base, responseId, { round:round + 1, status:'passed' })).status, 409, 'a verdict for another round is refused'); + const verdict = round === 1 ? { round, status:'failed', feedback:'lint failed: README must say fixed', failed_checks:['lint'] } : { round, status:'passed' }; + assert.equal((await postVerdict(base, responseId, verdict)).status, 202); + } }); + const final = events.at(-1); + assert.equal(final.type, 'response.completed', JSON.stringify(final)); + assert.deepEqual(seen, [{ round:1, complete:true, errors:0, readme:'first\n' }, { round:2, complete:true, errors:0, readme:'fixed\n' }]); + assert.deepEqual(final.response.metadata.worker_gate, { max_rounds:2, state:'finished', rounds:[{ round:1, status:'failed', failed_checks:['lint'] }, { round:2, status:'passed' }] }); + assert.equal(final.response.session_id, 'gate-session'); + assert.deepEqual(final.response.usage, { input_tokens:14, output_tokens:6, input_tokens_details:{ cached_tokens:4 } }); + const work = join(env.LOCAL_CLI_UHP_WORK, workspaceId); + const log = (await readFile(join(work, '.gate-log'), 'utf8')).trim().split('\n').map(line => JSON.parse(line)); + assert.deepEqual(log, [{ resume:'', kept:false, feedback:false }, { resume:'gate-session', kept:true, feedback:true }]); + assert.equal(await exists(`${work}.worker-session`), false, 'the kept session is removed when the task ends'); + const snapshot = await (await fetch(`${base}/extensions/foreman-workspace/v1/workspaces/${workspaceId}/snapshot`)).json(); + assert.equal(snapshot.complete, true); +}); + +test('check gate ends the Worker turn as it is on a passed first verdict, a timeout, or a cancellation', async t => { + const claudeBody = `import {writeFileSync,appendFileSync} from 'node:fs'; process.stdin.resume(); process.stdin.on('end',()=>{appendFileSync('.gate-count','x'); writeFileSync('README.md','edited\\n'); const model='claude-actual'; console.log(JSON.stringify({type:'system',subtype:'init',model,session_id:'s'})); console.log(JSON.stringify({type:'result',subtype:'success',is_error:false,result:'ok',model,session_id:'s',usage:{input_tokens:1,output_tokens:1}}));});`; + const { base, baseCommit, env } = await setup(t, { claudeBody, extraEnv:{ LOCAL_CLI_UHP_WORKER_GATE_TIMEOUT_MS:'300' } }); + const passedWs = await seedWorkspace(base, baseCommit); + const passed = await streamGatedWorker(base, { harness:'claude-code', model:'claude-requested', key:'gate-pass', workspaceId:passedWs, onGate: async (round, id) => { assert.equal((await postVerdict(base, id, { round, status:'passed' })).status, 202); } }); + assert.equal(passed.at(-1).type, 'response.completed'); + assert.deepEqual(passed.at(-1).response.metadata.worker_gate.rounds, [{ round:1, status:'passed' }]); + assert.equal(await readFile(join(env.LOCAL_CLI_UHP_WORK, passedWs, '.gate-count'), 'utf8'), 'x'); + const timedWs = await seedWorkspace(base, baseCommit); + const timed = await streamGatedWorker(base, { harness:'claude-code', model:'claude-requested', key:'gate-timeout', workspaceId:timedWs, onGate: async () => {} }); + assert.equal(timed.at(-1).type, 'response.completed'); + assert.deepEqual(timed.at(-1).response.metadata.worker_gate.rounds, [{ round:1, status:'skipped', reason:'verdict_timeout' }]); + const cancelWs = await seedWorkspace(base, baseCommit); + const cancelled = await streamGatedWorker(base, { harness:'claude-code', model:'claude-requested', key:'gate-cancel', workspaceId:cancelWs, onGate: async (round, id) => { assert.equal((await fetch(`${base}/v1/responses/${id}/cancel`, { method:'POST' })).status, 200); } }); + assert.equal(cancelled.at(-1).type, 'response.cancelled', JSON.stringify(cancelled.at(-1))); +}); + +test('check gate requests are validated and refused for Reviewer and role sessions', async t => { + const { base, baseCommit } = await setup(t); + const workspaceId = await seedWorkspace(base, baseCommit); + const post = (key, metadata) => fetch(`${base}/v1/responses`, { method:'POST', headers:{ 'Content-Type':'application/json', 'UHP-Version':'2026-09-12', 'Idempotency-Key':key }, body:JSON.stringify({ input:'x', model:'claude-requested', metadata:{ harness_id:'claude-code', ...metadata }, stream:true, timeout_seconds:5, max_step:1 }) }); + assert.equal((await (await post('gate-bad-rounds', { workspace_id:workspaceId, foreman_worker_gate:{ max_rounds:9 } })).json()).error.code, 'worker_gate_invalid'); + assert.equal((await (await post('gate-planner', { foreman_role_id:'planner', foreman_run_id:'run_1', foreman_worker_gate:{ max_rounds:1 } })).json()).error.code, 'worker_gate_role_unsupported'); + const discovery = await (await fetch(`${base}/v1/uhp`)).json(); + assert.deepEqual(discovery.capabilities.extensions.foreman_worker_gate_v1, { version:1, max_rounds:3 }); + assert.equal((await postVerdict(base, 'resp_missing', { round:1, status:'passed' })).status, 404); +}); + +test('Foreman answers check rounds so a Codex Worker fixes a failing check in its kept session within one attempt', async t => { + const fixture = await createWorkspaceFixture(); t.after(fixture.cleanup); + const parent = await mkdtemp(join(tmpdir(),'foreman-worker-gate-')); t.after(()=>rm(parent,{recursive:true,force:true})); + const codexBody = `import {readFileSync,writeFileSync,existsSync,appendFileSync,mkdirSync} from 'node:fs'; let prompt='';process.stdin.setEncoding('utf8');process.stdin.on('data',c=>prompt+=c);process.stdin.on('end',()=>{const i=process.argv.indexOf('resume'); const resume=i>=0?process.argv[i+1]:''; const marker=process.env.CODEX_HOME+'/sessions/gate'; const kept=existsSync(marker); mkdirSync(process.env.CODEX_HOME+'/sessions',{recursive:true}); writeFileSync(marker,'x'); appendFileSync('.gate-log',JSON.stringify({resume,kept,ephemeral:process.argv.includes('--ephemeral'),feedback:prompt.includes('assert fixed README')&&prompt.includes('README is not fixed yet')})+'\\n'); writeFileSync('README.md',resume?'# Fixture\\n\\nfixed\\n':'# Fixture\\n\\nfirst try\\n'); console.log(JSON.stringify({type:'thread.started',thread_id:'codex-gate-thread'})); console.log(JSON.stringify({type:'item.completed',item:{type:'agent_message',text:'Edited README.md.'}})); console.log(JSON.stringify({type:'turn.completed',usage:{input_tokens:11,output_tokens:7}}));});`; + const claudeBody=`let prompt='';process.stdin.setEncoding('utf8');process.stdin.on('data',chunk=>prompt+=chunk);process.stdin.on('end',()=>{const model='claude-actual';const result=prompt.includes('"workerTask"')&&prompt.includes('"targetFiles"')?JSON.stringify({workerTask:'Change README.md with one short sentence.',targetFiles:['README.md']}):'Planner recommends a concise README note.';console.log(JSON.stringify({type:'system',subtype:'init',model,session_id:'judgment-session'}));console.log(JSON.stringify({type:'result',subtype:'success',result,model,session_id:'judgment-session',usage:{input_tokens:5,output_tokens:2}}));});`; + const {base,env}=await setup(t,{sourceRepo:fixture.repo,baseCommit:fixture.baseCommit,codexBody,claudeBody}); + const uhp = new UhpClient({baseUrl:base,timeoutMs:20_000}); + const controller = new Controller(new JsonStore(join(parent,'foreman-state.json')),uhp,false,true); + const check = `const fs=require('node:fs');if(!fs.readFileSync('README.md','utf8').includes('fixed')){console.log('README is not fixed yet');process.exit(1)}`; + controller.configureVerifiedWorkspace({repoPath:fixture.repo,allowedScope:['README.md','.gate-log'],commands:[{name:'assert fixed README',command:process.execPath,args:['-e',check]}],bridgeBaseUrl:base,timeoutMs:10_000,maxOutputBytes:2_000,workerCheckRounds:2}); + await controller.refreshDiscovery(); + const project=await controller.createProject('Worker check gate fixture'); + const task=await controller.createTask(project.id,'Edit the assigned README'); + const run=await controller.createRun(task.id); + for (const role of ['planner','orchestrator']) await controller.selectRoleConfig(role,{harnessId:'claude-code',model:'claude-requested',options:{timeoutSeconds:5,maxStep:1}},undefined,run.id); + await controller.selectRoleConfig('worker',{harnessId:'codex-cli',model:'codex-requested',options:{timeoutSeconds:5,maxStep:1}},undefined,run.id); + await controller.prepareWorkerWorkspace(run.id,fixture.baseCommit); + await controller.addGuidance(run.id,'Keep the README change to one short note.'); + const orchestration=await controller.orchestrate(run.id,'Implement the requested README note.'); + assert.ok(orchestration.proposal, JSON.stringify(orchestration.assignment)); + const assignment=await controller.dispatchWorkerProposal(run.id,orchestration.proposal.id); + assert.equal(assignment.status,'succeeded',JSON.stringify(assignment)); + assert.equal(assignment.usage.inputTokens,22,'both Codex turns are counted'); + assert.deepEqual(assignment.cliInvocation?.args.slice(-3),['--model','codex-requested','-']); + const response=await uhp.retrieve(assignment.responseId); + assert.deepEqual(response.metadata.worker_gate.rounds,[{round:1,status:'failed',failed_checks:['assert fixed README']},{round:2,status:'passed'}]); + const log=(await readFile(join(env.LOCAL_CLI_UHP_WORK,response.metadata.workspace_id,'.gate-log'),'utf8')).trim().split('\n').map(line=>JSON.parse(line)); + assert.deepEqual(log,[{resume:'',kept:false,ephemeral:false,feedback:false},{resume:'codex-gate-thread',kept:true,ephemeral:false,feedback:true}]); + const state=await controller.state(); + assert.deepEqual(state.events.filter(e=>e.type==='worker.check_round').map(e=>[e.data.round,e.data.status,e.data.failedChecks]),[[1,'failed',['assert fixed README']],[2,'passed',[]]]); + const verified=await controller.verifyWorkerOutput(run.id,assignment.id); + assert.equal(verified.validation.status,'passed',JSON.stringify(verified.validation)); + assert.equal((await controller.state()).projects[0].tasks[0].runs[0].assignments.filter(a=>a.roleId==='worker').length,1,'the fix used no second Worker attempt'); +}); + +test('check gate resumes an Antigravity Worker in its kept conversation', async t => { + const agyBody = `import {writeFileSync,existsSync,appendFileSync} from 'node:fs'; if(process.argv[2]==='models'){console.log('gemini-3.8-flash-medium\\tGemini 3.8 Flash (Medium)');process.exit(0)} const ci=process.argv.indexOf('--conversation'); const resume=ci>=0?process.argv[ci+1]:''; const prompt=process.argv[process.argv.indexOf('-p')+1]||''; const marker=process.env.HOME+'/.gemini/antigravity-cli/conversations/gate'; const kept=existsSync(marker); writeFileSync(marker,'x'); appendFileSync('.gate-log',JSON.stringify({resume,kept,feedback:prompt.includes('typecheck failed here')})+'\\n'); writeFileSync('README.md',resume?'fixed\\n':'first\\n'); const conversation_id='agy-gate-conversation'; const model='gemini-3.8-flash-medium'; console.log(JSON.stringify({event:'init',conversation_id,agent:'foreman-worker',init:{cwd:'/workspace',model,tools:['view_file','write_to_file','finish']}})); console.log(JSON.stringify({event:'step_update',step_update:{conversation_id,step_index:0,state:'DONE',step_type:'tool',tool_name:'write_to_file',tool_info:{name:'write_to_file'}}})); console.log(JSON.stringify({event:'result',result:{conversation_id,status:'SUCCESS',response:'Edited README.',model,usage:{input_tokens:resume?30:10,output_tokens:resume?8:4}}}));`; + const { base, baseCommit, env } = await setup(t, { agyEnabled:true, agyBody }); + const workspaceId = await seedWorkspace(base, baseCommit); + const events = await streamGatedWorker(base, { harness:'antigravity-cli', model:'gemini-3.8-flash-medium', key:'gate-agy', workspaceId, rounds:1, onGate: async (round, id) => { + assert.equal((await postVerdict(base, id, { round, status:'failed', feedback:'typecheck failed here', failed_checks:['typecheck'] })).status, 202); + } }); + const final = events.at(-1); + assert.equal(final.type, 'response.completed', JSON.stringify(final)); + assert.deepEqual(final.response.metadata.worker_gate.rounds, [{ round:1, status:'failed', failed_checks:['typecheck'] }]); + assert.equal(final.response.usage.input_tokens, 30, 'Antigravity reports cumulative conversation usage, so the last turn is the total'); + const log = (await readFile(join(env.LOCAL_CLI_UHP_WORK, workspaceId, '.gate-log'), 'utf8')).trim().split('\n').map(line => JSON.parse(line)); + assert.deepEqual(log, [{ resume:'', kept:false, feedback:false }, { resume:'agy-gate-conversation', kept:true, feedback:true }]); +}); diff --git a/investigations/local-cli-uhp/worker-gate.mjs b/investigations/local-cli-uhp/worker-gate.mjs new file mode 100644 index 0000000..652fd49 --- /dev/null +++ b/investigations/local-cli-uhp/worker-gate.mjs @@ -0,0 +1,44 @@ +// Worker check gate: after a Worker turn completes, the bridge pauses the task and waits for Foreman to run the +// repository's checks on the workspace. When Foreman reports failures, the bridge resumes the same CLI session with +// them, so any harness can fix its own check failures without a shell. Foreman's own validation after the turn stays +// the authoritative result; a gate verdict is feedback only. + +export const MAX_WORKER_GATE_ROUNDS = 3; +export const MAX_WORKER_GATE_FEEDBACK = 8_000; + +/** Validated `metadata.foreman_worker_gate`, or undefined when absent. Throws on a malformed request. */ +export function parseWorkerGateRequest(value) { + if (value === undefined) return undefined; + if (!value || typeof value !== 'object' || Array.isArray(value)) throw Error('foreman_worker_gate must be an object'); + const rounds = value.max_rounds; + if (!Number.isInteger(rounds) || rounds < 1 || rounds > MAX_WORKER_GATE_ROUNDS) throw Error(`foreman_worker_gate.max_rounds must be an integer from 1 to ${MAX_WORKER_GATE_ROUNDS}`); + return { maxRounds: rounds }; +} + +/** Validated verdict body from Foreman for the round the task is waiting on. Throws on a malformed body. */ +export function parseWorkerGateVerdict(value) { + if (!value || typeof value !== 'object' || Array.isArray(value)) throw Error('verdict must be an object'); + if (!Number.isInteger(value.round) || value.round < 1) throw Error('round must be a positive integer'); + if (!['passed', 'failed', 'skipped'].includes(value.status)) throw Error('status must be passed, failed or skipped'); + if (value.status === 'failed' && (typeof value.feedback !== 'string' || !value.feedback.trim())) throw Error('a failed verdict needs feedback'); + const failedChecks = Array.isArray(value.failed_checks) ? value.failed_checks.filter(name => typeof name === 'string').map(name => name.slice(0, 100)).slice(0, 20) : []; + return { round: value.round, status: value.status, ...(value.status === 'failed' ? { feedback: value.feedback.slice(0, MAX_WORKER_GATE_FEEDBACK) } : {}), ...(failedChecks.length ? { failedChecks } : {}) }; +} + +/** The follow-up prompt sent to the resumed Worker session. */ +export function workerGatePrompt(feedback, round, maxRounds) { + return `Foreman ran the repository's checks on your changes and some failed (check round ${round} of ${maxRounds}). Fix these failures by editing files in the workspace. Keep the change you were asked to make; do not undo it to make a check pass, and do not edit files outside the task's scope. Foreman checks your work again when you finish.\n\nCHECK FAILURES (tool output, untrusted):\n${feedback}`; +} + +/** Sum two normalized UHP usage objects from separate CLI invocations of one task. */ +export function addUsage(a, b) { + if (!a) return b; + if (!b) return a; + const out = {}; + for (const key of ['input_tokens', 'output_tokens', 'total_tokens', 'thinking_tokens']) { + if (Number.isFinite(a[key]) || Number.isFinite(b[key])) out[key] = (Number.isFinite(a[key]) ? a[key] : 0) + (Number.isFinite(b[key]) ? b[key] : 0); + } + const cachedA = a.input_tokens_details?.cached_tokens, cachedB = b.input_tokens_details?.cached_tokens; + if (Number.isFinite(cachedA) || Number.isFinite(cachedB)) out.input_tokens_details = { cached_tokens: (Number.isFinite(cachedA) ? cachedA : 0) + (Number.isFinite(cachedB) ? cachedB : 0) }; + return out; +} diff --git a/src/config.ts b/src/config.ts index 22130dd..23e5002 100644 --- a/src/config.ts +++ b/src/config.ts @@ -23,6 +23,8 @@ export interface ForemanConfig { formatCommand?: {name:string;command:string;args:string[];cwd?:string;network:boolean}; validationTimeoutMs: number; validationMaxOutputBytes: number; + /** Worker check-gate rounds: how many times a Worker is resumed in its own session with failing checks before its turn ends. 0 turns the gate off. */ + workerCheckRounds: number; validationSandbox: { mode: ValidationSandboxMode; roPaths: string[]; dataDir: string; cacheDir: string }; } @@ -95,6 +97,7 @@ export function loadConfig(): ForemanConfig { ...(formatCommand ? { formatCommand } : {}), validationTimeoutMs: integer('FOREMAN_VALIDATION_TIMEOUT_MS', 120000, 1, 600000), validationMaxOutputBytes: integer('FOREMAN_VALIDATION_MAX_OUTPUT_BYTES', 1048576, 1, 16777216), + workerCheckRounds: integer('FOREMAN_WORKER_CHECK_ROUNDS', 2, 0, 3), validationSandbox: { mode: sandboxMode, roPaths: sandboxRoPaths.map(path => resolve(path)), dataDir, cacheDir: resolve(dataDir, 'validation-cache') } }; } diff --git a/src/controller.ts b/src/controller.ts index 9943b72..5803c52 100644 --- a/src/controller.ts +++ b/src/controller.ts @@ -13,7 +13,7 @@ const AUTOMATIC_RUN=Symbol('Foreman automatic run'); // Keep headroom for any adapter-side framing and check at every submit edge. const UHP_PROMPT_SAFE_LIMIT=15_000; function assertUhpPromptWithinLimit(prompt:string,label='UHP prompt'):void{const bytes=Buffer.byteLength(prompt,'utf8');if(bytes>UHP_PROMPT_SAFE_LIMIT)throw Object.assign(new Error(`${label} is ${bytes} UTF-8 bytes; the local UHP bridge safe limit is ${UHP_PROMPT_SAFE_LIMIT} bytes (bridge hard limit: 16,000 characters). Shorten the human request or reduce supplied context and retry. No provider request was sent.`),{statusCode:413});} -import { verifyWorkerSnapshot, reconstructRecordedSnapshot, validateWorkerOutput, fetchBridgeSnapshot, seedBridgeWorkspace, overlayBridgeWorkspace, formatReviewDiff, type OverlayEntry, type VerifiedWorkerWorkspace, type ValidationCommand, type ChangeEvidence, type ValidationCheckCallbacks, type ValidationObservation } from './verified-workspace.js'; +import { verifyWorkerSnapshot, reconstructRecordedSnapshot, validateWorkerOutput, fetchBridgeSnapshot, postWorkerGateVerdict, seedBridgeWorkspace, overlayBridgeWorkspace, formatReviewDiff, type OverlayEntry, type VerifiedWorkerWorkspace, type ValidationCommand, type ChangeEvidence, type ValidationCheckCallbacks, type ValidationObservation } from './verified-workspace.js'; import { snapshotGitCommit } from './git-workspace.js'; import { ensureBaseline, markFailsOnBase, baseFailureStopReason, baseFailureNote, workerCausedFailures } from './baseline-validation.js'; import type { WorkerEvidence, ValidationCheck } from './domain.js'; @@ -21,6 +21,7 @@ import { buildRepoDigest, extractKeywords } from './repo-digest.js'; import { promoteSnapshotToGit } from './git-promotion.js'; import { parseCiScripts, ciChecksNotConfigured as computeCiChecksNotConfigured } from './repository-inspector.js'; import { defaultNetworkAccess, validationSandboxStatus, type ValidationSandboxConfig } from './validation-sandbox.js'; +import { runWorkerCheckRound, type WorkerCheckRound } from './worker-check-gate.js'; const execFileAsync=promisify(execFile); export interface UhpAdapter { @@ -35,17 +36,19 @@ const must = (item: T | undefined, what: string): T => { if (!item) throw Obj const canonicalJson=(value:unknown):string=>JSON.stringify(value&&typeof value==='object'&&!Array.isArray(value)?Object.fromEntries(Object.entries(value as Record).filter(([,v])=>v!==undefined).sort(([a],[b])=>a.localeCompare(b)).map(([k,v])=>[k,JSON.parse(canonicalJson(v))])):Array.isArray(value)?value.map(v=>JSON.parse(canonicalJson(v))):value); const sha256=(value:unknown):string=>createHash('sha256').update(canonicalJson(value)).digest('hex'); const taskSpecDigest=(task:Task):string=>sha256({title:task.title,goal:task.goal??task.title,suggestedAllowedPaths:task.suggestedAllowedPaths??[],validationCriteria:task.validationCriteria??[],dependsOn:task.dependsOn??[]}); -/** `formats` says Foreman runs the configured formatter over the Worker's changed files afterwards, so the Worker need not hand-format. */ -function workerCapabilityStatement(harnessId:string,formats=false):string{ +/** `formats` says Foreman runs the configured formatter over the Worker's changed files afterwards, so the Worker need not hand-format. `checks` says the Worker check gate is on. */ +function workerCapabilityStatement(harnessId:string,formats=false,checks=false):string{ const note=formats?' Foreman also runs the repository formatter on the files you changed after you finish, so do not spend effort hand-formatting them.':''; - return workerHarnessCapabilities(harnessId)+note; + const gate=checks?' When the Worker finishes, Foreman runs the configured checks on its changes and sends any failures back to the same Worker session to fix, so do not ask the Worker to run or imitate checks.':''; + return workerHarnessCapabilities(harnessId)+note+gate; } function workerHarnessCapabilities(harnessId:string):string{ // claude-code bridge args: --tools Read,Edit,Write --permission-mode acceptEdits (no Bash) if(harnessId==='antigravity-cli')return 'The Worker can only read and edit files (view_file, write_to_file, replace_file_content, multi_replace_file_content). It cannot run shell commands, formatters, tests, or package scripts. Foreman runs the configured checks afterwards.'; if(harnessId==='claude-code')return 'The Worker can only read and edit files (Read, Edit, Write). It cannot run shell commands, formatters, tests, or package scripts. Foreman runs the configured checks afterwards.'; // codex-cli bridge args: --sandbox workspace-write — shell execution is permitted inside Codex's sandboxed workspace - if(harnessId==='codex-cli')return 'The Worker can read and edit files and run shell commands inside an isolated sandboxed workspace (Codex workspace-write sandbox). It cannot access the network or system paths outside the workspace. Foreman runs the configured checks afterwards.'; + // Its workspace is the commit's files only (no installed dependencies) and its sandbox has no network, so the repository's own tools are not available to it either. + if(harnessId==='codex-cli')return 'The Worker can read and edit files and run shell commands inside an isolated sandboxed workspace (Codex workspace-write sandbox). It cannot access the network, system paths outside the workspace, or installed dependencies, so the repository\'s formatters, linters and tests are not available to it. Foreman runs the configured checks afterwards.'; return `The Worker uses the ${harnessId} harness.`; } function buildOverlayEntries(changes:import('./workspace-snapshot.js').SnapshotChange[]):OverlayEntry[]{ @@ -62,7 +65,9 @@ export class Controller { private readonly automaticRuns = new Set(); private readonly taskStarts = new Set(); private readonly projectPlannerTurns = new Set(); - private verifiedWorkspaceConfig?:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];formatCommand?:ValidationCommand;bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig;bridgeToken?:string}; + private verifiedWorkspaceConfig?:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];formatCommand?:ValidationCommand;bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig;bridgeToken?:string;workerCheckRounds?:number}; + /** Worker check-gate rounds already answered or in progress, keyed by bridge response and round. */ + private readonly workerGateRounds = new Set(); private discovery?: { version:string; capabilities:Record; harnesses:Array<{id:string;models?:Array<{id:string;available?:boolean}>}>; models?:Array<{id:string;harnessId?:string;available?:boolean}> }; private discoveryError?: string; private memoryServiceState: 'ready'|'degraded'|'unavailable'|'not_configured'; @@ -88,8 +93,10 @@ export class Controller { return this.taskTimeoutSeconds; } async state(): Promise { const state=await this.store.load();for(const project of state.projects){const projectAssignments=[...(project.plannerAssignments??[]),...project.tasks.flatMap(t=>t.runs.flatMap(r=>r.assignments))];project.usage=aggregate(projectAssignments.map(a=>a.usage));project.usageByHarnessModel=groupUsage(projectAssignments);for(const task of project.tasks){const latest=task.runs.at(-1);if(latest?.controller?.active||latest?.status==='awaiting_approval'||latest?.status==='completed'&&latest.promotion?.status!=='applied')task.status='active';else if(latest?.promotion?.status==='applied')task.status='completed';else if(latest?.status==='failed')task.status='blocked';for(const run of task.runs){run.usage=aggregate(run.assignments.map(a=>a.usage));run.usageByRole=Object.fromEntries([...new Set(run.assignments.map(a=>a.roleId))].map(role=>[role,aggregate(run.assignments.filter(a=>a.roleId===role).map(a=>a.usage))]).filter(([,v])=>!!v));run.usageByHarnessModel=groupUsage(run.assignments);}}}for(const role of state.roles){const assignments=[...state.projects.flatMap(p=>p.tasks.flatMap(t=>t.runs.flatMap(r=>r.assignments))),...state.projects.flatMap(p=>p.plannerAssignments??[])].filter(a=>a.roleId===role.id);role.usage=aggregate(assignments.map(a=>a.usage));role.usageByHarnessModel=groupUsage(assignments);}return state; } - configureVerifiedWorkspace(config:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];formatCommand?:ValidationCommand;bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig;bridgeToken?:string}):void { if(!config.repoPath||!config.allowedScope.length||!config.commands.length)throw new Error('Workspace repository, allowed scope, and validation commands are required');this.verifiedWorkspaceConfig=structuredClone(config); } + configureVerifiedWorkspace(config:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];formatCommand?:ValidationCommand;bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig;bridgeToken?:string;workerCheckRounds?:number}):void { if(!config.repoPath||!config.allowedScope.length||!config.commands.length)throw new Error('Workspace repository, allowed scope, and validation commands are required');this.verifiedWorkspaceConfig=structuredClone(config); } private runFormatsWorkerFiles(run:Run):boolean{return Boolean(run.validationCommands?run.formatCommand:this.verifiedWorkspaceConfig?.formatCommand);} + /** Check-gate rounds for this run's Worker turns: the configured count, or 0 when there is no bridge workspace or automatic validation correction is off. */ + private workerCheckRounds(run:Run):number{const config=this.verifiedWorkspaceConfig;return config?.bridgeBaseUrl&&run.workspaceId&&run.autoValidationCorrection!==false?Math.max(0,Math.min(3,config.workerCheckRounds??0)):0;} private async workspaceConfigForRun(runId:string){const base=this.verifiedWorkspaceConfig;if(!base)return undefined;const run=this.findRun(await this.store.load(),runId);return {...base,allowedScope:run.allowedScope??base.allowedScope,commands:(run.validationCommands as ValidationCommand[]|undefined)??base.commands,formatCommand:run.validationCommands?(run.formatCommand as ValidationCommand|undefined):base.formatCommand};} async setValidationSettings(runId:string,input:{autoValidationCorrection:boolean}):Promise{if(typeof input?.autoValidationCorrection!=='boolean')throw Object.assign(new Error('autoValidationCorrection must be a boolean'),{statusCode:422});return this.store.mutate(s=>{const run=this.findRun(s,runId);run.autoValidationCorrection=input.autoValidationCorrection;s.events.push(event('validation.settings_changed','run',runId,{autoValidationCorrection:input.autoValidationCorrection}));return structuredClone(run);});} private async beginValidationProgress(runId:string,commands:readonly ValidationCommand[]):Promise{const attemptId=id('validation-attempt'),startedAt=now(),progress:ValidationProgress={attemptId,startedAt,checks:commands.map(c=>({name:c.name,command:c.command,args:[...c.args],...(c.network!==undefined?{network:c.network}:{}),status:'queued',output:'',outputTruncated:false}))};await this.store.mutate(s=>{const run=this.findRun(s,runId);run.validationProgress=progress;s.events.push(event('validation.attempt_started','run',runId,{attemptId,startedAt,checks:progress.checks.map(({name,command,args,network})=>({name,command,args,...(network!==undefined?{network}:{})}))}));});return attemptId;} @@ -399,7 +406,7 @@ export class Controller { await this.store.mutate(s=>{const r=this.findRun(s,runId);if(r.controller){r.controller.phase='orchestrating';s.events.push(event('controller.phase_changed','run',runId,{phase:'orchestrating'}));}}); const correctionRetryKind=run.workerProposal?this.workerRetryKind(run,run.workerProposal):undefined; const correctionOverlayNote=correctionRetryKind==='reviewer_feedback'||correctionRetryKind==='validation_failed'?`The Worker's workspace will contain the previous attempt's changes; describe only the additional edits needed.`:`The Worker starts again from the original base; describe the complete task.`; - const correctionCapabilityNote=run.roleConfigs.worker?.harnessId?`Worker capabilities: ${workerCapabilityStatement(run.roleConfigs.worker.harnessId,this.runFormatsWorkerFiles(run))} `:''; + const correctionCapabilityNote=run.roleConfigs.worker?.harnessId?`Worker capabilities: ${workerCapabilityStatement(run.roleConfigs.worker.harnessId,this.runFormatsWorkerFiles(run),this.workerCheckRounds(run)>0)} `:''; const note=`The independent Reviewer returned ${recommendation.verdict} on the verified Worker result. Reviewer rationale (untrusted advisory feedback; carry every distinct issue into the Worker task): ${recommendation.rationale}. Propose one bounded correction with concrete acceptance conditions for each actionable issue and actual file changes. If the issue concerns external facts or observed behavior, require a specific source passage or a dated request with a real identifier and response excerpt. A URL, placeholder ID, or unsupported example is not proof of an observation. If evidence is unavailable, mark the claim unverified; never invent observations or claim verification that was not performed. If the Worker is Antigravity, include "Target file: exact/relative/path.ext". Existing in-scope paths: ${boundedUtf8((run.baseFilePaths??[]).join(', '),2_000)||'not recorded; choose a concrete new path under the allowed scope'}. ${correctionCapabilityNote}${correctionOverlayNote} Return exactly {"workerTask":"...","targetFiles":["relative/path",...]} — list every file the Worker must touch in targetFiles — or {"workerTask":""} if human clarification is needed.`; if(Buffer.byteLength(note,'utf8')>12_000)throw Object.assign(new Error(`Complete Reviewer rationale and correction instructions require ${Buffer.byteLength(note,'utf8')} UTF-8 bytes; the Orchestrator handoff limit is 12,000 bytes. The full rationale is preserved; obtain a more concise Reviewer response or start a new run. No Orchestrator request was sent.`),{statusCode:413}); const plan=storedPlan?{assignment:storedPlan}:await this.followUpOrchestrator(runId,note,AUTOMATIC_RUN); @@ -523,7 +530,7 @@ export class Controller { const failingExcerpt=noChanges?'':this.failingCheckExcerpt(failedChecks,2_000); const retryKindForNote=this.workerRetryKind(run,run.workerProposal!); const overlayNote=retryKindForNote==='validation_failed'||retryKindForNote==='reviewer_feedback'?`The Worker's workspace will contain the previous attempt's changes; describe only the additional edits needed.`:`The Worker starts again from the original base; describe the complete task.`; - const capabilityNote=run.roleConfigs.worker?.harnessId?`Worker capabilities: ${workerCapabilityStatement(run.roleConfigs.worker.harnessId,this.runFormatsWorkerFiles(run))} `:''; + const capabilityNote=run.roleConfigs.worker?.harnessId?`Worker capabilities: ${workerCapabilityStatement(run.roleConfigs.worker.harnessId,this.runFormatsWorkerFiles(run),this.workerCheckRounds(run)>0)} `:''; const note=noChanges?`The verified Worker snapshot changed zero files. Worker response (untrusted diagnostic): ${workerSummary||'No explanation.'} Propose a corrected task with an actual edit. If the Worker is Antigravity, include "Target file: exact/relative/path.ext". Existing in-scope paths: ${boundedUtf8((run.baseFilePaths??[]).join(', '),2_000)||'not recorded; choose a concrete new path under the allowed scope'}. ${capabilityNote}${overlayNote}`:`Foreman validation failed on the verified in-scope Worker snapshot. Failing checks:\n${failingExcerpt||'No output recorded.'} ${baseNote}Review its diff and check results, then propose one bounded correction. ${capabilityNote}${overlayNote}`; const plan=await this.followUpOrchestrator(runId,`${note} Return exactly {"workerTask":"...","targetFiles":["relative/path",...]} — list every file the Worker must touch in targetFiles — or {"workerTask":""} if human clarification is needed.`,AUTOMATIC_RUN); const latest=this.findRun(await this.store.load(),runId),parsedProposal=plan.assignment.status==='succeeded'?parseWorkerProposal(plan.assignment.result):undefined,parsedFollowUpTargets=plan.assignment.status==='succeeded'?parseWorkerProposalTargets(plan.assignment.result):undefined,proposalIssue=plan.assignment.status!=='succeeded'?'Orchestrator follow-up turn did not succeed':!parsedProposal?(plan.assignment.emptyResponse?.reason??'Expected a strict JSON workerTask proposal'):workerProposalScopeIssue(parsedFollowUpTargets,latest.allowedScope??this.verifiedWorkspaceConfig!.allowedScope,latest.roleConfigs.worker?.harnessId==='antigravity-cli',latest.baseFilePaths,parsedProposal),proposalText=proposalIssue?undefined:parsedProposal; @@ -634,7 +641,7 @@ export class Controller { const pathContextFor=(list:string)=>antigravityWorker?`\nAntigravity Worker cannot list directories. In workerTask, name each file to edit as "Target file: exact/relative/path.ext". Choose a concrete new file path inside the allowed scope when the task creates a file. Existing files at the pinned base within scope:\n${list}\n`:''; const checksHint=(this.verifiedWorkspaceConfig?.commands??[]).map(c=>`${c.name} (${c.command}${c.args.length?` ${c.args.join(' ')}`:''})` ).join(', ')||'none configured'; const workerConfig=run.roleConfigs.worker; - const capabilityHint=workerConfig?.harnessId?`\nWorker capabilities: ${workerCapabilityStatement(workerConfig.harnessId,this.runFormatsWorkerFiles(run))}\n`:''; + const capabilityHint=workerConfig?.harnessId?`\nWorker capabilities: ${workerCapabilityStatement(workerConfig.harnessId,this.runFormatsWorkerFiles(run),this.workerCheckRounds(run)>0)}\n`:''; let orchRepoAccess:RepoAccessRecord|undefined,orchReadOnlyWorkspaceId:string|undefined,orchRepoAccessText='No tools, file access, or shell commands are available in this context. Respond with text only.',orchDigestRequest:DigestRequest|undefined;{const isAgy=config.harnessId==='antigravity-cli';const _ws=this.verifiedWorkspaceConfig;if(_ws?.bridgeBaseUrl&&!isAgy&&run.pinnedBaseCommit){const pinnedCommit=run.pinnedBaseCommit;try{const _cached=this.snapshotWorkspaceCache.get(pinnedCommit);if(_cached){orchReadOnlyWorkspaceId=_cached;orchRepoAccess={mode:'snapshot',commit:pinnedCommit};orchRepoAccessText=`You can read (not modify) a snapshot of the repository at commit ${pinnedCommit} in your working directory using Read, Grep and Glob. Inspect relevant code before proposing the Worker task; cite files you read.`;}else{const seeded=await seedBridgeWorkspace(_ws.bridgeBaseUrl,pinnedCommit,{token:_ws.bridgeToken});this.snapshotWorkspaceCache.set(pinnedCommit,seeded.workspaceId);orchReadOnlyWorkspaceId=seeded.workspaceId;orchRepoAccess={mode:'snapshot',commit:pinnedCommit};orchRepoAccessText=`You can read (not modify) a snapshot of the repository at commit ${pinnedCommit} in your working directory using Read, Grep and Glob. Inspect relevant code before proposing the Worker task; cite files you read.`;}}catch(err){const reason=err instanceof Error?err.message:String(err);orchDigestRequest={repoPath:_ws.repoPath,allowedScope:run.allowedScope??_ws.allowedScope,commit:pinnedCommit,keywords:extractKeywords(clean),reason,fallback:{mode:'digest',commit:pinnedCommit,reason}};}}else if(isAgy&&_ws&&run.pinnedBaseCommit){const pinnedCommit=run.pinnedBaseCommit;const reason='antigravity-cli does not support read-only workspace tools';orchDigestRequest={repoPath:_ws.repoPath,allowedScope:run.allowedScope??_ws.allowedScope,commit:pinnedCommit,keywords:extractKeywords(clean),reason};}} // Foreman-supplied context is trimmed to fit; only the instructions, the operator note and the task title are mandatory. let goal=task.goal??task.title,criteria=(task.validationCriteria??[]).join('; '),planner=plannerContext,paths=pathList,checks=checksHint,scopeText=allowedScope.join(', '),guidanceCount=guidanceLines.length; @@ -797,10 +804,10 @@ export class Controller { const snapshot=await this.store.load(); const {a}=this.findAssignment(snapshot,assignmentId);previousResponseRecovery=previousResponseRecovery||a.previousResponseRecoveryAttempted===true;if(a.status!=='submitting'&&a.status!=='cancel_requested') return a; if(retrying&&this.discovery&&this.discovery.capabilities.idempotency!==true){await this.store.mutate(s=>{s.events.push(event('assignment.recovery_blocked','assignment',assignmentId,{reason:'UHP idempotency is unavailable; controller will not repeat the submission'}));});return a;} try { - const context=await this.store.load();const {run}=this.findAssignment(context,assignmentId);const uhpSession=a.roleId==='planner'||a.roleId==='orchestrator'?run.sessions[a.roleId]:undefined;let previousResponseId=uhpSession?.responseId;if(uhpSession&&previousResponseId){const pointed=run.assignments.find(item=>item.roleId===a.roleId&&item.responseId===previousResponseId);if(pointed&&pointed.status!=='succeeded'){const latestSucceeded=run.assignments.filter(item=>item.roleId===a.roleId&&item.status==='succeeded'&&!!item.responseId&&canonicalJson(item.requestedConfig)===canonicalJson(a.requestedConfig)&&(!pointed.sessionId||item.sessionId===pointed.sessionId)).at(-1);previousResponseId=latestSucceeded?.responseId;await this.store.mutate(state=>{const current=this.findRun(state,runId),session=current.sessions[a.roleId as 'planner'|'orchestrator'];session.responseId=previousResponseId;if(latestSucceeded?.sessionId)session.uhpSessionId=latestSucceeded.sessionId;state.events.push(event('session.failed_pointer_repaired','session',session.localId,{runId,roleId:a.roleId,failedAssignmentId:pointed.id,previousResponseId}));});}}const reviewEvidence=a.roleId==='reviewer'?this.reviewerEvidencePackage(run):undefined;if(reviewEvidence?.orchestratorInboxId&&(a.reviewerInboxId!==reviewEvidence.orchestratorInboxId||a.reviewerEvidenceDigest!==reviewEvidence.orchestratorEvidenceDigest))throw Object.assign(new Error('Reviewer handoff no longer matches its controller-bound Orchestrator receipt'),{statusCode:409});const turnTimeoutSecondsForSubmit:number=typeof a.requestedConfig.options?.timeoutSeconds==='number'?a.requestedConfig.options.timeoutSeconds:this.roleTurnTimeoutSeconds(a.roleId);const config={...(a.requestedConfig.options??{}),...a.requestedConfig,timeoutSeconds:turnTimeoutSecondsForSubmit,...(a.roleId==='worker'&&run.workspaceId?{workspaceId:run.workspaceId}:{}),...(previousResponseId?{previousResponseId}:{}),...(a.roleId==='reviewer'?{reviewMode:'read_only',reviewEvidence}:{})}; + const context=await this.store.load();const {run}=this.findAssignment(context,assignmentId);const uhpSession=a.roleId==='planner'||a.roleId==='orchestrator'?run.sessions[a.roleId]:undefined;let previousResponseId=uhpSession?.responseId;if(uhpSession&&previousResponseId){const pointed=run.assignments.find(item=>item.roleId===a.roleId&&item.responseId===previousResponseId);if(pointed&&pointed.status!=='succeeded'){const latestSucceeded=run.assignments.filter(item=>item.roleId===a.roleId&&item.status==='succeeded'&&!!item.responseId&&canonicalJson(item.requestedConfig)===canonicalJson(a.requestedConfig)&&(!pointed.sessionId||item.sessionId===pointed.sessionId)).at(-1);previousResponseId=latestSucceeded?.responseId;await this.store.mutate(state=>{const current=this.findRun(state,runId),session=current.sessions[a.roleId as 'planner'|'orchestrator'];session.responseId=previousResponseId;if(latestSucceeded?.sessionId)session.uhpSessionId=latestSucceeded.sessionId;state.events.push(event('session.failed_pointer_repaired','session',session.localId,{runId,roleId:a.roleId,failedAssignmentId:pointed.id,previousResponseId}));});}}const reviewEvidence=a.roleId==='reviewer'?this.reviewerEvidencePackage(run):undefined;if(reviewEvidence?.orchestratorInboxId&&(a.reviewerInboxId!==reviewEvidence.orchestratorInboxId||a.reviewerEvidenceDigest!==reviewEvidence.orchestratorEvidenceDigest))throw Object.assign(new Error('Reviewer handoff no longer matches its controller-bound Orchestrator receipt'),{statusCode:409});const turnTimeoutSecondsForSubmit:number=typeof a.requestedConfig.options?.timeoutSeconds==='number'?a.requestedConfig.options.timeoutSeconds:this.roleTurnTimeoutSeconds(a.roleId);const config={...(a.requestedConfig.options??{}),...a.requestedConfig,timeoutSeconds:turnTimeoutSecondsForSubmit,...(a.roleId==='worker'&&run.workspaceId?{workspaceId:run.workspaceId}:{}),...(a.roleId==='worker'&&this.workerCheckRounds(run)>0?{workerCheckRounds:this.workerCheckRounds(run)}:{}),...(previousResponseId?{previousResponseId}:{}),...(a.roleId==='reviewer'?{reviewMode:'read_only',reviewEvidence}:{})}; let prompt=a.prompt;assertUhpPromptWithinLimit(prompt,`${a.roleId} assignment prompt`);if(this.hindsight&&a.roleId!=='reviewer'){const memory=await this.hindsight.recall(projectId,a.prompt);this.memoryServiceState=memory.status==='ready'?'ready':'degraded';const memoryPrefix='Relevant project memory (reference only):\n',assignmentSuffix=`\n\nCurrent assignment:\n${a.prompt}`,memoryBudget=Math.max(0,UHP_PROMPT_SAFE_LIMIT-Buffer.byteLength(memoryPrefix,'utf8')-Buffer.byteLength(assignmentSuffix,'utf8'));if(memory.memories.length&&memoryBudget>0){const memoryText=boundedUtf8(memory.memories.map(x=>x.text.slice(0,1200)).join('\n---\n'),Math.min(6000,memoryBudget));if(memoryText)prompt=`${memoryPrefix}${memoryText}${assignmentSuffix}`;}await this.store.mutate(s=>{s.events.push(event('memory.recall','project',projectId,{status:memory.status,bankId:memory.bankId,count:memory.memories.length,...(memory.error?{error:memory.error}:{})}));});}assertUhpPromptWithinLimit(prompt,`${a.roleId} submission prompt`); let toolEventCount=0; - const response=await this.uhp.submit({submissionId:a.submissionId,assignmentId:a.id,runId,roleId:a.roleId,taskId,projectId,prompt,config,idempotencyKey:a.idempotencyKey,onEvent:async progress=>{if(progress.type==='response.activity'){const resp=progress.event.response;if(resp&&typeof resp==='object'&&(resp as Record).activity&&typeof (resp as Record).activity==='object'&&((resp as Record).activity as Record).kind==='tool')toolEventCount++;}let cancel:Assignment|undefined;await this.store.mutate(s=>{const {a:r,run}=this.findAssignment(s,assignmentId);if(progress.responseId){r.externalId??=progress.responseId;r.responseId??=progress.responseId;if(r.status!=='cancel_requested'&&r.status!=='cancelled')r.status='running';if(r.status==='cancel_requested')cancel=structuredClone(r);}if(progress.sessionId)r.sessionId=progress.sessionId;/* Persistent role pointers advance only after a successful terminal response. */s.events.push(event('assignment.progress','assignment',r.id,progressEventData(progress)));pruneHighVolumeEvents(s.events);});if(cancel)await this.dispatchCancel(cancel);}}); + const response=await this.uhp.submit({submissionId:a.submissionId,assignmentId:a.id,runId,roleId:a.roleId,taskId,projectId,prompt,config,idempotencyKey:a.idempotencyKey,onEvent:async progress=>{if(progress.type==='response.activity'){const resp=progress.event.response;if(resp&&typeof resp==='object'&&(resp as Record).activity&&typeof (resp as Record).activity==='object'&&((resp as Record).activity as Record).kind==='tool')toolEventCount++;}let cancel:Assignment|undefined;await this.store.mutate(s=>{const {a:r,run}=this.findAssignment(s,assignmentId);if(progress.responseId){r.externalId??=progress.responseId;r.responseId??=progress.responseId;if(r.status!=='cancel_requested'&&r.status!=='cancelled')r.status='running';if(r.status==='cancel_requested')cancel=structuredClone(r);}if(progress.sessionId)r.sessionId=progress.sessionId;/* Persistent role pointers advance only after a successful terminal response. */s.events.push(event('assignment.progress','assignment',r.id,progressEventData(progress)));pruneHighVolumeEvents(s.events);});if(cancel)await this.dispatchCancel(cancel);const gateRound=a.roleId==='worker'&&progress.type==='response.activity'?workerGateRoundFrom(progress.event):undefined;if(gateRound&&progress.responseId)void this.answerWorkerGate(runId,a.id,progress.responseId,gateRound);}}); const turnTimeoutMsForDetection=turnTimeoutSecondsForSubmit*1000; const isTimeoutKill=response.cliFailureCategory==='terminated_sigterm'||(typeof response.runtimeMs==='number'&&response.runtimeMs>=turnTimeoutMsForDetection*0.95); const submitted=await this.store.mutate(s=>{const {a:r,run}=this.findAssignment(s,assignmentId); r.externalId=response.externalId; r.responseId=response.responseId; r.sessionId=response.sessionId; r.requestedModel=response.requestedModel??r.requestedConfig.model;r.actualModelStatus=response.actualModelStatus??(response.actualModel?'observed':undefined);r.selectedHarnessId=response.selectedHarnessId;r.reportedHarnessId=response.reportedHarnessId;r.modelFallback=response.modelFallback;r.cliInvocation=response.cliInvocation?structuredClone(response.cliInvocation):undefined;if(response.actualModel&&response.selectedHarnessId)r.actualConfig={harnessId:response.selectedHarnessId,model:response.actualModel};r.configOutcome=r.actualConfig?((response.modelFallback||response.actualModel!==r.requestedConfig.model||response.selectedHarnessId!==r.requestedConfig.harnessId)?'substituted':'confirmed'):'unavailable';r.configNotes={...(response.boundsApplied!==undefined?{boundsApplied:response.boundsApplied}:{}),...(response.ignoredFields?.length?{ignoredFields:[...response.ignoredFields]}:{})};if(response.usage&&Object.keys(response.usage).length)r.usage={...response.usage,measured:true,...(response.runtimeMs!==undefined?{runtimeMs:response.runtimeMs}:{})};else if(response.runtimeMs!==undefined)r.usage={runtimeMs:response.runtimeMs,measured:true};const terminal=normalizeStatus(response.status);const reviewPackage=r.roleId==='reviewer'?this.reviewerEvidencePackage(run):undefined;const reviewerInvalid=r.roleId==='reviewer'&&!this.reviewerProvesReadOnlyEvidence(run,reviewPackage,{terminal,responseId:response.responseId,sessionId:response.sessionId,model:response,mode:response.reviewerExecution?.mode,mutationAttempted:response.reviewerExecution?.mutationAttempted,validation:response.reviewerExecution?.validation},r.requestedConfig);r.reviewerExecution=response.reviewerExecution?.mode?{mode:response.reviewerExecution.mode,mutationAttempted:response.reviewerExecution.mutationAttempted===true,validation:response.reviewerExecution.validation}:undefined;if(reviewerInvalid){r.status='failed';r.error='Reviewer response did not prove a complete, distinct, read-only review of the controller evidence package';s.events.push(event('reviewer.execution_rejected','assignment',r.id,{responseId:r.responseId,sessionId:r.sessionId,actualModel:response.actualModel,mode:response.reviewerExecution?.mode,mutationAttempted:response.reviewerExecution?.mutationAttempted}));}else if(r.status==='cancel_requested'&&['succeeded','failed','needs_revision','cancelled'].includes(terminal)){r.status=terminal;if(terminal!=='cancelled')s.events.push(event('assignment.cancel_raced_terminal','assignment',r.id,{status:terminal,source:'submit-terminal'}));}else if(r.status!=='cancelled')r.status=terminal;r.result=terminal==='failed'||terminal==='needs_revision'?response.result??response.outputText??response.output:response.output??response.outputText??response.result; @@ -858,6 +865,27 @@ export class Controller { const envelope=await fetchBridgeSnapshot(config.bridgeBaseUrl,run.workspaceId,run.pinnedBaseCommit,{token:config.bridgeToken}); return this.processWorkerOutput(runId,workerAssignmentId,envelope,'bridge_snapshot'); } + /** + * Answer one Worker check-gate round: the bridge paused the Worker after a completed turn and waits for this verdict. + * Runs the configured checks on the paused workspace and posts the verdict; a failed verdict resumes the same Worker session with the failures. + * Never throws: any problem answers `skipped`, which ends the Worker turn as it is, and Foreman's validation afterwards still decides. + */ + private async answerWorkerGate(runId:string,assignmentId:string,responseId:string,round:number):Promise{ + const key=`${responseId}:${round}`;if(this.workerGateRounds.has(key))return undefined;this.workerGateRounds.add(key); + const config=await this.workspaceConfigForRun(runId).catch(()=>undefined);if(!config?.bridgeBaseUrl)return undefined; + let verdict:WorkerCheckRound; + try{ + const run=this.findRun(await this.store.load(),runId); + if(!run.workspaceId||!run.pinnedBaseCommit)throw new Error('The run has no pinned bridge workspace'); + const envelope=await fetchBridgeSnapshot(config.bridgeBaseUrl,run.workspaceId,run.pinnedBaseCommit,{token:config.bridgeToken}); + verdict=await runWorkerCheckRound({repoPath:config.repoPath,pinnedBaseCommit:run.pinnedBaseCommit,allowedScope:config.allowedScope,commands:config.commands,formatCommand:config.formatCommand,envelope,baseline:run.baselineValidation,timeoutMs:config.timeoutMs,maxOutputBytes:config.maxOutputBytes,sandbox:config.sandbox}); + }catch(error){verdict={status:'skipped',failedChecks:[],reason:error instanceof Error?error.message:String(error)};} + let delivery:string|undefined; + try{await postWorkerGateVerdict(config.bridgeBaseUrl,responseId,{round,status:verdict.status,...(verdict.feedback?{feedback:verdict.feedback}:{}),failedChecks:verdict.failedChecks},{token:config.bridgeToken});} + catch(error){delivery=error instanceof Error?error.message:String(error);} + await this.store.mutate(s=>{s.events.push(event('worker.check_round','assignment',assignmentId,{runId,responseId,round,status:verdict.status,failedChecks:verdict.failedChecks,...(verdict.reason?{reason:verdict.reason.slice(0,500)}:{}),...(delivery?{deliveryError:delivery.slice(0,500)}:{})}));}).catch(()=>undefined); + return verdict; + } /** Deterministic fixture path. Never exposed by the production HTTP API. */ async replayRecordedWorkerOutput(runId:string,workerAssignmentId:string,recordedEvidence:unknown):Promise{return this.processWorkerOutput(runId,workerAssignmentId,recordedEvidence,'recorded_replay');} /** Import a checked-in Worker response for controller verification without submitting a Worker task. */ @@ -1408,6 +1436,8 @@ async function fitRepoDigest(request:DigestRequest,room:number):Promise0?cut.slice(0,end):cut;} /** Bounded record of a streamed provider event for persisted state: identifiers, plus the scalar response fields and activity kind/summary (cut to 300 chars) the Live Response panel reads. Never the raw event. */ +/** The check round a bridge `response.activity` event announces when it pauses a Worker for Foreman's checks. */ +function workerGateRoundFrom(event:Record):number|undefined{const response=event.response,activity=response&&typeof response==='object'?(response as Record).activity:undefined;if(!activity||typeof activity!=='object')return undefined;const {kind,gate_round:round}=activity as Record;return kind==='worker_gate'&&Number.isInteger(round)&&(round as number)>0?round as number:undefined;} function progressEventData(progress:{type?:string;responseId?:string;sessionId?:string;event:Record}):Record{ const cut=(value:unknown,max:number):string|undefined=>{if(typeof value!=='string')return undefined;const head=value.slice(0,max);return /[\ud800-\udbff]$/.test(head)?head.slice(0,-1):head;}; const response=progress.event.response&&typeof progress.event.response==='object'?progress.event.response as Record:undefined,activity=response?.activity&&typeof response.activity==='object'?response.activity as Record:undefined; diff --git a/src/server.ts b/src/server.ts index aeb785d..335cd52 100644 --- a/src/server.ts +++ b/src/server.ts @@ -30,7 +30,7 @@ const uhp=config.uhpBaseUrl ? new UhpClient({baseUrl:config.uhpBaseUrl,...(uhpTo }; const hindsight=config.hindsightBaseUrl?new HindsightClient({baseUrl:config.hindsightBaseUrl,token:process.env.HINDSIGHT_TOKEN}):undefined; const controller=new Controller(store,uhp,!!config.hindsightBaseUrl,!!config.uhpBaseUrl,config.uhpHarnessId&&config.uhpModel?{harnessId:config.uhpHarnessId,model:config.uhpModel}:undefined,hindsight,Math.ceil(config.taskTimeoutMs/1000),Math.ceil(config.workerTimeoutMs/1000),300); -if(config.workspaceSourceRepo&&config.workspaceAllowedScope.length&&config.validationCommands.length)controller.configureVerifiedWorkspace({repoPath:config.workspaceSourceRepo,allowedScope:config.workspaceAllowedScope,commands:config.validationCommands,formatCommand:config.formatCommand,bridgeBaseUrl:config.workspaceBridgeUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox,bridgeToken:config.workspaceBridgeToken}); +if(config.workspaceSourceRepo&&config.workspaceAllowedScope.length&&config.validationCommands.length)controller.configureVerifiedWorkspace({repoPath:config.workspaceSourceRepo,allowedScope:config.workspaceAllowedScope,commands:config.validationCommands,formatCommand:config.formatCommand,bridgeBaseUrl:config.workspaceBridgeUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox,bridgeToken:config.workspaceBridgeToken,workerCheckRounds:config.workerCheckRounds}); const projectControllers=new Map>}>(); const createProjectRuntime=async(projectId:string,workspace:Awaited>)=>{ const bridge=new LocalBridge({dataDir:resolve(config.dataDir,'local-bridges'),onHealthChange:health=>process.stderr.write(`Local bridge for ${projectId} is ${health.state}${health.message?`: ${health.message}`:''}\n`)}); @@ -38,7 +38,7 @@ const createProjectRuntime=async(projectId:string,workspace:Awaited0?{foreman_worker_gate:{max_rounds:input.config.workerCheckRounds}}:{}), ...(input.roleId==='reviewer'?{foreman_review_mode:'read_only',review_evidence:reviewEvidence}:{}), ...((input.roleId==='planner'||input.roleId==='orchestrator')&&typeof input.config.readOnlyWorkspaceId==='string'?{foreman_read_only_workspace_id:input.config.readOnlyWorkspaceId}:{}) }, stream: true, store: true, timeout_seconds: timeoutSeconds, diff --git a/src/verified-workspace.ts b/src/verified-workspace.ts index b320bcb..77ed470 100644 --- a/src/verified-workspace.ts +++ b/src/verified-workspace.ts @@ -59,6 +59,14 @@ export async function overlayBridgeWorkspace(baseUrl:string,workspaceId:string,e if(typeof body.workspace_id!=='string'||typeof body.applied!=='number')throw new Error('Workspace bridge returned an invalid overlay result'); return {workspaceId:body.workspace_id,applied:body.applied}; } +/** Report Foreman's check verdict for a Worker the bridge paused after its turn (the Worker check gate). */ +export async function postWorkerGateVerdict(baseUrl:string,responseId:string,verdict:{round:number;status:'passed'|'failed'|'skipped';feedback?:string;failedChecks?:string[]},options:number|BridgeRequestOptions=15_000):Promise{ + if(!responseId||responseId.includes('/')||responseId.includes('\\'))throw new Error('Invalid bridge response ID'); + const endpoint=bridgeEndpoint(baseUrl,`/extensions/foreman-workspace/v1/responses/${encodeURIComponent(responseId)}/worker-gate`); + const body={round:verdict.round,status:verdict.status,...(verdict.feedback?{feedback:verdict.feedback}:{}),...(verdict.failedChecks?.length?{failed_checks:verdict.failedChecks}:{})}; + const response=await fetch(endpoint,{method:'POST',headers:bridgeHeaders(options,{'content-type':'application/json'}),body:JSON.stringify(body),signal:AbortSignal.timeout(bridgeTimeout(options))}); + if(!response.ok)throw bridgeFailure('check verdict request failed',response.status); +} export async function seedBridgeWorkspace(baseUrl:string,pinnedBaseCommit:string,options:number|BridgeRequestOptions=15_000):Promise<{workspaceId:string;baseCommit:string}>{ if(!SHA.test(pinnedBaseCommit))throw new Error('A full pinned base commit SHA is required'); const timeoutMs=bridgeTimeout(options); diff --git a/src/worker-check-gate.ts b/src/worker-check-gate.ts new file mode 100644 index 0000000..7fe044d --- /dev/null +++ b/src/worker-check-gate.ts @@ -0,0 +1,72 @@ +import { baselineCommandDigest } from './baseline-validation.js'; +import type { BaselineValidation } from './domain.js'; +import { formatWorkerSnapshot, formatSetupCommands } from './format-step.js'; +import type { ValidationSandboxConfig } from './validation-sandbox.js'; +import { validateWorkerOutput, verifyWorkerSnapshot, type BridgeSnapshotEnvelope, type ValidationCommand, type ValidationObservation } from './verified-workspace.js'; + +/** + * The Worker check gate. The bridge pauses a Worker after each completed turn; Foreman runs the configured checks on + * the paused workspace and returns a verdict. On `failed`, the bridge resumes the same Worker session with `feedback`, + * so every harness can fix its own check failures without a shell. The verdict is feedback only: Foreman's validation + * after the turn is still the authoritative result. + */ +export interface WorkerCheckRound { status: 'passed'|'failed'|'skipped'; failedChecks: string[]; feedback?: string; reason?: string } + +/** Upper bound on the feedback text sent back to the Worker. */ +export const WORKER_CHECK_FEEDBACK_BYTES = 7_000; +const MIN_PER_CHECK_BYTES = 600; + +const stripAnsi = (text:string):string => text.replace(/\u001b(?:\[[0-9;?]*[ -/]*[@-~]|\][^\u0007]*\u0007|[PX^_][^\u001b]*\u001b\\)/g, ''); +const passed = (o:ValidationObservation):boolean => o.exitCode === 0 && !o.timedOut && !o.outputTruncated; + +/** UTF-8-safe excerpt of at most `maxBytes`, keeping the start and (more of) the end, where tools print their summary. */ +export function headTailExcerpt(text:string, maxBytes:number):string { + const bytes = Buffer.from(text, 'utf8'); + if (bytes.length <= maxBytes) return text; + const marker = '\n[... output trimmed ...]\n', room = Math.max(0, maxBytes - Buffer.byteLength(marker)); + const head = Math.floor(room / 3), tail = room - head; + const start = bytes.subarray(0, head).toString('utf8').replace(/�$/, ''); + const end = bytes.subarray(bytes.length - tail).toString('utf8').replace(/^�/, ''); + return `${start}${marker}${end}`; +} + +/** Readable failure report for the Worker: each failed check's command, exit status and a head-and-tail output excerpt. */ +export function workerCheckFeedback(failures:readonly ValidationObservation[], maxBytes = WORKER_CHECK_FEEDBACK_BYTES):string { + if (!failures.length) return ''; + const perCheck = Math.max(MIN_PER_CHECK_BYTES, Math.floor(maxBytes / failures.length) - 200); + const parts = failures.map(f => { + const status = f.timedOut ? 'timed out' : f.outputTruncated ? `exit ${f.exitCode ?? 'none'}, output exceeded the limit` : `exit ${f.exitCode ?? 'none'}${f.signal ? `, ${f.signal}` : ''}`; + const output = stripAnsi(f.output ?? '').trim(); + return `## ${f.name} failed (${status})\n$ ${[f.command, ...f.args].join(' ')}\n${output ? headTailExcerpt(output, perCheck) : '(no output)'}`; + }); + return headTailExcerpt(parts.join('\n\n'), maxBytes); +} + +/** + * One gate round against a paused Worker's bridge snapshot: verify it against the pinned base and allowed scope, apply the + * configured formatter the same way the final validation does, run the checks, and report the failures the Worker could have caused. + * A scope violation is the Worker's to fix and fails the round; any other problem skips it, so the Worker is never sent + * feedback about Foreman's own infrastructure. + */ +export async function runWorkerCheckRound(input:{ + repoPath:string; pinnedBaseCommit:string; allowedScope:readonly string[]; commands:readonly ValidationCommand[]; formatCommand?:ValidationCommand; + envelope:BridgeSnapshotEnvelope; baseline?:BaselineValidation; timeoutMs?:number; maxOutputBytes?:number; sandbox?:ValidationSandboxConfig; +}):Promise { + let snapshot; + try { snapshot = await verifyWorkerSnapshot({ repoPath:input.repoPath, pinnedBaseCommit:input.pinnedBaseCommit, envelope:input.envelope, allowedScope:input.allowedScope }); } + catch (error) { + const message = error instanceof Error ? error.message : String(error); + if (/outside the allowed scope/.test(message)) return { status:'failed', failedChecks:['allowed scope'], feedback:`${message}. Only change files inside the allowed scope (${input.allowedScope.join(', ')}); restore any other file you changed.` }; + return { status:'skipped', failedChecks:[], reason:message }; + } + if (!snapshot.changes.length) return { status:'skipped', failedChecks:[], reason:'The Worker workspace has no changes to check' }; + const verified = input.formatCommand + ? (await formatWorkerSnapshot({ repoPath:input.repoPath, verified:snapshot, setupCommands:formatSetupCommands(input.commands), formatCommand:input.formatCommand, allowedScope:input.allowedScope, timeoutMs:input.timeoutMs, maxOutputBytes:input.maxOutputBytes, sandbox:input.sandbox })).verified + : snapshot; + const observed = await validateWorkerOutput({ repoPath:input.repoPath, evidence:verified, commands:input.commands, timeoutMs:input.timeoutMs, maxOutputBytes:input.maxOutputBytes, sandbox:input.sandbox }); + const baseline = input.baseline?.pinnedBaseCommit === input.pinnedBaseCommit && input.baseline.commandDigest === baselineCommandDigest(input.commands) ? input.baseline : undefined; + const baseFailing = new Set((baseline?.checks ?? []).filter(c => !c.passed).map(c => c.name)); + const failures = observed.checks.filter(c => !passed(c) && !baseFailing.has(c.name)); + if (!failures.length) return { status:'passed', failedChecks:[] }; + return { status:'failed', failedChecks:failures.map(f => f.name), feedback:workerCheckFeedback(failures) }; +} diff --git a/tests/worker-check-gate.test.ts b/tests/worker-check-gate.test.ts new file mode 100644 index 0000000..1770668 --- /dev/null +++ b/tests/worker-check-gate.test.ts @@ -0,0 +1,84 @@ +import { execFileSync } from 'node:child_process'; +import { createHash } from 'node:crypto'; +import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { dirname, join } from 'node:path'; +import { afterEach, describe, expect, it } from 'vitest'; +import { baselineCommandDigest } from '../src/baseline-validation.js'; +import { snapshotGitCommit } from '../src/git-workspace.js'; +import type { BridgeSnapshotEnvelope, ValidationCommand, ValidationObservation } from '../src/verified-workspace.js'; +import { headTailExcerpt, runWorkerCheckRound, workerCheckFeedback } from '../src/worker-check-gate.js'; + +const dirs: string[] = []; +afterEach(async () => { await Promise.all(dirs.splice(0).map(dir => rm(dir, { recursive: true, force: true }))); }); +const git = (cwd: string, ...args: string[]) => execFileSync('git', ['-C', cwd, ...args], { encoding: 'utf8', stdio: ['ignore', 'pipe', 'pipe'] }).trim(); + +async function fixture() { + const dir = await mkdtemp(join(tmpdir(), 'foreman-check-gate-')); + dirs.push(dir); + git(dir, 'init', '-q'); git(dir, 'config', 'user.name', 'Fixture'); git(dir, 'config', 'user.email', 'fixture@example.invalid'); + const files: Record = { 'src/a.txt': 'alpha\n', 'docs/outside.txt': 'outside\n' }; + for (const [path, content] of Object.entries(files)) { await mkdir(dirname(join(dir, path)), { recursive: true }); await writeFile(join(dir, path), content); } + git(dir, 'add', '-A'); git(dir, 'commit', '-q', '-m', 'base'); + return { dir, sha: git(dir, 'rev-parse', 'HEAD') }; +} +async function envelope(dir: string, sha: string, edits: Record): Promise { + const base = await snapshotGitCommit(dir, sha); + const entries = base.entries.map(entry => { + const bytes = Buffer.from(edits[entry.path] ?? Buffer.from(entry.contentBase64, 'base64').toString('utf8')); + return { path: entry.path, kind: 'file' as const, mode: '100644' as const, size: bytes.length, sha256: createHash('sha256').update(bytes).digest('hex'), contentBase64: bytes.toString('base64') }; + }); + return { complete: true, base_commit: sha, entries, errors: [] }; +} +const node = (name: string, source: string): ValidationCommand => ({ name, command: process.execPath, args: ['-e', source], network: false }); +const requireFixed = node('content check', `if(!require('node:fs').readFileSync('src/a.txt','utf8').includes('fixed')){console.log('src/a.txt: expected fixed');process.exit(1)}`); +const alwaysRed = node('red on base', `console.log('broken on base');process.exit(2)`); +const observation = (name: string, output: string, patch: Partial = {}): ValidationObservation => ({ name, command: 'pnpm', args: ['run', name], exitCode: 1, timedOut: false, output, outputTruncated: false, startedAt: '', finishedAt: '', sandbox: 'none', network: false, ...patch }); + +describe('Worker check gate', () => { + it('keeps the head and more of the tail when output is too long', () => { + const text = `${'a'.repeat(500)}MIDDLE${'z'.repeat(500)}`; + const excerpt = headTailExcerpt(text, 300); + expect(Buffer.byteLength(excerpt)).toBeLessThanOrEqual(300); + expect(excerpt.startsWith('aaa')).toBe(true); + expect(excerpt.endsWith('zzz')).toBe(true); + expect(excerpt).toContain('[... output trimmed ...]'); + expect(excerpt).not.toContain('MIDDLE'); + expect(headTailExcerpt('short', 300)).toBe('short'); + }); + + it('reports each failed check with its command, status and output without terminal colors', () => { + const feedback = workerCheckFeedback([observation('lint', '\u001b[31merror\u001b[0m src/a.ts:3 no-unused-vars'), observation('test', '', { exitCode: null, timedOut: true })]); + expect(feedback).toContain('## lint failed (exit 1)\n$ pnpm run lint\nerror src/a.ts:3 no-unused-vars'); + expect(feedback).toContain('## test failed (timed out)'); + expect(feedback).toContain('(no output)'); + expect(feedback).not.toContain('\u001b'); + }); + + it('fails a round on a Worker-caused check failure, leaves out checks red on the base, and passes once fixed', async () => { + const { dir, sha } = await fixture(); + const commands = [requireFixed, alwaysRed]; + const baseline = { pinnedBaseCommit: sha, commandDigest: baselineCommandDigest(commands), ranAt: '', checks: [{ name: 'red on base', passed: false, exitCode: 2, timedOut: false }] }; + const common = { repoPath: dir, pinnedBaseCommit: sha, allowedScope: ['src/'], commands, baseline, sandbox: { mode: 'none' as const } }; + const failed = await runWorkerCheckRound({ ...common, envelope: await envelope(dir, sha, { 'src/a.txt': 'first try\n' }) }); + expect(failed.status).toBe('failed'); + expect(failed.failedChecks).toEqual(['content check']); + expect(failed.feedback).toContain('src/a.txt: expected fixed'); + expect(failed.feedback).not.toContain('broken on base'); + const passed = await runWorkerCheckRound({ ...common, envelope: await envelope(dir, sha, { 'src/a.txt': 'fixed\n' }) }); + expect(passed).toEqual({ status: 'passed', failedChecks: [] }); + }); + + it('sends a scope violation back to the Worker but skips rounds it cannot fix', async () => { + const { dir, sha } = await fixture(); + const common = { repoPath: dir, pinnedBaseCommit: sha, allowedScope: ['src/'], commands: [requireFixed], sandbox: { mode: 'none' as const } }; + const outside = await runWorkerCheckRound({ ...common, envelope: await envelope(dir, sha, { 'src/a.txt': 'fixed\n', 'docs/outside.txt': 'changed\n' }) }); + expect(outside.status).toBe('failed'); + expect(outside.failedChecks).toEqual(['allowed scope']); + expect(outside.feedback).toContain('docs/outside.txt'); + const unchanged = await runWorkerCheckRound({ ...common, envelope: await envelope(dir, sha, {}) }); + expect(unchanged).toMatchObject({ status: 'skipped', reason: 'The Worker workspace has no changes to check' }); + const incomplete = await runWorkerCheckRound({ ...common, envelope: { ...(await envelope(dir, sha, { 'src/a.txt': 'x\n' })), complete: false } }); + expect(incomplete).toMatchObject({ status: 'skipped', reason: 'Bridge workspace snapshot is incomplete' }); + }); +});