diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 2371c3b52..8ba7752a3 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 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/`) - Builds a chat-mode agent using the `ChatDirector` diff --git a/src/permission/reactor-authorize.ts b/src/permission/reactor-authorize.ts index bef1e0fe2..acdc7dd51 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, + fingerprintDeniedCall, + formatWorkerDenyWithGrantId, + getProcessWorkerGrantStore, + 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,160 @@ 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; +} + +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 } + : {}), + }), + ); + return { + effect: "deny", + reason: formatWorkerDenyWithGrantId(baseReason, envelope.requestId), + }; +} + +/** 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, + 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) + 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); + // 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 { 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 +240,70 @@ 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) + 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); + // 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( + workerCallIdentity(grant.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 +355,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..102f3230d --- /dev/null +++ b/src/permission/worker-grant-flow.test.ts @@ -0,0 +1,393 @@ +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 { + WORKER_GRANT_TTL_MS, + 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("pending"); + + // 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 () => { + 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 pending", 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("pending"); + }); + + 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). + 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); + + // 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 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", () => { + 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"); + }); + + 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(); + 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, + }); + 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. Exactly-once is enforced at execution. + const [retryA, retryB] = await Promise.all([ + workerGate.executionVerdict(shellCall("c2", COMMAND)), + workerGate.executionVerdict(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("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 new file mode 100644 index 000000000..a951cd09a --- /dev/null +++ b/src/permission/worker-grant.test.ts @@ -0,0 +1,473 @@ +import { describe, test, expect } from "bun:test"; + +import { WORKER_CANNOT_COMPLETE_APPROVAL } from "./decline-markers.js"; +import { + WORKER_GRANT_TTL_MS, + WorkerGrantStore, + 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 cwd fails closed; sibling session falls through to its own round", () => { + 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, 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(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 [ + { + 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("another session's envelope never vetoes this session's round", () => { + 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).toEqual({ ok: true }); + }); + + 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; + // 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" })), + ); + expect(store.sweepExpired(second.expiresAt + 1)).toBe(1); + 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("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())); + // 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.peek(envelope.requestId)).toBe(envelope); + expect( + store.attachToAsk("worker-9", "ask-8", undefined, 1_000_001), + ).toBeUndefined(); + expect(envelope.audit.map((event) => event.event)).toEqual([ + "denied", + "asked", + ]); + }); + + 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.peek(envelopeB.requestId)?.questionId).toBe("ask-B"); + }); + + 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(); + }); + + 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())); + store.attachToAsk("worker-1", "ask-1", undefined, 1_000_001); + store.consumeOnAllow({ + sessionId: "worker-1", + canonicalTool: "run_shell", + args: { ...ARGS }, + cwd: CWD, + now: 1_000_001, + }); + expect(envelope.audit.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("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 new file mode 100644 index 000000000..581964835 --- /dev/null +++ b/src/permission/worker-grant.ts @@ -0,0 +1,486 @@ +// 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"; + +/** 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 = + | "pending" + | "consumed" + | "expired" + | "interrupted"; + +export interface WorkerGrantAuditEvent { + event: "denied" | "asked" | "consumed" | "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.`; +} + +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 "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(); + /** 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); + return envelope; + } + + peek(requestId: string): WorkerDeniedCallEnvelope | undefined { + return this.envelopes.get(requestId); + } + + pendingForSession( + sessionId: string, + now: number = Date.now(), + ): WorkerDeniedCallEnvelope | undefined { + this.sweepExpired(now); + 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). + * 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 ( + 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. 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 { + 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); + return envelope; + } + + /** + * Execution backstop pre-check for a worker call: fail closed when the exact + * 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 + * 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(); + this.sweepExpired(now); + const fingerprint = fingerprintOf(identity); + 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 + ) { + 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 === "expired") continue; + if (envelope.status !== "pending") { + return { + ok: false, + blocker: terminalBlocker(envelope, identity.canonicalTool), + }; + } + if (envelope.expiresAt <= now) { + envelope.status = "expired"; + audit(envelope, "expired"); + continue; + } + return { ok: true }; + } + // 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 }; + } + + /** + * 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; + } + + /** 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; + } + + /** 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()) { + 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; + } + + /** 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/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/run-ask-director-grant.test.ts b/src/subagent/run-ask-director-grant.test.ts new file mode 100644 index 000000000..6c10ff701 --- /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 "../../testkit/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 e35b395b5..736fe8598 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: { @@ -532,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"], }, @@ -598,11 +606,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(); @@ -851,6 +864,7 @@ async function runSubAgentInner( try { return await handleAskDirector({ question: rawArgs.question, + grantRequestId: rawArgs.grant_request_id, state: askDirectorState, port, signal, @@ -1144,6 +1158,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 +1193,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..530b46a81 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"]); @@ -301,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; }, @@ -308,7 +315,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 +626,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 +645,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 +718,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 { @@ -1891,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; }, @@ -1899,7 +1942,29 @@ 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. 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, + ask.grantRequestId, + now(), + ) ?? 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 +1986,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 +2012,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 +2286,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 +2328,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 +2373,11 @@ export function createSubAgentSessionStore( runInFlight.clear(); revisions.clear(); snapshotCache.clear(); + try { + getProcessWorkerGrantStore().invalidateAll(reason); + } catch { + // Grant invalidation must not throw out of teardown. + } notify(); }, }; 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; };