From 40c8bfdf9f139aca8230894f45bb1951a6b96a3e Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 28 Sep 2026 20:44:31 -0700 Subject: [PATCH 1/7] feat(permission): add harness-owned worker grant envelope A worker deny-on-ask registers a denied-call envelope keyed by worker session with single-parent-turn expiry. The deny reason names only the envelope requestId and the parent retries via resume_agent with exact args. --- docs/ARCHITECTURE.md | 4 + src/permission/reactor-authorize.ts | 170 +++++++- src/permission/worker-grant-flow.test.ts | 209 ++++++++++ src/permission/worker-grant.test.ts | 324 +++++++++++++++ src/permission/worker-grant.ts | 477 +++++++++++++++++++++++ src/subagent/run.ts | 27 +- src/subagent/session-store.ts | 98 ++++- 7 files changed, 1289 insertions(+), 20 deletions(-) create mode 100644 src/permission/worker-grant-flow.test.ts create mode 100644 src/permission/worker-grant.test.ts create mode 100644 src/permission/worker-grant.ts diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 2371c3b52..12434e63e 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -101,6 +101,10 @@ OAuth stores serialize refresh writes and treat the value observed under the sto The interactive TUI may offer an explicit alternate-provider/model selector only after the same-profile retry also ends in a terminal credential failure. That choice is generation-scoped and consumed once. It may continue the original operator message only when the failed attempt emitted no committing inference event; after any commitment it switches the live provider without replaying the message. `/connect` replaces or adds credentials and `/model` switches the live source explicitly. Exec and fleet workers use the same one-shot same-profile recovery but never open an auth or alternate-provider prompt; an exhausted failure terminates that attempt with a sanitized recovery diagnostic. +### Worker permission grants + +A worker deny-on-ask registers a harness-owned denied-call envelope (`src/permission/worker-grant.ts`): the exact denied ToolCall plus a stable path-aware fingerprint, keyed by worker session with a single-parent-turn expiry. The deny reason names only the envelope's requestId; the worker's `ask_director` prose carries no authority and plain `send_input` text can never mint, consume, or extend an envelope. The parent observes the envelope on the ask record, replays the exact call through its own gate operator path, and retries the retained worker via the existing `resume_agent` with exact args plus the questionId ref. The first covering-grant retry consumes the envelope; replays, tampered args/cwd/tool/session, expiry, decline, or interrupt fail closed with a truthful blocker. There is no dedicated grant verb — single-use and exactness are enforced by the envelope sidecar, not by tool plumbing. + ### TUI Runner (`src/tui/runner/`) - Builds a chat-mode agent using the `ChatDirector` diff --git a/src/permission/reactor-authorize.ts b/src/permission/reactor-authorize.ts index bef1e0fe2..b194b95ef 100644 --- a/src/permission/reactor-authorize.ts +++ b/src/permission/reactor-authorize.ts @@ -21,7 +21,15 @@ import type { AuthzCallResult } from "@intx/inference"; import { WORKER_CANNOT_COMPLETE_APPROVAL } from "./decline-markers.js"; import type { AuthorizeVerdict, GateVerdict, PermissionGate } from "./gate.js"; import type { PermissionRequest } from "./types.js"; +import { + createDeniedCallEnvelope, + formatWorkerDenyWithGrantId, + getProcessWorkerGrantStore, + type WorkerDeniedCallEnvelope, + type WorkerGrantStore, +} from "./worker-grant.js"; import { canonicalToolName } from "../agent/canonical-tool-name.js"; +import { getSubAgentIdentity } from "../subagent/identity-context.js"; import { FLEET_VERBS, ORCHESTRATOR_ONLY_FLEET_VERBS, @@ -71,21 +79,138 @@ function readAuthorizeToolCall( return call satisfies ToolCallType; } +export interface WorkerGrantOptions { + /** Defaults to the process-shared sidecar (worker deny side registers, + * parent observes/consumes). Tests pass an isolated store. */ + store?: WorkerGrantStore; + /** Owning worker session; absent means no envelope is minted or matched. */ + sessionId?: string | (() => string | undefined); + workspaceRoot?: string; + /** Observability hook for tests/telemetry; never authoritative. */ + onDeniedCall?: (envelope: WorkerDeniedCallEnvelope) => void; +} + +function resolveWorkerSessionId( + sessionId: WorkerGrantOptions["sessionId"], +): string | undefined { + return typeof sessionId === "function" ? sessionId() : sessionId; +} + +function resolveWorkerCwd(cwd: string | undefined): string { + return cwd ?? getSubAgentIdentity()?.cwd ?? process.cwd(); +} + +function workerCallIdentity( + sessionId: string, + call: ToolCallType, + cwd: string, +): { + sessionId: string; + canonicalTool: string; + args: Record; + cwd: string; +} { + return { + sessionId, + canonicalTool: canonicalToolName(call.name), + args: (call.arguments ?? {}) as Record, + cwd, + }; +} + +/** + * Harness-owned denied-call sidecar (CL-9475 Phase 1): on a worker + * deny-on-ask, register the exact denied call and name only its requestId in + * the deny reason. Reuses a still-pending envelope for the same exact call + * instead of minting duplicates across reactor retries with fresh call ids. + */ +function denyWorkerCallWithEnvelope( + options: WorkerGrantOptions | undefined, + call: ToolCallType, + request: PermissionRequest, + cwd: string, + baseReason: string, +): { effect: "deny"; reason: string } { + const sessionId = + options !== undefined + ? resolveWorkerSessionId(options.sessionId) + : undefined; + if (sessionId === undefined) return { effect: "deny", reason: baseReason }; + const store = options?.store ?? getProcessWorkerGrantStore(); + const identity = workerCallIdentity(sessionId, call, cwd); + const existing = store.pendingMatch( + identity.sessionId, + identity.canonicalTool, + identity.args, + identity.cwd, + ); + const envelope = + existing ?? + store.register( + createDeniedCallEnvelope({ + callId: call.id, + tool: call.name, + action: request.action, + subject: request.subject, + args: identity.args, + cwd: identity.cwd, + workerSessionId: identity.sessionId, + ...(options?.workspaceRoot !== undefined + ? { workspaceRoot: options.workspaceRoot } + : {}), + }), + ); + try { + options?.onDeniedCall?.(envelope); + } catch { + // Observability must not throw into the deny path. + } + return { + effect: "deny", + reason: formatWorkerDenyWithGrantId(baseReason, envelope.requestId), + }; +} + async function authorizeWorkerCall( gate: PermissionGate, call: ToolCallType, + grantOptions?: WorkerGrantOptions, + cwd?: string, ): Promise<{ effect: "allow" } | { effect: "deny"; reason: string }> { if (isWorkerControlPlaneTool(call.name)) return { effect: "allow" }; + const workerCwd = resolveWorkerCwd(cwd); + const sessionId = resolveWorkerSessionId(grantOptions?.sessionId); + if (sessionId !== undefined) { + const store = grantOptions?.store ?? getProcessWorkerGrantStore(); + const precheck = store.precheck( + workerCallIdentity(sessionId, call, workerCwd), + ); + if (!precheck.ok) return { effect: "deny", reason: precheck.blocker }; + } const verdict = await gate.authorizeCall(call); + if (verdict.effect === "allow") { + if (sessionId !== undefined) { + const store = grantOptions?.store ?? getProcessWorkerGrantStore(); + store.consumeOnAllow(workerCallIdentity(sessionId, call, workerCwd)); + } + return verdict; + } if (verdict.effect !== "ask") return verdict; - return { effect: "deny", reason: workerUnresolvedAskReason(verdict.request) }; + return denyWorkerCallWithEnvelope( + grantOptions, + call, + verdict.request, + workerCwd, + workerUnresolvedAskReason(verdict.request), + ); } async function evaluateWorkerCall( gate: PermissionGate, call: ToolCallType, + grantOptions?: WorkerGrantOptions, ): Promise { - const verdict = await authorizeWorkerCall(gate, call); + const verdict = await authorizeWorkerCall(gate, call, grantOptions); if (verdict.effect === "allow") return { allowed: true }; return { allowed: false, reason: verdict.reason }; } @@ -93,21 +218,49 @@ async function evaluateWorkerCall( async function executionVerdictWorkerCall( gate: PermissionGate, call: ToolCallType, + grantOptions?: WorkerGrantOptions, + cwd?: string, ): Promise { if (isWorkerControlPlaneTool(call.name)) return { effect: "allow" }; + const workerCwd = resolveWorkerCwd(cwd); + const sessionId = resolveWorkerSessionId(grantOptions?.sessionId); + if (sessionId !== undefined) { + const store = grantOptions?.store ?? getProcessWorkerGrantStore(); + const precheck = store.precheck( + workerCallIdentity(sessionId, call, workerCwd), + ); + if (!precheck.ok) return { effect: "deny", reason: precheck.blocker }; + } const verdict = await gate.executionVerdict(call); + if (verdict.effect === "allow") { + if (sessionId !== undefined) { + const store = grantOptions?.store ?? getProcessWorkerGrantStore(); + store.consumeOnAllow(workerCallIdentity(sessionId, call, workerCwd)); + } + return verdict; + } if (verdict.effect !== "ask") return verdict; - return { effect: "deny", reason: workerUnresolvedAskReason(verdict.request) }; + return denyWorkerCallWithEnvelope( + grantOptions, + call, + verdict.request, + workerCwd, + workerUnresolvedAskReason(verdict.request), + ); } // Shared-policy view for worker posix/MCP plugins and reactor authz: live // grants stay on the parent, reactor-gated middleware is `isReactorGated()`, // and authorizeCall never emits ask. -export function workerPermissionGate(gate: PermissionGate): PermissionGate { +export function workerPermissionGate( + gate: PermissionGate, + grantOptions?: WorkerGrantOptions, +): PermissionGate { return { - evaluate: (call) => evaluateWorkerCall(gate, call), - authorizeCall: (call) => authorizeWorkerCall(gate, call), - executionVerdict: (call) => executionVerdictWorkerCall(gate, call), + evaluate: (call) => evaluateWorkerCall(gate, call, grantOptions), + authorizeCall: (call) => authorizeWorkerCall(gate, call, grantOptions), + executionVerdict: (call) => + executionVerdictWorkerCall(gate, call, grantOptions), resolveSuspended: (request) => gate.resolveSuspended(request), isReactorGated: () => true, getApprovals: () => gate.getApprovals(), @@ -159,12 +312,13 @@ export function createReactorAuthorize( export function createWorkerAuthorize( gate: PermissionGate, + grantOptions?: WorkerGrantOptions, ): ( resource: string, action: string, context: unknown, ) => Promise { - const workerGate = workerPermissionGate(gate); + const workerGate = workerPermissionGate(gate, grantOptions); return async (resource, action, context) => { const call = readAuthorizeToolCall(resource, action, context); const verdict = await workerGate.authorizeCall(call); diff --git a/src/permission/worker-grant-flow.test.ts b/src/permission/worker-grant-flow.test.ts new file mode 100644 index 000000000..8593b3dfe --- /dev/null +++ b/src/permission/worker-grant-flow.test.ts @@ -0,0 +1,209 @@ +import { describe, test, expect } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import type { ToolCall } from "@intx/types/runtime"; +import { createPermissionGate } from "./gate.js"; +import { workerPermissionGate } from "./reactor-authorize.js"; +import { + WorkerGrantStore, + createDeniedCallEnvelope, + getProcessWorkerGrantStore, +} from "./worker-grant.js"; +import { createSubAgentSessionStore } from "../subagent/session-store.js"; + +const COMMAND = "npm test"; +const shellCall = (id: string, command: string): ToolCall => ({ + id, + name: "run_shell", + arguments: { command }, +}); + +function makeParentGate(cwd: string, approve: boolean) { + return createPermissionGate({ + approvals: [], + cwd, + skipPermissions: false, + interactive: true, + reactorGated: false, + // Operator surface: approve the offered exact scope as a session grant. + requestApproval: async (request) => { + if (!approve) return { allow: false }; + const exact = request.scopes[0]; + if (exact === undefined) return { allow: true }; + return { allow: true, persist: { ...exact, grant: "session" as const } }; + }, + }); +} + +function extractRequestId(reason: string): string { + const match = /grant request ([0-9a-f-]{36})/.exec(reason); + if (match?.[1] === undefined) { + throw new Error(`deny reason carries no grant request id: ${reason}`); + } + return match[1]; +} + +describe("worker grant-request flow: deny → parent replay grant → one retry", () => { + test("retained deny → grant → exactly one success, replay fails closed", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-flow-")); + const parentGate = makeParentGate(cwd, true); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(parentGate, { + sessionId: "worker-1", + store, + }); + + const denied = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + expect(denied.effect).toBe("deny"); + if (denied.effect !== "deny") throw new Error("expected worker deny"); + expect(denied.reason).toContain("deny"); + expect(denied.reason).toContain( + "workers cannot complete operator approval", + ); + const requestId = extractRequestId(denied.reason); + expect(store.peek(requestId)?.status).toBe("pending"); + + const replay = await parentGate.evaluate(shellCall("parent-1", COMMAND)); + expect(replay.allowed).toBe(true); + expect(parentGate.getSessionApprovals().length).toBeGreaterThan(0); + + const retry = await workerGate.authorizeCall(shellCall("c2", COMMAND)); + expect(retry).toEqual({ effect: "allow" }); + expect(store.peek(requestId)?.status).toBe("consumed"); + + const replayAgain = await workerGate.authorizeCall( + shellCall("c3", COMMAND), + ); + expect(replayAgain.effect).toBe("deny"); + if (replayAgain.effect !== "deny") throw new Error("expected replay deny"); + expect(replayAgain.reason).toContain("already consumed"); + expect(replayAgain.reason).toContain(requestId); + }); + + test("tampered args fall through to a fresh deny; original stays consumed", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-tamper-")); + const parentGate = makeParentGate(cwd, true); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(parentGate, { + sessionId: "worker-1", + store, + }); + + const denied = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + if (denied.effect !== "deny") throw new Error("expected worker deny"); + const requestId = extractRequestId(denied.reason); + expect(await parentGate.evaluate(shellCall("parent-1", COMMAND))).toEqual({ + allowed: true, + }); + expect(await workerGate.authorizeCall(shellCall("c2", COMMAND))).toEqual({ + effect: "allow", + }); + + const tampered = await workerGate.authorizeCall( + shellCall("c3", "npm run evil"), + ); + expect(tampered.effect).toBe("deny"); + if (tampered.effect !== "deny") throw new Error("expected tamper deny"); + expect(tampered.reason).not.toContain("already consumed"); + expect(store.peek(requestId)?.status).toBe("consumed"); + }); + + test("headless parent cannot grant: replay denies, auto-decline fails closed", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-headless-")); + // Interactive worker gate so the deny registers an envelope; the PARENT + // gate is headless (no operator seam). + const denyGate = createPermissionGate({ + approvals: [], + cwd, + skipPermissions: false, + interactive: true, + reactorGated: false, + // Seam wired (so decide preserves ask) but never consulted on the + // authorizeCall path — the worker deny registers its envelope here. + requestApproval: async () => ({ allow: false }), + }); + const headlessGate = createPermissionGate({ + approvals: [], + cwd, + skipPermissions: false, + interactive: false, + reactorGated: false, + }); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(denyGate, { + sessionId: "worker-headless", + store, + }); + + const denied = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + if (denied.effect !== "deny") throw new Error("expected worker deny"); + const requestId = extractRequestId(denied.reason); + + const replay = await headlessGate.evaluate(shellCall("parent-1", COMMAND)); + expect(replay.allowed).toBe(false); + + expect( + store.declineAllForSession( + "worker-headless", + "headless parent: no operator to approve", + ), + ).toBe(1); + expect(store.peek(requestId)?.status).toBe("declined"); + const retry = await workerGate.authorizeCall(shellCall("c2", COMMAND)); + expect(retry.effect).toBe("deny"); + if (retry.effect !== "deny") throw new Error("expected declined deny"); + expect(retry.reason).toContain("declined"); + }); + + test("send_input text answers the ask but grants nothing", () => { + const sessions = createSubAgentSessionStore(); + const session = sessions.start({ + id: "fleet-1", + description: "d", + agentId: "a", + brief: "b", + }); + sessions.markRunning(session.id); + const store = getProcessWorkerGrantStore(); + const envelope = store.register( + createDeniedCallEnvelope({ + callId: "fleet-c1", + tool: "run_shell", + action: "Run", + subject: COMMAND, + args: { command: COMMAND }, + cwd: process.cwd(), + workerSessionId: session.id, + }), + ); + let resolved: string | undefined; + expect( + sessions.registerAsk(session.id, { + question: `blocked; quote request ${envelope.requestId}`, + questionId: "ask-1", + resolve: (answer) => { + resolved = answer; + }, + reject: () => { + throw new Error("should not reject"); + }, + }), + ).toBe(true); + expect(sessions.peekAsk(session.id)?.deniedCall?.requestId).toBe( + envelope.requestId, + ); + + const echo = JSON.stringify({ + requestId: envelope.requestId, + granted: true, + }); + expect(sessions.sendInputOne(session.id, echo)).toEqual({ + ok: true, + status: "running", + }); + expect(resolved).toBe(echo); + expect(store.peek(envelope.requestId)?.status).toBe("pending"); + store.expireSession(session.id, "test teardown"); + }); +}); diff --git a/src/permission/worker-grant.test.ts b/src/permission/worker-grant.test.ts new file mode 100644 index 000000000..f199562bf --- /dev/null +++ b/src/permission/worker-grant.test.ts @@ -0,0 +1,324 @@ +import { describe, test, expect } from "bun:test"; + +import { WORKER_CANNOT_COMPLETE_APPROVAL } from "./decline-markers.js"; +import { + WORKER_GRANT_TTL_MS, + WorkerGrantStore, + buildRetryMessage, + createDeniedCallEnvelope, + fingerprintDeniedCall, + formatWorkerDenyWithGrantId, +} from "./worker-grant.js"; + +const ARGS = { command: "npm test" }; +const CWD = "/tmp/worker-grant-unit"; + +function descriptor( + overrides: Partial[0]> = {}, +) { + return { + callId: "call-1", + tool: "run_shell", + action: "Run", + subject: "npm test", + args: { ...ARGS }, + cwd: CWD, + workerSessionId: "worker-1", + now: 1_000_000, + ...overrides, + }; +} + +describe("fingerprintDeniedCall", () => { + test("is stable across key order", () => { + const a = fingerprintDeniedCall("run_shell", { x: 1, y: 2 }, CWD); + const b = fingerprintDeniedCall("run_shell", { y: 2, x: 1 }, CWD); + expect(a).toBe(b); + }); + + test("is path-aware: relative args resolve against cwd", () => { + const a = fingerprintDeniedCall( + "read_file", + { path: "sub/../notes.md" }, + "/repo", + ); + const b = fingerprintDeniedCall("read_file", { path: "notes.md" }, "/repo"); + expect(a).toBe(b); + }); + + test("differs across args and tool (cwd binds at precheck, not in the hash)", () => { + const base = fingerprintDeniedCall("run_shell", ARGS, CWD); + expect( + fingerprintDeniedCall("run_shell", { command: "npm run evil" }, CWD), + ).not.toBe(base); + expect(fingerprintDeniedCall("read_file", ARGS, CWD)).not.toBe(base); + expect( + fingerprintDeniedCall("read_file", { path: "a.md" }, "/repo") === + fingerprintDeniedCall("read_file", { path: "b.md" }, "/repo"), + ).toBe(false); + }); +}); + +describe("formatWorkerDenyWithGrantId", () => { + test("keeps deny + approval text, names only the requestId", () => { + const requestId = "11111111-2222-4333-8444-555555555555"; + const reason = formatWorkerDenyWithGrantId( + `deny: Run (npm test) requires a parent permission grant; ${WORKER_CANNOT_COMPLETE_APPROVAL}`, + requestId, + ); + expect(reason).toContain("deny"); + expect(reason).toContain(WORKER_CANNOT_COMPLETE_APPROVAL); + expect(reason).toContain(requestId); + expect(reason).not.toContain("npm test --secret=hunter2"); + }); +}); + +describe("WorkerGrantStore lifecycle", () => { + test("register → pending with ~10min expiry and denied audit", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + expect(envelope.status).toBe("pending"); + expect(envelope.expiresAt - envelope.createdAt).toBe(WORKER_GRANT_TTL_MS); + expect(envelope.audit.map((event) => event.event)).toEqual(["denied"]); + expect(store.peek(envelope.requestId)).toBe(envelope); + }); + + test("pending own envelope prechecks ok; unknown calls pass through", () => { + const store = new WorkerGrantStore(); + store.register(createDeniedCallEnvelope(descriptor())); + expect( + store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }), + ).toEqual({ ok: true }); + expect( + store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { command: "unrelated" }, + cwd: CWD, + }), + ).toEqual({ ok: true }); + }); + + test("consumeOnAllow consumes once; replay fails closed as consumed", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + const identity = { + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }; + expect(store.consumeOnAllow(identity)?.requestId).toBe(envelope.requestId); + expect(envelope.status).toBe("consumed"); + expect(store.consumeOnAllow(identity)).toBeUndefined(); + const replay = store.precheck(identity); + expect(replay.ok).toBe(false); + if (replay.ok) throw new Error("expected blocker"); + expect(replay.blocker).toContain(envelope.requestId); + expect(replay.blocker).toContain("already consumed"); + }); + + test("tampered session/cwd fail closed; tampered args fall through to the gate", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + store.consumeOnAllow({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }); + // Same fingerprint, owning session, wrong directory: the covering session + // grant is not cwd-scoped, so the backstop must refuse the ride. + const cwdTamper = store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: "/elsewhere", + now: 1_000_002, + }); + expect(cwdTamper.ok).toBe(false); + if (cwdTamper.ok) throw new Error("expected blocker"); + expect(cwdTamper.blocker).toContain(CWD); + // Same fingerprint, different session: must replay from the owning session. + const sessionTamper = store.precheck({ + sessionId: "worker-2", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_002, + }); + expect(sessionTamper.ok).toBe(false); + if (sessionTamper.ok) throw new Error("expected blocker"); + expect(sessionTamper.blocker).toContain("worker-1"); + // Different fingerprint (tampered args or tool): not this envelope — the + // gate denies downstream with a fresh deny since no grant covers it. + for (const identity of [ + { + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { command: "npm run evil" }, + cwd: CWD, + now: 1_000_002, + }, + { + sessionId: "worker-1", + canonicalTool: "read_file", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_002, + }, + ]) { + expect(store.precheck(identity)).toEqual({ ok: true }); + } + expect(envelope.status).toBe("consumed"); + }); + + test("cross-session replay of a pending envelope fails closed", () => { + const store = new WorkerGrantStore(); + store.register(createDeniedCallEnvelope(descriptor())); + const result = store.precheck({ + sessionId: "worker-2", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }); + expect(result.ok).toBe(false); + if (result.ok) throw new Error("expected blocker"); + expect(result.blocker).toContain("worker-1"); + }); + + test("expiry fails closed, lazily and via sweep", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + const after = envelope.expiresAt + 1; + const result = store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: after, + }); + expect(result.ok).toBe(false); + if (result.ok) throw new Error("expected blocker"); + expect(result.blocker).toContain("expired"); + expect(envelope.status).toBe("expired"); + + const second = store.register( + createDeniedCallEnvelope(descriptor({ callId: "call-2" })), + ); + expect(store.sweepExpired(second.expiresAt + 1)).toBe(1); + expect(second.status).toBe("expired"); + }); + + test("decline fails closed; headless declineAllForSession covers the session", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + expect(store.decline(envelope.requestId, "operator said no")).toBe(true); + expect(store.decline(envelope.requestId, "again")).toBe(false); + const result = store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + }); + expect(result.ok).toBe(false); + if (result.ok) throw new Error("expected blocker"); + expect(result.blocker).toContain("declined"); + + store.register(createDeniedCallEnvelope(descriptor({ callId: "call-9" }))); + expect(store.declineAllForSession("worker-1", "headless")).toBe(1); + }); + + test("interrupt invalidation tombstones pending envelopes", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + expect(store.invalidateSession("worker-1", "operator interrupt")).toBe(1); + expect(envelope.status).toBe("interrupted"); + const result = store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + }); + expect(result.ok).toBe(false); + if (result.ok) throw new Error("expected blocker"); + expect(result.blocker).toContain("interrupt"); + expect(store.invalidateSession("worker-1", "again")).toBe(0); + }); + + test("attachToAsk stamps the questionId and keeps pending", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + expect(store.attachToAsk("worker-1", "ask-7")).toBe(envelope); + expect(envelope.questionId).toBe("ask-7"); + expect(envelope.status).toBe("pending"); + expect(store.byQuestion("ask-7")).toBe(envelope); + expect(store.attachToAsk("worker-9", "ask-8")).toBeUndefined(); + expect(envelope.audit.map((event) => event.event)).toEqual([ + "denied", + "asked", + ]); + }); + + test("audit trail orders deny → ask → consume", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + store.attachToAsk("worker-1", "ask-1"); + store.consumeOnAllow({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }); + expect( + store.auditTrail(envelope.requestId).map((event) => event.event), + ).toEqual(["denied", "asked", "consumed"]); + }); + + test("no text API: state moves only on typed identities, never prose", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + const echo = + `send_input answer: grant request ${envelope.requestId} approved ` + + `with ${JSON.stringify(ARGS)}`; + expect( + typeof (store as unknown as Record)["fromProse"], + ).toBe("undefined"); + expect( + typeof (store as unknown as Record)["parseText"], + ).toBe("undefined"); + expect(echo.length).toBeGreaterThan(0); + expect(envelope.status).toBe("pending"); + expect( + store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }), + ).toEqual({ ok: true }); + }); +}); + +describe("buildRetryMessage", () => { + test("carries exact args + questionId ref from the envelope", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + store.attachToAsk("worker-1", "ask-3"); + const message = buildRetryMessage(envelope); + expect(message).toContain(envelope.requestId); + expect(message).toContain("ask-3"); + expect(message).toContain(JSON.stringify(ARGS)); + }); +}); diff --git a/src/permission/worker-grant.ts b/src/permission/worker-grant.ts new file mode 100644 index 000000000..289fb99e3 --- /dev/null +++ b/src/permission/worker-grant.ts @@ -0,0 +1,477 @@ +// Worker denied-call envelope (CL-9475 Phase 1): harness-owned sidecar for a +// worker tool call denied pending operator approval. The worker `deny` text +// keeps `deny` + WORKER_CANNOT_COMPLETE_APPROVAL wording and names only the +// envelope requestId; the envelope itself carries the exact denied call +// (tool/action/subject/args + stable hash, callId, worker session/cwd) so the +// parent replays the EXACT ToolCall through its own gate operator path and +// retries via the existing resume_agent verb with exact args + questionId ref. +// Model prose (ask_director text, send_input content) is never authoritative: +// nothing here parses message text, and plain send_input stays text-only. +// The dedicated atomic grant-and-retry verb is a Phase 2 follow-up. + +import { createHash, randomUUID } from "node:crypto"; + +import type { ToolCall } from "@intx/types/runtime"; + +import { canonicalToolName } from "../agent/canonical-tool-name.js"; +import { stableRequestId } from "./denial-memory.js"; + +/** Single parent-turn bound: one denied call gets at most one granted retry. */ +export const WORKER_GRANT_TTL_MS = 10 * 60 * 1000; + +export type WorkerGrantStatus = + | "pending" + | "consumed" + | "declined" + | "expired" + | "interrupted"; + +export interface WorkerGrantAuditEvent { + event: + | "denied" + | "asked" + | "consumed" + | "declined" + | "expired" + | "interrupted"; + at: number; + detail?: string; +} + +export interface WorkerDeniedCallEnvelope { + requestId: string; + workerSessionId: string; + turnId?: string; + deniedCallId: string; + canonicalTool: string; + toolAction?: string; + permissionSubject: string; + /** Exact denied arguments, JSON-cloned at deny time. */ + args: Record; + /** Stable path-aware hash of canonical tool + normalized args + cwd. */ + argsFingerprint: string; + cwd: string; + workspaceRoot?: string; + /** Attached when the worker's ask_director parks on this denial. */ + questionId?: string; + createdAt: number; + expiresAt: number; + status: WorkerGrantStatus; + audit: WorkerGrantAuditEvent[]; +} + +export interface DeniedCallDescriptor { + callId: string; + tool: string; + action?: string; + subject: string; + args: Record; + cwd: string; + workerSessionId: string; + turnId?: string; + workspaceRoot?: string; + now?: number; +} + +/** Stable path-aware fingerprint: same normalization the gate's denial memory + * uses (relative path args resolve against cwd), hashed to a fixed id. + * Object keys sort first so a retried call with reordered keys still matches + * the exact denied call. */ +export function fingerprintDeniedCall( + canonicalTool: string, + args: Record, + cwd: string, +): string { + const stable = stableRequestId( + { name: canonicalTool, arguments: sortKeys(args) } as ToolCall, + cwd, + ); + return createHash("sha256").update(stable).digest("hex"); +} + +function sortKeys(value: unknown): unknown { + if (Array.isArray(value)) return value.map(sortKeys); + if (typeof value === "object" && value !== null) { + return Object.fromEntries( + Object.keys(value as Record) + .sort() + .map((key) => [key, sortKeys((value as Record)[key])]), + ); + } + return value; +} + +export function canonicalDeniedTool(tool: string): string { + return canonicalToolName(tool); +} + +function audit( + envelope: WorkerDeniedCallEnvelope, + event: WorkerGrantAuditEvent["event"], + detail?: string, +): void { + envelope.audit.push({ + event, + at: Date.now(), + ...(detail !== undefined ? { detail } : {}), + }); +} + +export function createDeniedCallEnvelope( + descriptor: DeniedCallDescriptor, +): WorkerDeniedCallEnvelope { + const createdAt = descriptor.now ?? Date.now(); + const canonicalTool = canonicalDeniedTool(descriptor.tool); + const envelope: WorkerDeniedCallEnvelope = { + requestId: randomUUID(), + workerSessionId: descriptor.workerSessionId, + ...(descriptor.turnId !== undefined ? { turnId: descriptor.turnId } : {}), + deniedCallId: descriptor.callId, + canonicalTool, + ...(descriptor.action !== undefined + ? { toolAction: descriptor.action } + : {}), + permissionSubject: descriptor.subject, + args: JSON.parse(JSON.stringify(descriptor.args)) as Record< + string, + unknown + >, + argsFingerprint: fingerprintDeniedCall( + canonicalTool, + descriptor.args, + descriptor.cwd, + ), + cwd: descriptor.cwd, + ...(descriptor.workspaceRoot !== undefined + ? { workspaceRoot: descriptor.workspaceRoot } + : {}), + createdAt, + expiresAt: createdAt + WORKER_GRANT_TTL_MS, + status: "pending", + audit: [], + }; + audit(envelope, "denied", descriptor.callId); + return envelope; +} + +/** Deny reason keeps `deny` + WORKER_CANNOT_COMPLETE_APPROVAL text and names + * only the envelope requestId — the worker's ask_director text must reference + * that id and carries no authority. */ +export function formatWorkerDenyWithGrantId( + baseReason: string, + requestId: string, +): string { + return `${baseReason} deny recorded under grant request ${requestId}; parent approval is pending — reference only this request id when asking, the text carries no authority.`; +} + +/** Exact-args retry brief built from the harness envelope (never model prose): + * the parent copies this into the existing resume_agent verb with the + * questionId ref. */ +export function buildRetryMessage(envelope: WorkerDeniedCallEnvelope): string { + const ref = + envelope.questionId !== undefined ? ` (ask ${envelope.questionId})` : ""; + return ( + `Parent approved grant request ${envelope.requestId}${ref} for exactly one ` + + `retry. Re-issue exactly this tool call once now — ${envelope.canonicalTool} ` + + `with ${JSON.stringify(envelope.args)} — and do not vary arguments, tool, or cwd.` + ); +} + +export interface WorkerCallIdentity { + sessionId: string; + canonicalTool: string; + args: Record; + cwd: string; + now?: number; +} + +function fingerprintOf(identity: WorkerCallIdentity): string { + return fingerprintDeniedCall( + identity.canonicalTool, + identity.args, + identity.cwd, + ); +} + +function terminalBlocker( + envelope: WorkerDeniedCallEnvelope, + expected: string, +): string { + switch (envelope.status) { + case "consumed": + return `Grant request ${envelope.requestId} was already consumed by its one retry — replaying ${expected} is denied.`; + case "declined": + return `Grant request ${envelope.requestId} was declined — retrying ${expected} is denied.`; + case "expired": + return `Grant request ${envelope.requestId} expired — retrying ${expected} is denied. Re-ask instead.`; + case "interrupted": + return `Grant request ${envelope.requestId} was invalidated by interrupt — retrying ${expected} is denied. Re-ask instead.`; + default: + return `Grant request ${envelope.requestId} is not usable — retrying ${expected} is denied.`; + } +} + +export class WorkerGrantStore { + private readonly envelopes = new Map(); + + register(envelope: WorkerDeniedCallEnvelope): WorkerDeniedCallEnvelope { + this.envelopes.set(envelope.requestId, envelope); + return envelope; + } + + peek(requestId: string): WorkerDeniedCallEnvelope | undefined { + return this.envelopes.get(requestId); + } + + pendingForSession(sessionId: string): WorkerDeniedCallEnvelope | undefined { + for (const envelope of this.envelopes.values()) { + if ( + envelope.workerSessionId === sessionId && + envelope.status === "pending" + ) + return envelope; + } + return undefined; + } + + /** Still-pending envelope for the same exact denied call (dedupes reactor + * retries that mint fresh call ids for the same tool + normalized args). */ + pendingMatch( + sessionId: string, + canonicalTool: string, + args: Record, + cwd: string, + ): WorkerDeniedCallEnvelope | undefined { + const fingerprint = fingerprintDeniedCall(canonicalTool, args, cwd); + for (const envelope of this.envelopes.values()) { + if ( + envelope.workerSessionId === sessionId && + envelope.canonicalTool === canonicalTool && + envelope.argsFingerprint === fingerprint && + envelope.cwd === cwd && + envelope.status === "pending" + ) + return envelope; + } + return undefined; + } + + /** Attach the harness envelope to the worker's ask_director record: stamps + * the questionId so the parent's retry references it. */ + attachToAsk( + sessionId: string, + questionId: string, + ): WorkerDeniedCallEnvelope | undefined { + const envelope = this.pendingForSession(sessionId); + if (envelope === undefined) return undefined; + envelope.questionId = questionId; + audit(envelope, "asked", questionId); + return envelope; + } + + byQuestion(questionId: string): WorkerDeniedCallEnvelope | undefined { + for (const envelope of this.envelopes.values()) { + if (envelope.questionId === questionId) return envelope; + } + return undefined; + } + + /** + * Execution backstop pre-check for a worker call: fail closed when the exact + * call identity matches a non-pending envelope (replay, decline, expiry, + * interrupt), a tampered cwd, or another session's envelope. Pending own + * envelopes and unknown calls return ok and continue down the normal gate + * path — prose and send_input text never reach this check as authority. + * The cwd anchor binds the retry to the denied cwd: a covering session grant + * is not cwd-scoped, so without this check the same args from another + * directory would ride the parent's approval. + */ + precheck( + identity: WorkerCallIdentity, + ): { ok: true } | { ok: false; blocker: string } { + const now = identity.now ?? Date.now(); + const fingerprint = fingerprintOf(identity); + let crossSession: WorkerDeniedCallEnvelope | undefined; + for (const envelope of this.envelopes.values()) { + if (envelope.argsFingerprint !== fingerprint) continue; + if ( + envelope.canonicalTool !== identity.canonicalTool || + envelope.workerSessionId !== identity.sessionId + ) { + if (envelope.status === "pending" || envelope.status === "consumed") + crossSession = envelope; + continue; + } + if (envelope.cwd !== identity.cwd) { + if (envelope.status === "pending" || envelope.status === "consumed") + return { + ok: false, + blocker: + `Grant request ${envelope.requestId} covers ${identity.canonicalTool} ` + + `in ${envelope.cwd}, not ${identity.cwd}: retry from the denied ` + + `directory, or re-ask.`, + }; + return { + ok: false, + blocker: terminalBlocker(envelope, identity.canonicalTool), + }; + } + if (envelope.status !== "pending") { + return { + ok: false, + blocker: terminalBlocker(envelope, identity.canonicalTool), + }; + } + if (envelope.expiresAt <= now) { + envelope.status = "expired"; + audit(envelope, "expired"); + return { + ok: false, + blocker: terminalBlocker(envelope, identity.canonicalTool), + }; + } + return { ok: true }; + } + if (crossSession !== undefined) { + const expected = `${identity.canonicalTool} in ${identity.cwd}`; + return { + ok: false, + blocker: + `Grant request ${crossSession.requestId} does not cover ${expected}: ` + + `replay the exact denied tool from the owning session ${crossSession.workerSessionId}, or re-ask.`, + }; + } + return { ok: true }; + } + + /** + * Consume the pending own-session envelope when the exact call is allowed + * (the parent's operator grant now covers it): exactly one retry per + * requestId. Returns the consumed envelope, or undefined when no pending + * envelope matches (normal allow, nothing to consume). + */ + consumeOnAllow( + identity: WorkerCallIdentity, + ): WorkerDeniedCallEnvelope | undefined { + const now = identity.now ?? Date.now(); + const fingerprint = fingerprintOf(identity); + for (const envelope of this.envelopes.values()) { + if ( + envelope.argsFingerprint !== fingerprint || + envelope.cwd !== identity.cwd || + envelope.canonicalTool !== identity.canonicalTool || + envelope.workerSessionId !== identity.sessionId || + envelope.status !== "pending" || + envelope.expiresAt <= now + ) + continue; + envelope.status = "consumed"; + audit(envelope, "consumed", identity.canonicalTool); + return envelope; + } + return undefined; + } + + decline(requestId: string, reason: string): boolean { + const envelope = this.envelopes.get(requestId); + if (envelope === undefined || envelope.status !== "pending") return false; + envelope.status = "declined"; + audit(envelope, "declined", reason); + return true; + } + + /** Headless parent: decline every pending envelope (single truthful + * blocker each) instead of parking them for an operator who never comes. */ + declineAllForSession(sessionId: string, reason: string): number { + let declined = 0; + for (const envelope of this.envelopes.values()) { + if ( + envelope.workerSessionId !== sessionId || + envelope.status !== "pending" + ) + continue; + envelope.status = "declined"; + audit(envelope, "declined", reason); + declined += 1; + } + return declined; + } + + /** Interrupt invalidation: tombstone the session's envelopes so a later + * replay fails closed with a truthful blocker instead of falling through + * to whatever grant the gate now holds. Retained completion keeps pending + * envelopes — the retained resume_agent retry is the Phase 1 retry path. */ + invalidateSession(sessionId: string, reason: string): number { + let invalidated = 0; + for (const envelope of this.envelopes.values()) { + if ( + envelope.workerSessionId !== sessionId || + envelope.status !== "pending" + ) + continue; + envelope.status = "interrupted"; + audit(envelope, "interrupted", reason); + invalidated += 1; + } + return invalidated; + } + + sweepExpired(now = Date.now()): number { + let expired = 0; + for (const envelope of this.envelopes.values()) { + if (envelope.status !== "pending" || envelope.expiresAt > now) continue; + envelope.status = "expired"; + audit(envelope, "expired"); + expired += 1; + } + return expired; + } + + /** Ask-deadline expiry: the parent never answered this session's ask, so + * its pending envelopes fail closed as expired with a truthful reason. */ + expireSession(sessionId: string, reason: string): number { + let expired = 0; + for (const envelope of this.envelopes.values()) { + if ( + envelope.workerSessionId !== sessionId || + envelope.status !== "pending" + ) + continue; + envelope.status = "expired"; + audit(envelope, "expired", reason); + expired += 1; + } + return expired; + } + + /** Store teardown: invalidate every pending envelope so nothing granted + * before the teardown can be replayed afterwards. */ + invalidateAll(reason: string): number { + let invalidated = 0; + for (const envelope of this.envelopes.values()) { + if (envelope.status !== "pending") continue; + envelope.status = "interrupted"; + audit(envelope, "interrupted", reason); + invalidated += 1; + } + return invalidated; + } + + auditTrail(requestId: string): readonly WorkerGrantAuditEvent[] { + return this.envelopes.get(requestId)?.audit ?? []; + } + + /** Test-only reset: the process store is shared by parent and worker sides. */ + clear(): void { + this.envelopes.clear(); + } +} + +let processStore: WorkerGrantStore | undefined; + +/** Process-shared sidecar: the worker deny side registers, the parent + * observes/consumes — same single-use truth, no message-text authority. */ +export function getProcessWorkerGrantStore(): WorkerGrantStore { + if (processStore === undefined) processStore = new WorkerGrantStore(); + return processStore; +} diff --git a/src/subagent/run.ts b/src/subagent/run.ts index e35b395b5..23ca67e5f 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -524,7 +524,10 @@ const askDirectorDefinition: ToolDefinition = { "Ask the spawning director when the dispatch brief is genuinely ambiguous. " + "You cannot reach the operator. One pending question at a time; " + `at most ${ASK_DIRECTOR_MAX_QUESTIONS} questions per turn; ` + - `${ASK_DIRECTOR_MAX_BYTES} byte cap. The director answers with send_input (soft).`, + `${ASK_DIRECTOR_MAX_BYTES} byte cap. The director answers with send_input (soft). ` + + "For a denied tool call, reference only the grant requestId from the deny " + + "message — the harness-owned envelope carries the exact call, so repeating " + + "tool arguments here grants nothing.", inputSchema: { type: "object", properties: { @@ -598,11 +601,16 @@ async function runSubAgentInner( ): Promise { const inferenceDeps = await assembleInferenceBase(); - const permissionGate = workerPermissionGate(params.permissionGate); - // Identifies this dispatch to submit_result so a submission survives - // only for the turn it was spawned under — a stale call from a redirected - // orchestrator (echoing an old token) is rejected. Steering (followup) - // rotates it: the old token dies with the superseded turn. + // CL-9475: the fleet session id the parent observes (params.id); falls + // back to the local session id below when run without a fleet caller. + // Harness-owned denied-call envelopes are keyed by this id. + let workerGrantSessionId: string | undefined = + params.id !== undefined && /^[A-Za-z0-9_-]+$/.test(params.id) + ? params.id + : undefined; + const permissionGate = workerPermissionGate(params.permissionGate, { + sessionId: () => workerGrantSessionId, + }); let turnToken = params.tier === "leaf" ? generateSessionId() : undefined; const submitResultState = createSubmitResultState(); const askDirectorState = createAskDirectorState(); @@ -1144,6 +1152,9 @@ async function runSubAgentInner( ? params.id : undefined; const sessionId = safeRequestedId ?? generateSessionId(); + // CL-9475: fleet-less runs mint their own id — keep the grant sidecar + // keyed to the same session the parent would observe. + if (workerGrantSessionId === undefined) workerGrantSessionId = sessionId; const workdir = join(params.workdirBase, "subagents", sessionId); await mkdir(workdir, { recursive: true }); childContextDir = workdir; @@ -1176,7 +1187,9 @@ async function runSubAgentInner( const { storage, audit } = await createSessionStores(workdir); childBlobWriter = (key, bytes, contentType) => storage.writeBlob(key, bytes, contentType); - const authorize = createWorkerAuthorize(params.permissionGate); + const authorize = createWorkerAuthorize(params.permissionGate, { + sessionId: () => workerGrantSessionId, + }); const head = { provider: params.provider.providerName, diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index 0542ec481..a19108ec7 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -20,6 +20,10 @@ import { } from "./lifecycle.js"; import type { ForcedStopReason } from "./stop-policy.js"; import { toolCallPreview } from "./tool-preview.js"; +import { + getProcessWorkerGrantStore, + type WorkerDeniedCallEnvelope, +} from "../permission/worker-grant.js"; import type { AdmissionQueue, AdmissionStatus } from "./admission.js"; const log = getLogger([LOG_NAMESPACE_ROOT, "subagent", "session-store"]); @@ -308,7 +312,13 @@ export interface SubAgentSessionStore { resolveAsk(id: string, answer: string): boolean; cancelAsk(id: string, reason?: string): boolean; hasPendingAsk(id: string): boolean; - peekAsk(id: string): { question: string; questionId: string } | undefined; + peekAsk(id: string): + | { + question: string; + questionId: string; + deniedCall?: WorkerDeniedCallEnvelope; + } + | undefined; /** * Ask deadline (CL-8016): reject every pending ask older than `maxAgeMs` * with an explicit timeout error naming its question and session, so a @@ -613,6 +623,10 @@ export function createSubAgentSessionStore( // longer than the bound settles via `expireStaleAsks` instead of // waiting on a wake turn that may never land. askedAt: number; + // Harness-owned denied-call envelope (CL-9475 Phase 1): the exact + // denied ToolCall + hash attached at registerAsk. Parent replays from + // this record — model prose is never authoritative. + deniedCall?: WorkerDeniedCallEnvelope; } >(); const listeners = new Set<() => void>(); @@ -628,6 +642,16 @@ export function createSubAgentSessionStore( lifecycleStatus: projectLifecycleStatus(session.lifecycle), hint: EVICTED_RETENTION_HINT, }); + // CL-9475: an evicted session's denied-call envelopes fail closed — a + // later replay names the eviction instead of riding a lingering grant. + try { + getProcessWorkerGrantStore().invalidateSession( + session.id, + EVICTED_RETENTION_HINT, + ); + } catch { + // Grant invalidation must not throw out of eviction. + } if (evicted.size > MAX_EVICTED_TOMBSTONES) { const oldest = evicted.keys().next().value; if (oldest !== undefined) evicted.delete(oldest); @@ -691,10 +715,23 @@ export function createSubAgentSessionStore( id: string, reason: string, silent = false, + keepGrants = false, ): boolean => { const pending = pendingAsks.get(id); if (pending === undefined) return false; pendingAsks.delete(id); + if (!keepGrants) { + // CL-9475: interrupt/cancel invalidation — tombstone the session's + // denied-call envelopes so a later replay fails closed with a truthful + // blocker instead of falling through to whatever grant the gate holds. + // Retained run-settle keeps them (keepGrants): the retained + // resume_agent retry is the Phase 1 retry path. + try { + getProcessWorkerGrantStore().invalidateSession(id, reason); + } catch { + // Grant invalidation must not throw into settle/interrupt paths. + } + } try { pending.reject(new Error(reason)); } catch { @@ -1899,7 +1936,23 @@ export function createSubAgentSessionStore( if (session === undefined) return false; if (session.lifecycle.state !== "running") return false; if (pendingAsks.has(id)) return false; - pendingAsks.set(id, { ...ask, askedAt: now() }); + // CL-9475: attach the harness-owned denied-call envelope (exact + // ToolCall + hash) to the ask record; stamps its questionId for the + // parent's retry ref. Absent when the ask is not grant-backed. + let deniedCall: WorkerDeniedCallEnvelope | undefined; + try { + deniedCall = + getProcessWorkerGrantStore().attachToAsk(id, ask.questionId) ?? + undefined; + } catch { + // Envelope attach must not fail ask registration. + } + pendingAsks.set( + id, + deniedCall !== undefined + ? { ...ask, askedAt: now(), deniedCall } + : { ...ask, askedAt: now() }, + ); mutate(id, () => undefined); return true; }, @@ -1921,10 +1974,22 @@ export function createSubAgentSessionStore( return pendingAsks.has(id); }, - peekAsk(id: string): { question: string; questionId: string } | undefined { + peekAsk(id: string): + | { + question: string; + questionId: string; + deniedCall?: WorkerDeniedCallEnvelope; + } + | undefined { const pending = pendingAsks.get(id); if (pending === undefined) return undefined; - return { question: pending.question, questionId: pending.questionId }; + return { + question: pending.question, + questionId: pending.questionId, + ...(pending.deniedCall !== undefined + ? { deniedCall: pending.deniedCall } + : {}), + }; }, expireStaleAsks( @@ -1935,6 +2000,17 @@ export function createSubAgentSessionStore( for (const [id, pending] of pendingAsks) { if (pending.askedAt > cutoff) continue; expired.push({ sessionId: id, questionId: pending.questionId }); + // The parent never answered: this session's denied-call envelopes + // fail closed as expired rather than lingering for a later retry. + // Expire first — the cancel below invalidates only still-pending ones. + try { + getProcessWorkerGrantStore().expireSession( + id, + `ask ${pending.questionId} expired without an answer`, + ); + } catch { + // Expiry marking must not throw into the expiry sweep. + } cancelAskInternal( id, `ask_director question ${pending.questionId} for session ${id} expired without an answer after ${maxAgeMs}ms — reply with send_input before the deadline, or not at all`, @@ -2198,7 +2274,9 @@ export function createSubAgentSessionStore( }, settleRun(id: string): void { - cancelAskInternal(id, "run settled"); + // Retained run-settle keeps denied-call envelopes (keepGrants): the + // retained resume_agent retry is the Phase 1 retry path. + cancelAskInternal(id, "run settled", false, true); // CL-7988: the run settled with steers still queued — surface them. dropStashedFollowups(id, "run settled"); if (!runInFlight.delete(id)) return; @@ -2238,6 +2316,11 @@ export function createSubAgentSessionStore( revisions.clear(); snapshotCache.clear(); evicted.clear(); + try { + getProcessWorkerGrantStore().invalidateAll("store cleared"); + } catch { + // Grant invalidation must not throw out of clear. + } notify(); }, @@ -2278,6 +2361,11 @@ export function createSubAgentSessionStore( runInFlight.clear(); revisions.clear(); snapshotCache.clear(); + try { + getProcessWorkerGrantStore().invalidateAll(reason); + } catch { + // Grant invalidation must not throw out of teardown. + } notify(); }, }; From cdc80d428ca0cfb614c3daa2dc5ded2bc1b08537 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 28 Sep 2026 21:06:09 -0700 Subject: [PATCH 2/7] fix(permission): bind ask to its denial, expire open, consume once Ask names its denial via grant_request_id; unnamed keeps first-pending. Expiry yields a fresh gate round that mints a fresh envelope. Precheck-to-consume serializes per fingerprint. --- docs/ARCHITECTURE.md | 2 +- src/permission/reactor-authorize.ts | 76 +++++++++--- src/permission/worker-grant-flow.test.ts | 67 +++++++++++ src/permission/worker-grant.test.ts | 147 ++++++++++++++++++++--- src/permission/worker-grant.ts | 97 +++++++++++++-- src/subagent/agent-fleet.ts | 3 +- src/subagent/ask-director.ts | 5 + src/subagent/session-store.ts | 17 ++- src/subagent/types.ts | 1 + 9 files changed, 368 insertions(+), 47 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 12434e63e..8ba7752a3 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -103,7 +103,7 @@ The interactive TUI may offer an explicit alternate-provider/model selector only ### Worker permission grants -A worker deny-on-ask registers a harness-owned denied-call envelope (`src/permission/worker-grant.ts`): the exact denied ToolCall plus a stable path-aware fingerprint, keyed by worker session with a single-parent-turn expiry. The deny reason names only the envelope's requestId; the worker's `ask_director` prose carries no authority and plain `send_input` text can never mint, consume, or extend an envelope. The parent observes the envelope on the ask record, replays the exact call through its own gate operator path, and retries the retained worker via the existing `resume_agent` with exact args plus the questionId ref. The first covering-grant retry consumes the envelope; replays, tampered args/cwd/tool/session, expiry, decline, or interrupt fail closed with a truthful blocker. There is no dedicated grant verb — single-use and exactness are enforced by the envelope sidecar, not by tool plumbing. +A worker deny-on-ask registers a harness-owned denied-call envelope (`src/permission/worker-grant.ts`): the exact denied ToolCall plus a stable path-aware fingerprint, keyed by worker session with a ten-minute wall-clock TTL. The `turnId` on the envelope is metadata only — no turn enforcement exists, so expiry is purely `expiresAt`-driven. The deny reason names only the envelope's requestId; the worker's `ask_director` prose carries no authority and plain `send_input` text can never mint, consume, or extend an envelope. The worker quotes that requestId back as `grant_request_id`, binding the ask to its own denial at register time; an unnamed ask keeps the legacy first-pending bind. The parent observes the envelope on the ask record via the session-store subscription, replays the exact call through its own gate operator path, and retries the retained worker via the existing `resume_agent` with exact args plus the questionId ref. The first covering-grant retry consumes the envelope; replays, tampered args/cwd/tool/session, decline, or interrupt fail closed with a truthful blocker. Expiry deliberately does not fail closed: the lapsed envelope is marked, skipped, and the retry falls through to a fresh gate round that denies fresh with a new grant id, so a lapsed window can never strand the worker. Expiry is lazy — every pendency read sweeps overdue envelopes, and no periodic scheduler exists by decision, not omission. Concurrent identical retries serialize on a per-fingerprint mutex across the precheck-to-consume gap, so one envelope allows exactly once. There is no dedicated grant verb — single-use and exactness are enforced by the envelope sidecar, not by tool plumbing. Phase 2 (by design, not here): a broad parent session grant can cover more than the exact denied subject; binding the grant to the envelope args needs the atomic grant-and-retry verb. ### TUI Runner (`src/tui/runner/`) diff --git a/src/permission/reactor-authorize.ts b/src/permission/reactor-authorize.ts index b194b95ef..3f9f933fe 100644 --- a/src/permission/reactor-authorize.ts +++ b/src/permission/reactor-authorize.ts @@ -23,6 +23,7 @@ import type { AuthorizeVerdict, GateVerdict, PermissionGate } from "./gate.js"; import type { PermissionRequest } from "./types.js"; import { createDeniedCallEnvelope, + fingerprintDeniedCall, formatWorkerDenyWithGrantId, getProcessWorkerGrantStore, type WorkerDeniedCallEnvelope, @@ -171,6 +172,17 @@ function denyWorkerCallWithEnvelope( }; } +/** Mutex key covering the envelope match: concurrent identical worker calls + * serialize through precheck-to-consume so one envelope allows exactly once. */ +function workerGrantKey( + sessionId: string, + call: ToolCallType, + cwd: string, +): string { + const identity = workerCallIdentity(sessionId, call, cwd); + return `${sessionId}\n${fingerprintDeniedCall(identity.canonicalTool, identity.args, identity.cwd)}`; +} + async function authorizeWorkerCall( gate: PermissionGate, call: ToolCallType, @@ -180,18 +192,36 @@ async function authorizeWorkerCall( if (isWorkerControlPlaneTool(call.name)) return { effect: "allow" }; const workerCwd = resolveWorkerCwd(cwd); const sessionId = resolveWorkerSessionId(grantOptions?.sessionId); - if (sessionId !== undefined) { - const store = grantOptions?.store ?? getProcessWorkerGrantStore(); - const precheck = store.precheck( - workerCallIdentity(sessionId, call, workerCwd), + if (sessionId === undefined) + return authorizeWorkerCallInner(gate, call, grantOptions, workerCwd); + const store = grantOptions?.store ?? getProcessWorkerGrantStore(); + return store.runExclusive(workerGrantKey(sessionId, call, workerCwd), () => + authorizeWorkerCallInner(gate, call, grantOptions, workerCwd, { + sessionId, + store, + }), + ); +} + +async function authorizeWorkerCallInner( + gate: PermissionGate, + call: ToolCallType, + grantOptions: WorkerGrantOptions | undefined, + workerCwd: string, + grant?: { sessionId: string; store: WorkerGrantStore }, +): Promise<{ effect: "allow" } | { effect: "deny"; reason: string }> { + if (grant !== undefined) { + const precheck = grant.store.precheck( + workerCallIdentity(grant.sessionId, call, workerCwd), ); if (!precheck.ok) return { effect: "deny", reason: precheck.blocker }; } const verdict = await gate.authorizeCall(call); if (verdict.effect === "allow") { - if (sessionId !== undefined) { - const store = grantOptions?.store ?? getProcessWorkerGrantStore(); - store.consumeOnAllow(workerCallIdentity(sessionId, call, workerCwd)); + if (grant !== undefined) { + grant.store.consumeOnAllow( + workerCallIdentity(grant.sessionId, call, workerCwd), + ); } return verdict; } @@ -224,18 +254,36 @@ async function executionVerdictWorkerCall( if (isWorkerControlPlaneTool(call.name)) return { effect: "allow" }; const workerCwd = resolveWorkerCwd(cwd); const sessionId = resolveWorkerSessionId(grantOptions?.sessionId); - if (sessionId !== undefined) { - const store = grantOptions?.store ?? getProcessWorkerGrantStore(); - const precheck = store.precheck( - workerCallIdentity(sessionId, call, workerCwd), + if (sessionId === undefined) + return executionVerdictWorkerCallInner(gate, call, grantOptions, workerCwd); + const store = grantOptions?.store ?? getProcessWorkerGrantStore(); + return store.runExclusive(workerGrantKey(sessionId, call, workerCwd), () => + executionVerdictWorkerCallInner(gate, call, grantOptions, workerCwd, { + sessionId, + store, + }), + ); +} + +async function executionVerdictWorkerCallInner( + gate: PermissionGate, + call: ToolCallType, + grantOptions: WorkerGrantOptions | undefined, + workerCwd: string, + grant?: { sessionId: string; store: WorkerGrantStore }, +): Promise { + if (grant !== undefined) { + const precheck = grant.store.precheck( + workerCallIdentity(grant.sessionId, call, workerCwd), ); if (!precheck.ok) return { effect: "deny", reason: precheck.blocker }; } const verdict = await gate.executionVerdict(call); if (verdict.effect === "allow") { - if (sessionId !== undefined) { - const store = grantOptions?.store ?? getProcessWorkerGrantStore(); - store.consumeOnAllow(workerCallIdentity(sessionId, call, workerCwd)); + if (grant !== undefined) { + grant.store.consumeOnAllow( + workerCallIdentity(grant.sessionId, call, workerCwd), + ); } return verdict; } diff --git a/src/permission/worker-grant-flow.test.ts b/src/permission/worker-grant-flow.test.ts index 8593b3dfe..0c0d40140 100644 --- a/src/permission/worker-grant-flow.test.ts +++ b/src/permission/worker-grant-flow.test.ts @@ -6,6 +6,7 @@ import type { ToolCall } from "@intx/types/runtime"; import { createPermissionGate } from "./gate.js"; import { workerPermissionGate } from "./reactor-authorize.js"; import { + WORKER_GRANT_TTL_MS, WorkerGrantStore, createDeniedCallEnvelope, getProcessWorkerGrantStore, @@ -206,4 +207,70 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" expect(store.peek(envelope.requestId)?.status).toBe("pending"); store.expireSession(session.id, "test teardown"); }); + + test("concurrent identical retries allow exactly once", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-race-")); + const parentGate = makeParentGate(cwd, true); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(parentGate, { + sessionId: "worker-race", + store, + }); + + const denied = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + if (denied.effect !== "deny") throw new Error("expected worker deny"); + const requestId = extractRequestId(denied.reason); + expect(await parentGate.evaluate(shellCall("parent-1", COMMAND))).toEqual({ + allowed: true, + }); + + // Two in-flight copies of the exact call (fresh call ids, same args): + // without serialization both pass precheck before either consumes and + // the envelope allows twice. + const [retryA, retryB] = await Promise.all([ + workerGate.authorizeCall(shellCall("c2", COMMAND)), + workerGate.authorizeCall(shellCall("c3", COMMAND)), + ]); + const effects = [retryA.effect, retryB.effect].sort(); + expect(effects).toEqual(["allow", "deny"]); + const loser = retryA.effect === "deny" ? retryA : retryB; + if (loser.effect !== "deny") throw new Error("expected replay deny"); + expect(loser.reason).toContain("already consumed"); + expect(loser.reason).toContain(requestId); + expect(store.peek(requestId)?.status).toBe("consumed"); + }); + + test("post-expiry retry round mints a fresh envelope, then succeeds", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-reissue-")); + const parentGate = makeParentGate(cwd, true); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(parentGate, { + sessionId: "worker-reissue", + store, + }); + + const denied = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + if (denied.effect !== "deny") throw new Error("expected worker deny"); + const lapsedId = extractRequestId(denied.reason); + expect(store.sweepExpired(Date.now() + WORKER_GRANT_TTL_MS + 1)).toBe(1); + expect(store.peek(lapsedId)?.status).toBe("expired"); + + // The lapsed window is not a blackhole: the retry falls through to the + // gate, which denies fresh with a NEW grant id the worker can ask out of. + const reissue = await workerGate.authorizeCall(shellCall("c2", COMMAND)); + expect(reissue.effect).toBe("deny"); + if (reissue.effect !== "deny") throw new Error("expected re-issue deny"); + const freshId = extractRequestId(reissue.reason); + expect(freshId).not.toBe(lapsedId); + expect(store.peek(freshId)?.status).toBe("pending"); + expect(store.peek(lapsedId)?.status).toBe("expired"); + + expect(await parentGate.evaluate(shellCall("parent-1", COMMAND))).toEqual({ + allowed: true, + }); + expect(await workerGate.authorizeCall(shellCall("c3", COMMAND))).toEqual({ + effect: "allow", + }); + expect(store.peek(freshId)?.status).toBe("consumed"); + }); }); diff --git a/src/permission/worker-grant.test.ts b/src/permission/worker-grant.test.ts index f199562bf..da733912f 100644 --- a/src/permission/worker-grant.test.ts +++ b/src/permission/worker-grant.test.ts @@ -196,21 +196,27 @@ describe("WorkerGrantStore lifecycle", () => { expect(result.blocker).toContain("worker-1"); }); - test("expiry fails closed, lazily and via sweep", () => { + test("expiry yields a fresh gate round, never a blackhole", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); const after = envelope.expiresAt + 1; - const result = store.precheck({ - sessionId: "worker-1", - canonicalTool: "run_shell", - args: { ...ARGS }, - cwd: CWD, - now: after, - }); - expect(result.ok).toBe(false); - if (result.ok) throw new Error("expected blocker"); - expect(result.blocker).toContain("expired"); + // The lapsed window passes through instead of denying: the retry falls to + // the gate, which denies fresh and mints a fresh envelope the worker can + // ask out of. + expect( + store.precheck({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: after, + }), + ).toEqual({ ok: true }); expect(envelope.status).toBe("expired"); + expect(envelope.audit.map((event) => event.event)).toEqual([ + "denied", + "expired", + ]); const second = store.register( createDeniedCallEnvelope(descriptor({ callId: "call-2" })), @@ -219,6 +225,53 @@ describe("WorkerGrantStore lifecycle", () => { expect(second.status).toBe("expired"); }); + test("post-expiry re-issue mints a fresh envelope for the same exact call", () => { + const store = new WorkerGrantStore(); + const first = store.register(createDeniedCallEnvelope(descriptor())); + const after = first.expiresAt + 1; + expect(store.sweepExpired(after)).toBe(1); + expect(first.status).toBe("expired"); + // The retry round denies fresh: a new envelope with a new grant id. + const reissue = store.register( + createDeniedCallEnvelope( + descriptor({ callId: "call-2", now: after + 1_000 }), + ), + ); + expect(reissue.requestId).not.toBe(first.requestId); + expect(reissue.status).toBe("pending"); + const identity = { + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: after + 1_001, + }; + expect(store.precheck(identity)).toEqual({ ok: true }); + expect(store.consumeOnAllow(identity)?.requestId).toBe(reissue.requestId); + expect(reissue.status).toBe("consumed"); + expect(first.status).toBe("expired"); + }); + + test("read sites never surface an expired denial without an explicit sweep", () => { + const store = new WorkerGrantStore(); + const envelope = store.register(createDeniedCallEnvelope(descriptor())); + const after = envelope.expiresAt + 1; + // No sweepExpired call here: each read enforces the TTL itself. + expect( + store.pendingMatch("worker-1", "run_shell", { ...ARGS }, CWD, after), + ).toBeUndefined(); + expect( + store.attachToAsk("worker-1", "ask-late", undefined, after), + ).toBeUndefined(); + expect(store.pendingForSession("worker-1", after)).toBeUndefined(); + expect(envelope.status).toBe("expired"); + expect(envelope.audit.map((event) => event.event)).toEqual([ + "denied", + "expired", + ]); + expect(envelope.questionId).toBeUndefined(); + }); + test("decline fails closed; headless declineAllForSession covers the session", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); @@ -258,11 +311,17 @@ describe("WorkerGrantStore lifecycle", () => { test("attachToAsk stamps the questionId and keeps pending", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); - expect(store.attachToAsk("worker-1", "ask-7")).toBe(envelope); + // Reads pin the envelope's clock: the fixture now (1_000_000) predates + // wall time, so an unpinned read would sweep it as expired. + expect(store.attachToAsk("worker-1", "ask-7", undefined, 1_000_001)).toBe( + envelope, + ); expect(envelope.questionId).toBe("ask-7"); expect(envelope.status).toBe("pending"); expect(store.byQuestion("ask-7")).toBe(envelope); - expect(store.attachToAsk("worker-9", "ask-8")).toBeUndefined(); + expect( + store.attachToAsk("worker-9", "ask-8", undefined, 1_000_001), + ).toBeUndefined(); expect(envelope.audit.map((event) => event.event)).toEqual([ "denied", "asked", @@ -272,7 +331,7 @@ describe("WorkerGrantStore lifecycle", () => { test("audit trail orders deny → ask → consume", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); - store.attachToAsk("worker-1", "ask-1"); + store.attachToAsk("worker-1", "ask-1", undefined, 1_000_001); store.consumeOnAllow({ sessionId: "worker-1", canonicalTool: "run_shell", @@ -315,10 +374,68 @@ describe("buildRetryMessage", () => { test("carries exact args + questionId ref from the envelope", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); - store.attachToAsk("worker-1", "ask-3"); + store.attachToAsk("worker-1", "ask-3", undefined, 1_000_001); const message = buildRetryMessage(envelope); expect(message).toContain(envelope.requestId); expect(message).toContain("ask-3"); expect(message).toContain(JSON.stringify(ARGS)); }); }); + +describe("WorkerGrantStore.runExclusive", () => { + test("serializes same-key holders: no overlap, FIFO, error still releases", async () => { + const store = new WorkerGrantStore(); + const events: string[] = []; + let releaseFirst!: () => void; + const firstGate = new Promise((resolve) => { + releaseFirst = resolve; + }); + const first = store.runExclusive("k", async () => { + events.push("first-start"); + await firstGate; + events.push("first-end"); + return "first"; + }); + const second = store.runExclusive("k", async () => { + events.push("second-start"); + events.push("second-end"); + return "second"; + }); + // Let both register; the second holder must not start while the first + // holds the turn across the awaited gate. + await new Promise((resolve) => setTimeout(resolve, 5)); + expect(events).toEqual(["first-start"]); + releaseFirst(); + await expect(first).resolves.toBe("first"); + await expect(second).resolves.toBe("second"); + expect(events).toEqual([ + "first-start", + "first-end", + "second-start", + "second-end", + ]); + + // A throwing holder still releases: the next waiter proceeds. + const failing = store.runExclusive("k", async () => { + throw new Error("boom"); + }); + const after = store.runExclusive("k", async () => "after"); + await expect(failing).rejects.toThrow("boom"); + await expect(after).resolves.toBe("after"); + }); + + test("different keys do not block each other", async () => { + const store = new WorkerGrantStore(); + const order: string[] = []; + await Promise.all([ + store.runExclusive("a", async () => { + await new Promise((resolve) => setTimeout(resolve, 5)); + order.push("a"); + }), + store.runExclusive("b", async () => { + order.push("b"); + }), + ]); + expect(order).toEqual(["b", "a"]); + }); +}); diff --git a/src/permission/worker-grant.ts b/src/permission/worker-grant.ts index 289fb99e3..13375836a 100644 --- a/src/permission/worker-grant.ts +++ b/src/permission/worker-grant.ts @@ -16,7 +16,9 @@ import type { ToolCall } from "@intx/types/runtime"; import { canonicalToolName } from "../agent/canonical-tool-name.js"; import { stableRequestId } from "./denial-memory.js"; -/** Single parent-turn bound: one denied call gets at most one granted retry. */ +/** Wall-clock retry window: one denied call gets at most one granted retry + * within ten minutes of the deny. turnId is metadata only — no turn + * enforcement exists, so expiry is purely expiresAt-driven. */ export const WORKER_GRANT_TTL_MS = 10 * 60 * 1000; export type WorkerGrantStatus = @@ -213,6 +215,32 @@ function terminalBlocker( export class WorkerGrantStore { private readonly envelopes = new Map(); + /** Per-key async mutex chains: each entry resolves when its holder's turn + * ends, so waiters FIFO through the precheck-to-consume gap. */ + private readonly turns = new Map>(); + + /** Serialize concurrent identical worker retries across the + * precheck-to-consume gap: without this, two in-flight copies of the exact + * call both pass precheck before either consumes, and the gate allows both + * — two executions for one envelope. The key must cover the envelope match + * (session + tool/args/cwd fingerprint). Non-reentrant: fn must not call + * runExclusive with the same key. */ + async runExclusive(key: string, fn: () => Promise): Promise { + const prev = this.turns.get(key) ?? Promise.resolve(); + let release!: () => void; + const mine = new Promise((resolve) => { + release = resolve; + }); + const next = prev.then(() => mine); + this.turns.set(key, next); + await prev.catch(() => undefined); + try { + return await fn(); + } finally { + release(); + if (this.turns.get(key) === next) this.turns.delete(key); + } + } register(envelope: WorkerDeniedCallEnvelope): WorkerDeniedCallEnvelope { this.envelopes.set(envelope.requestId, envelope); @@ -223,7 +251,11 @@ export class WorkerGrantStore { return this.envelopes.get(requestId); } - pendingForSession(sessionId: string): WorkerDeniedCallEnvelope | undefined { + pendingForSession( + sessionId: string, + now: number = Date.now(), + ): WorkerDeniedCallEnvelope | undefined { + this.sweepExpired(now); for (const envelope of this.envelopes.values()) { if ( envelope.workerSessionId === sessionId && @@ -235,13 +267,16 @@ export class WorkerGrantStore { } /** Still-pending envelope for the same exact denied call (dedupes reactor - * retries that mint fresh call ids for the same tool + normalized args). */ + * retries that mint fresh call ids for the same tool + normalized args). + * Sweeps overdue envelopes first so a lapse can never read as pending. */ pendingMatch( sessionId: string, canonicalTool: string, args: Record, cwd: string, + now: number = Date.now(), ): WorkerDeniedCallEnvelope | undefined { + this.sweepExpired(now); const fingerprint = fingerprintDeniedCall(canonicalTool, args, cwd); for (const envelope of this.envelopes.values()) { if ( @@ -257,12 +292,36 @@ export class WorkerGrantStore { } /** Attach the harness envelope to the worker's ask_director record: stamps - * the questionId so the parent's retry references it. */ + * the questionId so the parent's retry references it. When the ask names + * its denial (the grant requestId quoted from the deny reason), bind that + * exact envelope — first-pending-wins would join question-about-B to + * exact-call-A when two denies share a session, poisoning the audit and + * (Phase 2) replaying the wrong call. A named id that resolves to no + * pending own-session envelope fails closed with no attach (never falls + * back to another denial). An unnamed ask keeps the legacy first-pending + * bind for the single-deny case. Sweeps overdue envelopes first: an ask + * can never join an expired denial. */ attachToAsk( sessionId: string, questionId: string, + requestId?: string, + now: number = Date.now(), ): WorkerDeniedCallEnvelope | undefined { - const envelope = this.pendingForSession(sessionId); + this.sweepExpired(now); + const named = requestId?.trim() || undefined; + if (named !== undefined) { + const envelope = this.envelopes.get(named); + if ( + envelope === undefined || + envelope.workerSessionId !== sessionId || + envelope.status !== "pending" + ) + return undefined; + envelope.questionId = questionId; + audit(envelope, "asked", questionId); + return envelope; + } + const envelope = this.pendingForSession(sessionId, now); if (envelope === undefined) return undefined; envelope.questionId = questionId; audit(envelope, "asked", questionId); @@ -278,10 +337,13 @@ export class WorkerGrantStore { /** * Execution backstop pre-check for a worker call: fail closed when the exact - * call identity matches a non-pending envelope (replay, decline, expiry, - * interrupt), a tampered cwd, or another session's envelope. Pending own - * envelopes and unknown calls return ok and continue down the normal gate - * path — prose and send_input text never reach this check as authority. + * call identity matches a terminal envelope (replay, decline, interrupt), + * a tampered cwd, or another session's envelope. An EXPIRED envelope is + * marked and skipped instead of denying: the lapsed window must yield a + * fresh gate round that mints a fresh envelope, never a blackhole the + * worker can never re-ask out of. Pending own envelopes and unknown calls + * return ok and continue down the normal gate path — prose and send_input + * text never reach this check as authority. * The cwd anchor binds the retry to the denied cwd: a covering session grant * is not cwd-scoped, so without this check the same args from another * directory would ride the parent's approval. @@ -290,6 +352,7 @@ export class WorkerGrantStore { identity: WorkerCallIdentity, ): { ok: true } | { ok: false; blocker: string } { const now = identity.now ?? Date.now(); + this.sweepExpired(now); const fingerprint = fingerprintOf(identity); let crossSession: WorkerDeniedCallEnvelope | undefined; for (const envelope of this.envelopes.values()) { @@ -316,6 +379,7 @@ export class WorkerGrantStore { blocker: terminalBlocker(envelope, identity.canonicalTool), }; } + if (envelope.status === "expired") continue; if (envelope.status !== "pending") { return { ok: false, @@ -325,10 +389,7 @@ export class WorkerGrantStore { if (envelope.expiresAt <= now) { envelope.status = "expired"; audit(envelope, "expired"); - return { - ok: false, - blocker: terminalBlocker(envelope, identity.canonicalTool), - }; + continue; } return { ok: true }; } @@ -341,6 +402,11 @@ export class WorkerGrantStore { `replay the exact denied tool from the owning session ${crossSession.workerSessionId}, or re-ask.`, }; } + // Phase 2 (by design, do not tighten here): a call whose fingerprint + // matches no envelope falls through to the normal gate path, where a + // broad parent session grant can cover more than the exact denied + // subject. Binding the grant to the envelope args needs the dedicated + // atomic grant-and-retry verb. return { ok: true }; } @@ -416,6 +482,11 @@ export class WorkerGrantStore { return invalidated; } + /** Lazy-expiry engine: marks overdue pending envelopes expired with an + * audit event. Every pendency read (precheck, pendingMatch, + * pendingForSession, attachToAsk) sweeps through here so a lapsed window + * can never read as pending; no periodic scheduler exists, and none is + * needed while every read enforces the TTL. */ sweepExpired(now = Date.now()): number { let expired = 0; for (const envelope of this.envelopes.values()) { diff --git a/src/subagent/agent-fleet.ts b/src/subagent/agent-fleet.ts index 2da7c9b91..766a2bfef 100644 --- a/src/subagent/agent-fleet.ts +++ b/src/subagent/agent-fleet.ts @@ -1449,7 +1449,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool { : {}), persist: deps.persist !== false, askDirectorPort: { - register: ({ question, questionId }) => { + register: ({ question, questionId, grantRequestId }) => { const hold: { resolve?: (answer: string) => void; reject?: (reason: unknown) => void; @@ -1469,6 +1469,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool { const ok = deps.sessions.registerAsk(session.id, { question, questionId, + ...(grantRequestId !== undefined ? { grantRequestId } : {}), resolve: hold.resolve, reject: hold.reject, }); diff --git a/src/subagent/ask-director.ts b/src/subagent/ask-director.ts index 895e06d7c..dc3f11631 100644 --- a/src/subagent/ask-director.ts +++ b/src/subagent/ask-director.ts @@ -24,6 +24,7 @@ export interface AskDirectorPort { register: (input: { question: string; questionId: string; + grantRequestId?: string; }) => Promise; cancel: (reason: string) => void; } @@ -117,6 +118,7 @@ export function createDeferredContinuation(): { export async function handleAskDirector(args: { question: unknown; + grantRequestId?: unknown; state: AskDirectorState; port: AskDirectorPort; signal: AbortSignal; @@ -138,10 +140,13 @@ export async function handleAskDirector(args: { return "Error: ask_director was cancelled."; } const questionId = `ask-${args.state.questions + 1}`; + const grantRequestId = + typeof args.grantRequestId === "string" ? args.grantRequestId : undefined; try { const answerP = args.port.register({ question: outcome.question, questionId, + ...(grantRequestId !== undefined ? { grantRequestId } : {}), }); if (args.signal.aborted) { onAbort(); diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index a19108ec7..f2b17db73 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -305,6 +305,9 @@ export interface SubAgentSessionStore { ask: { question: string; questionId: string; + /** Grant requestId quoted from the deny reason: binds this ask to + * its own denial instead of the session's first pending envelope. */ + grantRequestId?: string; resolve: (answer: string) => void; reject: (reason: unknown) => void; }, @@ -1928,6 +1931,9 @@ export function createSubAgentSessionStore( ask: { question: string; questionId: string; + /** Grant requestId quoted from the deny reason: binds this ask to + * its own denial instead of the session's first pending envelope. */ + grantRequestId?: string; resolve: (answer: string) => void; reject: (reason: unknown) => void; }, @@ -1938,12 +1944,17 @@ export function createSubAgentSessionStore( if (pendingAsks.has(id)) return false; // CL-9475: attach the harness-owned denied-call envelope (exact // ToolCall + hash) to the ask record; stamps its questionId for the - // parent's retry ref. Absent when the ask is not grant-backed. + // parent's retry ref. A named grantRequestId binds that exact denial; + // an unnamed ask falls back to the session's first pending envelope. + // Absent when the ask is not grant-backed. let deniedCall: WorkerDeniedCallEnvelope | undefined; try { deniedCall = - getProcessWorkerGrantStore().attachToAsk(id, ask.questionId) ?? - undefined; + getProcessWorkerGrantStore().attachToAsk( + id, + ask.questionId, + ask.grantRequestId, + ) ?? undefined; } catch { // Envelope attach must not fail ask registration. } diff --git a/src/subagent/types.ts b/src/subagent/types.ts index bb17fd35a..29d57309d 100644 --- a/src/subagent/types.ts +++ b/src/subagent/types.ts @@ -215,6 +215,7 @@ export type RunSubAgentParams = { register: (input: { question: string; questionId: string; + grantRequestId?: string; }) => Promise; cancel: (reason: string) => void; }; From 64605cccd9fba0440b28ead0544f6741f596b759 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 28 Sep 2026 21:21:53 -0700 Subject: [PATCH 3/7] fix(subagent): thread grant_request_id through ask_director leaf --- src/permission/worker-grant.test.ts | 65 +++++++ src/subagent/run-ask-director-grant.test.ts | 178 ++++++++++++++++++++ src/subagent/run.ts | 6 + 3 files changed, 249 insertions(+) create mode 100644 src/subagent/run-ask-director-grant.test.ts diff --git a/src/permission/worker-grant.test.ts b/src/permission/worker-grant.test.ts index da733912f..39d34e13c 100644 --- a/src/permission/worker-grant.test.ts +++ b/src/permission/worker-grant.test.ts @@ -328,6 +328,71 @@ describe("WorkerGrantStore lifecycle", () => { ]); }); + describe("two-denies correlation", () => { + function twoDenies() { + const store = new WorkerGrantStore(); + const envelopeA = store.register( + createDeniedCallEnvelope( + descriptor({ callId: "call-A", args: { command: "npm test" } }), + ), + ); + const envelopeB = store.register( + createDeniedCallEnvelope( + descriptor({ callId: "call-B", args: { command: "npm run lint" } }), + ), + ); + return { store, envelopeA, envelopeB }; + } + + test("named ask binds its exact denial, not the first pending", () => { + const { store, envelopeA, envelopeB } = twoDenies(); + expect( + store.attachToAsk("worker-1", "ask-B", envelopeB.requestId, 1_000_001), + ).toBe(envelopeB); + expect(envelopeB.questionId).toBe("ask-B"); + expect(envelopeA.questionId).toBeUndefined(); + expect(store.byQuestion("ask-B")).toBe(envelopeB); + }); + + test("unnamed ask keeps the legacy first-pending bind", () => { + const { store, envelopeA, envelopeB } = twoDenies(); + expect(store.attachToAsk("worker-1", "ask-x", undefined, 1_000_001)).toBe( + envelopeA, + ); + expect(envelopeA.questionId).toBe("ask-x"); + expect(envelopeB.questionId).toBeUndefined(); + }); + + test("bogus id fails closed with no fallback to another denial", () => { + const { store, envelopeA, envelopeB } = twoDenies(); + expect( + store.attachToAsk("worker-1", "ask-x", "not-a-real-id", 1_000_001), + ).toBeUndefined(); + expect(envelopeA.questionId).toBeUndefined(); + expect(envelopeB.questionId).toBeUndefined(); + expect(store.byQuestion("ask-x")).toBeUndefined(); + }); + + test("cross-session id fails closed with no fallback", () => { + const { store, envelopeA, envelopeB } = twoDenies(); + const other = store.register( + createDeniedCallEnvelope( + descriptor({ + callId: "call-other", + args: { command: "npm run build" }, + workerSessionId: "worker-9", + }), + ), + ); + expect( + store.attachToAsk("worker-1", "ask-x", other.requestId, 1_000_001), + ).toBeUndefined(); + expect(envelopeA.questionId).toBeUndefined(); + expect(envelopeB.questionId).toBeUndefined(); + expect(other.questionId).toBeUndefined(); + }); + }); + test("audit trail orders deny → ask → consume", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); diff --git a/src/subagent/run-ask-director-grant.test.ts b/src/subagent/run-ask-director-grant.test.ts new file mode 100644 index 000000000..090585f5e --- /dev/null +++ b/src/subagent/run-ask-director-grant.test.ts @@ -0,0 +1,178 @@ +/** + * run.ts must thread grant_request_id from the ask_director tool call through + * the leaf handler into handleAskDirector and out the askDirectorPort. + * Dropping it at the schema or the handler silently forces every production + * ask down the legacy first-pending path, so this test drives the real + * handler path end to end — unit tests on handleAskDirector alone cannot + * see the two call sites this covers. + */ +import { describe, expect, test } from "bun:test"; +import { mkdtemp } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +import type { DirectorFactory } from "@intx/agent"; + +import { withMockedModuleDuring } from "../../tests/helpers/mock-module.js"; +import { createPermissionGate } from "../permission/gate.js"; +import type { RunSubAgentParams } from "./types.js"; + +const testPermissionGate = createPermissionGate({ + approvals: [], + interactive: false, + skipPermissions: true, + reactorGated: false, +}); + +function createHangingStubAgent() { + return { + async send(_content: string, optsSend?: { signal?: AbortSignal }) { + return await new Promise((_, reject) => { + if (optsSend?.signal?.aborted === true) { + reject( + optsSend.signal.reason instanceof Error + ? optsSend.signal.reason + : new Error("aborted"), + ); + return; + } + optsSend?.signal?.addEventListener( + "abort", + () => { + const reason = optsSend.signal?.reason; + reject(reason instanceof Error ? reason : new Error("aborted")); + }, + { once: true }, + ); + }); + }, + stream: () => + (async function* () { + yield* []; + })(), + deliver: () => undefined, + close: async () => undefined, + setSource: () => undefined, + setSources: () => undefined, + history: async () => [], + checkpoints: async () => [], + readAt: async () => [], + blobReader: {}, + }; +} + +describe("runSubAgent ask_director grant_request_id threading", () => { + test("leaf forwards grant_request_id to the port; absent stays legacy", async () => { + const cwd = await mkdtemp(join(tmpdir(), "cl9475-ask-grant-")); + const registered: { + question: string; + questionId: string; + grantRequestId?: string; + }[] = []; + let capturedAskHandler: + | (( + rawArgs: Record, + signal: AbortSignal, + ) => Promise | string) + | undefined; + let capturedFactory: DirectorFactory | undefined; + let handles: { close: (ms?: number) => Promise } | undefined; + + const outcome = await withMockedModuleDuring( + import.meta.resolve("@intx/agent"), + (real: typeof import("@intx/agent")) => ({ + ...real, + defineDirector: (opts: Parameters[0]) => { + capturedFactory = opts.factory; + return real.defineDirector(opts); + }, + stringTool: (args: Parameters[0]) => { + const tool = real.stringTool(args); + if (args.definition.name === "ask_director") { + capturedAskHandler = args.handler; + } + return tool; + }, + }), + async () => + await withMockedModuleDuring( + import.meta.resolve("../agent/live-tool-dispatch.js"), + (real: typeof import("../agent/live-tool-dispatch.js")) => ({ + ...real, + createAgentWithLiveToolDispatch: async () => + createHangingStubAgent() as unknown as Awaited< + ReturnType + >, + }), + async () => { + const { runSubAgent } = await import("./run.js"); + + const params: RunSubAgentParams = { + cwd, + workdirBase: join(cwd, ".ctx"), + permissionGate: testPermissionGate, + provider: { + providerName: "test", + baseURL: "http://localhost", + model: "test-model", + }, + description: "ask_director grant threading", + prompt: "hold for ask_director", + persist: true, + tier: "leaf", + askDirectorPort: { + register: (input) => { + registered.push({ ...input }); + return Promise.resolve(`answer:${registered.length}`); + }, + cancel: () => undefined, + }, + onAgentReady: (h) => { + handles = h; + }, + }; + + const runPromise = runSubAgent(params); + for (let i = 0; i < 500 && handles === undefined; i++) { + await new Promise((resolve) => setTimeout(resolve, 1)); + } + if (handles === undefined) + throw new Error("onAgentReady never fired"); + if (capturedFactory === undefined) + throw new Error("defineDirector factory not seen"); + if (capturedAskHandler === undefined) + throw new Error("ask_director tool was not mounted"); + + // Exercise the factory so the leaf tool wiring is live. + capturedFactory({}, {} as never, { + systemPrompt: "system", + toolDefinitions: [], + compactorNames: [], + }); + + const named = await capturedAskHandler( + { question: "may I retry?", grant_request_id: "grant-123" }, + new AbortController().signal, + ); + expect(named).toBe("answer:1"); + + const legacy = await capturedAskHandler( + { question: "plain question" }, + new AbortController().signal, + ); + expect(legacy).toBe("answer:2"); + + await handles.close().catch(() => undefined); + await runPromise.catch(() => undefined); + return { registered }; + }, + ), + ); + + expect(outcome.registered.length).toBe(2); + expect(outcome.registered[0]?.question).toBe("may I retry?"); + expect(outcome.registered[0]?.grantRequestId).toBe("grant-123"); + expect(outcome.registered[1]?.question).toBe("plain question"); + expect(outcome.registered[1]?.grantRequestId).toBeUndefined(); + }); +}); diff --git a/src/subagent/run.ts b/src/subagent/run.ts index 23ca67e5f..736fe8598 100644 --- a/src/subagent/run.ts +++ b/src/subagent/run.ts @@ -535,6 +535,11 @@ const askDirectorDefinition: ToolDefinition = { type: "string", description: "The question for the spawning director (non-empty).", }, + grant_request_id: { + type: "string", + description: + "Grant request id quoted from the deny reason; binds this ask to its own denial.", + }, }, required: ["question"], }, @@ -859,6 +864,7 @@ async function runSubAgentInner( try { return await handleAskDirector({ question: rawArgs.question, + grantRequestId: rawArgs.grant_request_id, state: askDirectorState, port, signal, From a0c81aff810b07ebff76d4f5d7d275769fab0fc9 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Tue, 29 Sep 2026 11:15:24 -0700 Subject: [PATCH 4/7] test(subagent): follow test-helper move to repo-root testkit origin/main relocated shared test helpers from tests/helpers to testkit; the grant threading test follows the mock-module import. --- src/subagent/run-ask-director-grant.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/subagent/run-ask-director-grant.test.ts b/src/subagent/run-ask-director-grant.test.ts index 090585f5e..6c10ff701 100644 --- a/src/subagent/run-ask-director-grant.test.ts +++ b/src/subagent/run-ask-director-grant.test.ts @@ -13,7 +13,7 @@ import { join } from "node:path"; import type { DirectorFactory } from "@intx/agent"; -import { withMockedModuleDuring } from "../../tests/helpers/mock-module.js"; +import { withMockedModuleDuring } from "../../testkit/mock-module.js"; import { createPermissionGate } from "../permission/gate.js"; import type { RunSubAgentParams } from "./types.js"; From 0f69c3c3c4cb18334cfa5c83c725b2db5677ab2f Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Tue, 29 Sep 2026 11:17:30 -0700 Subject: [PATCH 5/7] test(permission): cover envelope reuse across reactor retries Two denies of the same exact call under fresh call ids name the same grant request instead of prompting the parent twice. --- src/permission/worker-grant-flow.test.ts | 20 ++++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/src/permission/worker-grant-flow.test.ts b/src/permission/worker-grant-flow.test.ts index 0c0d40140..6baa3cb3f 100644 --- a/src/permission/worker-grant-flow.test.ts +++ b/src/permission/worker-grant-flow.test.ts @@ -82,6 +82,26 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" expect(replayAgain.reason).toContain(requestId); }); + test("reactor retry of the same exact call reuses the pending envelope", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-dedupe-")); + const parentGate = makeParentGate(cwd, true); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(parentGate, { + sessionId: "worker-dedupe", + store, + }); + + const first = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + if (first.effect !== "deny") throw new Error("expected worker deny"); + // Reactor retries mint fresh call ids for the same tool + args: the + // deny must name the same grant request, not prompt the parent twice. + const second = await workerGate.authorizeCall(shellCall("c2", COMMAND)); + if (second.effect !== "deny") throw new Error("expected worker deny"); + expect(extractRequestId(second.reason)).toBe( + extractRequestId(first.reason), + ); + }); + test("tampered args fall through to a fresh deny; original stays consumed", async () => { const cwd = mkdtempSync(join(tmpdir(), "worker-grant-tamper-")); const parentGate = makeParentGate(cwd, true); From 9dfe7e357e4ea219e78ea21fc11d7d4c8d634168 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Tue, 29 Sep 2026 12:42:45 -0700 Subject: [PATCH 6/7] fix(permission): spend worker grant at execution, session-scope precheck Authorize leaves the envelope pending so the granted retry survives the execution backstop, which now spends it exactly once; precheck no longer vetoes sibling sessions, and unwired helpers are removed. --- src/permission/reactor-authorize.ts | 27 ++++------ src/permission/worker-grant.ts | 84 ++++------------------------- src/subagent/session-store.ts | 1 + 3 files changed, 23 insertions(+), 89 deletions(-) diff --git a/src/permission/reactor-authorize.ts b/src/permission/reactor-authorize.ts index 3f9f933fe..acdc7dd51 100644 --- a/src/permission/reactor-authorize.ts +++ b/src/permission/reactor-authorize.ts @@ -26,7 +26,6 @@ import { fingerprintDeniedCall, formatWorkerDenyWithGrantId, getProcessWorkerGrantStore, - type WorkerDeniedCallEnvelope, type WorkerGrantStore, } from "./worker-grant.js"; import { canonicalToolName } from "../agent/canonical-tool-name.js"; @@ -87,8 +86,6 @@ export interface WorkerGrantOptions { /** Owning worker session; absent means no envelope is minted or matched. */ sessionId?: string | (() => string | undefined); workspaceRoot?: string; - /** Observability hook for tests/telemetry; never authoritative. */ - onDeniedCall?: (envelope: WorkerDeniedCallEnvelope) => void; } function resolveWorkerSessionId( @@ -161,11 +158,6 @@ function denyWorkerCallWithEnvelope( : {}), }), ); - try { - options?.onDeniedCall?.(envelope); - } catch { - // Observability must not throw into the deny path. - } return { effect: "deny", reason: formatWorkerDenyWithGrantId(baseReason, envelope.requestId), @@ -217,14 +209,14 @@ async function authorizeWorkerCallInner( if (!precheck.ok) return { effect: "deny", reason: precheck.blocker }; } const verdict = await gate.authorizeCall(call); - if (verdict.effect === "allow") { - if (grant !== undefined) { - grant.store.consumeOnAllow( - workerCallIdentity(grant.sessionId, call, workerCwd), - ); - } - return verdict; - } + // Authorize never consumes: it only permits the call to proceed to the + // tool-runner middleware, which enforces executionVerdict. Consuming here + // would spend the envelope before execution, so the execution backstop's + // precheck would deny the granted retry as "already consumed". The single + // consumption happens in executionVerdictWorkerCallInner below — the stage + // that actually permits execution — where the same mutex still serializes + // concurrent identical retries to exactly one. + if (verdict.effect === "allow") return verdict; if (verdict.effect !== "ask") return verdict; return denyWorkerCallWithEnvelope( grantOptions, @@ -279,6 +271,9 @@ async function executionVerdictWorkerCallInner( if (!precheck.ok) return { effect: "deny", reason: precheck.blocker }; } const verdict = await gate.executionVerdict(call); + // Sole consumption point: the execution backstop spends the envelope when + // the exact call is allowed, so one grant request permits exactly one + // execution. Authorize deliberately does not consume (see above). if (verdict.effect === "allow") { if (grant !== undefined) { grant.store.consumeOnAllow( diff --git a/src/permission/worker-grant.ts b/src/permission/worker-grant.ts index 13375836a..581964835 100644 --- a/src/permission/worker-grant.ts +++ b/src/permission/worker-grant.ts @@ -24,18 +24,11 @@ export const WORKER_GRANT_TTL_MS = 10 * 60 * 1000; export type WorkerGrantStatus = | "pending" | "consumed" - | "declined" | "expired" | "interrupted"; export interface WorkerGrantAuditEvent { - event: - | "denied" - | "asked" - | "consumed" - | "declined" - | "expired" - | "interrupted"; + event: "denied" | "asked" | "consumed" | "expired" | "interrupted"; at: number; detail?: string; } @@ -166,19 +159,6 @@ export function formatWorkerDenyWithGrantId( return `${baseReason} deny recorded under grant request ${requestId}; parent approval is pending — reference only this request id when asking, the text carries no authority.`; } -/** Exact-args retry brief built from the harness envelope (never model prose): - * the parent copies this into the existing resume_agent verb with the - * questionId ref. */ -export function buildRetryMessage(envelope: WorkerDeniedCallEnvelope): string { - const ref = - envelope.questionId !== undefined ? ` (ask ${envelope.questionId})` : ""; - return ( - `Parent approved grant request ${envelope.requestId}${ref} for exactly one ` + - `retry. Re-issue exactly this tool call once now — ${envelope.canonicalTool} ` + - `with ${JSON.stringify(envelope.args)} — and do not vary arguments, tool, or cwd.` - ); -} - export interface WorkerCallIdentity { sessionId: string; canonicalTool: string; @@ -202,8 +182,6 @@ function terminalBlocker( switch (envelope.status) { case "consumed": return `Grant request ${envelope.requestId} was already consumed by its one retry — replaying ${expected} is denied.`; - case "declined": - return `Grant request ${envelope.requestId} was declined — retrying ${expected} is denied.`; case "expired": return `Grant request ${envelope.requestId} expired — retrying ${expected} is denied. Re-ask instead.`; case "interrupted": @@ -328,17 +306,12 @@ export class WorkerGrantStore { return envelope; } - byQuestion(questionId: string): WorkerDeniedCallEnvelope | undefined { - for (const envelope of this.envelopes.values()) { - if (envelope.questionId === questionId) return envelope; - } - return undefined; - } - /** * Execution backstop pre-check for a worker call: fail closed when the exact - * call identity matches a terminal envelope (replay, decline, interrupt), - * a tampered cwd, or another session's envelope. An EXPIRED envelope is + * call identity matches a terminal envelope (replay, interrupt) or a + * tampered cwd. Envelopes from other sessions never veto this session: + * each session mints and spends its own envelope through its own grant + * round. An EXPIRED envelope is * marked and skipped instead of denying: the lapsed window must yield a * fresh gate round that mints a fresh envelope, never a blackhole the * worker can never re-ask out of. Pending own envelopes and unknown calls @@ -354,15 +327,18 @@ export class WorkerGrantStore { const now = identity.now ?? Date.now(); this.sweepExpired(now); const fingerprint = fingerprintOf(identity); - let crossSession: WorkerDeniedCallEnvelope | undefined; for (const envelope of this.envelopes.values()) { if (envelope.argsFingerprint !== fingerprint) continue; + // Another session's envelope never vetoes this session: each session + // mints and spends its own envelope through its own grant round, so a + // sibling's pending or already-consumed envelope is irrelevant here. + // (Reactor retries share the worker's session, so consume-once within + // the session survives this scoping. Spending stays session-scoped in + // consumeOnAllow, and tampered args/cwd still fail closed below.) if ( envelope.canonicalTool !== identity.canonicalTool || envelope.workerSessionId !== identity.sessionId ) { - if (envelope.status === "pending" || envelope.status === "consumed") - crossSession = envelope; continue; } if (envelope.cwd !== identity.cwd) { @@ -393,15 +369,6 @@ export class WorkerGrantStore { } return { ok: true }; } - if (crossSession !== undefined) { - const expected = `${identity.canonicalTool} in ${identity.cwd}`; - return { - ok: false, - blocker: - `Grant request ${crossSession.requestId} does not cover ${expected}: ` + - `replay the exact denied tool from the owning session ${crossSession.workerSessionId}, or re-ask.`, - }; - } // Phase 2 (by design, do not tighten here): a call whose fingerprint // matches no envelope falls through to the normal gate path, where a // broad parent session grant can cover more than the exact denied @@ -438,31 +405,6 @@ export class WorkerGrantStore { return undefined; } - decline(requestId: string, reason: string): boolean { - const envelope = this.envelopes.get(requestId); - if (envelope === undefined || envelope.status !== "pending") return false; - envelope.status = "declined"; - audit(envelope, "declined", reason); - return true; - } - - /** Headless parent: decline every pending envelope (single truthful - * blocker each) instead of parking them for an operator who never comes. */ - declineAllForSession(sessionId: string, reason: string): number { - let declined = 0; - for (const envelope of this.envelopes.values()) { - if ( - envelope.workerSessionId !== sessionId || - envelope.status !== "pending" - ) - continue; - envelope.status = "declined"; - audit(envelope, "declined", reason); - declined += 1; - } - return declined; - } - /** Interrupt invalidation: tombstone the session's envelopes so a later * replay fails closed with a truthful blocker instead of falling through * to whatever grant the gate now holds. Retained completion keeps pending @@ -528,10 +470,6 @@ export class WorkerGrantStore { return invalidated; } - auditTrail(requestId: string): readonly WorkerGrantAuditEvent[] { - return this.envelopes.get(requestId)?.audit ?? []; - } - /** Test-only reset: the process store is shared by parent and worker sides. */ clear(): void { this.envelopes.clear(); diff --git a/src/subagent/session-store.ts b/src/subagent/session-store.ts index f2b17db73..530b46a81 100644 --- a/src/subagent/session-store.ts +++ b/src/subagent/session-store.ts @@ -1954,6 +1954,7 @@ export function createSubAgentSessionStore( id, ask.questionId, ask.grantRequestId, + now(), ) ?? undefined; } catch { // Envelope attach must not fail ask registration. From a36b57efeab1c9824461818123c4ad4b7ec35c9e Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Tue, 29 Sep 2026 12:42:47 -0700 Subject: [PATCH 7/7] test(permission): keep two-stage execute and per-session envelopes --- src/permission/worker-grant-flow.test.ts | 147 +++++++++++++++++++---- src/permission/worker-grant.test.ts | 63 +++------- 2 files changed, 137 insertions(+), 73 deletions(-) diff --git a/src/permission/worker-grant-flow.test.ts b/src/permission/worker-grant-flow.test.ts index 6baa3cb3f..102f3230d 100644 --- a/src/permission/worker-grant-flow.test.ts +++ b/src/permission/worker-grant-flow.test.ts @@ -71,15 +71,14 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" const retry = await workerGate.authorizeCall(shellCall("c2", COMMAND)); expect(retry).toEqual({ effect: "allow" }); - expect(store.peek(requestId)?.status).toBe("consumed"); + expect(store.peek(requestId)?.status).toBe("pending"); - const replayAgain = await workerGate.authorizeCall( - shellCall("c3", COMMAND), - ); - expect(replayAgain.effect).toBe("deny"); - if (replayAgain.effect !== "deny") throw new Error("expected replay deny"); - expect(replayAgain.reason).toContain("already consumed"); - expect(replayAgain.reason).toContain(requestId); + // Authorize is not execution: a second authorize of the same call still + // allows without spending the envelope. Exactly-once is enforced when the + // call executes (see the two-stage KEEPER below). + const retryAgain = await workerGate.authorizeCall(shellCall("c3", COMMAND)); + expect(retryAgain).toEqual({ effect: "allow" }); + expect(store.peek(requestId)?.status).toBe("pending"); }); test("reactor retry of the same exact call reuses the pending envelope", async () => { @@ -102,7 +101,7 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" ); }); - test("tampered args fall through to a fresh deny; original stays consumed", async () => { + test("tampered args fall through to a fresh deny; original stays pending", async () => { const cwd = mkdtempSync(join(tmpdir(), "worker-grant-tamper-")); const parentGate = makeParentGate(cwd, true); const store = new WorkerGrantStore(); @@ -127,10 +126,10 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" expect(tampered.effect).toBe("deny"); if (tampered.effect !== "deny") throw new Error("expected tamper deny"); expect(tampered.reason).not.toContain("already consumed"); - expect(store.peek(requestId)?.status).toBe("consumed"); + expect(store.peek(requestId)?.status).toBe("pending"); }); - test("headless parent cannot grant: replay denies, auto-decline fails closed", async () => { + test("headless parent cannot grant: retry denies, envelope stays pending", async () => { const cwd = mkdtempSync(join(tmpdir(), "worker-grant-headless-")); // Interactive worker gate so the deny registers an envelope; the PARENT // gate is headless (no operator seam). @@ -164,17 +163,14 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" const replay = await headlessGate.evaluate(shellCall("parent-1", COMMAND)); expect(replay.allowed).toBe(false); - expect( - store.declineAllForSession( - "worker-headless", - "headless parent: no operator to approve", - ), - ).toBe(1); - expect(store.peek(requestId)?.status).toBe("declined"); + // No operator, no approval, no decline verb: the retry keeps denying + // against the same pending envelope — nothing is spent, nothing is + // granted, and the deny still names the original grant request. const retry = await workerGate.authorizeCall(shellCall("c2", COMMAND)); expect(retry.effect).toBe("deny"); - if (retry.effect !== "deny") throw new Error("expected declined deny"); - expect(retry.reason).toContain("declined"); + if (retry.effect !== "deny") throw new Error("expected headless deny"); + expect(extractRequestId(retry.reason)).toBe(requestId); + expect(store.peek(requestId)?.status).toBe("pending"); }); test("send_input text answers the ask but grants nothing", () => { @@ -228,7 +224,7 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" store.expireSession(session.id, "test teardown"); }); - test("concurrent identical retries allow exactly once", async () => { + test("concurrent identical executions allow exactly once", async () => { const cwd = mkdtempSync(join(tmpdir(), "worker-grant-race-")); const parentGate = makeParentGate(cwd, true); const store = new WorkerGrantStore(); @@ -243,13 +239,17 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" expect(await parentGate.evaluate(shellCall("parent-1", COMMAND))).toEqual({ allowed: true, }); + expect(await workerGate.authorizeCall(shellCall("c2", COMMAND))).toEqual({ + effect: "allow", + }); + expect(store.peek(requestId)?.status).toBe("pending"); // Two in-flight copies of the exact call (fresh call ids, same args): // without serialization both pass precheck before either consumes and - // the envelope allows twice. + // the envelope allows twice. Exactly-once is enforced at execution. const [retryA, retryB] = await Promise.all([ - workerGate.authorizeCall(shellCall("c2", COMMAND)), - workerGate.authorizeCall(shellCall("c3", COMMAND)), + workerGate.executionVerdict(shellCall("c2", COMMAND)), + workerGate.executionVerdict(shellCall("c3", COMMAND)), ]); const effects = [retryA.effect, retryB.effect].sort(); expect(effects).toEqual(["allow", "deny"]); @@ -291,6 +291,103 @@ describe("worker grant-request flow: deny → parent replay grant → one retry" expect(await workerGate.authorizeCall(shellCall("c3", COMMAND))).toEqual({ effect: "allow", }); - expect(store.peek(freshId)?.status).toBe("consumed"); + expect(store.peek(freshId)?.status).toBe("pending"); + }); + + test("KEEPER production two-stage: authorize allows, execution consumes, replay denied", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-two-stage-")); + const parentGate = makeParentGate(cwd, true); + const store = new WorkerGrantStore(); + const workerGate = workerPermissionGate(parentGate, { + sessionId: "worker-1", + store, + }); + + const denied = await workerGate.authorizeCall(shellCall("c1", COMMAND)); + if (denied.effect !== "deny") throw new Error("expected worker deny"); + const requestId = extractRequestId(denied.reason); + expect(await parentGate.evaluate(shellCall("parent-1", COMMAND))).toEqual({ + allowed: true, + }); + + // Stage 1 (reactor authorize hook): the granted retry authorizes but the + // envelope stays pending — nothing has executed yet. + const retry = await workerGate.authorizeCall(shellCall("c2", COMMAND)); + expect(retry).toEqual({ effect: "allow" }); + expect(store.peek(requestId)?.status).toBe("pending"); + + // Stage 2 (tool-runner middleware executionVerdict on the identical call): + // allows and consumes the single granted execution. + const executed = await workerGate.executionVerdict( + shellCall("c2", COMMAND), + ); + expect(executed).toEqual({ effect: "allow" }); + expect(store.peek(requestId)?.status).toBe("consumed"); + + // Exactly-once: a second execution of the same call fails closed, as does + // a second authorize. + const replayExecution = await workerGate.executionVerdict( + shellCall("c2", COMMAND), + ); + expect(replayExecution.effect).toBe("deny"); + if (replayExecution.effect !== "deny") + throw new Error("expected replay deny"); + expect(replayExecution.reason).toContain("already consumed"); + expect(replayExecution.reason).toContain(requestId); + const replayAuthorize = await workerGate.authorizeCall( + shellCall("c2", COMMAND), + ); + expect(replayAuthorize.effect).toBe("deny"); + }); + + test("KEEPER second session gets its own envelope: pending and consumed variants", async () => { + const cwd = mkdtempSync(join(tmpdir(), "worker-grant-sessions-")); + const store = new WorkerGrantStore(); + const gateA = workerPermissionGate(makeParentGate(cwd, true), { + sessionId: "worker-A", + store, + }); + const gateB = workerPermissionGate(makeParentGate(cwd, true), { + sessionId: "worker-B", + store, + }); + + // Pending variant: A's pending envelope must not veto B's own deny+ask. + const deniedA = await gateA.authorizeCall(shellCall("a1", COMMAND)); + if (deniedA.effect !== "deny") throw new Error("expected A deny"); + const idA = extractRequestId(deniedA.reason); + const deniedB = await gateB.authorizeCall(shellCall("b1", COMMAND)); + if (deniedB.effect !== "deny") throw new Error("expected B deny"); + const idB = extractRequestId(deniedB.reason); + expect(idB).not.toBe(idA); + expect(store.peek(idB)?.workerSessionId).toBe("worker-B"); + expect(store.peek(idB)?.status).toBe("pending"); + + // Consumed variant: A completes its grant + two-stage retry; a fresh + // session's identical call still gets its own envelope, never A's veto. + const parentA = makeParentGate(cwd, true); + const gateA2 = workerPermissionGate(parentA, { + sessionId: "worker-A", + store, + }); + expect(await parentA.evaluate(shellCall("parent-A", COMMAND))).toEqual({ + allowed: true, + }); + expect(await gateA2.authorizeCall(shellCall("a2", COMMAND))).toEqual({ + effect: "allow", + }); + expect(await gateA2.executionVerdict(shellCall("a2", COMMAND))).toEqual({ + effect: "allow", + }); + expect(store.peek(idA)?.status).toBe("consumed"); + const gateC = workerPermissionGate(makeParentGate(cwd, true), { + sessionId: "worker-C", + store, + }); + const deniedC = await gateC.authorizeCall(shellCall("c1", COMMAND)); + if (deniedC.effect !== "deny") throw new Error("expected C deny"); + const idC = extractRequestId(deniedC.reason); + expect(idC).not.toBe(idA); + expect(store.peek(idC)?.workerSessionId).toBe("worker-C"); }); }); diff --git a/src/permission/worker-grant.test.ts b/src/permission/worker-grant.test.ts index 39d34e13c..a951cd09a 100644 --- a/src/permission/worker-grant.test.ts +++ b/src/permission/worker-grant.test.ts @@ -4,7 +4,6 @@ import { WORKER_CANNOT_COMPLETE_APPROVAL } from "./decline-markers.js"; import { WORKER_GRANT_TTL_MS, WorkerGrantStore, - buildRetryMessage, createDeniedCallEnvelope, fingerprintDeniedCall, formatWorkerDenyWithGrantId, @@ -125,7 +124,7 @@ describe("WorkerGrantStore lifecycle", () => { expect(replay.blocker).toContain("already consumed"); }); - test("tampered session/cwd fail closed; tampered args fall through to the gate", () => { + test("tampered cwd fails closed; sibling session falls through to its own round", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); store.consumeOnAllow({ @@ -147,17 +146,17 @@ describe("WorkerGrantStore lifecycle", () => { expect(cwdTamper.ok).toBe(false); if (cwdTamper.ok) throw new Error("expected blocker"); expect(cwdTamper.blocker).toContain(CWD); - // Same fingerprint, different session: must replay from the owning session. - const sessionTamper = store.precheck({ + // Same fingerprint, another session: never vetoed by this envelope — the + // sibling falls through to its own gate round and mints its own envelope + // there instead of riding this session's grant. + const sessionFallthrough = store.precheck({ sessionId: "worker-2", canonicalTool: "run_shell", args: { ...ARGS }, cwd: CWD, now: 1_000_002, }); - expect(sessionTamper.ok).toBe(false); - if (sessionTamper.ok) throw new Error("expected blocker"); - expect(sessionTamper.blocker).toContain("worker-1"); + expect(sessionFallthrough).toEqual({ ok: true }); // Different fingerprint (tampered args or tool): not this envelope — the // gate denies downstream with a fresh deny since no grant covers it. for (const identity of [ @@ -181,7 +180,7 @@ describe("WorkerGrantStore lifecycle", () => { expect(envelope.status).toBe("consumed"); }); - test("cross-session replay of a pending envelope fails closed", () => { + test("another session's envelope never vetoes this session's round", () => { const store = new WorkerGrantStore(); store.register(createDeniedCallEnvelope(descriptor())); const result = store.precheck({ @@ -191,9 +190,7 @@ describe("WorkerGrantStore lifecycle", () => { cwd: CWD, now: 1_000_001, }); - expect(result.ok).toBe(false); - if (result.ok) throw new Error("expected blocker"); - expect(result.blocker).toContain("worker-1"); + expect(result).toEqual({ ok: true }); }); test("expiry yields a fresh gate round, never a blackhole", () => { @@ -272,25 +269,6 @@ describe("WorkerGrantStore lifecycle", () => { expect(envelope.questionId).toBeUndefined(); }); - test("decline fails closed; headless declineAllForSession covers the session", () => { - const store = new WorkerGrantStore(); - const envelope = store.register(createDeniedCallEnvelope(descriptor())); - expect(store.decline(envelope.requestId, "operator said no")).toBe(true); - expect(store.decline(envelope.requestId, "again")).toBe(false); - const result = store.precheck({ - sessionId: "worker-1", - canonicalTool: "run_shell", - args: { ...ARGS }, - cwd: CWD, - }); - expect(result.ok).toBe(false); - if (result.ok) throw new Error("expected blocker"); - expect(result.blocker).toContain("declined"); - - store.register(createDeniedCallEnvelope(descriptor({ callId: "call-9" }))); - expect(store.declineAllForSession("worker-1", "headless")).toBe(1); - }); - test("interrupt invalidation tombstones pending envelopes", () => { const store = new WorkerGrantStore(); const envelope = store.register(createDeniedCallEnvelope(descriptor())); @@ -318,7 +296,7 @@ describe("WorkerGrantStore lifecycle", () => { ); expect(envelope.questionId).toBe("ask-7"); expect(envelope.status).toBe("pending"); - expect(store.byQuestion("ask-7")).toBe(envelope); + expect(store.peek(envelope.requestId)).toBe(envelope); expect( store.attachToAsk("worker-9", "ask-8", undefined, 1_000_001), ).toBeUndefined(); @@ -351,7 +329,7 @@ describe("WorkerGrantStore lifecycle", () => { ).toBe(envelopeB); expect(envelopeB.questionId).toBe("ask-B"); expect(envelopeA.questionId).toBeUndefined(); - expect(store.byQuestion("ask-B")).toBe(envelopeB); + expect(store.peek(envelopeB.requestId)?.questionId).toBe("ask-B"); }); test("unnamed ask keeps the legacy first-pending bind", () => { @@ -370,7 +348,6 @@ describe("WorkerGrantStore lifecycle", () => { ).toBeUndefined(); expect(envelopeA.questionId).toBeUndefined(); expect(envelopeB.questionId).toBeUndefined(); - expect(store.byQuestion("ask-x")).toBeUndefined(); }); test("cross-session id fails closed with no fallback", () => { @@ -404,9 +381,11 @@ describe("WorkerGrantStore lifecycle", () => { cwd: CWD, now: 1_000_001, }); - expect( - store.auditTrail(envelope.requestId).map((event) => event.event), - ).toEqual(["denied", "asked", "consumed"]); + expect(envelope.audit.map((event) => event.event)).toEqual([ + "denied", + "asked", + "consumed", + ]); }); test("no text API: state moves only on typed identities, never prose", () => { @@ -435,18 +414,6 @@ describe("WorkerGrantStore lifecycle", () => { }); }); -describe("buildRetryMessage", () => { - test("carries exact args + questionId ref from the envelope", () => { - const store = new WorkerGrantStore(); - const envelope = store.register(createDeniedCallEnvelope(descriptor())); - store.attachToAsk("worker-1", "ask-3", undefined, 1_000_001); - const message = buildRetryMessage(envelope); - expect(message).toContain(envelope.requestId); - expect(message).toContain("ask-3"); - expect(message).toContain(JSON.stringify(ARGS)); - }); -}); - describe("WorkerGrantStore.runExclusive", () => { test("serializes same-key holders: no overlap, FIFO, error still releases", async () => { const store = new WorkerGrantStore();