Dewey's round 3 candidate, manifest
agents/dewey/work/queue-40/candidate-manifest-r3.sha256 (d0aa0ded,
27 files, checked OK in the canonical tree).
- WebUI inbox, tasks, agents and trail views, read-only over /api/bus.
The README says the bus proof ends at the Console process.
- CHAT-03 seal: the engine command is fixed, the engine environment is
explicit, SEAL_FLAGS has --no-approve, escalating is cleared on throw.
- Terminal input typed after Ctrl-T or Ctrl-O is held. Only the run whose
own parse set held drains it (T1), and #run catches errors per action.
- DEFERRED keeps N2 and moves F2 to done, citing T1.
Reviews: Filbert approve (comment 27011, rev 260), Darkwing approve
(27013, rev 264). Landing gate on 8cad7722 plus the candidate: webui 22,
conversation 161, control-board 124, every scripts/test-*.sh green,
test-task 98/0. Mutant Mr survives; its flows test is the first
follow-up row.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
1681 lines
81 KiB
JavaScript
1681 lines
81 KiB
JavaScript
// 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, engineEnv } 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"]);
|
||
|
||
// The engine a test runs instead of Pi: `{ command, preArgs, env }`. A
|
||
// symbol key, so no JSON configuration can carry it; the plain `engine`
|
||
// option takes only extraArgs, cwd and envKeys (I3, both reviewers on #1507).
|
||
export const TEST_ENGINE = Symbol("conversation.test-engine");
|
||
const ENGINE_KEYS = new Set(["extraArgs", "cwd", "envKeys"]);
|
||
|
||
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 <project>/.pi/state/<seat>/sessions/<name>.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));
|
||
if (!engine || typeof engine !== "object" || Array.isArray(engine)) throw new ControlRefusal(UNSEALED_ENGINE, "engine is not an object");
|
||
const extra = Object.keys(engine).find((k) => !ENGINE_KEYS.has(k));
|
||
if (extra !== undefined) throw new ControlRefusal(UNSEALED_ENGINE, `engine.${extra.slice(0, 40)} can't be configured; the controller launches the pinned Pi with its own environment`);
|
||
const test = opts[TEST_ENGINE] ?? null;
|
||
this.engine = {
|
||
sealed: test === null,
|
||
command: test ? test.command : process.execPath,
|
||
preArgs: test ? test.preArgs : [join(pinRoot, PI_BIN)],
|
||
extraArgs: engine.extraArgs ?? [],
|
||
cwd: engine.cwd ?? projectRoot,
|
||
env: test ? test.env : engineEnv(engine.envKeys),
|
||
};
|
||
for (const k of ["preArgs", "extraArgs"]) {
|
||
if (!Array.isArray(this.engine[k])) throw new ControlRefusal(UNSEALED_ENGINE, `engine.${k} is not a list`);
|
||
}
|
||
if (typeof this.engine.command !== "string" || !this.engine.command) throw new ControlRefusal(UNSEALED_ENGINE, "the engine command is not a path");
|
||
this.piArgs = buildPiArgs({ sessionFile: this.paths.sessionFile, extraArgs: this.engine.extraArgs });
|
||
this.pinRoot = pinRoot;
|
||
this.#checkSeal();
|
||
this.launcher = launcher;
|
||
this.verifier = verifier;
|
||
this.units = units;
|
||
this.barrier = barrier;
|
||
this.now = now;
|
||
this.T = { ...TIMEOUTS, ...timeouts };
|
||
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);
|
||
}
|
||
|
||
// The seal covers the command: unless a test engine was given, the launch
|
||
// is this Node running the pinned Pi's bin, and nothing else.
|
||
#checkSeal() {
|
||
if (this.engine.sealed && (this.engine.command !== process.execPath || this.engine.preArgs.length !== 1 || this.engine.preArgs[0] !== join(this.pinRoot, PI_BIN))) {
|
||
throw new ControlRefusal(UNSEALED_ENGINE, "the engine command is not the pinned Pi");
|
||
}
|
||
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;
|
||
try {
|
||
this.closers.add("force-stop");
|
||
this.#admission();
|
||
this.exec?.link?.poison("force-stop");
|
||
} catch (err) {
|
||
// #handle drops `after` on a throw, so #forceStop never runs to clear
|
||
// the flag; without this every later force stop is refused `fenced`
|
||
// until restart (Darkwing F2 on #1507).
|
||
if (this.escalating === s.id) this.escalating = null;
|
||
throw err;
|
||
}
|
||
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 : "<invalid>";
|
||
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)]);
|
||
}
|
||
}
|