// The CHAT-03 I1 controller (#1507, CHAT-03 §§1–3, §6). // // One controller process per execution owns the engine's stdin. It serves a // Unix socket in a 0700 directory; every connection starts as an observer. // Commands are CHAT-01 v2 client requests and are evaluated in check.mjs's // order, with CHAT-03's narrowings: no broker queue (`busy`), one pending // dispatch slot, the incarnation token (`stale-incarnation`), and Interrupt's // `no-turn`. The controller never writes a session file and never opens one // through Pi. // // Fixture-only in code: the claim root, socket directory and session file // are constructor arguments checked by the live-session guard at // construction and again at bind (§2). // // Proofs go to the injected verifier, the fixture's trusted digest registry // as in CHAT-01. With no verifier nothing verifies, and every stop ends // `uncertain`. import { lstatSync, mkdirSync, readFileSync, unlinkSync } from "node:fs"; import { createServer } from "node:net"; import { basename, dirname, join, resolve } from "node:path"; import { fileURLToPath } from "node:url"; import { randomBytes } from "node:crypto"; import { processStart } from "../../discord/src/journal.mjs"; import { ALREADY_ACTIVE, ClaimStore, FOREIGN_HOST, UNSAFE_REPLACEMENT, machineId, seatKey, sessionKey } from "./claim.mjs"; import { AUTHORITY, PgroupLauncher, bootProof, cohortProof, cohortRefOf, effectReport, forceStopCohort, systemdUnits } from "./cohort.mjs"; import { EngineLink } from "./engine.mjs"; import { DIALOG_METHODS, KNOWN_UNSHOWN, NOTIFY_METHODS, deltaBlocks, isFoldedUpdate, messageBlocks, partsOf, roleOf, toolResultText } from "./events.mjs"; import { LineSplitter, encodeLine, parseLine } from "./framing.mjs"; import { LiveSessionGuard, realPath } from "./guard.mjs"; import { ID, fragments, safeId } from "./parts.mjs"; import { parseSnapshot } from "./pi.mjs"; import { ENGINE_PIN_MISMATCH, PI_BIN, UNSEALED_ENGINE, argvDigest, buildPiArgs, checkEnginePin, checkSeal } from "./pi-pin.mjs"; import { ACTOR, conversationId, createReader, rootsFromSpecs } from "./reader.mjs"; import { clone, equal, hash, newId, receiptAllows, record, scopeMatch, sealProof, sha256, targetOf } from "./records.mjs"; import { ControlRefusal, Refusal } from "./safe-fs.mjs"; import { TEXT_POLICY, textPolicy } from "./text-policy.mjs"; import { DONE_STOPS, Tracker } from "./turns.mjs"; export const REPO_ROOT = resolve(dirname(fileURLToPath(import.meta.url)), "..", "..", ".."); // CHAT-03 refusal names (§"New names"), plus the two Sage approved for this // build (malformed, eligibility). export const BUSY = "busy"; export const STALE_INCARNATION = "stale-incarnation"; export const TRANSPORT_UNKNOWN = "transport-unknown"; export const HANDLED_WITHOUT_RUN = "handled-without-run"; export const ACK_WITHOUT_START = "ack-without-start"; export const INTERRUPTED = "interrupted"; export const NO_TURN = "no-turn"; export const RUN_OVERLAP = "run-overlap"; export const MALFORMED = "malformed"; export const ELIGIBILITY = "eligibility"; // check.mjs `caps`. export const CAPS = Object.freeze({ observe: "observe", prompt: "send", takeover: "take-control", "edit-queued": "send", "cancel-queued": "send", approval: "approve", interrupt: "interrupt", "force-stop": "force-stop", recover: "recover", "acquire-recovery-control": "recover-control", "list-drafts": "draft", "create-draft": "draft", "update-draft": "draft", "discard-draft": "draft", "begin-upload": "upload", "append-upload": "upload", "complete-upload": "upload", "discard-upload": "upload", "issue-confirmation": "confirm", "answer-confirmation": "confirm", }); export const ALL_CAPABILITIES = Object.freeze(["observe", "send", "take-control", "approve", "interrupt", "force-stop", "recover", "recover-control", "draft", "upload", "confirm"]); // Operations I1 implements. Drafts, uploads, queue edits and approvals are // CHAT-04 or I4 and refuse `unsupported-capability`. export const VERIFIED_OPERATIONS = Object.freeze(["observe", "prompt", "takeover", "acquire-recovery-control", "interrupt", "force-stop", "recover", "issue-confirmation", "answer-confirmation"]); export const TIMEOUTS = Object.freeze({ ack: 5000, state: 5000, start: 5000, clear: 5000, abort: 10000, settle: 10000, grace: 1000, write: 5000, maxRounds: 3 }); const FINAL = new Set(["finished", "failed", "dispatch-refused", "delivery-unknown"]); const STOP_IN_PROGRESS = new Set(["fenced", "cancelling", "stopping"]); const STOP_NEXT = { fenced: ["cancelling", "stopping", "uncertain"], cancelling: ["stopping", "uncertain"], stopping: ["uncertain"], uncertain: [] }; const SEAL_BASIS = "seal: --no-extensions and no --extension argument (pi-pin.mjs SEAL_FLAGS)"; const DIALOG_REASON = "Pi dialogs are not answered in CHAT-03 (lead decision 30); shown disabled"; const refused = (code, extra = {}) => ({ outcome: `refused:${code}`, refusal: code, ...extra }); const isGen = (v) => Number.isInteger(v) && v >= 1 && v <= Number.MAX_SAFE_INTEGER; const isId = (v) => typeof v === "string" && ID.test(v); const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); // The client request envelope, per operation (CHAT-01 schema `command`). // Unknown operations pass the shape check and are refused by evaluate. const SHAPES = { observe: { cursor: (v) => v === null || isId(v), limit: (v) => Number.isInteger(v) && v >= 1 && v <= 100 }, prompt: { draft: isId, draftRevision: isGen }, takeover: {}, interrupt: {}, "force-stop": { confirmation: isId }, recover: { stop: isId, confirmation: isId }, "acquire-recovery-control": { confirmation: isId }, "issue-confirmation": { operationToConfirm: (v) => ["force-stop", "recover", "acquire-recovery-control"].includes(v) }, "answer-confirmation": { confirmation: isId, answer: (v) => v === "confirm" || v === "cancel" }, }; export function malformedRequest(r) { if (!r || typeof r !== "object" || Array.isArray(r)) return "the request is not an object"; const keys = Object.keys(r).sort().join(","); if (keys !== "command,connection,id,kind,target,version") return "the request has missing or extra fields"; if (r.version !== 2 || r.kind !== "clientRequest" || !isId(r.id) || !isId(r.connection)) return "bad version, kind, id or connection"; const t = r.target; if (!t || typeof t !== "object" || Object.keys(t).sort().join(",") !== "branch,controllerGeneration,conversation,execution") return "bad target"; if (!isId(t.conversation) || !isId(t.branch) || !isId(t.execution) || !isGen(t.controllerGeneration)) return "bad target"; const cmd = r.command; if (!cmd || typeof cmd !== "object" || Array.isArray(cmd) || typeof cmd.operation !== "string" || !CAPS[cmd.operation]) return "bad command"; const shape = SHAPES[cmd.operation]; if (!shape) return null; const fields = Object.keys(cmd).filter((k) => k !== "operation"); if (fields.length !== Object.keys(shape).length || fields.some((k) => !shape[k] || !shape[k](cmd[k]))) return `bad ${cmd.operation} fields`; return null; } class Lock { #tail = Promise.resolve(); run(fn) { const p = this.#tail.then(() => fn()); this.#tail = p.then(() => {}, () => {}); return p; } } // Tool calls in the session file, for effects when no stream was observed // (an orphan, a boot proof). A call with no result entry is open. export function sessionTools(parsed) { const tools = new Map(); for (const { entry } of parsed.entries) { if (entry.type !== "message") continue; const m = entry.message; if (m?.role === "assistant" && Array.isArray(m.content)) { for (const c of m.content) if (c?.type === "toolCall" && typeof c.id === "string") tools.set(c.id, { call: safeId(c.id), start: safeId(`s.${entry.id}`), end: null }); } if (m?.role === "toolResult" && tools.has(m.toolCallId)) tools.get(m.toolCallId).end = safeId(`s.${entry.id}`); } return tools; } export class Controller { constructor(opts = {}) { const { fixtureRoot, claimRoot, socketDir, sessionFile, seat, workspace = "default", engine = {}, launcher = new PgroupLauncher(), verifier = null, host, identity, barrier = null, units = systemdUnits, grants = null, guardOptions = {}, timeouts = {}, now = () => new Date(), policyRevision = "policy-fixture-1", sourceRootRef = "source-fixture-1", approvedMappings = null, pinRoot = REPO_ROOT, auditWritable = true, hostId = null, } = opts; for (const [k, v] of Object.entries({ claimRoot, socketDir, sessionFile, seat })) { if (typeof v !== "string" || !v) throw new ControlRefusal("configuration", `${k} is required; CHAT-03 has no default`); } this.guard = new LiveSessionGuard({ fixtureRoot, ...guardOptions }); this.paths = { claimRoot: resolve(claimRoot), socketDir: resolve(socketDir), sessionFile: resolve(sessionFile) }; this.guard.check(this.paths, "construction"); const dir = dirname(this.paths.sessionFile); const projectRoot = resolve(dir, "..", "..", "..", ".."); const project = basename(projectRoot); const roots = rootsFromSpecs([{ sessionsDir: dir, agent: seat, project }]); if (roots.length !== 1) throw new ControlRefusal("configuration", "the session file must be /.pi/state//sessions/.jsonl"); this.root = roots[0]; this.reader = createReader({ roots, now: () => now().getTime() }); this.seat = seat; this.project = project; this.workspace = workspace; this.conversation = conversationId(this.root, basename(this.paths.sessionFile)); this.engine = { command: engine.command ?? process.execPath, preArgs: engine.preArgs ?? [join(pinRoot, PI_BIN)], extraArgs: engine.extraArgs ?? [], cwd: engine.cwd ?? projectRoot, env: engine.env ?? process.env, }; for (const k of ["preArgs", "extraArgs"]) { if (!Array.isArray(this.engine[k])) throw new ControlRefusal(UNSEALED_ENGINE, `engine.${k} is not a list`); } this.piArgs = buildPiArgs({ sessionFile: this.paths.sessionFile, extraArgs: this.engine.extraArgs }); this.#checkSeal(); this.launcher = launcher; this.verifier = verifier; this.units = units; this.barrier = barrier; this.now = now; this.T = { ...TIMEOUTS, ...timeouts }; this.pinRoot = pinRoot; this.policyRevision = policyRevision; this.sourceRootRef = sourceRootRef; this.approvedMappings = approvedMappings ?? [sourceRootRef]; this.auditWritable = auditWritable; this.hostId = hostId ?? `host-${machineId() ?? "unknown"}`; this.grantSpecs = grants; this.store = new ClaimStore({ root: this.paths.claimRoot, host, identity, barrier, now }); this.seatK = seatKey({ seat, project, workspace }); // The session key is the Pi header ID (BRIEF §2 D1), not the // conversation ID: two paths to one session file, a hard link or a copy, // share it. The conversation itself needs no key of its own: it names one // file under one seat, which the seat key already covers. this.nativeSession = this.#readSession().nativeSession; this.sessionK = sessionKey({ harness: "pi", nativeSession: this.nativeSession }); this.incarnation = randomBytes(16).toString("hex"); this.lock = new Lock(); this.b = null; this.exec = null; this.claim = null; this.claimChain = Promise.resolve(); this.closers = new Set(["starting"]); this.stops = new Map(); this.receipts = new Map(); this.requests = new Map(); this.connections = new Map(); this.channels = new Set(); this.confirmations = new Map(); this.grants = new Map(); this.dedup = new Map(); this.eligibility = new Map(); this.outcomeUnknown = new Set(); this.escalating = null; this.abortWritten = new Set(); this.preflightOk = false; this.server = null; this.sockets = new Set(); this.events = []; this.evidence = { dropped: { lines: 0, bytes: 0 }, unknownEvents: {}, unshown: {}, folded: 0, dialogs: [], notices: 0, gaps: [], overlaps: [], uncertain: [], stops: [], refusedRevisions: [], internal: [], stderrTail: "" }; } get binding() { return this.b; } get socketPath() { return join(this.paths.socketDir, "controller.sock"); } async #pause(name, detail = {}) { if (this.barrier) await this.barrier(name, detail); } #checkSeal() { const bad = this.engine.preArgs.find((a) => typeof a !== "string" || a === "-e" || a === "--extension" || a.startsWith("--extension=")); if (bad !== undefined) throw new ControlRefusal(UNSEALED_ENGINE, `engine pre-arguments carry ${bad}`); checkSeal(this.piArgs); } #pins() { const pin = checkEnginePin(this.pinRoot); return { engineVersion: pin.version, enginePin: pin.pin, argvDigest: argvDigest(this.engine.command, [...this.engine.preArgs, ...this.piArgs]) }; } // Every read after construction refuses a header ID other than the one the // claim is keyed on. #readSession() { let text; try { text = readFileSync(this.paths.sessionFile, "utf8"); } catch (err) { throw new ControlRefusal("configuration", `the session file is unreadable (${err.code ?? "error"})`); } const parsed = parseSnapshot(text); if (!isId(parsed.header.id)) throw new ControlRefusal("configuration", "the session header id is not a CHAT-01 id"); if (this.nativeSession !== undefined && parsed.header.id !== this.nativeSession) throw new ControlRefusal("target", "the session header id changed since construction"); return { nativeSession: parsed.header.id, branch: parsed.defaultBranch, leaf: parsed.defaultLeaf?.entry.id ?? null, snapshotDigest: sha256(text), parsed }; } #guardCheck(when) { this.guard.check(this.paths, when); } #socketDir() { const dir = this.paths.socketDir; mkdirSync(dir, { recursive: true, mode: 0o700 }); const st = lstatSync(dir); if (!st.isDirectory() || st.isSymbolicLink() || st.uid !== process.getuid() || (st.mode & 0o077) !== 0) { throw new ControlRefusal("channel", "the socket directory must be a real 0700 directory owned by this user"); } let sock = null; try { sock = lstatSync(this.socketPath); } catch { sock = null; } if (sock) { if (!sock.isSocket()) throw new ControlRefusal("channel", "controller.sock exists and is not a socket"); unlinkSync(this.socketPath); } } #grantRecords(scope) { const specs = this.grantSpecs ?? [{ id: "grant-local", revision: "rev-1", actor: ACTOR, capabilities: ALL_CAPABILITIES }]; for (const g of specs) { this.grants.set(g.id, record("grant", { id: g.id, revision: g.revision ?? "rev-1", actor: g.actor ?? ACTOR, scope: clone(g.scope ?? scope), capabilities: [...(g.capabilities ?? ALL_CAPABILITIES)], state: g.state ?? "active", expiresAt: g.expiresAt ?? new Date(this.now().getTime() + 24 * 3600 * 1000).toISOString(), })); } } #newBinding({ execution, generation, session, cohortRef, state, pins }) { const scope = { host: this.hostId, seat: this.seat, project: this.project, workspace: this.workspace, conversation: this.conversation }; if (this.grants.size === 0) this.#grantRecords(scope); return record("binding", { id: newId("binding"), scope, harness: "pi", nativeSession: session.nativeSession, branch: session.branch, execution, engineDigest: hash({ engineVersion: pins.engineVersion, enginePin: pins.enginePin }), configDigest: hash(pins), policyRevision: this.policyRevision, sourceRootRef: this.sourceRootRef, snapshotDigest: session.snapshotDigest, cohortRef, controllerGeneration: generation, controllerConnection: null, state, admission: "closed", createdAt: this.now().toISOString(), stop: null, }); } // ---- start ------------------------------------------------------------- // Bind: guard, seal and pins, then classify the claim pair. Returns // { launched, classified }. async start() { this.#guardCheck("bind"); this.#checkSeal(); const pins = this.#pins(); const session = this.#readSession(); const classified = await this.store.classify(this.seatK, this.sessionK, { units: this.units, adopt: this.store.owner(this.incarnation), stoppedPatch: { leafAtProof: session.leaf, branchAtProof: session.branch }, bootProof: (latest) => this.#bootProof(latest, session), }); if (classified.refusal === FOREIGN_HOST || classified.damaged) throw new ControlRefusal(classified.refusal, `the claim pair is held (${classified.refusal})`); if (classified.refusal === ALREADY_ACTIVE || classified.state === "owned") throw new ControlRefusal(ALREADY_ACTIVE, "the writer claim for this seat or session is held by a live controller"); if (classified.state === "stopped") return { launched: false, classified: { state: "stopped", proofKind: classified.proofKind } }; if (classified.state === "uncertain" || classified.state === "stopping") { if (!classified.claim) throw new ControlRefusal(classified.refusal ?? UNSAFE_REPLACEMENT, `the claim pair is uncertain (${classified.needs ?? "held"})`); await this.#serveOrphan(classified, session, pins); return { launched: false, classified: { state: classified.state, unit: classified.unit?.state ?? null } }; } // Free. A resume of this session checks the last proven stop. const prior = this.store.head(this.sessionK); let generation = 1; let priorRef = null; if (prior.n > 0 && !prior.damaged && prior.record.state === "stopped") { const p = prior.record; if (!equal(p.pins, pins)) throw new ControlRefusal(ENGINE_PIN_MISMATCH, "the engine pin or launch argv changed since the proven stop"); const leaf = p.leafAtProof !== undefined ? p.leafAtProof : p.leaf; const branch = p.branchAtProof ?? p.branch; if (leaf !== session.leaf || branch !== session.branch) throw new ControlRefusal("target", "the session branch or leaf changed since the proven stop"); generation = p.generation + 1; priorRef = { claimId: p.claimId, stop: p.stop?.id ?? null }; } const execution = newId("exec"); this.b = this.#newBinding({ execution, generation, session, cohortRef: `pending-${execution}`, state: "reserved", pins }); const claim = await this.store.acquire(this.seatK, this.sessionK, { bindingId: this.b.id, harness: "pi", conversation: this.conversation, branch: session.branch, leaf: session.leaf, pins, owner: this.store.owner(this.incarnation), generation, prior: priorRef, execution, cohortRef: null, }); // The socket directory is prepared only once this controller holds the // pair: a refused contender never touches a live controller's socket. this.#socketDir(); await this.#listen(); await this.#launchInto(claim, session); return { launched: true, classified: { state: "free" } }; } async #bootProof(latest, session) { if (!this.verifier) return null; const binding = { scope: { conversation: latest.conversation }, execution: latest.execution ?? `boot-${latest.claimId}`, cohortRef: latest.cohortRef ?? `cohort-boot-${latest.claimId}` }; const stop = latest.stop?.id ?? null; const proof = bootProof({ claim: latest, conversation: binding.scope.conversation, execution: binding.execution, cohortRef: binding.cohortRef, stop, now: this.now }); const effects = effectReport({ binding, stop, tools: sessionTools(session.parsed), observedAt: proof.observedAt }); this.verifier.post(proof); this.verifier.post(effects); const ok = this.verifier.cohort(proof, effects, { binding, stop, now: this.now(), epoch: proof.membershipEpoch }); return ok ? { proof, effects } : null; } // An orphan or an interrupted force stop: no engine pipe, the binding is // `uncertain` (or `stopping`), and only a confirmed force stop moves it on. async #serveOrphan(classified, session, pins) { const rec = classified.claim.record; this.claim = classified.claim; const execution = rec.execution ?? newId("exec"); this.b = this.#newBinding({ execution, generation: rec.generation, session, cohortRef: rec.cohortRef ?? `pending-${execution}`, state: rec.state === "stopping" ? "stopping" : "uncertain", pins }); this.b.id = rec.bindingId ?? this.b.id; this.closers.add("uncertain"); this.closers.delete("starting"); if (rec.stop?.record) { this.stops.set(rec.stop.record.id, clone(rec.stop.record)); this.b.stop = rec.stop.record.id; if (rec.stop.record.mode === "force-stop") this.closers.add("force-stop"); } this.orphanTools = sessionTools(session.parsed); this.#admission(); this.#socketDir(); await this.#listen(); this.#evidenceAdd("uncertain", { reason: "orphan", unit: classified.unit?.state ?? null, adopted: classified.adopted }); if (classified.resumeStop) { const stop = this.stops.get(classified.resumeStop.id) ?? null; if (stop) void this.#forceStop(stop, { confirmation: classified.resumeStop.confirmation ?? null, resumed: true }).catch((e) => this.#internal("force-stop", e)); } } // ---- launch ------------------------------------------------------------ #newExec(execution) { const exec = { execution, link: null, proc: null, tools: new Map(), msgSerial: 0, toolSerial: 0, message: null, seq: 0, streamEpoch: `${execution}.${this.incarnation.slice(0, 12)}`, responseIdx: new Map(), clears: new Map(), waiters: new Set(), dialogs: [], closing: false, promptSlots: new Map(), }; exec.tracker = new Tracker({ onOverlap: (signal, entry) => this.#onOverlap(exec, signal, entry), onWorking: (slot) => this.#revise(slot.receipt, "working", null), onSlotSettled: (slot, result) => this.#onSlotSettled(exec, slot, result), // Lead decision 34: only the controller's abort produces `aborted`, so // it links to the stop in progress only once an abort was written in // that stop's chain (the stop or one it superseded). stopLink: () => { const s = this.stops.get(this.b?.stop); if (!s || !STOP_IN_PROGRESS.has(s.state)) return null; for (let p = s; p; p = this.stops.get(p.supersedes)) if (this.abortWritten.has(p.id)) return s.id; return null; }, }); return exec; } async #launchInto(claim, session) { this.claim = claim; const b = this.b; await this.#pause("reserved", { claimId: claim.claimId }); await this.#claimAdvance({ spawnMarker: true }); await this.#pause("spawn-marker", { claimId: claim.claimId }); const exec = this.#newExec(b.execution); this.exec = exec; let proc; try { proc = await this.launcher.launch({ unitName: claim.record.unitName, socketPath: join(this.paths.socketDir, `shim-${claim.claimId.slice(0, 12)}.sock`), command: this.engine.command, args: [...this.engine.preArgs, ...this.piArgs], cwd: this.engine.cwd, env: this.engine.env, }); } catch (err) { this.#uncertain("launch-failed", { error: String(err.code ?? err.message).slice(0, 200) }); return; } exec.proc = proc; const cohortRef = cohortRefOf({ machineId: claim.record.host.machineId, bootId: claim.record.host.bootId, unitName: claim.record.unitName, invocationId: proc.invocationId ?? `pgroup-${proc.pid}-${proc.start}` }); b.cohortRef = cohortRef; proc.stderr?.on("data", (c) => { this.evidence.stderrTail = (this.evidence.stderrTail + c.toString("utf8")).slice(-65536); }); exec.link = new EngineLink({ execution: b.execution, stdin: proc.stdin, stdout: proc.stdout, writeTimeoutMs: this.T.write, onLine: (link, value, bytes) => this.#onLine(exec, link, value, bytes), onResponse: (link, value, entry) => this.#onResponse(exec, link, value, entry), onGap: (link, kind, detail) => this.#onGap(exec, link, kind, detail), onEnd: (link) => this.#onEnd(exec, link), }); await this.#claimAdvance({ engine: { kind: proc.kind, pid: proc.pid, start: proc.start }, invocationId: proc.invocationId, shim: proc.shimSocket, execution: b.execution, cohortRef }); await this.#pause("spawned", { claimId: claim.claimId, pid: proc.pid }); this.#push({ kind: "binding", binding: b }); // K8: the engine must have loaded this session file at this leaf. const why = await this.#checkLoaded(exec, session); if (this.exec !== exec) return; if (why) { this.#evidenceAdd("uncertain", { reason: "loaded-session", detail: why }); this.#uncertain("target", { detail: why }); return; } if (exec.link.poisoned || exec.tracker.overlapped || this.closers.has("uncertain")) return; try { await this.#claimAdvance({ state: "active" }); } catch (err) { this.#uncertain("claim", { error: String(err.code ?? err.message) }); return; } this.preflightOk = true; b.state = "active"; this.closers.delete("starting"); this.#admission(); await this.#pause("active", { claimId: claim.claimId }); } async #checkLoaded(exec, session) { const st = await this.#ask(exec, "get_state", {}, this.T.state); if (!st.response) return `get_state: ${st.why}`; const d = st.response.data ?? {}; if (!st.response.success || typeof d.sessionFile !== "string" || realPath(d.sessionFile) !== realPath(this.paths.sessionFile)) return "the engine loaded another session file"; if (d.sessionId !== session.nativeSession) return "the engine reports another session id"; const tree = await this.#ask(exec, "get_tree", {}, this.T.state); if (!tree.response) return `get_tree: ${tree.why}`; if (!tree.response.success || (tree.response.data?.leafId ?? null) !== session.leaf) return "the engine loaded another leaf"; return null; } // One request with a bounded response. Transport failures poison the link. async #ask(exec, type, fields, timeoutMs, beforeWrite = null) { if (!exec.link || exec.link.poisoned) return { why: "poisoned" }; const req = exec.link.request(type, fields, { timeoutMs, beforeWrite }); const res = await req.response; if (res.unsent) { this.#transport(exec, `${type}-write-${res.unsent.reason ?? res.unsent.outcome}`); return { why: "unsent", id: req.id }; } if (res.timeout) { this.#transport(exec, `${type}-timeout`); return { why: "timeout", id: req.id }; } return { response: res.response, idx: exec.responseIdx.get(req.id), id: req.id }; } // ---- socket ------------------------------------------------------------ async #listen() { if (this.server) return; this.server = createServer((sock) => this.#accept(sock)); await new Promise((res, rej) => { this.server.once("error", rej); this.server.listen(this.socketPath, () => { this.server.off("error", rej); res(); }); }); } #accept(sock) { this.sockets.add(sock); const state = { conn: null }; const send = (v) => { if (!sock.destroyed) sock.write(encodeLine(v)); }; state.send = send; const splitter = new LineSplitter((line) => this.#onClientLine(state, line), { maxBytes: 4 * 1024 * 1024, onOverflow: () => sock.destroy() }); sock.on("data", (c) => splitter.push(c)); sock.on("error", () => {}); sock.on("close", () => { this.sockets.delete(sock); const c = state.conn ? this.connections.get(state.conn) : null; if (c && c.rec.state === "connected") { // Disconnect never transfers control and never changes a claim (H11). c.rec.state = "disconnected"; c.send = null; this.#push({ kind: "connection", connection: c.rec }); } }); } #onClientLine(state, line) { const parsed = parseLine(line); if (parsed.error) return state.send({ type: "refused", refusal: MALFORMED, detail: parsed.error }); const msg = parsed.value; if (!state.conn) { if (msg?.type !== "hello") return state.send({ type: "refused", refusal: MALFORMED, detail: "hello first" }); const g = typeof msg.grant === "string" ? this.grants.get(msg.grant) : undefined; if (!g || !this.b) return state.send({ type: "refused", refusal: "grant" }); const rec = record("connection", { id: newId("conn"), actor: g.actor, grant: g.id, grantRevision: g.revision, conversation: this.b.scope.conversation, mode: "observer", transport: "local-terminal", state: "connected", authenticatedChannelRef: newId("channel"), createdAt: this.now().toISOString(), generation: 1, }); this.channels.add(rec.authenticatedChannelRef); this.connections.set(rec.id, { rec, send: state.send }); state.conn = rec.id; return state.send({ type: "welcome", incarnation: this.incarnation, connection: rec, binding: this.b, target: targetOf(this.b), streamEpoch: this.exec?.streamEpoch ?? null }); } if (msg?.type !== "request") return state.send({ type: "refused", refusal: MALFORMED, detail: "unknown message type" }); const id = isId(msg.request?.id) ? msg.request.id : null; const reply = (res) => state.send({ type: "reply", id, ...res }); if (msg.incarnation !== this.incarnation) return reply(refused(STALE_INCARNATION)); const bad = malformedRequest(msg.request); if (bad) return reply(refused(MALFORMED, { detail: bad })); if (msg.request.connection !== state.conn) return reply(refused("channel")); if (msg.text !== undefined && typeof msg.text !== "string") return reply(refused(MALFORMED, { detail: "text must be a string" })); void this.#handle(state, msg.request, msg.text, reply); return undefined; } async #handle(state, r, text, reply) { const c = this.connections.get(state.conn)?.rec ?? null; let res; try { if (r.command.operation === "interrupt") { // Interrupt sets its fence before it takes the dispatch lock (§3). const pre = this.#evaluate(c, r, text); if (!pre.fence) res = pre; else { await this.#pause("interrupt-fenced", { request: r.id }); res = await this.lock.run(async () => { const out = await this.#interruptLocked(c, r, pre); if (!out.outcome.startsWith("refused:")) this.dedup.set(pre.dedupKey, { actor: c.actor, digest: pre.digest, outcome: out.outcome, receipt: null, request: null }); return out; }); } } else { res = await this.lock.run(() => this.#evaluate(c, r, text)); } } catch (err) { if (err instanceof Refusal) res = refused(err.code, { detail: err.message }); else { this.#internal("evaluate", err); res = { outcome: "error", detail: "internal error; recorded in controller evidence" }; } } const { after, fence, closer, fenceIdx, dedupKey, digest, ...wire } = res; reply(wire); if (after) void after().catch((e) => this.#internal("after", e)); } // ---- evaluate (check.mjs order) ---------------------------------------- #grantOk(c) { const g = this.grants.get(c.grant); if (!g || g.actor !== c.actor || g.revision !== c.grantRevision || g.state !== "active" || Date.parse(g.expiresAt) <= this.now().getTime()) return null; return g; } #evaluate(c, r, text) { const b = this.b, t = r.target, cmd = r.command, op = cmd.operation; if (!c || c.state !== "connected" || !this.channels.has(c.authenticatedChannelRef)) return refused("channel"); if (c.transport === "private-host") return refused("unsupported-capability"); const g = this.#grantOk(c); if (!g) return refused("grant"); if (!equal(g.scope, b.scope) || c.conversation !== b.scope.conversation) return refused("scope"); if (!this.approvedMappings.includes(b.sourceRootRef)) return refused("mapping"); if (!this.auditWritable) return refused("audit"); if (!g.capabilities.includes(CAPS[op])) return refused("capability"); if (!VERIFIED_OPERATIONS.includes(op)) return refused("unsupported-capability"); if (t.conversation !== b.scope.conversation || t.execution !== b.execution) return refused("target"); if (op === "observe") return this.#observe(c, r); if (t.branch !== b.branch) return refused("target"); const key = `${c.actor}\0${t.conversation}\0${r.id}`; const digest = hash({ target: t, command: cmd, payloadDigest: sha256(text ?? "") }); const prior = this.dedup.get(key); if (prior) { if (prior.actor !== c.actor || prior.digest !== digest) return refused("conflicting-request"); const receipt = prior.receipt ? this.receipts.get(prior.receipt) : null; return { outcome: `existing:${receipt?.state ?? prior.outcome}`, receipt: receipt ?? null, request: prior.request ? this.requests.get(prior.request) : null }; } const finish = (res) => { if (res?.fence) return { ...res, dedupKey: key, digest }; if (res && !res.outcome.startsWith("refused:")) { this.dedup.set(key, { actor: c.actor, digest, outcome: res.outcome, receipt: res.receipt?.id ?? null, request: res.request?.id ?? null }); } return res; }; // `recover` is asynchronous (it acquires a claim); everything else, // including Interrupt's fence, returns synchronously. const res = this.#evaluateOp(c, g, r, text); return typeof res?.then === "function" ? res.then(finish) : finish(res); } #evaluateOp(c, g, r, text) { const b = this.b, t = r.target, cmd = r.command, op = cmd.operation; if (t.controllerGeneration !== b.controllerGeneration) return refused("generation"); if (op === "issue-confirmation") { if (!g.capabilities.includes(CAPS[cmd.operationToConfirm])) return refused("capability"); const x = record("confirmation", { id: newId("confirmation"), actor: c.actor, connection: c.id, target: clone(t), operation: cmd.operationToConfirm, intentDigest: hash({ target: t, operation: cmd.operationToConfirm, stop: b.stop }), expiresAt: new Date(this.now().getTime() + 60000).toISOString(), state: "pending", connectionGeneration: c.generation, stop: b.stop, }); this.confirmations.set(x.id, x); this.#pushTo(c.id, { kind: "confirmation", confirmation: x }); return { outcome: "confirmation-issued", data: { confirmation: x } }; } if (op === "answer-confirmation") { const x = this.confirmations.get(cmd.confirmation); if (!x || x.state !== "pending" || x.actor !== c.actor || x.connection !== c.id || x.connectionGeneration !== c.generation || !equal(x.target, t) || Date.parse(x.expiresAt) <= this.now().getTime()) return refused("confirmation"); x.state = cmd.answer === "confirm" ? "confirmed" : "cancelled"; this.#pushTo(c.id, { kind: "confirmation", confirmation: x }); return { outcome: `confirmation-${x.state}`, data: { confirmation: x } }; } if (op === "takeover") { if (c.id === b.controllerConnection) return refused("already-controller"); if (!g.capabilities.includes("observe")) return refused("capability"); if (b.state !== "active" || b.admission !== "open") return refused("fenced"); return { outcome: this.#transfer(c, false) }; } if (op === "acquire-recovery-control") { if (!g.capabilities.includes("observe")) return refused("capability"); if (b.admission !== "closed") return refused("fenced"); const old = this.connections.get(b.controllerConnection)?.rec; if (old?.state === "connected") return refused("controller-present"); if (!this.#checkConfirmation(r, c, op)) return refused("confirmation"); return { outcome: this.#transfer(c, true) }; } if (c.id !== b.controllerConnection || c.mode !== "controller") return refused("controller"); if (op === "force-stop") { // One escalation at a time: a second force stop waits until the first // ends, and can be retried once it ends `uncertain`. if (!["active", "stopping", "uncertain"].includes(b.state) || this.escalating) return refused("fenced"); if (!this.#checkConfirmation(r, c, op)) return refused("confirmation"); const s = this.#startStop("force-stop", { requestId: r.id, connection: c.id, target: t }); this.escalating = s.id; this.closers.add("force-stop"); this.#admission(); this.exec?.link?.poison("force-stop"); return { outcome: "force-stop-fenced", stop: s, after: () => this.#forceStop(s, { confirmation: cmd.confirmation }) }; } if (op === "recover") return this.#recover(c, r); if (!this.preflightOk) return refused("preflight"); if (b.state !== "active" || b.admission !== "open") return refused("fenced"); if (op === "interrupt") { // The fence, set synchronously; the rest runs under the lock. const closer = `interrupt:${r.id}`; this.closers.add(closer); this.#admission(); return { fence: true, closer, fenceIdx: this.exec?.tracker.idx ?? 0 }; } if (op === "prompt") return this.#admitPrompt(c, g, r, text); return refused("unsupported-capability"); } #checkConfirmation(r, c, op) { const b = this.b, x = this.confirmations.get(r.command.confirmation); if (!x || x.state !== "confirmed" || x.actor !== c.actor || x.connection !== c.id || x.connectionGeneration !== c.generation || x.operation !== op || !equal(x.target, r.target) || x.stop !== b.stop || (op === "recover" && x.stop !== r.command.stop) || x.intentDigest !== hash({ target: r.target, operation: op, stop: x.stop }) || Date.parse(x.expiresAt) <= this.now().getTime()) return false; x.state = "consumed"; this.#pushTo(c.id, { kind: "confirmation", confirmation: x }); return true; } #transfer(c, recovery) { const b = this.b; const old = this.connections.get(b.controllerConnection)?.rec; if (old) old.mode = "observer"; b.controllerGeneration += 1; b.controllerConnection = c.id; c.mode = "controller"; this.#emit("control-transferred", { stop: b.stop }); this.#push({ kind: "binding", binding: b }); for (const x of [old, c].filter(Boolean)) this.#push({ kind: "connection", connection: x }); this.#claimGeneration(); return recovery ? "recovery-control-acquired" : "transferred"; } #claimGeneration() { if (!this.claim || this.claim.record.state === "stopped") return; const generation = this.b.controllerGeneration; this.#claimAdvance({ generation }).catch((err) => this.#uncertain("claim", { error: String(err.code ?? err.message) })); } #observe(c, r) { const t = r.target, cmd = r.command; const res = cmd.cursor ? this.reader.next({ cursor: cmd.cursor, conversation: t.conversation, branch: t.branch, actor: c.actor }) : this.reader.open({ conversation: t.conversation, branch: t.branch, actor: c.actor }); if (!res.ok) { const code = res.refusal.code; if (code === "unknown-branch" || code === "unknown-conversation") return refused("target"); if (code.startsWith("cursor-") || code === "source-replaced") return refused("cursor", { detail: code }); return refused(code, { detail: res.refusal.message }); } const exec = this.exec; return { outcome: "observing", data: { page: res.page, cursor: res.cursor, follow: res.follow, view: res.view, incarnation: this.incarnation, // The seam: no replay in CHAT-03. Events from fromSequence on are // live; the page and the stream are not deduplicated by ID (E4). // `quiet` says no run was visible and no prompt held the slot when // the page was read, so every earlier message had settled and Pi had // persisted it (agent-session.js 386–398 persists on message_end, // before agent_settled). Otherwise the cut is not atomic, and the // client marks the seam and re-reads after the run settles. seam: { replay: "unavailable", streamEpoch: exec?.streamEpoch ?? null, fromSequence: (exec?.seq ?? 0) + 1, reconcile: true, quiet: !exec?.tracker.current && !exec?.tracker.slot }, limitAdvisory: true, }, }; } // ---- prompt ------------------------------------------------------------ #admitPrompt(c, g, r, text) { const b = this.b, t = r.target, cmd = r.command; if (typeof text !== "string" || text.length > 262144) return refused("draft"); if (textPolicy(text)) return refused(TEXT_POLICY); const exec = this.exec; if (exec?.tracker.slot) return refused(BUSY); if (!exec?.link || exec.link.poisoned) return refused("fenced"); const now = this.now().toISOString(); const request = record("request", { id: newId("request"), clientRequest: r.id, connection: c.id, actor: c.actor, grant: g.id, grantRevision: g.revision, target: clone(t), operationDigest: hash({ target: t, command: cmd, payloadDigest: sha256(text) }), command: clone(cmd), admittedAt: now, frozenPayload: null, }); const receipt = record("receipt", { id: newId("receipt"), request: request.id, target: clone(t), state: "admitted", revision: 1, reasonCode: null, event: null, createdAt: now }); this.requests.set(request.id, request); this.receipts.set(receipt.id, receipt); const slot = { request: request.id, receipt, text, connection: c.id, actor: c.actor, generation: t.controllerGeneration, grant: g.id, incarnation: this.incarnation, exec, writeStarted: false, written: false, refused: false, released: false, nativeError: null }; exec.tracker.attach(slot); this.#push({ kind: "receipt", receipt }); return { outcome: "admitted", receipt, request, after: () => this.#dispatch(slot) }; } #recheck(slot) { const b = this.b, exec = slot.exec; if (b.state !== "active" || b.admission !== "open") return "fenced"; const c = this.connections.get(slot.connection)?.rec; if (!c || c.state === "revoked") return "channel"; if (b.controllerConnection !== c.id || c.mode !== "controller") return "controller"; if (b.controllerGeneration !== slot.generation) return "generation"; if (!this.#grantOk(c) || c.grant !== slot.grant) return "grant"; if (exec !== this.exec || exec.execution !== b.execution || !exec.link || exec.link.poisoned) return "fenced"; if (slot.incarnation !== this.incarnation) return STALE_INCARNATION; if (textPolicy(slot.text)) return TEXT_POLICY; return null; } async #dispatch(slot) { const exec = slot.exec; await this.#pause("admitted", { request: slot.request }); let req = null; await this.lock.run(async () => { if (slot.refused || slot.released) return; await this.#pause("before-recheck", { request: slot.request }); const why = this.#recheck(slot); if (why) return this.#dispatchRefused(slot, why); slot.writeStarted = true; req = exec.link.request("prompt", { message: slot.text }, { timeoutMs: this.T.ack }); exec.promptSlots.set(req.id, slot); const w = await req.written; if (w.outcome === "written") this.#markWritten(slot); else this.#transport(exec, `prompt-write-${w.reason ?? w.outcome}`); return undefined; }); if (!slot.written) return; await this.#pause("written", { request: slot.request }); const res = await req.response; if (res.timeout && slot.ackIdx === null && !this.#final(slot.receipt)) this.#transport(exec, "no-prompt-response"); } // The write callback can fire after the engine has read the line and // answered it; a response for the prompt shows the write completed first. #markWritten(slot) { if (slot.written) return; slot.written = true; this.#revise(slot.receipt, "dispatched", null); } #dispatchRefused(slot, code) { slot.refused = true; this.#revise(slot.receipt, "dispatch-refused", code); this.#releaseSlot(slot); } #releaseSlot(slot) { const tr = slot.exec.tracker; if (tr.slot === slot) tr.release(); slot.released = true; this.#wake(slot.exec); } #final(receipt) { return FINAL.has(receipt.state) || this.outcomeUnknown.has(receipt.id); } #revise(receipt, state, reasonCode, extra = {}) { if (!receiptAllows(receipt.state, state)) { this.evidence.refusedRevisions.push({ receipt: receipt.id, from: receipt.state, to: state, reasonCode }); return false; } receipt.state = state; receipt.revision += 1; receipt.reasonCode = reasonCode; receipt.createdAt = this.now().toISOString(); this.#push({ kind: "receipt", receipt, ...extra }); return true; } // A `working` receipt never moves back; it is shown as outcome unknown. #markOutcomeUnknown(slot, reason) { if (this.outcomeUnknown.has(slot.receipt.id)) return; this.outcomeUnknown.add(slot.receipt.id); this.#push({ kind: "receipt", receipt: slot.receipt, outcomeUnknown: true, reason }); } #settleFailure(slot, kind) { if (!slot || !slot.writeStarted || this.#final(slot.receipt)) return; if (slot.receipt.state === "working") this.#markOutcomeUnknown(slot, kind); else this.#revise(slot.receipt, "delivery-unknown", kind); } // ---- engine output ----------------------------------------------------- #stale(exec, link) { return exec !== this.exec || link !== exec.link; } #onResponse(exec, link, value, entry) { if (this.#stale(exec, link)) { this.evidence.dropped.lines += 1; return; } const tr = exec.tracker; const idx = tr.read(); exec.responseIdx.set(value.id, idx); if (entry.type === "clear_queue") this.#closeClearWindow(exec, value.id); if (entry.type === "get_state" && value.success && Number(value.data?.pendingMessageCount) > 0) { tr.overlap("O5", idx, { cause: "pending-messages", count: Number(value.data.pendingMessageCount) }); } if (entry.type === "prompt") { const slot = exec.promptSlots.get(value.id); if (slot) { this.#markWritten(slot); if (value.success) { if (tr.slot === slot) tr.acked(idx); slot.ackIdx = idx; this.#revise(slot.receipt, "acknowledged", null); if (!slot.released && tr.slot === slot) void this.#postAck(slot).catch((e) => this.#internal("post-ack", e)); } else { slot.nativeError = typeof value.error === "string" ? value.error.slice(0, 2000) : null; if (!this.#final(slot.receipt) && this.#revise(slot.receipt, "failed", null, { nativeError: slot.nativeError })) this.#releaseSlot(slot); } } } this.#wake(exec); } async #postAck(slot) { const exec = slot.exec, tr = exec.tracker; const ackIdx = slot.ackIdx; const st = await this.#ask(exec, "get_state", {}, this.T.state); if (slot.released || this.#final(slot.receipt)) return; if (!st.response) return; const replyIdx = st.idx; const between = (i) => i !== null && i > ackIdx && i < replyIdx; const quiet = !tr.runs.some((r) => between(r.startIdx) || between(r.settleIdx)) && !tr.looseSettles.some(between); if (st.response.data?.isStreaming === false && quiet && slot.runId === null) { if (tr.overlapped || tr.gap) return; this.#revise(slot.receipt, "delivery-unknown", HANDLED_WITHOUT_RUN); this.#releaseSlot(slot); return; } if (slot.runId === null) { const ok = await this.#waitFor(exec, () => slot.runId !== null || this.#final(slot.receipt) || slot.released, this.T.start); if (!ok) this.#transport(exec, "no-agent-start"); } } #onSlotSettled(exec, slot, result) { if (result.outcomeUnknown) { this.#markOutcomeUnknown(slot, result.outcomeUnknown); this.#uncertain(result.outcomeUnknown, {}); return; } const extra = {}; if (result.nativeError !== undefined) extra.nativeError = result.nativeError; if (result.stop) extra.stop = result.stop; this.#revise(slot.receipt, result.state, result.reason ?? null, extra); this.#releaseSlot(slot); } #onOverlap(exec, signal, entry) { this.closers.add("overlap"); this.evidence.overlaps.push(entry); this.#settleFailure(exec.tracker.slot, RUN_OVERLAP); this.#uncertain(RUN_OVERLAP, { signal }); this.#push({ kind: "evidence", evidence: { overlap: entry } }); } #onGap(exec, link, kind, detail) { if (this.#stale(exec, link)) { this.evidence.dropped.lines += 1; return; } if (exec.closing) return; this.evidence.gaps.push({ kind, detail, idx: exec.tracker.idx }); this.#transport(exec, kind); } #onEnd(exec, link) { if (this.#stale(exec, link) || exec.closing) return; this.#transport(exec, "eof"); } #transport(exec, reason) { if (exec !== this.exec) return; exec.link?.poison(reason); const tr = exec.tracker; if (!tr.gap) tr.gap = { reason, idx: tr.idx }; this.#settleFailure(tr.slot, TRANSPORT_UNKNOWN); this.#uncertain(TRANSPORT_UNKNOWN, { reason }); this.#wake(exec); } #uncertain(reason, detail = {}) { const b = this.b; this.evidence.uncertain.push({ reason, ...detail, at: this.now().toISOString() }); const first = !this.closers.has("uncertain"); this.closers.add("uncertain"); if (b.state !== "stopping" && b.state !== "stopped") b.state = "uncertain"; this.#admission(); if (first) this.#emit("uncertain", { stop: b.stop }); if (this.claim && ["active", "reserved"].includes(this.claim.record.state)) { this.#claimAdvance({ state: "uncertain" }).catch((err) => this.evidence.uncertain.push({ reason: "claim", error: String(err.code ?? err.message) })); } } #onLine(exec, link, value, bytes) { if (this.#stale(exec, link)) { // H14: late output of a replaced engine is dropped and counted. this.evidence.dropped.lines += 1; this.evidence.dropped.bytes += bytes; return; } const tr = exec.tracker; const idx = tr.read(); this.#map(exec, value, bytes); tr.event(value, idx); this.#wake(exec); } #slotRequest(exec) { const tr = exec.tracker, s = tr.slot; return s && s.runId && tr.current?.id === s.runId && !tr.overlapped ? s.request : null; } #map(exec, ev, bytes) { const type = typeof ev?.type === "string" ? ev.type : null; switch (type) { case "message_start": { exec.message = { id: `${exec.execution}.m${++exec.msgSerial}`, role: roleOf(ev.message) }; return this.#emit("message-start", { message: exec.message.id, role: exec.message.role, request: this.#slotRequest(exec) }); } case "message_update": { const u = ev.assistantMessageEvent; const kind = u?.type === "text_delta" ? "text" : u?.type === "thinking_delta" ? "thinking" : null; if (!kind) { if (isFoldedUpdate(u?.type)) this.evidence.folded += 1; else this.#unknown(`message_update:${u?.type}`, bytes); return undefined; } const m = exec.message ?? (exec.message = { id: `${exec.execution}.m${++exec.msgSerial}`, role: "assistant" }); const ci = Number.isInteger(u.contentIndex) && u.contentIndex >= 0 && u.contentIndex <= 63 ? u.contentIndex : null; return this.#emit(kind === "text" ? "text-delta" : "thinking-delta", { message: m.id, role: m.role, contentIndex: ci, updateMode: "append", content: deltaBlocks(kind, u.delta, u.contentIndex), request: this.#slotRequest(exec) }); } case "message_end": { const m = exec.message ?? { id: `${exec.execution}.m${++exec.msgSerial}`, role: roleOf(ev.message) }; exec.message = null; const parts = partsOf(messageBlocks(ev.message)); parts.forEach((content, i) => this.#emit("message-end", { message: m.id, role: roleOf(ev.message), entry: `${m.id}.e`, part: i, lastPart: i === parts.length - 1, updateMode: "replace", content, request: this.#slotRequest(exec) })); return undefined; } case "tool_execution_start": case "tool_execution_update": case "tool_execution_end": return this.#mapTool(exec, type, ev); case "agent_settled": return this.#emit("run-settled", { request: this.#slotRequest(exec) }); case "extension_ui_request": { if (DIALOG_METHODS.has(ev.method)) { // P3: shown disabled with a reason; never answered. const dialog = { id: safeId(ev.id), method: ev.method, disabled: true, reason: DIALOG_REASON }; exec.dialogs.push(dialog); this.evidence.dialogs.push(dialog); this.#push({ kind: "dialog", dialog }); } else if (NOTIFY_METHODS.has(ev.method)) { this.evidence.notices += 1; } else this.#unknown(`extension_ui_request:${ev.method}`, bytes); return undefined; } default: if (type && KNOWN_UNSHOWN.has(type)) { this.evidence.unshown[type] = (this.evidence.unshown[type] ?? 0) + 1; return undefined; } return this.#unknown(type, bytes); } } #mapTool(exec, type, ev) { const call = safeId(ev.toolCallId); let t = exec.tools.get(call); if (!t) { t = { call, start: null, end: null, message: `${exec.execution}.t${++exec.toolSerial}` }; exec.tools.set(call, t); } const request = this.#slotRequest(exec); if (type === "tool_execution_start") { const parts = fragments(JSON.stringify(ev.args ?? {})); const content = parts.map((a, i) => ({ type: "tool-call", call, name: safeId(ev.toolName), argumentsText: a, block: 0, fragment: i, lastFragment: i === parts.length - 1 })).slice(0, 64); t.start = this.#emit("tool-start", { message: t.message, role: "tool", content, updateMode: "none", request }).id; return undefined; } const text = toolResultText(type === "tool_execution_end" ? ev.result : ev.partialResult); const parts = fragments(text); const content = parts.map((x, i) => ({ type: "tool-result", call, text: x, isError: ev.isError === true, block: 0, fragment: i, lastFragment: i === parts.length - 1 })).slice(0, 64); const e = this.#emit(type === "tool_execution_end" ? "tool-end" : "tool-update", { message: t.message, role: "tool", content, updateMode: "replace", request }); if (type === "tool_execution_end") t.end = e.id; return undefined; } #unknown(type, bytes) { const name = typeof type === "string" && /^[A-Za-z0-9_:.-]{1,80}$/.test(type) ? type : ""; const u = (this.evidence.unknownEvents[name] ??= { count: 0, bytes: 0 }); u.count += 1; u.bytes += bytes; const total = Object.values(this.evidence.unknownEvents).reduce((n, x) => n + x.count, 0); this.#push({ kind: "unknown", count: total, nativeType: name }); } #emit(type, fields = {}) { const exec = this.exec; const seq = exec ? ++exec.seq : this.events.length + 1; const e = record("event", { id: exec ? `${exec.execution}.${seq}` : `${this.b.execution}.c${seq}`, target: targetOf(this.b), sequence: seq, streamEpoch: exec?.streamEpoch ?? `${this.b.execution}.${this.incarnation.slice(0, 12)}`, request: fields.request ?? null, entry: fields.entry ?? null, type, contentIndex: fields.contentIndex ?? null, updateMode: fields.updateMode ?? "none", content: fields.content ?? [], visibility: "permitted-visible", createdAt: this.now().toISOString(), stop: fields.stop ?? null, message: fields.message ?? null, part: fields.part ?? null, lastPart: fields.lastPart ?? null, role: fields.role ?? null, }); this.events.push(e); this.#push({ kind: "event", event: e }); return e; } // ---- admission, pushes, waits ----------------------------------------- // Recomputes admission and pushes the binding when its admission or state // differs from what clients last saw. #admission() { const b = this.b; if (!b) return; const next = b.state === "active" && this.closers.size === 0 ? "open" : "closed"; if (b.admission !== next || this.pushed?.id !== b.id || this.pushed.state !== b.state) { b.admission = next; this.#push({ kind: "binding", binding: b }); } } #push(msg) { if (msg.kind === "binding") this.pushed = { id: msg.binding.id, state: msg.binding.state, admission: msg.binding.admission }; for (const c of this.connections.values()) if (c.send && c.rec.state === "connected") c.send({ ...msg, type: "push" }); } #pushTo(id, msg) { const c = this.connections.get(id); if (c?.send && c.rec.state === "connected") c.send({ ...msg, type: "push" }); } #wake(exec) { for (const w of [...exec.waiters]) w(); } #waitFor(exec, pred, ms) { if (pred()) return Promise.resolve(true); return new Promise((resolve) => { const w = () => { if (!pred()) return; exec.waiters.delete(w); clearTimeout(timer); resolve(true); }; const timer = setTimeout(() => { exec.waiters.delete(w); resolve(pred()); }, ms); exec.waiters.add(w); }); } #evidenceAdd(kind, detail) { const entry = { kind, ...detail, at: this.now().toISOString() }; if (kind === "uncertain") this.evidence.uncertain.push(entry); else this.evidence.stops.push(entry); this.#push({ kind: "evidence", evidence: entry }); } #internal(where, err) { this.evidence.internal.push({ where, error: String(err?.stack ?? err).slice(0, 2000) }); } #claimAdvance(patch) { const run = this.claimChain.then(() => { const p = typeof patch === "function" ? patch(this.claim.record) : patch; return this.store.advance(this.claim, p); }); this.claimChain = run.catch(() => {}); return run; } #claimFinish(patch) { const run = this.claimChain.then(() => this.store.finish(this.claim, patch)); this.claimChain = run.catch(() => {}); return run; } // ---- stops ------------------------------------------------------------- #startStop(mode, { requestId, connection, target }) { const b = this.b; const pred = this.stops.get(b.stop); if (pred && !["stopped", "turn-interrupted", "superseded"].includes(pred.state)) { pred.state = "superseded"; if (pred.mode === "interrupt") this.closers.delete(`interrupt:${pred.id}`); this.#push({ kind: "stop", stop: pred }); } const s = record("stop", { id: newId("stop"), request: requestId, target: clone(target), mode, state: "fenced", queueDrafts: [], cohortRef: b.cohortRef, supervisorEvidence: null, effectsEvidence: null, externalEffects: "uncertain", createdAt: this.now().toISOString(), nativeQueue: "pending", approvalDisposition: "pending", turnEvidence: null, supersedes: b.stop, queueFailures: [], revokedConnection: mode === "revocation" ? connection : null, }); this.stops.set(s.id, s); b.stop = s.id; if (mode === "force-stop") b.state = "stopping"; this.#admission(); if (mode !== "revocation") this.#emit("stopping", { stop: s.id }); this.#push({ kind: "stop", stop: s }); this.#push({ kind: "binding", binding: b }); return s; } #advanceStop(stop, state) { if (!STOP_NEXT[stop.state]?.includes(state)) return false; stop.state = state; if (stop.mode === "force-stop") this.b.state = state === "uncertain" ? "uncertain" : "stopping"; this.#push({ kind: "stop", stop }); this.#admission(); return true; } // Interrupt, under the dispatch lock, after its fence was set. async #interruptLocked(c, r, pre) { const exec = this.exec, tr = exec?.tracker; const dropFence = () => { this.closers.delete(pre.closer); this.#admission(); }; let dispatchRefused = null; const slot = tr?.slot ?? null; if (slot && !slot.writeStarted) { this.#dispatchRefused(slot, "fenced"); dispatchRefused = slot.receipt.id; } const written = tr?.slot && tr.slot.writeStarted ? tr.slot : null; if (!written && !tr?.current) { // H10: lifts only its own fence. dropFence(); return refused(NO_TURN, { effect: { dispatchRefused } }); } const b = this.b; const others = [...this.closers].filter((k) => k !== pre.closer); if (b.state !== "active" || others.length > 0 || exec.link?.poisoned) { dropFence(); return refused("fenced", { effect: { dispatchRefused } }); } if (b.controllerConnection !== c.id || c.mode !== "controller") { dropFence(); return refused("controller", { effect: { dispatchRefused } }); } if (b.controllerGeneration !== r.target.controllerGeneration) { dropFence(); return refused("generation", { effect: { dispatchRefused } }); } const s = this.#startStop("interrupt", { requestId: r.id, connection: c.id, target: r.target }); this.closers.delete(pre.closer); this.closers.add(`interrupt:${s.id}`); this.#admission(); const ctx = { exec, stop: s, fenceIdx: pre.fenceIdx, slot: written, abortIdxs: [], lastAbortIdx: null, clears: [], reason: null }; return { outcome: "interrupt-fenced", stop: s, effect: { dispatchRefused }, after: () => this.#interruptLoop(ctx) }; } async #clear(exec, ctx) { const tr = exec.tracker; let id = null; const res = await this.#ask(exec, "clear_queue", {}, this.T.clear, (rid) => { id = rid; tr.clearWindow += 1; exec.clears.set(rid, { closed: false }); }); if (id) this.#closeClearWindow(exec, id); if (!res.response) return { ok: false, why: res.why }; if (!res.response.success) return { ok: false, why: "error-response" }; const d = res.response.data ?? {}; const items = [...(Array.isArray(d.steering) ? d.steering : []), ...(Array.isArray(d.followUp) ? d.followUp : [])]; const at = this.now().toISOString(); const entry = { at, idx: res.idx, empty: items.length === 0, removed: items.map((x) => ({ digest: sha256(typeof x === "string" ? x : JSON.stringify(x)), bytes: Buffer.byteLength(typeof x === "string" ? x : JSON.stringify(x)) })) }; ctx.clears.push(entry); return { ok: true, ...entry }; } #closeClearWindow(exec, id) { const w = exec.clears.get(id); if (!w || w.closed) return; w.closed = true; exec.tracker.clearWindow = Math.max(0, exec.tracker.clearWindow - 1); } async #interruptLoop(ctx) { const { exec, stop } = ctx, tr = exec.tracker; const gone = () => stop.state === "superseded" || this.exec !== exec; let unknown = false; for (let round = 1; ; round++) { if (gone()) return; if (tr.overlapped || tr.gap || exec.link.poisoned) { unknown = true; break; } const clear = await this.#clear(exec, ctx); if (gone()) return; if (round === 1) this.#advanceStop(stop, "cancelling"); if (!clear.ok) { // Rule 2: no abort, which would run whatever is queued. stop.nativeQueue = "unknown"; ctx.reason = `clear-${clear.why}`; unknown = true; break; } if (!clear.empty) { tr.overlap("O5", clear.idx, { cause: "non-empty-clear", removed: clear.removed }); unknown = true; break; } await this.#pause("before-abort", { stop: stop.id, round }); // An overlap read with the clear's response (O5 from a queue_update in // the same chunk) is recorded by now: no abort, which would run it. if (tr.overlapped || tr.gap || exec.link.poisoned) { unknown = true; break; } const ab = await this.#ask(exec, "abort", {}, this.T.abort, () => { this.abortWritten.add(stop.id); ctx.lastAbortIdx = tr.idx; ctx.abortIdxs.push(tr.idx); }); if (gone()) return; if (!ab.response) { ctx.reason = `abort-${ab.why}`; unknown = true; break; } const slot = ctx.slot; if (slot && !this.#final(slot.receipt)) { if (!(await this.#waitFor(exec, () => slot.ackIdx != null || this.#final(slot.receipt), this.T.ack))) { this.#transport(exec, "no-prompt-response"); unknown = true; break; } if (!(await this.#waitFor(exec, () => slot.runId !== null || this.#final(slot.receipt), this.T.start))) { this.#transport(exec, "no-agent-start"); unknown = true; break; } } const cur = tr.current; if (cur && cur.startIdx <= ctx.lastAbortIdx && !(await tr.waitSettle(cur, this.T.settle))) { // K9: an interrupt that never settles stays uncertain. ctx.reason = "no-settle"; unknown = true; break; } if (gone()) return; if (tr.startedAfter(ctx.lastAbortIdx).length === 0 && !tr.current) break; if (round >= this.T.maxRounds) { stop.nativeQueue = "unknown"; ctx.reason = "rounds-exhausted"; unknown = true; break; } } if (gone()) return; const outcome = unknown ? "unknown" : this.#classifyStop(ctx); if (outcome === "interrupted") return this.#reconcile(ctx); return this.#stopUncertain(ctx, outcome); } // §3 rule 5. Unknown's conditions first. #classifyStop(ctx) { const { exec } = ctx, tr = exec.tracker; if (ctx.lastAbortIdx === null || tr.overlapped || tr.gap || exec.link.poisoned) return "unknown"; if (tr.runs.some((r) => r.startIdx > ctx.lastAbortIdx)) return "unknown"; const active = (r, i) => r.startIdx <= i && (r.settleIdx === null || r.settleIdx > i); const inScope = tr.runs.filter((r) => active(r, ctx.fenceIdx) || (r.startIdx > ctx.fenceIdx && r.startIdx <= ctx.lastAbortIdx)); if (inScope.length > 1) return "unknown"; if (inScope.length === 1) { const run = inScope[0]; if (run.settleIdx === null) return "unknown"; if (run.lastStop === "aborted" && ctx.abortIdxs.some((i) => active(run, i))) return "interrupted"; if (DONE_STOPS.has(run.lastStop)) return "completed-first"; if (run.lastStop === "error") return "failed-on-its-own"; return "unknown"; } const r = ctx.slot?.receipt; if (r && ((r.state === "failed" && r.reasonCode === null) || (r.state === "delivery-unknown" && r.reasonCode === HANDLED_WITHOUT_RUN))) return "no-run"; return "unknown"; } #stopUncertain(ctx, outcome, why = null) { const { stop } = ctx; if (stop.state === "superseded") return; // Rule 2: nothing proved the native queue clear. if (stop.nativeQueue === "pending") stop.nativeQueue = "unknown"; this.#advanceStop(stop, "uncertain"); const entry = { stop: stop.id, mode: stop.mode, outcome, reason: why ?? ctx.reason, queueBasis: SEAL_BASIS, agentLevelQueues: "unobservable", clears: ctx.clears, lastAbortIdx: ctx.lastAbortIdx, fenceIdx: ctx.fenceIdx, }; this.#evidenceAdd("stop", entry); } // §3 rule 6. async #reconcile(ctx) { const { exec, stop } = ctx, tr = exec.tracker; const slot = ctx.slot; if (slot) { const { state, reasonCode } = slot.receipt; const settled = state === "finished" || state === "failed" || (state === "delivery-unknown" && [HANDLED_WITHOUT_RUN, ACK_WITHOUT_START].includes(reasonCode)); if (!settled || this.outcomeUnknown.has(slot.receipt.id)) return this.#stopUncertain(ctx, "interrupted", "receipt-unsettled"); } const st = await this.#ask(exec, "get_state", {}, this.T.state); if (stop.state === "superseded") return undefined; if (!st.response || !st.response.success || st.response.data?.isStreaming !== false || st.response.data?.pendingMessageCount !== 0) return this.#stopUncertain(ctx, "interrupted", "not-idle-after-abort"); const clear = await this.#clear(exec, ctx); if (stop.state === "superseded") return undefined; if (!clear.ok) { stop.nativeQueue = "unknown"; return this.#stopUncertain(ctx, "interrupted", `post-settle-clear-${clear.why}`); } if (!clear.empty) { tr.overlap("O5", clear.idx, { cause: "non-empty-post-settle-clear", removed: clear.removed }); return this.#stopUncertain(ctx, "interrupted", "post-settle-clear-not-empty"); } if (tr.overlapped || tr.gap || exec.link.poisoned) return this.#stopUncertain(ctx, "interrupted", "overlap-or-gap"); if (!this.verifier) return this.#stopUncertain(ctx, "interrupted", "no-verifier"); const b = this.b; const effects = effectReport({ binding: b, stop: stop.id, tools: exec.tools, observedAt: clear.at }); const turn = sealProof(record("turnProof", { id: newId("turn-proof"), authority: AUTHORITY, conversation: b.scope.conversation, execution: b.execution, stop: stop.id, cohortRef: b.cohortRef, nativeQueue: "cleared", nativePending: [], turnState: "interrupted", approvalDisposition: exec.dialogs.length ? "uncertain" : "resolved", effectsEvidence: effects.id, observedAt: clear.at, verificationDigest: "", decisionOutcomes: [], })); this.verifier.post(effects); this.verifier.post(turn); const vctx = { binding: b, stop: stop.id, now: this.now() }; if (!this.verifier.verify(turn, "turnProof", vctx) || !this.verifier.effects(effects, vctx) || turn.approvalDisposition === "uncertain") { return this.#stopUncertain(ctx, "interrupted", "proof-not-verified"); } const c = this.connections.get(b.controllerConnection)?.rec; const g = c ? this.#grantOk(c) : null; if (!c || c.state !== "connected" || !g || !equal(g.scope, b.scope) || !this.channels.has(c.authenticatedChannelRef) || !this.approvedMappings.includes(b.sourceRootRef) || !this.auditWritable) { return this.#stopUncertain(ctx, "interrupted", "controller-not-current"); } stop.state = "turn-interrupted"; stop.nativeQueue = "cleared"; stop.approvalDisposition = turn.approvalDisposition; stop.turnEvidence = turn.id; stop.effectsEvidence = effects.id; this.closers.delete(`interrupt:${stop.id}`); this.#push({ kind: "stop", stop }); this.#emit("reconciled", { stop: stop.id }); this.#evidenceAdd("stop", { stop: stop.id, mode: "interrupt", outcome: "interrupted", queueBasis: SEAL_BASIS, agentLevelQueues: "unobservable", clears: ctx.clears, proofs: { turn: turn.id, effects: effects.id } }); this.#admission(); return undefined; } // Holds `escalating` for the whole escalation, so no second one runs beside // it and writes its phases onto this stop's claim revision. async #forceStop(stop, opts) { this.escalating = stop.id; try { return await this.#escalate(stop, opts); } finally { if (this.escalating === stop.id) this.escalating = null; } } async #escalate(stop, { confirmation = null, resumed = false } = {}) { const b = this.b; const rec = () => this.claim.record; const fail = async (reason) => { if (stop.state !== "superseded") this.#advanceStop(stop, "uncertain"); b.state = "uncertain"; this.closers.add("uncertain"); this.#admission(); this.#emit("uncertain", { stop: stop.id }); this.#evidenceAdd("stop", { stop: stop.id, mode: "force-stop", outcome: "uncertain", reason, resumed }); try { await this.#claimAdvance({ state: "uncertain" }); } catch (err) { this.#evidenceAdd("uncertain", { reason: "claim", error: String(err.code ?? err.message) }); } }; try { await this.#claimAdvance((r) => ({ state: "stopping", stop: { id: stop.id, confirmation, phaseStarted: resumed ? r.stop?.phaseStarted ?? null : null, record: stop } })); } catch (err) { return fail(`claim: ${err.code ?? err.message}`); } await this.#pause("force-stop-recorded", { stop: stop.id }); const engine = rec().engine; if (!engine?.pid) return fail("no engine identity recorded"); if (engine.kind === "pgroup" && (!engine.start || processStart(engine.pid) !== engine.start)) { return fail("process identity can't be checked; no signal sent"); } if (this.exec) this.exec.closing = true; // A launcher may own how its cohort stops (the fixture launcher does); // the result still goes through the verifier. const stopCohort = typeof this.launcher.forceStop === "function" ? (a) => this.launcher.forceStop(a) : forceStopCohort; const result = await stopCohort({ kind: engine.kind, unitName: rec().unitName, invocationId: rec().invocationId, shimSocket: rec().shim, pid: engine.pid, graceMs: this.T.grace, onPhase: async (name) => { await this.#claimAdvance((r) => ({ stop: { ...r.stop, phaseStarted: name } })); await this.#pause(`phase-${name}`, { stop: stop.id }); }, }); if (stop.state === "superseded") return undefined; this.#advanceStop(stop, "stopping"); if (result.outcome !== "proven") return fail(result.reason ?? "cohort evidence unavailable"); if (!this.verifier) return fail("no verifier"); const tools = this.exec && this.exec.execution === b.execution ? this.exec.tools : this.orphanTools ?? new Map(); const proof = cohortProof({ binding: b, stop: stop.id, result }); const effects = effectReport({ binding: b, stop: stop.id, tools, observedAt: result.observedAt }); this.verifier.post(proof); this.verifier.post(effects); const ok = b.stop === stop.id && scopeMatch(stop.target, targetOf(b)) && this.verifier.cohort(proof, effects, { binding: b, stop: stop.id, now: this.now(), epoch: rec().invocationId }); if (!ok) return fail("proof not verified"); let session = null; try { session = this.#readSession(); } catch (err) { return fail(`session unreadable at proof: ${err.code ?? err.message}`); } try { await this.#claimFinish({ state: "stopped", proof: { kind: "cohortProof", ref: proof.id, effects: effects.id }, leafAtProof: session.leaf, branchAtProof: session.branch, stop: { ...rec().stop, record: stop } }); } catch (err) { return fail(`claim: ${err.code ?? err.message}`); } stop.state = "stopped"; stop.supervisorEvidence = proof.id; stop.effectsEvidence = effects.id; stop.externalEffects = effects.invocations.some((i) => i.disposition === "uncertain") ? "uncertain" : effects.invocations.length ? "completed" : "none"; b.state = "stopped"; this.stoppedProof = { proof, effects, epoch: rec().invocationId }; this.#admission(); this.#push({ kind: "stop", stop }); this.#push({ kind: "binding", binding: b }); this.#emit("stopped", { stop: stop.id }); this.#evidenceAdd("stop", { stop: stop.id, mode: "force-stop", outcome: "stopped", proofs: { cohort: proof.id, effects: effects.id }, resumed }); return undefined; } // check.mjs `stopped`. #stopped(s) { const b = this.b, p = this.stoppedProof; if (!s || s.state !== "stopped" || b.stop !== s.id || !scopeMatch(s.target, targetOf(b)) || !p || !this.verifier) return false; if (p.proof.id !== s.supervisorEvidence || p.effects.id !== s.effectsEvidence) return false; return this.verifier.cohort(p.proof, p.effects, { binding: b, stop: s.id, now: this.now(), epoch: p.epoch }); } async #recover(c, r) { const b = this.b, cmd = r.command; const s = this.stops.get(cmd.stop); if (b.state !== "stopped" || b.admission !== "closed" || !this.#stopped(s)) return refused("stop-proof"); let pins; try { pins = this.#pins(); } catch (err) { return refused(err.code ?? ENGINE_PIN_MISMATCH); } if (!equal(pins, this.claim.record.pins)) return refused(ENGINE_PIN_MISMATCH); const session = this.#readSession(); if (session.leaf !== this.claim.record.leafAtProof || session.branch !== this.claim.record.branchAtProof) return refused("target"); if (!this.#checkConfirmation(r, c, "recover")) return refused("confirmation"); const execution = newId("exec"); const generation = this.claim.record.generation + 1; const bindingId = newId("binding"); const claim = await this.store.acquire(this.seatK, this.sessionK, { bindingId, harness: "pi", conversation: this.conversation, branch: session.branch, leaf: session.leaf, pins, owner: this.store.owner(this.incarnation), generation, prior: { claimId: this.claim.claimId, stop: s.id }, execution, cohortRef: null, }); const e = { id: newId("eligibility"), claim, bindingId, execution, generation, stop: s.id, pins, leaf: session.leaf, branch: session.branch, incarnation: this.incarnation, used: false, launched: false }; this.eligibility.set(e.id, e); return { outcome: "recovery-eligible", data: { eligibility: e.id, claim: claim.claimId, generation } }; } // A launcher's call, never a client command (Q15). Single-use: the record // is consumed before any check, so a second call refuses (K17). async launch(id) { const e = this.eligibility.get(id); if (!e || e.used) throw new ControlRefusal(ELIGIBILITY, "no unused eligibility record with that id"); e.used = true; if (e.incarnation !== this.incarnation) throw new ControlRefusal(ELIGIBILITY, "the eligibility record belongs to another controller incarnation"); this.#guardCheck("launch"); for (const key of [this.seatK, this.sessionK]) { const h = this.store.head(key); if (h.damaged || h.n === 0 || h.record.claimId !== e.claim.claimId || h.record.state !== "reserved" || h.record.spawnMarker) { throw new ControlRefusal(ELIGIBILITY, "the reserved claim changed since eligibility"); } } const pins = this.#pins(); if (!equal(pins, e.pins)) throw new ControlRefusal(ENGINE_PIN_MISMATCH, "the engine pin or launch argv changed since eligibility"); const session = this.#readSession(); if (session.leaf !== e.leaf || session.branch !== e.branch) throw new ControlRefusal("target", "the session branch or leaf changed since eligibility (K18)"); e.launched = true; const old = this.exec; this.exec = null; if (old) old.closing = true; this.b = this.#newBinding({ execution: e.execution, generation: e.generation, session, cohortRef: `pending-${e.execution}`, state: "reserved", pins }); this.b.id = e.bindingId; this.closers = new Set(["starting"]); this.preflightOk = false; this.stoppedProof = null; this.orphanTools = null; for (const c of this.connections.values()) if (c.rec.mode === "controller") c.rec.mode = "observer"; this.#push({ kind: "binding", binding: this.b }); await this.#launchInto(e.claim, session); return { binding: this.b, claim: e.claim.claimId }; } // Releases an unlaunched reservation with a no-unit observation: nothing // was spawned under it. async release(id) { const e = this.eligibility.get(id); if (!e || e.launched) throw new ControlRefusal(ELIGIBILITY, "no unlaunched eligibility record with that id"); e.used = true; e.launched = true; await this.store.finish(e.claim, { state: "stopped", proof: { kind: "no-unit", ref: null }, leafAtProof: this.claim.record.leafAtProof ?? null, branchAtProof: this.claim.record.branchAtProof ?? null }); return { released: e.claim.claimId }; } // Server-origin (check.mjs `revoke-connection`). The revocation fence is // never lifted in CHAT-03: reconcile-revocation needs evidence I1 lacks. revokeConnection(id) { const entry = this.connections.get(id); if (!entry) throw new ControlRefusal("channel", "no such connection"); const c = entry.rec, b = this.b; c.state = "revoked"; c.mode = "observer"; if (b.controllerConnection === c.id) { if (b.state === "active" && b.admission === "open") this.#startStop("revocation", { requestId: null, connection: c.id, target: targetOf(b) }); b.controllerConnection = null; b.controllerGeneration += 1; this.closers.add("revocation"); this.#admission(); this.#emit("control-transferred", { stop: b.stop }); this.#claimGeneration(); } this.#push({ kind: "connection", connection: c }); this.#push({ kind: "binding", binding: b }); return "revoked"; } // Test and shutdown hook. Never a release: the claim is unchanged. async close({ killEngine = false } = {}) { const exec = this.exec; if (exec) exec.closing = true; if (killEngine && typeof exec?.proc?.kill === "function") exec.proc.kill(); else if (killEngine && exec?.proc?.pid) { try { process.kill(exec.proc.kind === "pgroup" ? -exec.proc.pid : exec.proc.pid, "SIGKILL"); } catch { // already gone } } for (const s of this.sockets) s.destroy(); if (this.server) await new Promise((r) => this.server.close(() => r())); this.server = null; await this.claimChain; if (exec?.proc && killEngine) await Promise.race([exec.proc.exited, sleep(2000)]); } }