feat(conversation): CHAT-03 I1, mediated control of a sealed headless Pi (#1507)
Controller, claim store, live-session guard, engine link and seal, turn tracker, cohort force stop and recovery, client library, transcript and mediated terminal, with the fake engine and tests. Fixtures only; no live cutover. Dewey built it. Darkwing (comment 26690) and Filbert (comment 26694) approved round 2. Manifest I1-r2-manifest.sha256 (2b48e333, 27 files). Suites on an export: conversation 152/152, control-board 124, webui 14, seat 19, chat-00/01/01c checks, and all nine scripts/test-*.sh green. Follow-ups for I3 are in DEFERRED. Gate E stays with Jason. Co-Authored-By: Claude Opus 5.5 <[email protected]>
This commit is contained in:
@@ -0,0 +1,354 @@
|
||||
// The writer-claim record (#1507, CHAT-03 §2, D1).
|
||||
//
|
||||
// A claim holds two keys, the seat tuple (seat, project, workspace) and the
|
||||
// native session identity. Each key is a directory under the claim root and
|
||||
// each revision a numbered file in it (r0000000001.json, ...). A revision is
|
||||
// published by writing a temporary file in the same directory, fsyncing it,
|
||||
// link()ing it to the next revision name and fsyncing the directory. link()
|
||||
// fails when the name exists, so publication is exclusive and a revision name
|
||||
// only ever points at complete contents. Revisions are never rewritten.
|
||||
//
|
||||
// Keys are taken seat first, then session. Every step of the lifecycle
|
||||
// publishes on the seat key and then on the session key. The pair's state is
|
||||
// the more conservative of the two: uncertain > stopping > active > reserved
|
||||
// > stopped. A key is held unless its highest revision is `stopped` with a
|
||||
// proof reference. A highest revision that does not parse holds the key as
|
||||
// `uncertain`; an older revision is never reused.
|
||||
//
|
||||
// The store is internal. It never crosses the wire, and it stores only the
|
||||
// CHAT-01 binding fields it needs.
|
||||
|
||||
import { closeSync, fsyncSync, linkSync, mkdirSync, openSync, readdirSync, readFileSync, unlinkSync, writeSync } from "node:fs";
|
||||
import { join } from "node:path";
|
||||
import { createHash, randomBytes } from "node:crypto";
|
||||
import { bootId as readBootId, identityOf, ownerState } from "../../discord/src/journal.mjs";
|
||||
import { ControlRefusal } from "./safe-fs.mjs";
|
||||
import { ID } from "./parts.mjs";
|
||||
|
||||
export const CLAIM_VERSION = 1;
|
||||
export const STATES = Object.freeze(["uncertain", "stopping", "active", "reserved", "stopped"]);
|
||||
export const PROOF_KINDS = Object.freeze(["cohortProof", "boot", "no-unit"]);
|
||||
export const ALREADY_ACTIVE = "already-active";
|
||||
export const UNSAFE_REPLACEMENT = "unsafe-replacement";
|
||||
export const FOREIGN_HOST = "foreign-host";
|
||||
|
||||
const REVISION = /^r(\d{10})\.json$/;
|
||||
const revName = (n) => `r${String(n).padStart(10, "0")}.json`;
|
||||
const rank = (state) => STATES.indexOf(state);
|
||||
|
||||
export function machineId() {
|
||||
try {
|
||||
const id = readFileSync("/etc/machine-id", "utf8").trim();
|
||||
return /^[0-9a-f]{32}$/.test(id) ? id : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
export const defaultHost = Object.freeze({ machineId, bootId: readBootId });
|
||||
|
||||
export function newClaimId() {
|
||||
return "c" + randomBytes(12).toString("hex");
|
||||
}
|
||||
|
||||
export function unitNameFor(claimId) {
|
||||
return `mosaic-chat-${claimId}`;
|
||||
}
|
||||
|
||||
// The more conservative of two states.
|
||||
export function conservative(a, b) {
|
||||
return rank(a) <= rank(b) ? a : b;
|
||||
}
|
||||
|
||||
const keyHash = (key) => createHash("sha256").update(JSON.stringify(key)).digest("hex").slice(0, 32);
|
||||
|
||||
export function seatKey({ seat, project, workspace }) {
|
||||
return { kind: "seat", seat, project, workspace };
|
||||
}
|
||||
|
||||
export function sessionKey(session) {
|
||||
return { kind: "session", session };
|
||||
}
|
||||
|
||||
function valid(rec, key, n) {
|
||||
return rec && typeof rec === "object" && rec.version === CLAIM_VERSION && rec.kind === "writer-claim" &&
|
||||
rec.revision === n && JSON.stringify(rec.key) === JSON.stringify(key) && typeof rec.claimId === "string" &&
|
||||
ID.test(rec.claimId) && STATES.includes(rec.state) && rec.host && typeof rec.host === "object" &&
|
||||
(rec.proof === null || (rec.proof && PROOF_KINDS.includes(rec.proof.kind)));
|
||||
}
|
||||
|
||||
// A key's highest revision is free only when it is `stopped` with a proof.
|
||||
export function held(head) {
|
||||
if (head.n === 0) return false;
|
||||
if (head.damaged) return true;
|
||||
return !(head.record.state === "stopped" && head.record.proof);
|
||||
}
|
||||
|
||||
export function headState(head) {
|
||||
if (head.n === 0) return "stopped";
|
||||
if (head.damaged) return "uncertain";
|
||||
if (head.record.state === "stopped" && !head.record.proof) return "uncertain";
|
||||
return head.record.state;
|
||||
}
|
||||
|
||||
export class ClaimStore {
|
||||
constructor({ root, host = defaultHost, identity = identityOf, barrier = null, now = () => new Date() }) {
|
||||
if (typeof root !== "string" || !root) throw new ControlRefusal("configuration", "a claim root is required; CHAT-03 has no default");
|
||||
this.root = root;
|
||||
this.host = host;
|
||||
this.identity = identity;
|
||||
this.barrier = barrier;
|
||||
this.now = now;
|
||||
}
|
||||
|
||||
dir(key) {
|
||||
return join(this.root, key.kind, keyHash(key));
|
||||
}
|
||||
|
||||
async pause(name, detail) {
|
||||
if (this.barrier) await this.barrier(name, detail);
|
||||
}
|
||||
|
||||
head(key) {
|
||||
const dir = this.dir(key);
|
||||
let names;
|
||||
try {
|
||||
names = readdirSync(dir);
|
||||
} catch (err) {
|
||||
if (err.code === "ENOENT") return { key, n: 0, record: null, damaged: false };
|
||||
throw err;
|
||||
}
|
||||
let n = 0;
|
||||
for (const name of names) {
|
||||
const m = REVISION.exec(name);
|
||||
if (m) n = Math.max(n, Number(m[1]));
|
||||
}
|
||||
if (n === 0) return { key, n: 0, record: null, damaged: false };
|
||||
let record = null;
|
||||
try {
|
||||
record = JSON.parse(readFileSync(join(dir, revName(n)), "utf8"));
|
||||
} catch {
|
||||
record = null;
|
||||
}
|
||||
if (!valid(record, key, n)) return { key, n, record: null, damaged: true };
|
||||
return { key, n, record, damaged: false };
|
||||
}
|
||||
|
||||
revisions(key) {
|
||||
const dir = this.dir(key);
|
||||
let names;
|
||||
try {
|
||||
names = readdirSync(dir);
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
return names.filter((x) => REVISION.test(x)).sort().map((name) => ({ name, text: readFileSync(join(dir, name), "utf8") }));
|
||||
}
|
||||
|
||||
// Publishes revision `n` on `key`. Returns null when another writer
|
||||
// published that revision first.
|
||||
async publish(key, n, body) {
|
||||
const dir = this.dir(key);
|
||||
mkdirSync(dir, { recursive: true, mode: 0o700 });
|
||||
const record = { ...body, key, revision: n, writtenAt: this.now().toISOString() };
|
||||
const tmp = join(dir, `.tmp-${randomBytes(8).toString("hex")}`);
|
||||
const fd = openSync(tmp, "wx", 0o400);
|
||||
try {
|
||||
writeSync(fd, JSON.stringify(record) + "\n");
|
||||
await this.pause("temp-written", { key, n });
|
||||
fsyncSync(fd);
|
||||
} finally {
|
||||
closeSync(fd);
|
||||
}
|
||||
await this.pause("temp-synced", { key, n });
|
||||
let won = true;
|
||||
try {
|
||||
linkSync(tmp, join(dir, revName(n)));
|
||||
} catch (err) {
|
||||
if (err.code !== "EEXIST") throw err;
|
||||
won = false;
|
||||
}
|
||||
unlinkSync(tmp);
|
||||
if (!won) return null;
|
||||
await this.pause("linked", { key, n });
|
||||
const dfd = openSync(dir, "r");
|
||||
try {
|
||||
fsyncSync(dfd);
|
||||
} finally {
|
||||
closeSync(dfd);
|
||||
}
|
||||
await this.pause("dir-synced", { key, n });
|
||||
return record;
|
||||
}
|
||||
|
||||
owner(incarnation) {
|
||||
const id = this.identity(process.pid);
|
||||
return { pid: process.pid, start: id.start, boot: id.boot, incarnation };
|
||||
}
|
||||
|
||||
// Why a held pair refuses a new claim.
|
||||
refusal(seat, session) {
|
||||
const heads = [seat, session].filter(held);
|
||||
const mine = this.host.machineId();
|
||||
if (heads.some((h) => !h.damaged && h.record.host.machineId !== mine)) return FOREIGN_HOST;
|
||||
if (heads.some((h) => h.damaged)) return UNSAFE_REPLACEMENT;
|
||||
if (heads.length === 2 && heads[0].record.claimId !== heads[1].record.claimId) return UNSAFE_REPLACEMENT;
|
||||
const state = heads.map(headState).reduce(conservative, "stopped");
|
||||
return state === "reserved" || state === "active" ? ALREADY_ACTIVE : UNSAFE_REPLACEMENT;
|
||||
}
|
||||
|
||||
refuse(seat, session) {
|
||||
const code = this.refusal(seat, session);
|
||||
return new ControlRefusal(code, `the writer claim for this seat or session is held (${code})`);
|
||||
}
|
||||
|
||||
// Reserves both keys under a new claim ID. `fields` holds the binding
|
||||
// fields the record stores. Refuses when either key is held.
|
||||
async acquire(seatK, sessionK, fields) {
|
||||
for (let attempt = 0; attempt < 4; attempt++) {
|
||||
const seat = this.head(seatK), session = this.head(sessionK);
|
||||
if (held(seat) || held(session)) throw this.refuse(seat, session);
|
||||
const claimId = fields.claimId ?? newClaimId();
|
||||
const body = {
|
||||
version: CLAIM_VERSION, kind: "writer-claim", claimId,
|
||||
bindingId: fields.bindingId, harness: fields.harness, conversation: fields.conversation,
|
||||
branch: fields.branch, leaf: fields.leaf, pins: fields.pins,
|
||||
host: { machineId: this.host.machineId(), bootId: this.host.bootId() },
|
||||
owner: fields.owner, unitName: unitNameFor(claimId), spawnMarker: false,
|
||||
invocationId: null, engine: null, shim: null, generation: fields.generation,
|
||||
state: "reserved", proof: null, stop: null, prior: fields.prior ?? null,
|
||||
execution: fields.execution ?? null, cohortRef: fields.cohortRef ?? null,
|
||||
};
|
||||
const first = await this.publish(seatK, seat.n + 1, body);
|
||||
if (!first) continue; // lost the seat key: re-read and re-decide
|
||||
await this.pause("between-keys", { claimId });
|
||||
const now = this.head(sessionK);
|
||||
const second = held(now) ? null : await this.publish(sessionK, now.n + 1, body);
|
||||
if (!second) {
|
||||
// The session key is held by someone else: close our seat revision
|
||||
// with a no-unit proof. Nothing was spawned.
|
||||
await this.publish(seatK, seat.n + 2, { ...body, state: "stopped", proof: { kind: "no-unit", ref: null } });
|
||||
throw this.refuse(this.head(seatK), this.head(sessionK));
|
||||
}
|
||||
return { claimId, seatKey: seatK, sessionKey: sessionK, n: { seat: first.revision, session: second.revision }, record: body };
|
||||
}
|
||||
throw new ControlRefusal(UNSAFE_REPLACEMENT, "the writer claim changed during every attempt");
|
||||
}
|
||||
|
||||
// Publishes `patch` on both keys, seat first. Both heads must still be this
|
||||
// claim's own last revisions.
|
||||
async advance(claim, patch) {
|
||||
const next = { ...claim.record, ...patch };
|
||||
const n = { ...claim.n };
|
||||
for (const kind of ["seat", "session"]) {
|
||||
const key = kind === "seat" ? claim.seatKey : claim.sessionKey;
|
||||
const head = this.head(key);
|
||||
if (head.damaged || head.n !== n[kind] || head.record.claimId !== claim.claimId) {
|
||||
throw new ControlRefusal(UNSAFE_REPLACEMENT, `the ${kind} key changed under claim ${claim.claimId}`);
|
||||
}
|
||||
const rec = await this.publish(key, head.n + 1, next);
|
||||
if (!rec) throw new ControlRefusal(UNSAFE_REPLACEMENT, `another writer published on the ${kind} key of claim ${claim.claimId}`);
|
||||
n[kind] = rec.revision;
|
||||
if (kind === "seat") await this.pause("between-keys", { claimId: claim.claimId, patch });
|
||||
}
|
||||
claim.record = next;
|
||||
claim.n = n;
|
||||
return claim;
|
||||
}
|
||||
|
||||
// Restart classification for a pair whose recorded owner may be gone. Acts
|
||||
// only when the owner is proven gone: a different boot, or no process with
|
||||
// the recorded pid and start time. `units.lookup(record)` reports the
|
||||
// intended unit as { state: "absent" | "alive" | "unknown" }. `bootProof`
|
||||
// issues the boot proof: it returns { proof, effects } once the verifier
|
||||
// accepted both, or null, which leaves the pair `uncertain`. `adopt` is the
|
||||
// owner record a recovery controller publishes when it takes over an
|
||||
// orphan; without it nothing is adopted. `stoppedPatch` adds fields (the
|
||||
// leaf at proof) to a `stopped` revision this classification publishes.
|
||||
async classify(seatK, sessionK, { units, bootProof = null, adopt = null, stoppedPatch = {} } = {}) {
|
||||
const seat = this.head(seatK), session = this.head(sessionK);
|
||||
if (!held(seat) && !held(session)) return { state: "free" };
|
||||
const code = this.refusal(seat, session);
|
||||
if (code === FOREIGN_HOST) return { state: "uncertain", refusal: FOREIGN_HOST };
|
||||
if (seat.damaged || session.damaged) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT, damaged: true };
|
||||
const heads = [seat, session].filter(held);
|
||||
if (heads.length === 2 && heads[0].record.claimId !== heads[1].record.claimId) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT };
|
||||
const claimId = heads[0].record.claimId;
|
||||
const ours = [seat, session].filter((h) => h.n > 0 && !h.damaged && h.record.claimId === claimId);
|
||||
const latest = ours.reduce((a, b) => (a.record.writtenAt >= b.record.writtenAt ? a : b)).record;
|
||||
const owner = latest.owner;
|
||||
if (adopt && owner?.incarnation === adopt.incarnation && owner?.pid === adopt.pid) return { state: "owned", claimId, record: latest };
|
||||
const status = ownerState(owner ? { pid: owner.pid, start: owner.start ?? null, boot: owner.boot ?? null } : { invalid: true }, { identity: this.identity });
|
||||
if (status === "live" || status === "unknown" || status === "invalid") return { state: headState(heads[0]), refusal: ALREADY_ACTIVE, claimId, record: latest };
|
||||
|
||||
const pair = heads.map(headState).reduce(conservative, "stopped");
|
||||
const marker = ours.some((h) => h.record.spawnMarker);
|
||||
const claim = this.claimFrom(seatK, sessionK, claimId, seat, session, latest);
|
||||
|
||||
// A different boot on the same host: every process of that boot is gone.
|
||||
if (latest.host.bootId !== this.host.bootId()) {
|
||||
if (!bootProof) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT, claimId, record: latest, needs: "boot-proof" };
|
||||
const issued = await bootProof(latest);
|
||||
if (!issued) return { state: "uncertain", refusal: UNSAFE_REPLACEMENT, claimId, record: latest, needs: "boot-proof" };
|
||||
const { proof, effects = null } = issued;
|
||||
await this.finish(claim, { ...stoppedPatch, state: "stopped", proof: { kind: "boot", ref: proof.id, effects: effects?.id ?? null }, spawnMarker: marker });
|
||||
return { state: "stopped", proofKind: "boot", proof, effects, claimId, record: claim.record };
|
||||
}
|
||||
|
||||
// Keys disagree because a release stopped between them: finish it under
|
||||
// the same claim ID and proof.
|
||||
const stoppedHead = ours.find((h) => h.record.state === "stopped" && h.record.proof);
|
||||
if (stoppedHead) {
|
||||
await this.finish(claim, { state: "stopped", proof: stoppedHead.record.proof, spawnMarker: marker || latest.spawnMarker, leafAtProof: stoppedHead.record.leafAtProof ?? null, branchAtProof: stoppedHead.record.branchAtProof ?? null });
|
||||
return { state: "stopped", proofKind: stoppedHead.record.proof.kind, claimId, record: claim.record, completed: true };
|
||||
}
|
||||
|
||||
const unit = await units.lookup(latest);
|
||||
if (pair === "reserved" && !marker && unit.state === "absent") {
|
||||
await this.finish(claim, { ...stoppedPatch, state: "stopped", proof: { kind: "no-unit", ref: null } });
|
||||
return { state: "stopped", proofKind: "no-unit", claimId, record: claim.record };
|
||||
}
|
||||
|
||||
// A marker with no unit, a unit that exists, or an engine that may run:
|
||||
// the pair is uncertain. A collected scope is an absent observation, not
|
||||
// proof that every process ended. The marker is copied, never dropped.
|
||||
const patch = { state: latest.state === "stopping" && latest.stop ? "stopping" : "uncertain", spawnMarker: marker };
|
||||
if (adopt) patch.owner = adopt;
|
||||
await this.finish(claim, patch);
|
||||
return { state: patch.state, claimId, record: claim.record, unit, adopted: Boolean(adopt), claim, resumeStop: latest.stop && latest.state === "stopping" ? latest.stop : null };
|
||||
}
|
||||
|
||||
claimFrom(seatK, sessionK, claimId, seat, session, latest) {
|
||||
return {
|
||||
claimId, seatKey: seatK, sessionKey: sessionK,
|
||||
n: { seat: seat.n, session: session.n },
|
||||
heads: { seat, session },
|
||||
record: { ...latest },
|
||||
};
|
||||
}
|
||||
|
||||
// Brings both keys of a claim to `patch`, seat first. A key whose head
|
||||
// belongs to an older claim (a half-done acquire) gains the revision too.
|
||||
async finish(claim, patch) {
|
||||
const next = { ...claim.record, ...patch };
|
||||
for (const kind of ["seat", "session"]) {
|
||||
const key = kind === "seat" ? claim.seatKey : claim.sessionKey;
|
||||
const head = this.head(key);
|
||||
if (head.damaged) throw new ControlRefusal(UNSAFE_REPLACEMENT, `the ${kind} key is unreadable`);
|
||||
if (head.n > 0 && head.record.claimId !== claim.claimId && held(head)) throw new ControlRefusal(UNSAFE_REPLACEMENT, `the ${kind} key belongs to another claim`);
|
||||
const same = head.n > 0 && head.record.claimId === claim.claimId && head.record.state === next.state &&
|
||||
JSON.stringify(head.record.proof) === JSON.stringify(next.proof) && head.record.spawnMarker === next.spawnMarker &&
|
||||
JSON.stringify(head.record.owner) === JSON.stringify(next.owner);
|
||||
if (same) {
|
||||
claim.n[kind] = head.n;
|
||||
continue;
|
||||
}
|
||||
const rec = await this.publish(key, head.n + 1, next);
|
||||
if (!rec) throw new ControlRefusal(UNSAFE_REPLACEMENT, `another writer published on the ${kind} key of claim ${claim.claimId}`);
|
||||
claim.n[kind] = rec.revision;
|
||||
if (kind === "seat") await this.pause("between-keys", { claimId: claim.claimId, patch });
|
||||
}
|
||||
claim.record = next;
|
||||
return claim;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,234 @@
|
||||
// The CHAT-03 client library (#1507, CHAT-03 §1, deviation V-1).
|
||||
//
|
||||
// It speaks the controller's line protocol over the Unix socket: `hello`
|
||||
// with a grant, then CHAT-01 v2 client requests, each carrying the
|
||||
// controller's incarnation token. The library never resends. A request still
|
||||
// pending when the connection drops, or one refused `stale-incarnation`, is
|
||||
// shown as "outcome unknown, check the transcript" (lead decision 25, H21–
|
||||
// H23). An actor may copy its text into a new request; the library doesn't
|
||||
// present that as a safe retry.
|
||||
//
|
||||
// Drafts are client-local in CHAT-03: the draft ID and revision name the
|
||||
// client's own buffer, and the text travels in the request envelope. CHAT-04
|
||||
// owns server-side drafts.
|
||||
|
||||
import { connect as netConnect } from "node:net";
|
||||
import { LineSplitter, encodeLine, parseLine } from "./framing.mjs";
|
||||
import { newId, targetOf } from "./records.mjs";
|
||||
|
||||
export const OUTCOME_UNKNOWN = "outcome unknown, check the transcript";
|
||||
const STALE = "stale-incarnation";
|
||||
const FINAL = new Set(["finished", "failed", "dispatch-refused", "delivery-unknown"]);
|
||||
|
||||
export class ConversationClient {
|
||||
constructor({ socketPath, grant = "grant-local", connect = netConnect }) {
|
||||
this.socketPath = socketPath;
|
||||
this.grant = grant;
|
||||
this.netConnect = connect;
|
||||
this.sock = null;
|
||||
this.incarnation = null;
|
||||
this.connection = null;
|
||||
this.binding = null;
|
||||
this.streamEpoch = null;
|
||||
this.pending = new Map();
|
||||
this.unknown = [];
|
||||
this.receipts = new Map();
|
||||
this.events = [];
|
||||
this.pushes = [];
|
||||
this.listeners = new Set();
|
||||
this.waiters = new Set();
|
||||
this.closed = true;
|
||||
}
|
||||
|
||||
get target() {
|
||||
return this.binding ? targetOf(this.binding) : null;
|
||||
}
|
||||
|
||||
get isController() {
|
||||
return !!this.connection && this.binding?.controllerConnection === this.connection.id;
|
||||
}
|
||||
|
||||
// `fn(message)` for every push, reply and status change.
|
||||
on(fn) {
|
||||
this.listeners.add(fn);
|
||||
return () => this.listeners.delete(fn);
|
||||
}
|
||||
|
||||
#notify(msg) {
|
||||
for (const fn of this.listeners) fn(msg);
|
||||
for (const w of [...this.waiters]) w();
|
||||
}
|
||||
|
||||
// Resolves with the welcome. A new incarnation turns every request still
|
||||
// pending under the old one into outcome-unknown; none is resent (H23).
|
||||
async connect() {
|
||||
this.#dropPending("disconnected");
|
||||
const sock = this.netConnect(this.socketPath);
|
||||
this.sock = sock;
|
||||
await new Promise((resolve, reject) => {
|
||||
sock.once("connect", resolve);
|
||||
sock.once("error", reject);
|
||||
});
|
||||
this.closed = false;
|
||||
const welcome = new Promise((resolve, reject) => {
|
||||
this.welcomeWaiter = { resolve, reject };
|
||||
});
|
||||
const splitter = new LineSplitter((line) => this.#onLine(line));
|
||||
sock.on("data", (c) => splitter.push(c));
|
||||
sock.on("error", () => {});
|
||||
sock.on("close", () => {
|
||||
if (this.sock !== sock) return;
|
||||
this.closed = true;
|
||||
this.welcomeWaiter?.reject(new Error("the controller closed the connection"));
|
||||
this.welcomeWaiter = null;
|
||||
this.#dropPending("disconnected");
|
||||
this.#notify({ type: "status", status: "disconnected" });
|
||||
});
|
||||
sock.write(encodeLine({ type: "hello", grant: this.grant }));
|
||||
return welcome;
|
||||
}
|
||||
|
||||
close() {
|
||||
this.sock?.destroy();
|
||||
}
|
||||
|
||||
#dropPending(reason) {
|
||||
for (const [id, p] of this.pending) {
|
||||
const entry = { request: id, operation: p.envelope.command.operation, incarnation: p.incarnation, reason, display: OUTCOME_UNKNOWN };
|
||||
this.unknown.push(entry);
|
||||
p.resolve({ outcome: "outcome-unknown", refusal: null, display: OUTCOME_UNKNOWN, reason, request: null, receipt: null });
|
||||
this.#notify({ type: "outcome-unknown", ...entry });
|
||||
}
|
||||
this.pending.clear();
|
||||
}
|
||||
|
||||
// A receipt last seen short of a final state belongs to the old controller;
|
||||
// the new one doesn't know it. Its outcome is unknown (H20, H23).
|
||||
#unsettled() {
|
||||
for (const r of this.receipts.values()) {
|
||||
if (FINAL.has(r.state) || this.unknown.some((u) => u.receipt === r.id)) continue;
|
||||
const entry = { request: null, receipt: r.id, operation: "prompt", incarnation: this.incarnation, reason: STALE, display: OUTCOME_UNKNOWN };
|
||||
this.unknown.push(entry);
|
||||
this.#notify({ type: "outcome-unknown", ...entry });
|
||||
}
|
||||
}
|
||||
|
||||
#onLine(line) {
|
||||
const parsed = parseLine(line);
|
||||
if (parsed.error) return;
|
||||
const msg = parsed.value;
|
||||
if (msg.type === "welcome") {
|
||||
if (this.incarnation !== null && this.incarnation !== msg.incarnation) {
|
||||
this.#dropPending(STALE);
|
||||
this.#unsettled();
|
||||
}
|
||||
this.incarnation = msg.incarnation;
|
||||
this.connection = msg.connection;
|
||||
this.binding = msg.binding;
|
||||
this.streamEpoch = msg.streamEpoch;
|
||||
this.welcomeWaiter?.resolve(msg);
|
||||
this.welcomeWaiter = null;
|
||||
this.#notify(msg);
|
||||
return;
|
||||
}
|
||||
if (msg.type === "refused" && this.welcomeWaiter) {
|
||||
this.welcomeWaiter.reject(Object.assign(new Error(`refused: ${msg.refusal}`), { refusal: msg.refusal }));
|
||||
this.welcomeWaiter = null;
|
||||
return;
|
||||
}
|
||||
if (msg.type === "reply") {
|
||||
const p = this.pending.get(msg.id);
|
||||
if (!p) return;
|
||||
this.pending.delete(msg.id);
|
||||
const reply = { ...msg };
|
||||
if (msg.refusal === STALE) {
|
||||
reply.display = OUTCOME_UNKNOWN;
|
||||
this.unknown.push({ request: msg.id, operation: p.envelope.command.operation, incarnation: p.incarnation, reason: STALE, display: OUTCOME_UNKNOWN });
|
||||
}
|
||||
if (msg.receipt) this.receipts.set(msg.receipt.id, msg.receipt);
|
||||
p.resolve(reply);
|
||||
this.#notify(reply);
|
||||
return;
|
||||
}
|
||||
if (msg.type === "push") {
|
||||
this.pushes.push(msg);
|
||||
if (msg.kind === "binding") this.binding = msg.binding;
|
||||
else if (msg.kind === "connection" && msg.connection.id === this.connection?.id) this.connection = msg.connection;
|
||||
else if (msg.kind === "receipt") this.receipts.set(msg.receipt.id, msg.receipt);
|
||||
else if (msg.kind === "event") this.events.push(msg.event);
|
||||
this.#notify(msg);
|
||||
}
|
||||
}
|
||||
|
||||
// The envelope for one operation on the current target.
|
||||
envelope(operation, fields = {}) {
|
||||
if (!this.connection || !this.binding) throw new Error("not connected");
|
||||
return { version: 2, kind: "clientRequest", id: newId("req"), connection: this.connection.id, target: this.target, command: { operation, ...fields } };
|
||||
}
|
||||
|
||||
// Sends one request and resolves with the controller's reply.
|
||||
request(operation, fields = {}, text) {
|
||||
return this.send(this.envelope(operation, fields), text);
|
||||
}
|
||||
|
||||
// Sends an envelope under the token this client holds now. An envelope is
|
||||
// sent once; the library has no retry path.
|
||||
send(envelope, text, { incarnation = this.incarnation } = {}) {
|
||||
if (this.closed) return Promise.resolve({ outcome: "refused:channel", refusal: "channel", display: "not connected" });
|
||||
return new Promise((resolve) => {
|
||||
this.pending.set(envelope.id, { envelope, incarnation, resolve });
|
||||
const msg = { type: "request", incarnation, request: envelope };
|
||||
if (text !== undefined) msg.text = text;
|
||||
this.sock.write(encodeLine(msg));
|
||||
});
|
||||
}
|
||||
|
||||
prompt(text, draft = { id: newId("draft"), revision: 1 }) {
|
||||
return this.request("prompt", { draft: draft.id, draftRevision: draft.revision }, text);
|
||||
}
|
||||
|
||||
observe({ cursor = null, limit = 50 } = {}) {
|
||||
return this.request("observe", { cursor, limit });
|
||||
}
|
||||
|
||||
takeover() {
|
||||
return this.request("takeover");
|
||||
}
|
||||
|
||||
interrupt() {
|
||||
return this.request("interrupt");
|
||||
}
|
||||
|
||||
// Issue, confirm and use a confirmation for a force stop, a recovery or
|
||||
// recovery control.
|
||||
async confirmed(operation, fields = {}) {
|
||||
const issued = await this.request("issue-confirmation", { operationToConfirm: operation });
|
||||
if (issued.outcome !== "confirmation-issued") return issued;
|
||||
const id = issued.data.confirmation.id;
|
||||
const answered = await this.request("answer-confirmation", { confirmation: id, answer: "confirm" });
|
||||
if (answered.outcome !== "confirmation-confirmed") return answered;
|
||||
return this.request(operation, { ...fields, confirmation: id });
|
||||
}
|
||||
|
||||
receipt(id) {
|
||||
return this.receipts.get(id) ?? null;
|
||||
}
|
||||
|
||||
// Resolves when `pred()` holds, or false after `ms`.
|
||||
waitFor(pred, ms = 5000) {
|
||||
if (pred()) return Promise.resolve(true);
|
||||
return new Promise((resolve) => {
|
||||
const w = () => {
|
||||
if (!pred()) return;
|
||||
this.waiters.delete(w);
|
||||
clearTimeout(t);
|
||||
resolve(true);
|
||||
};
|
||||
const t = setTimeout(() => {
|
||||
this.waiters.delete(w);
|
||||
resolve(false);
|
||||
}, ms);
|
||||
this.waiters.add(w);
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,235 @@
|
||||
// Engine launch, cohort observation and force stop (#1507, CHAT-03 §6).
|
||||
//
|
||||
// ScopeLauncher starts the supervisor shim (shim.mjs) in a delegated
|
||||
// systemd user scope named after the claim. The cohort is the scope's
|
||||
// `engine` cgroup, and the scope's invocation ID is the membership epoch. A
|
||||
// unit with the right name but another invocation ID is a different cohort:
|
||||
// it gets no signal, and its evidence is unavailable.
|
||||
//
|
||||
// PgroupLauncher is the fallback with no scope: the engine leads its own
|
||||
// process group. A process group can't be enumerated completely, so a force
|
||||
// stop on it always ends `uncertain`.
|
||||
//
|
||||
// Force stop: check the invocation ID against systemd and the shim; TERM
|
||||
// every member and wait a bounded grace; freeze `engine` and wait for
|
||||
// `frozen 1`; enumerate every member with pid and start time; write
|
||||
// `cgroup.kill`; wait for `populated 0`. Only when all of that succeeded is
|
||||
// membership complete. Anything unavailable ends the stop `uncertain`.
|
||||
|
||||
import { spawn, spawnSync } from "node:child_process";
|
||||
import { existsSync, readFileSync } from "node:fs";
|
||||
import { connect } from "node:net";
|
||||
import { dirname, join } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { processStart } from "../../discord/src/journal.mjs";
|
||||
import { LineSplitter, encodeLine, parseLine } from "./framing.mjs";
|
||||
import { hash, newId, record, sealProof } from "./records.mjs";
|
||||
|
||||
export const AUTHORITY = "mosaic-conversation-shim-fixture";
|
||||
export const SHIM_PATH = join(dirname(fileURLToPath(import.meta.url)), "shim.mjs");
|
||||
|
||||
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
|
||||
|
||||
export function cohortRefOf({ machineId, bootId, unitName, invocationId }) {
|
||||
return "cohort-" + hash({ machineId, bootId, unitName, invocationId }).slice(0, 40);
|
||||
}
|
||||
|
||||
// One request per connection keeps a dead shim from wedging a caller.
|
||||
export function shimRequest(socketPath, op, extra = {}, timeoutMs = 3000) {
|
||||
return new Promise((resolve) => {
|
||||
if (!socketPath || !existsSync(socketPath)) return resolve({ ok: false, unavailable: "the shim socket is gone" });
|
||||
const sock = connect(socketPath);
|
||||
let done = false;
|
||||
const finish = (v) => {
|
||||
if (done) return;
|
||||
done = true;
|
||||
clearTimeout(timer);
|
||||
sock.destroy();
|
||||
resolve(v);
|
||||
};
|
||||
const timer = setTimeout(() => finish({ ok: false, unavailable: `the shim did not answer ${op}` }), timeoutMs);
|
||||
const splitter = new LineSplitter((line) => {
|
||||
const p = parseLine(line);
|
||||
finish(p.error ? { ok: false, unavailable: `shim answer ${p.error}` } : p.value);
|
||||
});
|
||||
sock.on("data", (c) => splitter.push(c));
|
||||
sock.on("error", () => finish({ ok: false, unavailable: "the shim is unreachable" }));
|
||||
sock.on("close", () => finish({ ok: false, unavailable: "the shim closed the connection" }));
|
||||
sock.on("connect", () => sock.write(encodeLine({ id: 1, op, ...extra })));
|
||||
});
|
||||
}
|
||||
|
||||
export function systemctlShow(unitName) {
|
||||
const r = spawnSync("systemctl", ["--user", "show", "-p", "LoadState,ActiveState,InvocationID,ControlGroup", `${unitName}.scope`], { encoding: "utf8", timeout: 5000 });
|
||||
if (r.status !== 0) return null;
|
||||
const out = {};
|
||||
for (const line of r.stdout.split("\n")) {
|
||||
const i = line.indexOf("=");
|
||||
if (i > 0) out[line.slice(0, i)] = line.slice(i + 1);
|
||||
}
|
||||
return { loadState: out.LoadState ?? null, activeState: out.ActiveState ?? null, invocationId: out.InvocationID || null, controlGroup: out.ControlGroup || null };
|
||||
}
|
||||
|
||||
// claim.classify's unit lookup. A unit systemd has collected is absent; one
|
||||
// that is loaded and active is alive; anything unreadable is unknown.
|
||||
export const systemdUnits = Object.freeze({
|
||||
lookup(rec) {
|
||||
if (!rec?.unitName) return { state: "unknown" };
|
||||
const s = systemctlShow(rec.unitName);
|
||||
if (!s) return { state: "unknown" };
|
||||
if (s.loadState === "not-found" || (s.activeState === "inactive" && !s.controlGroup)) return { state: "absent" };
|
||||
if (["active", "activating", "deactivating", "reloading"].includes(s.activeState)) return { state: "alive", invocationId: s.invocationId };
|
||||
return { state: "unknown", detail: s };
|
||||
},
|
||||
});
|
||||
|
||||
export function scopeAvailable() {
|
||||
const r = spawnSync("systemd-run", ["--user", "--scope", "--quiet", "-p", "Delegate=yes", "--", "true"], { timeout: 10000 });
|
||||
return r.status === 0;
|
||||
}
|
||||
|
||||
export class ScopeLauncher {
|
||||
constructor({ shimPath = SHIM_PATH, startTimeoutMs = 10000 } = {}) {
|
||||
this.kind = "scope";
|
||||
this.shimPath = shimPath;
|
||||
this.startTimeoutMs = startTimeoutMs;
|
||||
}
|
||||
|
||||
async launch({ unitName, socketPath, command, args, cwd, env }) {
|
||||
const proc = spawn("systemd-run", ["--user", "--scope", "-p", "Delegate=yes", `--unit=${unitName}`, "--quiet", "--", process.execPath, this.shimPath, "--socket", socketPath, "--", command, ...args], {
|
||||
cwd, env, stdio: ["pipe", "pipe", "pipe"],
|
||||
});
|
||||
const exited = new Promise((r) => proc.on("exit", (code, signal) => r({ code, signal })));
|
||||
const end = Date.now() + this.startTimeoutMs;
|
||||
let hello = null;
|
||||
while (Date.now() < end) {
|
||||
const race = await Promise.race([exited.then((e) => ({ exited: e })), sleep(20).then(() => null)]);
|
||||
if (race?.exited) throw Object.assign(new Error(`the scope exited before the shim answered (${JSON.stringify(race.exited)})`), { code: "launch-failed" });
|
||||
if (existsSync(socketPath)) {
|
||||
hello = await shimRequest(socketPath, "hello");
|
||||
if (hello.ok) break;
|
||||
}
|
||||
}
|
||||
if (!hello?.ok) throw Object.assign(new Error("the shim never answered"), { code: "launch-failed", proc });
|
||||
const show = systemctlShow(unitName);
|
||||
return {
|
||||
kind: "scope", proc, stdin: proc.stdin, stdout: proc.stdout, stderr: proc.stderr,
|
||||
pid: hello.enginePid, start: hello.engineStart === null ? null : String(hello.engineStart),
|
||||
invocationId: hello.invocationId, scope: hello.scope, shimSocket: socketPath,
|
||||
systemd: show, exited,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
export class PgroupLauncher {
|
||||
constructor() {
|
||||
this.kind = "pgroup";
|
||||
}
|
||||
|
||||
async launch({ command, args, cwd, env }) {
|
||||
const proc = spawn(command, args, { cwd, env, stdio: ["pipe", "pipe", "pipe"], detached: true });
|
||||
const exited = new Promise((r) => proc.on("exit", (code, signal) => r({ code, signal })));
|
||||
await new Promise((resolve, reject) => {
|
||||
proc.once("spawn", resolve);
|
||||
proc.once("error", reject);
|
||||
});
|
||||
return { kind: "pgroup", proc, stdin: proc.stdin, stdout: proc.stdout, stderr: proc.stderr, pid: proc.pid, start: processStart(proc.pid), invocationId: null, scope: null, shimSocket: null, exited };
|
||||
}
|
||||
}
|
||||
|
||||
// Runs the force-stop escalation from the TERM phase. `onPhase(name)` is
|
||||
// awaited before each phase begins, so the caller records the phase before
|
||||
// any signal. Returns { outcome: "proven" | "unavailable", ... }. Nothing is
|
||||
// recorded as done unless it was observed.
|
||||
export async function forceStopCohort({ kind, unitName, invocationId, shimSocket, pid, graceMs = 1000, onPhase = async () => {} }) {
|
||||
if (kind === "pgroup") {
|
||||
await onPhase("term");
|
||||
try {
|
||||
process.kill(-pid, "SIGTERM");
|
||||
} catch {
|
||||
// the group is gone
|
||||
}
|
||||
await sleep(graceMs);
|
||||
await onPhase("kill");
|
||||
try {
|
||||
process.kill(-pid, "SIGKILL");
|
||||
} catch {
|
||||
// the group is gone
|
||||
}
|
||||
return { outcome: "unavailable", reason: "process-group fallback: membership can't be enumerated completely" };
|
||||
}
|
||||
const show = systemctlShow(unitName);
|
||||
if (!show || !invocationId || show.invocationId !== invocationId) {
|
||||
return { outcome: "unavailable", reason: `invocation ID mismatch or unreadable (recorded ${invocationId}, found ${show?.invocationId ?? "none"}); no signal sent` };
|
||||
}
|
||||
const hello = await shimRequest(shimSocket, "hello");
|
||||
if (!hello.ok) return { outcome: "unavailable", reason: hello.unavailable };
|
||||
if (hello.invocationId !== invocationId || hello.scope !== show.controlGroup) {
|
||||
return { outcome: "unavailable", reason: "the shim does not answer for the recorded scope; no signal sent" };
|
||||
}
|
||||
await onPhase("term");
|
||||
const term = await shimRequest(shimSocket, "term");
|
||||
if (!term.ok) return { outcome: "unavailable", reason: term.unavailable, phase: "term" };
|
||||
const end = Date.now() + graceMs;
|
||||
for (;;) {
|
||||
const e = await shimRequest(shimSocket, "events");
|
||||
if (!e.ok) return { outcome: "unavailable", reason: e.unavailable, phase: "term" };
|
||||
if (e.populated === 0 || Date.now() >= end) break;
|
||||
await sleep(20);
|
||||
}
|
||||
await onPhase("kill");
|
||||
const frozen = await shimRequest(shimSocket, "freeze", { timeoutMs: 3000 }, 6000);
|
||||
if (!frozen.ok) return { outcome: "unavailable", reason: frozen.unavailable, phase: "kill" };
|
||||
if (frozen.timedOut) return { outcome: "unavailable", reason: "engine never reported frozen 1", phase: "kill" };
|
||||
const listed = await shimRequest(shimSocket, "members");
|
||||
if (!listed.ok) return { outcome: "unavailable", reason: listed.unavailable, phase: "kill" };
|
||||
const killed = await shimRequest(shimSocket, "kill", { timeoutMs: 5000 }, 8000);
|
||||
if (!killed.ok) return { outcome: "unavailable", reason: killed.unavailable, phase: "kill" };
|
||||
if (killed.timedOut || killed.populated !== 0) return { outcome: "unavailable", reason: "engine never reported populated 0", phase: "kill" };
|
||||
const observedAt = new Date().toISOString();
|
||||
if (listed.members.some((m) => !Number.isInteger(m.startTicks) || m.startTicks < 1)) {
|
||||
return { outcome: "unavailable", reason: "a member's start time was unreadable at enumeration", phase: "kill" };
|
||||
}
|
||||
return {
|
||||
outcome: "proven", membershipComplete: true, epoch: invocationId, observedAt, boot: listed.boot,
|
||||
members: listed.members.map((m) => ({ pid: m.pid, boot: listed.boot, startTicks: m.startTicks, terminatedAt: observedAt })),
|
||||
};
|
||||
}
|
||||
|
||||
export function cohortProof({ binding, stop, result }) {
|
||||
return sealProof(record("cohortProof", {
|
||||
id: newId("cohort-proof"), authority: AUTHORITY, conversation: binding.scope.conversation, execution: binding.execution,
|
||||
cohortRef: binding.cohortRef, membershipEpoch: result.epoch, membershipComplete: result.membershipComplete === true,
|
||||
members: result.members, observedAt: result.observedAt, verificationDigest: "", stop,
|
||||
}));
|
||||
}
|
||||
|
||||
// A tool-start with no tool-end at stop time is an `uncertain` effect.
|
||||
// Killing never counts as rollback.
|
||||
export function effectReport({ binding, stop, tools, observedAt }) {
|
||||
return sealProof(record("effectReport", {
|
||||
id: newId("effects"), authority: AUTHORITY, conversation: binding.scope.conversation, execution: binding.execution,
|
||||
cohortRef: binding.cohortRef,
|
||||
invocations: [...tools.values()].map((t) => ({ id: t.call, disposition: t.end ? "completed" : "uncertain", evidence: t.end ?? t.start })),
|
||||
observedAt, verificationDigest: "", stop,
|
||||
}));
|
||||
}
|
||||
|
||||
// The current boot's start time, from /proc/stat btime. Every process of an
|
||||
// earlier boot ended before it.
|
||||
export function bootTime() {
|
||||
const line = readFileSync("/proc/stat", "utf8").split("\n").find((l) => l.startsWith("btime "));
|
||||
return new Date(Number(line.slice(6)) * 1000).toISOString();
|
||||
}
|
||||
|
||||
// A boot proof for a claim recorded under an earlier boot of this host. Its
|
||||
// evidence is the boot change; it goes through the same verifier.
|
||||
export function bootProof({ claim, conversation, execution, cohortRef, stop, now = () => new Date() }) {
|
||||
const terminatedAt = bootTime();
|
||||
const members = claim.engine?.pid && claim.engine?.start ? [{ pid: claim.engine.pid, boot: claim.host.bootId, startTicks: Number(claim.engine.start), terminatedAt }] : [];
|
||||
return sealProof(record("cohortProof", {
|
||||
id: newId("boot-proof"), authority: AUTHORITY, conversation, execution, cohortRef,
|
||||
membershipEpoch: claim.invocationId ?? `boot-${claim.host.bootId}`, membershipComplete: true, members,
|
||||
observedAt: now().toISOString(), verificationDigest: "", stop,
|
||||
}));
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,126 @@
|
||||
// The engine pipe (#1507, CHAT-03 §1).
|
||||
//
|
||||
// One EngineLink per execution owns the engine's stdin and reads its stdout.
|
||||
// Output is split on LF only (framing.mjs); every line reaches the controller
|
||||
// in order, tagged with the link's execution, so a replaced engine's late
|
||||
// output can be dropped by incarnation (H14).
|
||||
//
|
||||
// Write outcomes (§1):
|
||||
// written the whole line was accepted by the pipe (the write callback ran
|
||||
// without an error). That is not native consumption.
|
||||
// unknown the write returned an error (EPIPE), or its callback did not run
|
||||
// within the bound. A partial line may be in the pipe, so the link
|
||||
// is poisoned and never written again.
|
||||
// `acknowledged` is the response to a request, which `request()` reports
|
||||
// separately. A request whose response does not come within its bound is
|
||||
// reported as a timeout; the controller decides what that means.
|
||||
|
||||
import { LineSplitter, encodeLine, parseLine } from "./framing.mjs";
|
||||
|
||||
export const WRITE_TIMEOUT_MS = 5000;
|
||||
|
||||
export class EngineLink {
|
||||
// `onLine(link, value, bytes)` gets every parsed non-response line;
|
||||
// `onResponse(link, value, entry)` every response, before its promise
|
||||
// resolves; `onGap(link, kind, detail)` an unparseable line, an unmatched
|
||||
// response, an overlong line or EOF with a partial line.
|
||||
constructor({ execution, stdin, stdout, onLine, onResponse = () => {}, onGap, onEnd = () => {}, writeTimeoutMs = WRITE_TIMEOUT_MS }) {
|
||||
this.execution = execution;
|
||||
this.stdin = stdin;
|
||||
this.onLine = onLine;
|
||||
this.onResponse = onResponse;
|
||||
this.onGap = onGap;
|
||||
this.onEnd = onEnd;
|
||||
this.writeTimeoutMs = writeTimeoutMs;
|
||||
this.poisoned = null;
|
||||
this.pending = new Map();
|
||||
this.serial = 0;
|
||||
this.bytesWritten = 0;
|
||||
this.ended = false;
|
||||
this.writes = [];
|
||||
stdin.on("error", (err) => {
|
||||
this.stdinError = err.code ?? err.message;
|
||||
});
|
||||
this.splitter = new LineSplitter((line) => this.#line(line), { onOverflow: () => this.onGap(this, "overlong-line", null) });
|
||||
stdout.on("data", (chunk) => this.splitter.push(chunk));
|
||||
stdout.on("end", () => {
|
||||
this.ended = true;
|
||||
if (this.splitter.pending() > 0) this.onGap(this, "partial-line-at-eof", { bytes: this.splitter.pending() });
|
||||
this.onEnd(this);
|
||||
});
|
||||
stdout.on("error", () => {});
|
||||
}
|
||||
|
||||
#line(line) {
|
||||
const bytes = Buffer.byteLength(line, "utf8");
|
||||
const parsed = parseLine(line);
|
||||
if (parsed.error) return this.onGap(this, "unparseable-line", { bytes, error: parsed.error });
|
||||
const value = parsed.value;
|
||||
if (value.type === "response") {
|
||||
const entry = typeof value.id === "string" ? this.pending.get(value.id) : undefined;
|
||||
if (!entry) return this.onGap(this, "unmatched-response", { bytes, command: typeof value.command === "string" ? value.command.slice(0, 40) : null });
|
||||
this.pending.delete(value.id);
|
||||
entry.late = entry.timedOut;
|
||||
this.onResponse(this, value, entry);
|
||||
entry.resolve({ response: value, late: entry.late });
|
||||
return undefined;
|
||||
}
|
||||
return this.onLine(this, value, bytes);
|
||||
}
|
||||
|
||||
poison(reason) {
|
||||
if (!this.poisoned) this.poisoned = reason;
|
||||
}
|
||||
|
||||
// Writes one record. Never writes to a poisoned link.
|
||||
write(value) {
|
||||
if (this.poisoned) return Promise.resolve({ outcome: "refused", reason: "poisoned" });
|
||||
const line = encodeLine(value);
|
||||
const bytes = Buffer.byteLength(line, "utf8");
|
||||
return new Promise((resolve) => {
|
||||
let done = false;
|
||||
const finish = (o) => {
|
||||
if (done) return;
|
||||
done = true;
|
||||
clearTimeout(timer);
|
||||
if (o.outcome === "unknown") this.poison(o.reason);
|
||||
else this.bytesWritten += bytes;
|
||||
this.writes.push({ type: value.type, id: value.id, outcome: o.outcome });
|
||||
resolve(o);
|
||||
};
|
||||
const timer = setTimeout(() => finish({ outcome: "unknown", reason: "write-timeout" }), this.writeTimeoutMs);
|
||||
try {
|
||||
this.stdin.write(line, (err) => finish(err ? { outcome: "unknown", reason: err.code ?? "write-error" } : { outcome: "written" }));
|
||||
} catch (err) {
|
||||
finish({ outcome: "unknown", reason: err.code ?? "write-error" });
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// Sends a command. Returns { id, written, response }: `written` settles
|
||||
// with the write outcome, `response` with { response } or { timeout: true }
|
||||
// or { unsent: outcome }. The response is registered before the write, so
|
||||
// it can never be unmatched.
|
||||
request(type, fields = {}, { timeoutMs = 5000, beforeWrite = null } = {}) {
|
||||
const id = `${type}-${++this.serial}`;
|
||||
let resolve;
|
||||
const response = new Promise((r) => (resolve = r));
|
||||
const entry = { id, type, resolve, timedOut: false, late: false };
|
||||
this.pending.set(id, entry);
|
||||
beforeWrite?.(id);
|
||||
const written = this.write({ id, type, ...fields });
|
||||
written.then((w) => {
|
||||
if (w.outcome !== "written") {
|
||||
this.pending.delete(id);
|
||||
resolve({ unsent: w });
|
||||
return;
|
||||
}
|
||||
setTimeout(() => {
|
||||
if (!this.pending.has(id)) return;
|
||||
entry.timedOut = true; // kept, so a late answer is recognised and not a gap
|
||||
resolve({ timeout: true });
|
||||
}, timeoutMs);
|
||||
});
|
||||
return { id, written, response };
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,121 @@
|
||||
// Pi RPC events to CHAT-01 events (#1507, CHAT-03 §3 "Events").
|
||||
//
|
||||
// Pinned Pi 0.85.1 (docs/rpc.md, rpc-types.d.ts) emits JSON lines. Mapped:
|
||||
//
|
||||
// message_start -> message-start
|
||||
// message_update text_delta -> text-delta (append)
|
||||
// message_update thinking_delta -> thinking-delta (append)
|
||||
// tool_execution_start -> tool-start
|
||||
// tool_execution_update -> tool-update (replace)
|
||||
// tool_execution_end -> tool-end (replace)
|
||||
// message_end -> message-end (replace, one or more parts)
|
||||
// agent_settled -> run-settled (never cohort termination)
|
||||
//
|
||||
// Known and deliberately not shown: turn_start, turn_end, agent_start,
|
||||
// agent_end, queue_update, compaction_*, auto_retry_*, summarization_retry_*,
|
||||
// bash_execution_update, extension_error, extension_ui_request (P3 has its own notice), and the
|
||||
// message_update subtypes that message-end supersedes (text_start, text_end,
|
||||
// thinking_start, thinking_end, toolcall_*, start, done, error). The
|
||||
// controller counts them in its evidence. Anything else is an unknown event:
|
||||
// no client event, counted with its type and byte size, and never passed
|
||||
// through raw.
|
||||
//
|
||||
// Pi's stream carries no entry ID. `entry` on message-end is a stream-local
|
||||
// ID, not a session entry ID, so the seam between a history page and the
|
||||
// stream can't be deduplicated by ID (E4).
|
||||
|
||||
import { fragments, safeId, LIMITS } from "./parts.mjs";
|
||||
|
||||
export const KNOWN_UNSHOWN = Object.freeze(new Set([
|
||||
"turn_start", "turn_end", "agent_start", "agent_end", "queue_update", "compaction_start", "compaction_end",
|
||||
"auto_retry_start", "auto_retry_end", "summarization_retry_scheduled", "summarization_retry_attempt_start",
|
||||
"summarization_retry_finished", "bash_execution_update", "extension_error", "extension_ui_request", "response",
|
||||
]));
|
||||
export const MAPPED = Object.freeze(new Set(["message_start", "message_update", "message_end", "tool_execution_start", "tool_execution_update", "tool_execution_end", "agent_settled"]));
|
||||
const FOLDED_UPDATES = new Set(["start", "text_start", "text_end", "thinking_start", "thinking_end", "toolcall_start", "toolcall_delta", "toolcall_end", "done", "error"]);
|
||||
|
||||
export const DIALOG_METHODS = Object.freeze(new Set(["select", "confirm", "input", "editor"]));
|
||||
export const NOTIFY_METHODS = Object.freeze(new Set(["notify", "setStatus", "setWidget", "setTitle", "set_editor_text"]));
|
||||
|
||||
export function roleOf(message) {
|
||||
switch (message?.role) {
|
||||
case "user":
|
||||
return "user";
|
||||
case "assistant":
|
||||
return "assistant";
|
||||
case "toolResult":
|
||||
return "tool";
|
||||
case "compactionSummary":
|
||||
return "compaction";
|
||||
default:
|
||||
return "notice";
|
||||
}
|
||||
}
|
||||
|
||||
const textOf = (content) => {
|
||||
if (typeof content === "string") return content;
|
||||
if (!Array.isArray(content)) return "";
|
||||
return content.filter((c) => c && c.type === "text" && typeof c.text === "string").map((c) => c.text).join("");
|
||||
};
|
||||
|
||||
function pushText(out, type, text, block, extra = {}) {
|
||||
const parts = fragments(text);
|
||||
parts.forEach((t, i) => out.push({ type, ...extra, text: t, block, fragment: i, lastFragment: i === parts.length - 1 }));
|
||||
}
|
||||
|
||||
// Every content block of a finished native message, in CHAT-01 form.
|
||||
export function messageBlocks(message) {
|
||||
const out = [];
|
||||
const role = roleOf(message);
|
||||
if (role === "tool") {
|
||||
pushText(out, "tool-result", textOf(message.content), 0, { call: safeId(message.toolCallId), isError: message.isError === true });
|
||||
return out;
|
||||
}
|
||||
if (role === "compaction") {
|
||||
const parts = fragments(typeof message.summary === "string" ? message.summary : "");
|
||||
parts.forEach((t, i) => out.push({ type: "compaction", summary: t, nativeEntry: safeId(message.firstKeptEntryId ?? "compaction"), block: 0, fragment: i, lastFragment: i === parts.length - 1 }));
|
||||
return out;
|
||||
}
|
||||
const content = typeof message?.content === "string" ? [{ type: "text", text: message.content }] : Array.isArray(message?.content) ? message.content : [];
|
||||
content.forEach((c, block) => {
|
||||
if (!c || typeof c !== "object") return;
|
||||
if (c.type === "text") pushText(out, "text", String(c.text ?? ""), block);
|
||||
else if (c.type === "thinking") {
|
||||
if (c.redacted) out.push({ type: "thinking", text: "", visibility: "unavailable", block, fragment: 0, lastFragment: true });
|
||||
else pushText(out, "thinking", String(c.thinking ?? ""), block, { visibility: "permitted-visible" });
|
||||
} else if (c.type === "toolCall") {
|
||||
const parts = fragments(JSON.stringify(c.arguments ?? {}));
|
||||
parts.forEach((t, i) => out.push({ type: "tool-call", call: safeId(c.id), name: safeId(c.name), argumentsText: t, block, fragment: i, lastFragment: i === parts.length - 1 }));
|
||||
} else if (c.type === "image") {
|
||||
out.push({ type: "attachment", attachment: safeId(`image-${block}`), block, fragment: 0, lastFragment: true });
|
||||
} else {
|
||||
pushText(out, "text", `[${String(c.type).slice(0, 40)} block not shown]`, block);
|
||||
}
|
||||
});
|
||||
return out;
|
||||
}
|
||||
|
||||
// Splits blocks into message-end parts of at most LIMITS.blocks blocks.
|
||||
export function partsOf(blocks) {
|
||||
if (blocks.length === 0) return [[]];
|
||||
const out = [];
|
||||
for (let i = 0; i < blocks.length; i += LIMITS.blocks) out.push(blocks.slice(i, i + LIMITS.blocks));
|
||||
return out;
|
||||
}
|
||||
|
||||
export function deltaBlocks(kind, delta, contentIndex) {
|
||||
const out = [];
|
||||
const block = Number.isInteger(contentIndex) && contentIndex >= 0 ? contentIndex : 0;
|
||||
if (kind === "thinking") pushText(out, "thinking", String(delta ?? ""), block, { visibility: "permitted-visible" });
|
||||
else pushText(out, "text", String(delta ?? ""), block);
|
||||
return out;
|
||||
}
|
||||
|
||||
export function isFoldedUpdate(type) {
|
||||
return FOLDED_UPDATES.has(type);
|
||||
}
|
||||
|
||||
export function toolResultText(result) {
|
||||
if (result && typeof result === "object") return textOf(result.content);
|
||||
return typeof result === "string" ? result : "";
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
// JSONL framing for the engine pipe and the control socket (#1507, CHAT-03 §1).
|
||||
//
|
||||
// Records are split on LF (0x0A) only, and a trailing CR is stripped (rpc.md
|
||||
// lines 30–38). Node's readline is not used, because it also splits on U+2028
|
||||
// and U+2029, which JSON strings may carry raw (E2). Splitting happens on
|
||||
// bytes, before UTF-8 decoding, so a multibyte character split across two
|
||||
// chunks is never damaged.
|
||||
|
||||
export const MAX_LINE_BYTES = 64 * 1024 * 1024;
|
||||
|
||||
export class LineSplitter {
|
||||
constructor(onLine, { maxBytes = MAX_LINE_BYTES, onOverflow = null } = {}) {
|
||||
this.onLine = onLine;
|
||||
this.maxBytes = maxBytes;
|
||||
this.onOverflow = onOverflow;
|
||||
this.parts = [];
|
||||
this.size = 0;
|
||||
this.overflowed = false;
|
||||
}
|
||||
|
||||
push(chunk) {
|
||||
if (this.overflowed) return;
|
||||
let start = 0;
|
||||
for (;;) {
|
||||
const lf = chunk.indexOf(0x0a, start);
|
||||
if (lf === -1) break;
|
||||
this.parts.push(chunk.subarray(start, lf));
|
||||
const line = Buffer.concat(this.parts);
|
||||
this.parts = [];
|
||||
this.size = 0;
|
||||
const end = line.length > 0 && line[line.length - 1] === 0x0d ? line.length - 1 : line.length;
|
||||
this.onLine(line.subarray(0, end).toString("utf8"));
|
||||
start = lf + 1;
|
||||
}
|
||||
if (start < chunk.length) {
|
||||
const rest = chunk.subarray(start);
|
||||
this.size += rest.length;
|
||||
if (this.size > this.maxBytes) {
|
||||
this.overflowed = true;
|
||||
this.parts = [];
|
||||
this.onOverflow?.();
|
||||
return;
|
||||
}
|
||||
this.parts.push(Buffer.from(rest));
|
||||
}
|
||||
}
|
||||
|
||||
// Bytes left without a terminating LF when the stream ends. They are never
|
||||
// parsed as a record: a partial line is a transport gap, not data.
|
||||
pending() {
|
||||
return this.size;
|
||||
}
|
||||
}
|
||||
|
||||
export function encodeLine(value) {
|
||||
return JSON.stringify(value) + "\n";
|
||||
}
|
||||
|
||||
// Parses one line. Returns {value} or {error}; never throws.
|
||||
export function parseLine(line) {
|
||||
try {
|
||||
const value = JSON.parse(line);
|
||||
if (value === null || typeof value !== "object" || Array.isArray(value)) return { error: "not-an-object" };
|
||||
return { value };
|
||||
} catch {
|
||||
return { error: "unparseable" };
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,143 @@
|
||||
// The live-session guard (#1507, CHAT-03 §2).
|
||||
//
|
||||
// CHAT-03 is fixture-only in code. The claim root, the socket directory and
|
||||
// every session path are constructor arguments, with no default. At
|
||||
// construction and again at bind, their real paths must lie inside the
|
||||
// explicit fixture root and outside every protected location: the
|
||||
// repository's .pi/state/, ~/.pi, ~/.claude, the configured data root and any
|
||||
// path a seat registration names. A path is refused when it is inside a
|
||||
// protected location or contains one. The real protections always apply;
|
||||
// options only add to them, so a test can't switch them off.
|
||||
//
|
||||
// The guard comes out only at cutover (CHAT-07), as a reviewed data-map
|
||||
// change.
|
||||
|
||||
import { existsSync, readdirSync, readFileSync, realpathSync } from "node:fs";
|
||||
import { homedir } from "node:os";
|
||||
import { basename, dirname, isAbsolute, join, resolve, sep } from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { ControlRefusal } from "./safe-fs.mjs";
|
||||
|
||||
export const LIVE_SESSION_REFUSED = "live-session-refused";
|
||||
|
||||
const REPO_ROOT = resolve(dirname(fileURLToPath(import.meta.url)), "..", "..", "..");
|
||||
|
||||
// The real path of `p`, or of its nearest existing ancestor with the rest
|
||||
// appended when `p` doesn't exist yet.
|
||||
export function realPath(p) {
|
||||
const abs = resolve(p);
|
||||
const tail = [];
|
||||
let cur = abs;
|
||||
for (;;) {
|
||||
try {
|
||||
return join(realpathSync(cur), ...tail.reverse());
|
||||
} catch (err) {
|
||||
if (err.code !== "ENOENT" && err.code !== "ENOTDIR") throw err;
|
||||
const up = dirname(cur);
|
||||
if (up === cur) return abs;
|
||||
tail.push(basename(cur));
|
||||
cur = up;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const inside = (child, parent) => child === parent || child.startsWith(parent.endsWith(sep) ? parent : parent + sep);
|
||||
const overlaps = (a, b) => inside(a, b) || inside(b, a);
|
||||
|
||||
function configuredDataRoot(home) {
|
||||
try {
|
||||
const cfg = JSON.parse(readFileSync(join(home, ".config", "mosaic-dev", "config.json"), "utf8"));
|
||||
if (typeof cfg?.dataRoot === "string" && cfg.dataRoot) return cfg.dataRoot.replace(/^~(?=$|\/)/, home);
|
||||
} catch {
|
||||
// An unreadable config adds no location; the default data root still applies.
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
// Registrations under <dataRoot>/seats/<layout>/<seat>/registration.json. An
|
||||
// unreadable one adds nothing; the data root itself is protected anyway.
|
||||
function registrationsUnder(dataRoot) {
|
||||
const out = [];
|
||||
const seats = join(dataRoot, "seats");
|
||||
let layouts = [];
|
||||
try {
|
||||
layouts = readdirSync(seats);
|
||||
} catch {
|
||||
return out;
|
||||
}
|
||||
for (const layout of layouts) {
|
||||
let names = [];
|
||||
try {
|
||||
names = readdirSync(join(seats, layout));
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
for (const seat of names) {
|
||||
try {
|
||||
out.push(JSON.parse(readFileSync(join(seats, layout, seat, "registration.json"), "utf8")));
|
||||
} catch {
|
||||
// skipped
|
||||
}
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
function pathsIn(value, out) {
|
||||
if (typeof value === "string") {
|
||||
if (isAbsolute(value)) out.push(value);
|
||||
} else if (Array.isArray(value)) {
|
||||
for (const v of value) pathsIn(v, out);
|
||||
} else if (value && typeof value === "object") {
|
||||
for (const v of Object.values(value)) pathsIn(v, out);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
export class LiveSessionGuard {
|
||||
// `fixtureRoot` is required. `repoRoots`, `homes`, `dataRoots` and
|
||||
// `registrations` add protected locations to the real ones.
|
||||
constructor({ fixtureRoot, repoRoots = [], homes = [], dataRoots = [], registrations = [] } = {}) {
|
||||
if (typeof fixtureRoot !== "string" || !isAbsolute(fixtureRoot)) throw new ControlRefusal(LIVE_SESSION_REFUSED, "an absolute fixture root is required");
|
||||
if (!existsSync(fixtureRoot)) throw new ControlRefusal(LIVE_SESSION_REFUSED, "the fixture root does not exist");
|
||||
this.fixtureRoot = fixtureRoot;
|
||||
const home = homedir();
|
||||
const allHomes = [home, ...homes];
|
||||
const protectedPaths = [];
|
||||
for (const repo of [REPO_ROOT, ...repoRoots]) protectedPaths.push({ path: join(repo, ".pi", "state"), why: "a repository .pi/state" });
|
||||
for (const h of allHomes) {
|
||||
protectedPaths.push({ path: join(h, ".pi"), why: "~/.pi" });
|
||||
protectedPaths.push({ path: join(h, ".claude"), why: "~/.claude" });
|
||||
protectedPaths.push({ path: join(h, ".mosaic-dev"), why: "the default data root" });
|
||||
const configured = configuredDataRoot(h);
|
||||
if (configured) protectedPaths.push({ path: configured, why: "the configured data root" });
|
||||
}
|
||||
for (const d of dataRoots) protectedPaths.push({ path: d, why: "the data root" });
|
||||
const roots = protectedPaths.filter((p) => p.why.includes("data root")).map((p) => p.path);
|
||||
const found = roots.flatMap(registrationsUnder);
|
||||
for (const reg of [...found, ...registrations]) for (const p of pathsIn(reg, [])) protectedPaths.push({ path: p, why: "a seat registration" });
|
||||
this.protected = protectedPaths;
|
||||
}
|
||||
|
||||
// Refuses unless every path is inside the fixture root and clear of every
|
||||
// protected location, both as written and as real paths.
|
||||
check(paths, when) {
|
||||
const root = realPath(this.fixtureRoot);
|
||||
const guarded = this.protected.flatMap((p) => [{ ...p, path: resolve(p.path) }, { ...p, path: realPath(p.path) }]);
|
||||
for (const [name, p] of Object.entries(paths)) {
|
||||
if (typeof p !== "string" || !isAbsolute(p)) throw new ControlRefusal(LIVE_SESSION_REFUSED, `${name} must be an absolute path (${when})`);
|
||||
const written = resolve(p), real = realPath(p);
|
||||
if (!inside(written, root) && !inside(written, resolve(this.fixtureRoot))) {
|
||||
throw new ControlRefusal(LIVE_SESSION_REFUSED, `${name} is outside the fixture root (${when})`);
|
||||
}
|
||||
// The real path must be inside the real fixture root: a symlink out of
|
||||
// it is refused even when the link itself sits inside.
|
||||
if (!inside(real, root)) throw new ControlRefusal(LIVE_SESSION_REFUSED, `${name} resolves outside the fixture root (${when})`);
|
||||
for (const candidate of [written, real]) {
|
||||
const hit = guarded.find((g) => overlaps(candidate, g.path));
|
||||
if (hit) throw new ControlRefusal(LIVE_SESSION_REFUSED, `${name} overlaps ${hit.why} (${when})`);
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,92 @@
|
||||
// The Pi pin and the engine seal (#1507, CHAT-03 §3, lead decisions 31–32).
|
||||
//
|
||||
// Pin: package-lock.json and npm's installed record
|
||||
// (node_modules/.package-lock.json) must both name the pinned version with the
|
||||
// pinned integrity. That ties the install to the package through npm's record;
|
||||
// it is not a hash of the files on disk. `pi` runs dist/bundle/cli.js, the
|
||||
// package's bin, and the built-in llama.cpp extension ships inside it.
|
||||
//
|
||||
// Seal: the controller builds the launch argv. It always carries
|
||||
// --no-extensions, --no-prompt-templates and --no-themes, and never an
|
||||
// --extension argument (cli/args.js; usage.md 224 and 233–236). With
|
||||
// --no-extensions Pi loads only command-line extension paths
|
||||
// (resource-loader.js 316–318), so no explicit extension loads. Under the seal
|
||||
// the Mosaic prompt in the slot is the only thing that can start a run, which
|
||||
// is the basis for attributing a run to it by order.
|
||||
//
|
||||
// The seal is an allow-list. Pi's parser (cli/args.js) keeps the last --mode
|
||||
// and the last --session, reads a bare word as a prompt and an `@` word as a
|
||||
// file, so the argv must be exactly the controller's prefix followed by
|
||||
// ENGINE_OPTIONS pairs, each at most once with one plain value.
|
||||
|
||||
import { readFileSync } from "node:fs";
|
||||
import { isAbsolute, join } from "node:path";
|
||||
import { createHash } from "node:crypto";
|
||||
import { ControlRefusal } from "./safe-fs.mjs";
|
||||
|
||||
export const PI_PACKAGE = "@earendil-works/pi-coding-agent";
|
||||
export const PI_VERSION = "0.85.1";
|
||||
export const PI_INTEGRITY = "sha512-FGRN+OHbWaefBPGaTggAdLjrIHW+s2PzLyglz/5dfLzb9of7uuXMXYC0fJIeZTw+shS32o2cuQ9jF7YSDuL/oQ==";
|
||||
export const PI_BIN = join("node_modules", PI_PACKAGE, "dist", "bundle", "cli.js");
|
||||
export const SEAL_FLAGS = Object.freeze(["--no-extensions", "--no-prompt-templates", "--no-themes"]);
|
||||
export const ENGINE_OPTIONS = Object.freeze(["--model", "--provider", "--thinking"]);
|
||||
|
||||
export const ENGINE_PIN_MISMATCH = "engine-pin-mismatch";
|
||||
export const UNSEALED_ENGINE = "unsealed-engine";
|
||||
|
||||
function lockEntry(path) {
|
||||
let lock;
|
||||
try {
|
||||
lock = JSON.parse(readFileSync(path, "utf8"));
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
const entry = lock?.packages?.[`node_modules/${PI_PACKAGE}`];
|
||||
return entry && typeof entry === "object" ? entry : null;
|
||||
}
|
||||
|
||||
// `root` holds package-lock.json and node_modules/.package-lock.json.
|
||||
export function checkEnginePin(root) {
|
||||
for (const path of [join(root, "package-lock.json"), join(root, "node_modules", ".package-lock.json")]) {
|
||||
const entry = lockEntry(path);
|
||||
if (!entry || entry.version !== PI_VERSION || entry.integrity !== PI_INTEGRITY) {
|
||||
throw new ControlRefusal(ENGINE_PIN_MISMATCH, `${path} does not pin ${PI_PACKAGE} ${PI_VERSION} with the pinned integrity`);
|
||||
}
|
||||
}
|
||||
return { version: PI_VERSION, pin: PI_INTEGRITY };
|
||||
}
|
||||
|
||||
export function buildPiArgs({ sessionFile, extraArgs = [] }) {
|
||||
return ["--mode", "rpc", ...SEAL_FLAGS, "--session", sessionFile, ...extraArgs];
|
||||
}
|
||||
|
||||
// Refuses any argv that is not `--mode rpc`, the three --no-* flags and
|
||||
// `--session <absolute path>`, in that order, followed by ENGINE_OPTIONS
|
||||
// pairs. That covers --extension in either spelling, a second --mode or
|
||||
// --session, session and output flags (--no-session, --fork, --export, ...)
|
||||
// and stray prompt words.
|
||||
export function checkSeal(args) {
|
||||
if (!Array.isArray(args) || args.some((a) => typeof a !== "string")) throw new ControlRefusal(UNSEALED_ENGINE, "launch argv is not a list of strings");
|
||||
const extension = args.find((a) => a === "-e" || a === "--extension" || a.startsWith("--extension="));
|
||||
if (extension !== undefined) throw new ControlRefusal(UNSEALED_ENGINE, `launch argv carries ${extension}`);
|
||||
for (const flag of SEAL_FLAGS) {
|
||||
if (!args.includes(flag)) throw new ControlRefusal(UNSEALED_ENGINE, `launch argv lacks ${flag}`);
|
||||
}
|
||||
const prefix = ["--mode", "rpc", ...SEAL_FLAGS, "--session"];
|
||||
if (prefix.some((a, i) => args[i] !== a)) throw new ControlRefusal(UNSEALED_ENGINE, `launch argv does not start with ${prefix.join(" ")}`);
|
||||
const file = args[prefix.length];
|
||||
if (typeof file !== "string" || !isAbsolute(file)) throw new ControlRefusal(UNSEALED_ENGINE, "the --session value is not an absolute path");
|
||||
const seen = new Set();
|
||||
for (let i = prefix.length + 1; i < args.length; i += 2) {
|
||||
const flag = args[i], value = args[i + 1];
|
||||
if (!ENGINE_OPTIONS.includes(flag)) throw new ControlRefusal(UNSEALED_ENGINE, `launch argv carries ${flag}, which is not one of ${ENGINE_OPTIONS.join(", ")}`);
|
||||
if (seen.has(flag)) throw new ControlRefusal(UNSEALED_ENGINE, `launch argv repeats ${flag}`);
|
||||
if (typeof value !== "string" || !value || value.startsWith("-") || value.startsWith("@")) throw new ControlRefusal(UNSEALED_ENGINE, `${flag} needs one plain value`);
|
||||
seen.add(flag);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
export function argvDigest(command, args) {
|
||||
return createHash("sha256").update(JSON.stringify([command, ...args])).digest("hex");
|
||||
}
|
||||
@@ -0,0 +1,101 @@
|
||||
// CHAT-01 records for the live controller (#1507, CHAT-03).
|
||||
//
|
||||
// Every record the controller emits is a CHAT-01 v2 record (schema
|
||||
// docs/plans/chat-01/contracts.schema.json). The digest rule is CHAT-01's:
|
||||
// sha256 of JSON with sorted keys (check.mjs `hash`).
|
||||
//
|
||||
// The verifier is the fixture's trusted digest registry, as in CHAT-01
|
||||
// (check.mjs `proof` and `stopped`). A proof the producer posts to it is
|
||||
// self-posted, so no CHAT-03 proof is live authority (CHAT-01 lines 330–333).
|
||||
// Without a verifier no proof verifies, and every stop ends `uncertain`.
|
||||
|
||||
import { createHash, randomBytes } from "node:crypto";
|
||||
import { ID } from "./parts.mjs";
|
||||
|
||||
export const VERSION = 2;
|
||||
|
||||
const sortKeys = (v) =>
|
||||
Array.isArray(v) ? v.map(sortKeys) : v && typeof v === "object" ? Object.fromEntries(Object.keys(v).sort().map((k) => [k, sortKeys(v[k])])) : v;
|
||||
export const canonical = (v) => JSON.stringify(sortKeys(v));
|
||||
export const hash = (v) => createHash("sha256").update(canonical(v)).digest("hex");
|
||||
export const sha256 = (data) => createHash("sha256").update(data).digest("hex");
|
||||
export const without = (v, key) => Object.fromEntries(Object.entries(v).filter(([k]) => k !== key));
|
||||
export const equal = (a, b) => canonical(a) === canonical(b);
|
||||
export const clone = (v) => structuredClone(v);
|
||||
|
||||
export function newId(prefix) {
|
||||
const id = `${prefix}-${randomBytes(8).toString("hex")}`;
|
||||
if (!ID.test(id)) throw new Error(`bad id prefix ${prefix}`);
|
||||
return id;
|
||||
}
|
||||
|
||||
export const record = (kind, fields) => ({ version: VERSION, kind, ...fields });
|
||||
|
||||
export function targetOf(binding) {
|
||||
return {
|
||||
conversation: binding.scope.conversation,
|
||||
branch: binding.branch,
|
||||
execution: binding.execution,
|
||||
controllerGeneration: binding.controllerGeneration,
|
||||
};
|
||||
}
|
||||
|
||||
export const scopeMatch = (a, b) => a.conversation === b.conversation && a.branch === b.branch && a.execution === b.execution;
|
||||
|
||||
// Seals a proof: its verification digest covers every other field.
|
||||
export function sealProof(p) {
|
||||
const body = without(p, "verificationDigest");
|
||||
return { ...body, verificationDigest: hash(body) };
|
||||
}
|
||||
|
||||
export class FixtureVerifier {
|
||||
constructor({ authorities = [] } = {}) {
|
||||
this.authorities = new Set(authorities);
|
||||
this.trusted = new Map();
|
||||
}
|
||||
|
||||
// The fixture registry records a posted proof's digest (CHAT-01 fixtures'
|
||||
// `trustedProofs`). Only proofs from a trusted authority are recorded.
|
||||
post(p) {
|
||||
if (!p || !this.authorities.has(p.authority)) return false;
|
||||
const digest = hash(without(p, "verificationDigest"));
|
||||
if (digest !== p.verificationDigest) return false;
|
||||
this.trusted.set(p.id, digest);
|
||||
return true;
|
||||
}
|
||||
|
||||
// check.mjs `proof`: authority, scope, stop, time and digest.
|
||||
verify(p, kind, { binding, stop, now }) {
|
||||
if (!p || p.kind !== kind || !this.authorities.has(p.authority)) return null;
|
||||
if (p.conversation !== binding.scope.conversation || p.execution !== binding.execution || p.cohortRef !== binding.cohortRef || p.stop !== stop) return null;
|
||||
if (Date.parse(p.observedAt) > now.getTime()) return null;
|
||||
const digest = hash(without(p, "verificationDigest"));
|
||||
return digest === p.verificationDigest && this.trusted.get(p.id) === digest ? p : null;
|
||||
}
|
||||
|
||||
// check.mjs `effects`.
|
||||
effects(report, ctx) {
|
||||
const p = this.verify(report, "effectReport", ctx);
|
||||
return Boolean(p && p.invocations.every((i) => ["completed", "uncertain", "not-started"].includes(i.disposition) && (i.disposition === "not-started" || i.evidence)));
|
||||
}
|
||||
|
||||
// check.mjs `stopped`, less the stop-record checks the controller makes.
|
||||
cohort(proof, report, ctx) {
|
||||
const p = this.verify(proof, "cohortProof", ctx);
|
||||
if (!p || !p.membershipComplete || p.membershipEpoch !== ctx.epoch) return false;
|
||||
const unique = new Set(p.members.map((m) => `${m.boot}:${m.pid}:${m.startTicks}`)).size === p.members.length;
|
||||
const dead = p.members.every((m) => m.terminatedAt && Date.parse(m.terminatedAt) <= Date.parse(p.observedAt));
|
||||
return unique && dead && this.effects(report, ctx);
|
||||
}
|
||||
}
|
||||
|
||||
// Receipt order (§3 rule 8): admitted < dispatched < acknowledged < working <
|
||||
// finished | failed. dispatch-refused only before dispatched,
|
||||
// delivery-unknown only before working.
|
||||
const ORDER = { admitted: 0, dispatched: 1, acknowledged: 2, working: 3, finished: 4, failed: 4 };
|
||||
export function receiptAllows(from, to) {
|
||||
if (["finished", "failed", "dispatch-refused", "delivery-unknown"].includes(from)) return false;
|
||||
if (to === "dispatch-refused") return from === "admitted";
|
||||
if (to === "delivery-unknown") return ORDER[from] < ORDER.working;
|
||||
return ORDER[to] > ORDER[from];
|
||||
}
|
||||
@@ -26,6 +26,11 @@ export class Refusal extends Error {
|
||||
}
|
||||
}
|
||||
|
||||
// Refusals of the CHAT-03 control layer (controller, claims, guard, engine
|
||||
// pin). They reach a socket client, never the reader's HTTP routes, so the
|
||||
// control board's status map covers only the reader's own codes.
|
||||
export class ControlRefusal extends Refusal {}
|
||||
|
||||
const SESSION_NAME = /^[A-Za-z0-9][A-Za-z0-9._:-]*\.jsonl$/;
|
||||
|
||||
// Permission errors inside a root are one conversation's problem, not the
|
||||
|
||||
@@ -0,0 +1,189 @@
|
||||
#!/usr/bin/env node
|
||||
// The supervisor shim (#1507, CHAT-03 §6). Not a library: `cohort.mjs` starts
|
||||
// it as
|
||||
//
|
||||
// systemd-run --user --scope -p Delegate=yes --unit=<name> --quiet -- \
|
||||
// node shim.mjs --socket <path> -- <engine argv...>
|
||||
//
|
||||
// Inside the delegated scope it moves itself into a `supervisor` child cgroup
|
||||
// and starts the engine in an `engine` child cgroup. A small shell writes its
|
||||
// own pid into engine/cgroup.procs before it execs `unshare -U
|
||||
// --map-current-user --cgroup`, which execs the engine, so the engine is
|
||||
// contained from its first instruction and its cgroup namespace is rooted at
|
||||
// `engine`. The engine inherits the controller's stdin and stdout directly;
|
||||
// the shim closes its own copies. The shim is not a cohort member. It holds
|
||||
// the scope open, so systemd can't collect the cgroup before emptiness is
|
||||
// read, and it outlives its controller so a restarted controller can reach
|
||||
// the cohort again.
|
||||
//
|
||||
// Control is JSON lines over a Unix socket in the controller's 0700 socket
|
||||
// directory. Every request names an op; every answer is {ok, ...} or
|
||||
// {ok:false, unavailable}. Emptiness is a readable engine/cgroup.events with
|
||||
// `populated 0`. A missing or unreadable file is unavailable, never empty.
|
||||
|
||||
import { spawn } from "node:child_process";
|
||||
import { closeSync, mkdirSync, readdirSync, readFileSync, unlinkSync, writeFileSync } from "node:fs";
|
||||
import { createServer } from "node:net";
|
||||
import { join } from "node:path";
|
||||
import { LineSplitter, encodeLine, parseLine } from "./framing.mjs";
|
||||
|
||||
const args = process.argv.slice(2);
|
||||
const sep = args.indexOf("--");
|
||||
const opt = (name) => {
|
||||
const i = args.indexOf(name);
|
||||
return i >= 0 && i < sep ? args[i + 1] : null;
|
||||
};
|
||||
const socketPath = opt("--socket");
|
||||
const engineArgv = sep >= 0 ? args.slice(sep + 1) : [];
|
||||
if (!socketPath || engineArgv.length === 0) {
|
||||
process.stderr.write("shim: usage: shim.mjs --socket <path> -- <engine argv...>\n");
|
||||
process.exit(2);
|
||||
}
|
||||
|
||||
const own = readFileSync("/proc/self/cgroup", "utf8").split("\n").find((l) => l.startsWith("0::"));
|
||||
const scope = join("/sys/fs/cgroup", own.slice(3).trim());
|
||||
const supervisor = join(scope, "supervisor");
|
||||
const engine = join(scope, "engine");
|
||||
const invocationId = process.env.INVOCATION_ID ?? null;
|
||||
const bootId = readFileSync("/proc/sys/kernel/random/boot_id", "utf8").trim();
|
||||
|
||||
mkdirSync(supervisor, { recursive: true });
|
||||
mkdirSync(engine, { recursive: true });
|
||||
writeFileSync(join(supervisor, "cgroup.procs"), String(process.pid));
|
||||
|
||||
const startOf = (pid) => {
|
||||
try {
|
||||
const stat = readFileSync(`/proc/${pid}/stat`, "utf8");
|
||||
const fields = stat.slice(stat.lastIndexOf(")") + 2).split(" ");
|
||||
return Number(fields[19]);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
};
|
||||
|
||||
const child = spawn(
|
||||
"/bin/sh",
|
||||
["-c", 'echo $$ > "$1/cgroup.procs" || exit 97; shift; exec unshare -U --map-current-user --cgroup -- "$@"', "mosaic-engine", engine, ...engineArgv],
|
||||
{ stdio: [0, 1, 2] },
|
||||
);
|
||||
const enginePid = child.pid;
|
||||
let engineExit = null;
|
||||
child.on("exit", (code, signal) => {
|
||||
engineExit = { code, signal, at: new Date().toISOString() };
|
||||
});
|
||||
// The engine holds the controller's pipes; the shim's copies would hide EOF.
|
||||
closeSync(0);
|
||||
closeSync(1);
|
||||
|
||||
function events() {
|
||||
try {
|
||||
const text = readFileSync(join(engine, "cgroup.events"), "utf8");
|
||||
const out = {};
|
||||
for (const line of text.split("\n")) {
|
||||
const [k, v] = line.split(" ");
|
||||
if (k) out[k] = Number(v);
|
||||
}
|
||||
if (out.populated !== 0 && out.populated !== 1) return { ok: false, unavailable: "cgroup.events has no populated field" };
|
||||
return { ok: true, populated: out.populated, frozen: out.frozen ?? null };
|
||||
} catch (err) {
|
||||
return { ok: false, unavailable: `engine/cgroup.events unreadable (${err.code ?? err.message})` };
|
||||
}
|
||||
}
|
||||
|
||||
// Every member of `engine` and its descendants: the engine's namespace can
|
||||
// create child cgroups, so enumeration is recursive.
|
||||
function members() {
|
||||
const out = [];
|
||||
const walk = (dir) => {
|
||||
const procs = readFileSync(join(dir, "cgroup.procs"), "utf8").split("\n").filter(Boolean).map(Number);
|
||||
for (const pid of procs) out.push({ pid, startTicks: startOf(pid), cgroup: dir.slice(scope.length) || "/" });
|
||||
for (const name of readdirSync(dir, { withFileTypes: true })) if (name.isDirectory()) walk(join(dir, name.name));
|
||||
};
|
||||
try {
|
||||
walk(engine);
|
||||
return { ok: true, members: out, boot: bootId };
|
||||
} catch (err) {
|
||||
return { ok: false, unavailable: `engine enumeration failed (${err.code ?? err.message})` };
|
||||
}
|
||||
}
|
||||
|
||||
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
|
||||
|
||||
async function waitFor(pred, ms) {
|
||||
const end = Date.now() + ms;
|
||||
for (;;) {
|
||||
const e = events();
|
||||
if (!e.ok) return e;
|
||||
if (pred(e)) return e;
|
||||
if (Date.now() >= end) return { ...e, timedOut: true };
|
||||
await sleep(10);
|
||||
}
|
||||
}
|
||||
|
||||
async function handle(req) {
|
||||
switch (req.op) {
|
||||
case "hello":
|
||||
return { ok: true, invocationId, scope: scope.slice("/sys/fs/cgroup".length), enginePid, engineStart: startOf(enginePid), shimPid: process.pid, shimStart: startOf(process.pid), boot: bootId, engineExit };
|
||||
case "events":
|
||||
return events();
|
||||
case "members":
|
||||
return members();
|
||||
case "term": {
|
||||
const m = members();
|
||||
if (!m.ok) return m;
|
||||
const signalled = [];
|
||||
for (const { pid, startTicks } of m.members) {
|
||||
if (startOf(pid) !== startTicks) continue;
|
||||
try {
|
||||
process.kill(pid, "SIGTERM");
|
||||
signalled.push(pid);
|
||||
} catch {
|
||||
// gone already
|
||||
}
|
||||
}
|
||||
return { ok: true, signalled };
|
||||
}
|
||||
case "freeze":
|
||||
try {
|
||||
writeFileSync(join(engine, "cgroup.freeze"), "1");
|
||||
} catch (err) {
|
||||
return { ok: false, unavailable: `cgroup.freeze unwritable (${err.code ?? err.message})` };
|
||||
}
|
||||
return waitFor((e) => e.frozen === 1 || e.populated === 0, Number(req.timeoutMs) || 2000);
|
||||
case "kill":
|
||||
try {
|
||||
writeFileSync(join(engine, "cgroup.kill"), "1");
|
||||
} catch (err) {
|
||||
return { ok: false, unavailable: `cgroup.kill unwritable (${err.code ?? err.message})` };
|
||||
}
|
||||
return waitFor((e) => e.populated === 0, Number(req.timeoutMs) || 5000);
|
||||
case "release": {
|
||||
const e = events();
|
||||
if (!e.ok || e.populated !== 0) return { ok: false, unavailable: "the engine cgroup is not empty" };
|
||||
setTimeout(() => {
|
||||
try {
|
||||
unlinkSync(socketPath);
|
||||
} catch {
|
||||
// already gone
|
||||
}
|
||||
process.exit(0);
|
||||
}, 10);
|
||||
return { ok: true };
|
||||
}
|
||||
default:
|
||||
return { ok: false, unavailable: `unknown op ${String(req.op)}` };
|
||||
}
|
||||
}
|
||||
|
||||
const server = createServer((sock) => {
|
||||
const splitter = new LineSplitter(async (line) => {
|
||||
const parsed = parseLine(line);
|
||||
const answer = parsed.error ? { ok: false, unavailable: parsed.error } : await handle(parsed.value);
|
||||
if (!sock.destroyed) sock.write(encodeLine({ id: parsed.value?.id ?? null, ...answer }));
|
||||
}, { maxBytes: 65536 });
|
||||
sock.on("data", (chunk) => splitter.push(chunk));
|
||||
sock.on("error", () => {});
|
||||
});
|
||||
server.listen(socketPath);
|
||||
process.on("SIGTERM", () => {}); // a stray TERM never drops the scope's anchor
|
||||
process.on("SIGHUP", () => {});
|
||||
@@ -0,0 +1,280 @@
|
||||
// The mediated terminal (#1507, CHAT-03 §1, §4, §7, §9).
|
||||
//
|
||||
// node packages/conversation/src/terminal.mjs --socket <path> [--grant <id>]
|
||||
//
|
||||
// A thin view over the client library. It renders the same Transcript the
|
||||
// library offers (E7), so it shows what any other client shows. It holds no
|
||||
// engine-side state.
|
||||
//
|
||||
// The composer is a local buffer (§4). It is empty at start, cleared after
|
||||
// each submit and whenever the controller changes, and it submits only while
|
||||
// this connection is the controller. An observer's submit is refused here,
|
||||
// "not admitted: controller", and sends nothing (S5). A bracketed paste is
|
||||
// inserted literally, newlines included; it never submits by itself.
|
||||
//
|
||||
// Keys: Enter submits; Ctrl-J or Alt-Enter adds a newline; Ctrl-T takes
|
||||
// control; Ctrl-G interrupts; Ctrl-O reconnects if needed and re-reads the
|
||||
// page; PageUp and PageDown scroll; Ctrl-C or Ctrl-D quits.
|
||||
//
|
||||
// Engine text is shown with control characters made visible, so transcript
|
||||
// content can't drive the operator's terminal.
|
||||
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { ConversationClient, OUTCOME_UNKNOWN } from "./client.mjs";
|
||||
import { Transcript } from "./transcript.mjs";
|
||||
|
||||
export const NOT_CONTROLLER = "not admitted: controller";
|
||||
const PASTE_START = "\x1b[200~";
|
||||
const PASTE_END = "\x1b[201~";
|
||||
const KEYS = Object.freeze({ "\r": "submit", "\n": "newline", "\x7f": "backspace", "\b": "backspace", "\x14": "takeover", "\x07": "interrupt", "\x0f": "reload", "\x03": "quit", "\x04": "quit" });
|
||||
|
||||
// Control characters, line and paragraph separators, bidi controls, invisible
|
||||
// characters that can hide or spoof text (zero-width space, word joiner and
|
||||
// invisible operators, BOM, tag characters) shown as text. ZWJ and ZWNJ pass:
|
||||
// emoji sequences and joining scripts need them.
|
||||
export function visible(s) {
|
||||
return String(s).replace(/[\x00-\x08\x0b-\x1f\x7f-\x9f\u061c\u200b\u200e\u200f\u2028\u2029\u202a-\u202e\u2060-\u2064\u2066-\u2069\ufeff\u{e0000}-\u{e007f}]|\t/gu, (c) => {
|
||||
if (c === "\t") return " ";
|
||||
const n = c.codePointAt(0);
|
||||
if (n < 0x20) return "^" + String.fromCharCode(n + 64);
|
||||
if (n === 0x7f) return "^?";
|
||||
return `<U+${n.toString(16).toUpperCase().padStart(4, "0")}>`;
|
||||
});
|
||||
}
|
||||
|
||||
// For lines that must stay one line (head, status, notices): LF shown too.
|
||||
const oneLine = (s) => visible(s).replace(/\n/g, "^J");
|
||||
|
||||
export class Terminal {
|
||||
constructor({ client, write = () => {}, rows = 24, onQuit = () => {} }) {
|
||||
this.client = client;
|
||||
this.write = write;
|
||||
this.rows = rows;
|
||||
this.onQuit = onQuit;
|
||||
this.transcript = new Transcript({ client });
|
||||
this.composer = "";
|
||||
this.inPaste = false;
|
||||
this.carry = "";
|
||||
this.scroll = 0;
|
||||
this.status = "";
|
||||
this.unknownCount = 0;
|
||||
this.dialogs = [];
|
||||
this.notices = [];
|
||||
this.lastReceipt = null;
|
||||
this.controller = client.binding?.controllerConnection ?? null;
|
||||
this.frame = [];
|
||||
this.sent = 0;
|
||||
this.queue = Promise.resolve();
|
||||
client.on((m) => this.#onClient(m));
|
||||
this.transcript.on(() => this.render());
|
||||
}
|
||||
|
||||
#onClient(m) {
|
||||
if (m.type === "welcome" || (m.type === "push" && m.kind === "binding")) {
|
||||
const next = this.client.binding?.controllerConnection ?? null;
|
||||
// §4: the composer clears on every control transfer (mutant 10).
|
||||
if (next !== this.controller) this.composer = "";
|
||||
this.controller = next;
|
||||
} else if (m.type === "push" && m.kind === "unknown") {
|
||||
this.unknownCount = m.count;
|
||||
} else if (m.type === "push" && m.kind === "dialog") {
|
||||
this.dialogs.push(m.dialog);
|
||||
} else if (m.type === "push" && m.kind === "receipt") {
|
||||
if (m.outcomeUnknown) this.notices.push(`prompt ${m.receipt.id}: ${OUTCOME_UNKNOWN}`);
|
||||
if (m.receipt.id === this.lastReceipt?.id) this.lastReceipt = m.receipt;
|
||||
} else if (m.type === "outcome-unknown") {
|
||||
this.notices.push(`${m.operation}: ${OUTCOME_UNKNOWN}`);
|
||||
} else if (m.type === "status" && m.status === "disconnected") {
|
||||
this.status = "disconnected; Ctrl-O reconnects";
|
||||
}
|
||||
this.render();
|
||||
}
|
||||
|
||||
// Feeds raw terminal input. Resolves when the actions it started finish.
|
||||
key(data) {
|
||||
const actions = [];
|
||||
let s = this.carry + data;
|
||||
this.carry = "";
|
||||
let i = 0;
|
||||
while (i < s.length) {
|
||||
if (this.inPaste) {
|
||||
const end = s.indexOf(PASTE_END, i);
|
||||
if (end === -1) {
|
||||
const keep = partialSuffix(s.slice(i), PASTE_END);
|
||||
this.composer += s.slice(i, s.length - keep);
|
||||
this.carry = s.slice(s.length - keep);
|
||||
break;
|
||||
}
|
||||
this.composer += s.slice(i, end);
|
||||
this.inPaste = false;
|
||||
i = end + PASTE_END.length;
|
||||
continue;
|
||||
}
|
||||
if (s.startsWith(PASTE_START, i)) {
|
||||
this.inPaste = true;
|
||||
i += PASTE_START.length;
|
||||
continue;
|
||||
}
|
||||
if (s[i] === "\x1b") {
|
||||
// A trailing ESC, alone or with more of the paste-start marker, waits
|
||||
// for the next chunk. A lone Escape has no action here, so holding it
|
||||
// costs nothing.
|
||||
if (s.length - i < PASTE_START.length && PASTE_START.startsWith(s.slice(i))) {
|
||||
this.carry = s.slice(i);
|
||||
break;
|
||||
}
|
||||
const seq = /^\x1b(?:\[[0-9;?]*[ -/]*[@-~]|O.|[\s\S])?/.exec(s.slice(i))[0];
|
||||
if (seq === "\x1b[5~") this.scroll += Math.max(1, this.rows - 4);
|
||||
else if (seq === "\x1b[6~") this.scroll = Math.max(0, this.scroll - Math.max(1, this.rows - 4));
|
||||
else if (seq === "\x1b\r") this.composer += "\n";
|
||||
i += seq.length;
|
||||
continue;
|
||||
}
|
||||
const action = KEYS[s[i]];
|
||||
if (action === "newline") this.composer += "\n";
|
||||
else if (action === "backspace") this.composer = Array.from(this.composer).slice(0, -1).join("");
|
||||
else if (action === "submit") {
|
||||
// The composer is taken at the Enter, so text after it in the same
|
||||
// chunk starts the next message instead of joining this one.
|
||||
const text = this.#take();
|
||||
if (text !== null) actions.push(() => this.#send(text));
|
||||
} else if (action) actions.push(() => this.#act(action));
|
||||
else if (s[i] >= " ") this.composer += s[i];
|
||||
i += 1;
|
||||
}
|
||||
this.render();
|
||||
for (const run of actions) this.queue = this.queue.then(run);
|
||||
return this.queue;
|
||||
}
|
||||
|
||||
async #act(action) {
|
||||
if (action === "takeover") return this.#show("takeover", await this.client.takeover());
|
||||
if (action === "interrupt") return this.#show("interrupt", await this.client.interrupt());
|
||||
if (action === "reload") {
|
||||
if (this.client.closed) {
|
||||
try {
|
||||
await this.client.connect();
|
||||
} catch (err) {
|
||||
this.status = `connect: ${err.refusal ?? err.message}`;
|
||||
}
|
||||
} else await this.transcript.reload("manual");
|
||||
return this.render();
|
||||
}
|
||||
if (action === "quit") return this.onQuit();
|
||||
}
|
||||
|
||||
// Sends the composer as one prompt, only as the controller. A submit clears
|
||||
// the composer whatever its outcome; an observer's Enter is refused before
|
||||
// that and leaves the buffer for the operator.
|
||||
async submit() {
|
||||
const text = this.#take();
|
||||
if (text === null) {
|
||||
this.render();
|
||||
return { sent: false, refusal: "controller" };
|
||||
}
|
||||
return this.#send(text);
|
||||
}
|
||||
|
||||
// Empties the composer and returns its text, or refuses an observer and
|
||||
// returns null, leaving the buffer.
|
||||
#take() {
|
||||
if (!this.client.isController) {
|
||||
this.status = NOT_CONTROLLER;
|
||||
return null;
|
||||
}
|
||||
const text = this.composer;
|
||||
this.composer = "";
|
||||
return text;
|
||||
}
|
||||
|
||||
async #send(text) {
|
||||
if (!text) return { sent: false, refusal: null };
|
||||
this.sent += 1;
|
||||
const r = await this.client.prompt(text);
|
||||
if (r.receipt) this.lastReceipt = r.receipt;
|
||||
this.#show("prompt", r);
|
||||
return { sent: true, reply: r };
|
||||
}
|
||||
|
||||
#show(op, r) {
|
||||
if (r.outcome === "outcome-unknown") this.status = `${op}: ${OUTCOME_UNKNOWN}`;
|
||||
else if (r.refusal) this.status = `not admitted: ${r.refusal}`;
|
||||
else this.status = `${op}: ${r.outcome}`;
|
||||
this.render();
|
||||
}
|
||||
|
||||
// The frame: a header, the transcript window, dialogs, notices, the status
|
||||
// line and the composer.
|
||||
render() {
|
||||
const b = this.client.binding;
|
||||
const head = `${oneLine(b?.state ?? "disconnected")} | ${this.client.isController ? "controller" : "observer"} | unknown events: ${this.unknownCount}`;
|
||||
const body = [];
|
||||
for (const line of this.transcript.lines()) body.push(...visible(line.replace(/\r\n/g, "\n")).split("\n"));
|
||||
for (const d of this.dialogs) body.push(oneLine(`[dialog ${d.method} disabled: ${d.reason}]`));
|
||||
for (const n of this.notices) body.push(oneLine(n));
|
||||
const receipt = this.lastReceipt ? `last prompt: ${this.lastReceipt.state}${this.lastReceipt.reasonCode ? ` (${this.lastReceipt.reasonCode})` : ""}` : "";
|
||||
const status = oneLine([this.status, receipt].filter(Boolean).join(" | "));
|
||||
const composer = (this.composer ? this.composer.split("\n") : [""]).map((l, i) => (i ? " " : "> ") + visible(l));
|
||||
const room = Math.max(1, this.rows - 2 - composer.length);
|
||||
const end = Math.max(0, body.length - this.scroll);
|
||||
this.frame = [head, ...body.slice(Math.max(0, end - room), end), status, ...composer];
|
||||
this.write("\x1b[H\x1b[2J" + this.frame.join("\r\n"));
|
||||
return this.frame;
|
||||
}
|
||||
}
|
||||
|
||||
// How many trailing characters of `s` begin `marker`.
|
||||
function partialSuffix(s, marker) {
|
||||
for (let n = Math.min(s.length, marker.length - 1); n > 0; n--) if (marker.startsWith(s.slice(-n))) return n;
|
||||
return 0;
|
||||
}
|
||||
|
||||
function parseArgs(argv) {
|
||||
const out = { socket: null, grant: "grant-local" };
|
||||
for (let i = 0; i < argv.length; i++) {
|
||||
if (argv[i] === "--socket") out.socket = argv[++i];
|
||||
else if (argv[i] === "--grant") out.grant = argv[++i];
|
||||
else throw new Error(`unknown argument: ${argv[i]}`);
|
||||
}
|
||||
if (!out.socket) throw new Error("usage: terminal.mjs --socket <path> [--grant <id>]");
|
||||
return out;
|
||||
}
|
||||
|
||||
async function main() {
|
||||
let args;
|
||||
try {
|
||||
args = parseArgs(process.argv.slice(2));
|
||||
} catch (err) {
|
||||
process.stderr.write(err.message + "\n");
|
||||
process.exit(2);
|
||||
}
|
||||
const { stdin, stdout } = process;
|
||||
const client = new ConversationClient({ socketPath: args.socket, grant: args.grant });
|
||||
const restore = () => {
|
||||
stdout.write("\x1b[?2004l\r\n");
|
||||
if (stdin.isTTY) stdin.setRawMode(false);
|
||||
};
|
||||
const quit = () => {
|
||||
restore();
|
||||
client.close();
|
||||
process.exit(0);
|
||||
};
|
||||
const term = new Terminal({ client, write: (s) => stdout.write(s), rows: stdout.rows || 24, onQuit: quit });
|
||||
try {
|
||||
await client.connect();
|
||||
} catch (err) {
|
||||
process.stderr.write(`could not connect: ${err.refusal ?? err.message}\n`);
|
||||
process.exit(1);
|
||||
}
|
||||
if (stdin.isTTY) stdin.setRawMode(true);
|
||||
stdout.write("\x1b[?2004h");
|
||||
stdout.on("resize", () => {
|
||||
term.rows = stdout.rows || 24;
|
||||
term.render();
|
||||
});
|
||||
stdin.on("data", (c) => void term.key(c.toString("utf8")));
|
||||
stdin.on("end", quit);
|
||||
term.render();
|
||||
}
|
||||
|
||||
if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) main();
|
||||
@@ -0,0 +1,41 @@
|
||||
// The text policy for mediated prompts (#1507, CHAT-03 §4).
|
||||
//
|
||||
// CHAT-01 lines 181–183 refuse `text-policy` when the first non-whitespace
|
||||
// character is `/`. CHAT-01 lines 368–370 leave `!`, `@` and slashes on later
|
||||
// lines open; CHAT-03 settles them from the pinned Pi 0.85.1 source. The RPC
|
||||
// `prompt` command calls AgentSession.prompt(message) (rpc-mode.js 298–318),
|
||||
// and that path interprets text only through three checks, all on the exact
|
||||
// first character:
|
||||
//
|
||||
// - extension commands: `text.startsWith("/")` (agent-session.js 828);
|
||||
// - skills: `text.startsWith("/skill:")` (agent-session.js 984);
|
||||
// - prompt templates: `text.startsWith("/")` (prompt-templates.js 222).
|
||||
//
|
||||
// The bundle `pi` runs (dist/bundle/chunks/chunk-JVUZSMYM.js) carries the same
|
||||
// three checks. `!` and `!!` are shell shortcuts of the interactive editor only
|
||||
// (interactive-mode.js 2502), and `@` is a file argument of the command line
|
||||
// only (cli/args.js 214); neither is on the RPC prompt path. Pi never trims
|
||||
// before its checks, so a slash after leading whitespace or on a later line is
|
||||
// plain text to Pi. CHAT-01's rule still refuses the leading-whitespace case.
|
||||
//
|
||||
// Every prefix is therefore classified: a first non-whitespace `/` is refused,
|
||||
// and everything else is text. PREFIXES is the list the tests read (S2).
|
||||
|
||||
export const TEXT_POLICY = "text-policy";
|
||||
|
||||
export const PREFIXES = Object.freeze([
|
||||
Object.freeze({ prefix: "/", interpreted: true, example: "/goal x", basis: "agent-session.js 828: extension commands run immediately" }),
|
||||
Object.freeze({ prefix: "/skill:", interpreted: true, example: "/skill:ms-unslop", basis: "agent-session.js 984: skills expand" }),
|
||||
Object.freeze({ prefix: "/<template>", interpreted: true, example: "/review this", basis: "prompt-templates.js 222: templates expand" }),
|
||||
Object.freeze({ prefix: "!", interpreted: false, example: "!ls", basis: "interactive-mode.js 2502 only; not on the RPC prompt path" }),
|
||||
Object.freeze({ prefix: "!!", interpreted: false, example: "!!ls", basis: "interactive-mode.js 2503 only; not on the RPC prompt path" }),
|
||||
Object.freeze({ prefix: "@", interpreted: false, example: "@README.md", basis: "cli/args.js 214 only; not on the RPC prompt path" }),
|
||||
]);
|
||||
|
||||
// Later-line slashes are not interpreted: every check is at index 0 (S3).
|
||||
export const LATER_LINE_SLASH_INTERPRETED = false;
|
||||
|
||||
export function textPolicy(text) {
|
||||
if (typeof text !== "string") return TEXT_POLICY;
|
||||
return text.trimStart().startsWith("/") ? TEXT_POLICY : null;
|
||||
}
|
||||
@@ -0,0 +1,304 @@
|
||||
// The transcript view shared by the library and the terminal (#1507, CHAT-03
|
||||
// §9, E1–E7).
|
||||
//
|
||||
// It joins a history page with the live event stream. Replay is unavailable
|
||||
// in CHAT-03, and page entries and events carry different IDs, so overlap at
|
||||
// the seam can't be deduplicated (CHAT-01 lines 104–110). The rule is
|
||||
// reconciliation instead:
|
||||
//
|
||||
// - The controller's seam says where live events start (`fromSequence`) and
|
||||
// whether the cut was quiet: no run visible and no prompt in the slot, so
|
||||
// every earlier message had settled and was persisted.
|
||||
// - A cut that wasn't quiet puts a reconcile marker at the seam. Near it a
|
||||
// message may show twice or be missing. The view re-reads the page after
|
||||
// the next `run-settled` and clears the marker once a read is quiet.
|
||||
// - A sequence gap, a new stream epoch or a repeated event ID with different
|
||||
// bytes also marks the seam and re-reads at once. A gap is never
|
||||
// concatenated across.
|
||||
// - A repeated event ID with identical bytes is dropped (E3).
|
||||
//
|
||||
// Tool progress (`tool-start`, `tool-update`, `tool-end`) shows a call while
|
||||
// it runs. Once a finished tool message carries that call's result, from the
|
||||
// stream or the page, the progress item is hidden, so each message appears
|
||||
// once.
|
||||
|
||||
export const RECONCILE_MARKER = "-- reconciling: messages near here may repeat or be missing until the run settles --";
|
||||
|
||||
const MESSAGE_EVENTS = new Set(["message-start", "text-delta", "thinking-delta", "message-end"]);
|
||||
const TOOL_EVENTS = new Set(["tool-start", "tool-update", "tool-end"]);
|
||||
|
||||
export class Transcript {
|
||||
// With a client, the view follows it: it reads the page on every welcome
|
||||
// and re-reads when the stream asks for it.
|
||||
constructor({ client = null } = {}) {
|
||||
this.client = client;
|
||||
this.entries = [];
|
||||
this.events = [];
|
||||
this.held = [];
|
||||
this.seen = new Map();
|
||||
this.seam = null;
|
||||
this.expected = null;
|
||||
this.reconcile = false;
|
||||
this.reason = null;
|
||||
this.refusal = null;
|
||||
this.loads = 0;
|
||||
this.loading = null;
|
||||
this.again = false;
|
||||
this.listeners = new Set();
|
||||
if (client) client.on((m) => this.#onClient(m));
|
||||
// Attached to a client that is already connected: no welcome will come,
|
||||
// so read the page now.
|
||||
if (client && !client.closed && client.connection) void this.reload("attach");
|
||||
}
|
||||
|
||||
on(fn) {
|
||||
this.listeners.add(fn);
|
||||
return () => this.listeners.delete(fn);
|
||||
}
|
||||
|
||||
#changed() {
|
||||
for (const fn of this.listeners) fn();
|
||||
}
|
||||
|
||||
#onClient(m) {
|
||||
if (m.type === "welcome") {
|
||||
void this.reload("connect");
|
||||
return;
|
||||
}
|
||||
if (m.type !== "push" || m.kind !== "event") return;
|
||||
const why = this.event(m.event);
|
||||
if (why) void this.reload(why);
|
||||
this.#changed();
|
||||
}
|
||||
|
||||
// Installs a full read of the page (every entry, in order) and its seam.
|
||||
load(entries, seam) {
|
||||
this.entries = entries.slice();
|
||||
this.seam = { ...seam };
|
||||
this.reconcile = seam.quiet !== true;
|
||||
this.reason = this.reconcile ? "cut" : null;
|
||||
this.loads += 1;
|
||||
const kept = [...this.events, ...this.held].filter((e) => e.streamEpoch === seam.streamEpoch && e.sequence >= seam.fromSequence).sort((a, b) => a.sequence - b.sequence);
|
||||
this.events = [];
|
||||
this.held = [];
|
||||
this.expected = seam.fromSequence;
|
||||
for (const e of kept) {
|
||||
if (e.sequence < this.expected) continue;
|
||||
if (e.sequence > this.expected) {
|
||||
this.#mark("gap");
|
||||
this.held = kept.filter((x) => x.sequence > this.expected);
|
||||
// A run settled past the gap: the next read starts after it.
|
||||
if (this.held.some((x) => x.type === "run-settled")) this.again = true;
|
||||
break;
|
||||
}
|
||||
this.events.push(e);
|
||||
this.expected = e.sequence + 1;
|
||||
}
|
||||
}
|
||||
|
||||
// Takes one live event. Returns why the page must be re-read, or null.
|
||||
event(e) {
|
||||
const bytes = JSON.stringify(e);
|
||||
const prior = this.seen.get(e.id);
|
||||
if (prior !== undefined) return prior === bytes ? null : this.#mark("conflict");
|
||||
this.seen.set(e.id, bytes);
|
||||
// Events past a gap or in another epoch are held, not shown, until the
|
||||
// next read; the read keeps those after its seam.
|
||||
if (!this.seam) {
|
||||
this.held.push(e);
|
||||
return null;
|
||||
}
|
||||
// Once something is held a read is already due; a later held event asks
|
||||
// again only when its run settles, when the page has caught up.
|
||||
if (this.held.length) {
|
||||
this.held.push(e);
|
||||
return e.type === "run-settled" ? "settled" : null;
|
||||
}
|
||||
if (e.streamEpoch !== this.seam.streamEpoch) {
|
||||
this.held.push(e);
|
||||
return this.#mark("epoch");
|
||||
}
|
||||
if (e.sequence < this.expected) return null;
|
||||
if (e.sequence > this.expected) {
|
||||
this.held.push(e);
|
||||
return this.#mark("gap");
|
||||
}
|
||||
this.events.push(e);
|
||||
this.expected = e.sequence + 1;
|
||||
return e.type === "run-settled" && this.reconcile ? "settled" : null;
|
||||
}
|
||||
|
||||
#mark(why) {
|
||||
this.reconcile = true;
|
||||
this.reason = why;
|
||||
return why;
|
||||
}
|
||||
|
||||
// Reads every page and installs it. Overlapping calls coalesce; a call made
|
||||
// while a read is in flight runs one more read after it.
|
||||
reload(why = "manual") {
|
||||
if (!this.client) return Promise.resolve();
|
||||
if (this.loading) {
|
||||
this.again = true;
|
||||
return this.loading;
|
||||
}
|
||||
this.loading = (async () => {
|
||||
do {
|
||||
this.again = false;
|
||||
const r = await this.#readAll();
|
||||
if (!r.ok) {
|
||||
this.refusal = r.refusal;
|
||||
break;
|
||||
}
|
||||
this.refusal = null;
|
||||
this.load(r.entries, r.seam);
|
||||
} while (this.again);
|
||||
})().finally(() => {
|
||||
this.loading = null;
|
||||
this.#changed();
|
||||
});
|
||||
this.lastReload = why;
|
||||
return this.loading;
|
||||
}
|
||||
|
||||
async #readAll() {
|
||||
let r = await this.client.observe();
|
||||
if (r.outcome !== "observing") return { ok: false, refusal: r.refusal ?? r.outcome };
|
||||
const seam = r.data.seam;
|
||||
const entries = [...r.data.page.entries];
|
||||
while (r.data.page.hasMore) {
|
||||
r = await this.client.observe({ cursor: r.data.page.nextCursor });
|
||||
if (r.outcome !== "observing") return { ok: false, refusal: r.refusal ?? r.outcome };
|
||||
entries.push(...r.data.page.entries);
|
||||
}
|
||||
return { ok: true, entries, seam };
|
||||
}
|
||||
|
||||
// The view: page messages, the marker when the seam is unreconciled, then
|
||||
// live messages. Each item is {key, role, source, final, text} or
|
||||
// {marker: true, text, reason}.
|
||||
messages() {
|
||||
const items = [];
|
||||
for (const en of this.entries) {
|
||||
const last = items.at(-1);
|
||||
const it = last && last.key === en.message ? last : null;
|
||||
if (it) it.parts.set(en.part, en.content);
|
||||
else items.push({ key: en.message, role: en.role, source: "page", parts: new Map([[en.part, en.content]]), lastPart: null, exec: false });
|
||||
if (en.lastPart) (it ?? items.at(-1)).lastPart = en.part;
|
||||
}
|
||||
const pageCount = items.length;
|
||||
const live = new Map();
|
||||
const epoch = this.seam?.streamEpoch;
|
||||
for (const e of this.events) {
|
||||
if (!this.seam || e.streamEpoch !== epoch) continue;
|
||||
if (!MESSAGE_EVENTS.has(e.type) && !TOOL_EVENTS.has(e.type)) continue;
|
||||
let it = live.get(e.message);
|
||||
if (!it) {
|
||||
it = { key: e.message, role: e.role, source: "live", parts: new Map(), lastPart: null, stream: new Map(), exec: TOOL_EVENTS.has(e.type), status: null, call: null };
|
||||
live.set(e.message, it);
|
||||
items.push(it);
|
||||
}
|
||||
if (e.type === "text-delta" || e.type === "thinking-delta") {
|
||||
for (const b of e.content) {
|
||||
const s = it.stream.get(b.block) ?? { type: b.type, text: "", visibility: b.visibility };
|
||||
s.text += b.text ?? "";
|
||||
it.stream.set(b.block, s);
|
||||
}
|
||||
} else if (e.type === "message-end") {
|
||||
it.role = e.role;
|
||||
it.parts.set(e.part, e.content);
|
||||
if (e.lastPart) it.lastPart = e.part;
|
||||
} else if (TOOL_EVENTS.has(e.type)) {
|
||||
it.call = e.content[0]?.call ?? it.call;
|
||||
if (e.type === "tool-start") {
|
||||
it.status = "running";
|
||||
it.name = e.content[0]?.name ?? null;
|
||||
} else {
|
||||
it.status = e.type === "tool-end" ? "done" : "running";
|
||||
it.result = e.content;
|
||||
}
|
||||
}
|
||||
}
|
||||
const finals = new Set();
|
||||
for (const it of items) {
|
||||
if (!isFinal(it) || it.exec) continue;
|
||||
for (const b of blocksOf(it)) if (b.type === "tool-result") finals.add(b.call);
|
||||
}
|
||||
const out = [];
|
||||
items.forEach((it, i) => {
|
||||
if (i === pageCount && this.reconcile) out.push(this.#marker());
|
||||
if (it.exec && finals.has(it.call)) return;
|
||||
out.push({ key: it.key, role: it.role, source: it.source, final: isFinal(it), text: textOf(it) });
|
||||
});
|
||||
if (items.length === pageCount && this.reconcile) out.push(this.#marker());
|
||||
return out;
|
||||
}
|
||||
|
||||
#marker() {
|
||||
return { marker: true, text: RECONCILE_MARKER, reason: this.reason };
|
||||
}
|
||||
|
||||
// One line per message (a message may span several lines of text).
|
||||
lines() {
|
||||
return this.messages().map((m) => (m.marker ? m.text : `${m.role}: ${m.text}`));
|
||||
}
|
||||
}
|
||||
|
||||
function isFinal(it) {
|
||||
if (it.lastPart === null) return false;
|
||||
for (let p = 0; p <= it.lastPart; p++) if (!it.parts.has(p)) return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
function blocksOf(it) {
|
||||
const out = [];
|
||||
const n = it.lastPart ?? Math.max(-1, ...it.parts.keys());
|
||||
for (let p = 0; p <= n; p++) out.push(...(it.parts.get(p) ?? []));
|
||||
return out;
|
||||
}
|
||||
|
||||
// Joins fragments by block and renders each block.
|
||||
export function renderBlocks(blocks) {
|
||||
const joined = [];
|
||||
for (const b of blocks) {
|
||||
const last = joined.at(-1);
|
||||
if (last && last.block === b.block && last.type === b.type && (b.fragment ?? 0) > 0 && (last.call ?? null) === (b.call ?? null)) {
|
||||
last.text = (last.text ?? "") + (b.text ?? "");
|
||||
last.argumentsText = (last.argumentsText ?? "") + (b.argumentsText ?? "");
|
||||
last.summary = (last.summary ?? "") + (b.summary ?? "");
|
||||
} else joined.push({ ...b });
|
||||
}
|
||||
return joined.map(renderBlock).join("\n");
|
||||
}
|
||||
|
||||
function renderBlock(b) {
|
||||
switch (b.type) {
|
||||
case "text":
|
||||
return b.text ?? "";
|
||||
case "thinking":
|
||||
return b.visibility === "unavailable" ? "(thinking not shown)" : `(thinking) ${b.text ?? ""}`;
|
||||
case "tool-call":
|
||||
return `[tool-call ${b.name} ${b.argumentsText ?? ""}]`;
|
||||
case "tool-result":
|
||||
return `[tool-result${b.isError ? " error" : ""}] ${b.text ?? ""}`;
|
||||
case "compaction":
|
||||
return `[compaction] ${b.summary ?? ""}`;
|
||||
case "attachment":
|
||||
return `[attachment ${b.attachment}]`;
|
||||
default:
|
||||
return `[${b.type}]`;
|
||||
}
|
||||
}
|
||||
|
||||
function textOf(it) {
|
||||
if (isFinal(it)) {
|
||||
const blocks = blocksOf(it);
|
||||
return blocks.length ? renderBlocks(blocks) : "(empty)";
|
||||
}
|
||||
if (it.exec) {
|
||||
const head = `[tool ${it.name ?? it.call} ${it.status}]`;
|
||||
return it.result?.length ? `${head} ${renderBlocks(it.result)}` : head;
|
||||
}
|
||||
const parts = [...it.stream.entries()].sort(([a], [b]) => a - b).map(([, s]) => renderBlock(s));
|
||||
return parts.length ? parts.join("\n") : "…";
|
||||
}
|
||||
@@ -0,0 +1,225 @@
|
||||
// Run tracking, overlap signals and receipt settlement (#1507, CHAT-03 §3).
|
||||
//
|
||||
// Pinned Pi ties no event to a prompt. Under the seal (pi-pin.mjs) the Mosaic
|
||||
// prompt in the slot is the only thing that can start a run, so a run that
|
||||
// follows the slot's ack is attributed to it by order. The tracker watches
|
||||
// for signs that the seal failed (O1–O6) and for lead decision 34's
|
||||
// `aborted` with no stop in progress, and settles the slot's receipt only on
|
||||
// evidence tied to its own prompt (§3 rule 4).
|
||||
//
|
||||
// Every line read from the engine gets an observation index (`read()`), so
|
||||
// "before the ack", "since the fence" and "before the last abort" are exact
|
||||
// comparisons on one ordered observation.
|
||||
//
|
||||
// A run is one prompt's run, from its first `agent_start` to its
|
||||
// `agent_settled`; retries and compaction can add further `agent_start` …
|
||||
// `agent_end` pairs inside it (agent-session.js 787–810).
|
||||
|
||||
export const DONE_STOPS = Object.freeze(new Set(["stop", "length", "toolUse"]));
|
||||
|
||||
export class Tracker {
|
||||
// `onOverlap(signal, detail)`; `onWorking(slot)`; `onSlotSettled(slot,
|
||||
// result)` with result { state, reason, stop? } or { outcomeUnknown };
|
||||
// `stopLink()` returns the stop ID an `aborted` may be linked to, or null
|
||||
// (lead decision 34).
|
||||
constructor({ onOverlap, onWorking, onSlotSettled, stopLink }) {
|
||||
this.onOverlap = onOverlap;
|
||||
this.onWorking = onWorking;
|
||||
this.onSlotSettled = onSlotSettled;
|
||||
this.stopLink = stopLink;
|
||||
this.idx = 0;
|
||||
this.slot = null;
|
||||
this.lastSlot = null;
|
||||
this.runs = [];
|
||||
this.current = null;
|
||||
this.overlaps = [];
|
||||
this.gap = null;
|
||||
this.clearWindow = 0;
|
||||
this.runSerial = 0;
|
||||
this.waiters = new Set();
|
||||
this.ignoredAgentEnds = 0;
|
||||
this.looseSettles = [];
|
||||
}
|
||||
|
||||
get overlapped() {
|
||||
return this.overlaps.length > 0;
|
||||
}
|
||||
|
||||
read() {
|
||||
return ++this.idx;
|
||||
}
|
||||
|
||||
overlap(signal, idx, detail = {}) {
|
||||
const entry = { signal, idx, request: (this.slot ?? this.lastSlot)?.request ?? null, ...detail };
|
||||
this.overlaps.push(entry);
|
||||
this.onOverlap(signal, entry);
|
||||
}
|
||||
|
||||
attach(slot) {
|
||||
this.slot = slot;
|
||||
Object.assign(slot, { ackIdx: null, runId: null, working: false, settleSeen: false, failureNoRun: false });
|
||||
}
|
||||
|
||||
acked(idx) {
|
||||
if (this.slot) this.slot.ackIdx = idx;
|
||||
}
|
||||
|
||||
release() {
|
||||
if (this.slot) this.lastSlot = this.slot;
|
||||
this.slot = null;
|
||||
}
|
||||
|
||||
// Runs whose first agent_start was read after `idx`.
|
||||
startedAfter(idx) {
|
||||
return this.runs.filter((r) => r.startIdx > idx);
|
||||
}
|
||||
|
||||
settledAfter(idx) {
|
||||
return this.runs.some((r) => r.settleIdx !== null && r.settleIdx > idx) || this.looseSettles.some((i) => i > idx);
|
||||
}
|
||||
|
||||
waitSettle(run, ms) {
|
||||
if (run.settleIdx !== null) return Promise.resolve(true);
|
||||
return new Promise((resolve) => {
|
||||
const w = () => {
|
||||
if (run.settleIdx === null) return;
|
||||
this.waiters.delete(w);
|
||||
clearTimeout(t);
|
||||
resolve(true);
|
||||
};
|
||||
const t = setTimeout(() => {
|
||||
this.waiters.delete(w);
|
||||
resolve(false);
|
||||
}, ms);
|
||||
this.waiters.add(w);
|
||||
});
|
||||
}
|
||||
|
||||
// A run in the slot's window, or null.
|
||||
slotRun() {
|
||||
return this.slot?.runId ? this.runs.find((r) => r.id === this.slot.runId) ?? null : null;
|
||||
}
|
||||
|
||||
event(ev, idx) {
|
||||
switch (ev.type) {
|
||||
case "agent_start":
|
||||
return this.#agentStart(idx);
|
||||
case "agent_end":
|
||||
if (this.current && this.current.openStarts > 0) this.current.openStarts -= 1;
|
||||
else this.ignoredAgentEnds += 1;
|
||||
return undefined;
|
||||
case "message_start":
|
||||
if (ev.message?.role === "user") return this.#userStart(idx);
|
||||
return undefined;
|
||||
case "message_end":
|
||||
if (ev.message?.role === "assistant") return this.#assistantEnd(ev.message, idx);
|
||||
return undefined;
|
||||
case "agent_settled":
|
||||
return this.#settled(idx);
|
||||
case "queue_update": {
|
||||
const empty = (Array.isArray(ev.steering) ? ev.steering.length : 1) === 0 && (Array.isArray(ev.followUp) ? ev.followUp.length : 1) === 0;
|
||||
// Pi's own clear emits an empty queue_update before its response
|
||||
// (agent-session.js 1201, rpc-mode.js 334): that one is caused (N25).
|
||||
if (empty && this.clearWindow > 0) return undefined;
|
||||
return this.overlap("O5", idx, { cause: empty ? "uncaused-queue-update" : "queue-update", items: queueDigestCount(ev) });
|
||||
}
|
||||
default:
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
#agentStart(idx) {
|
||||
if (this.current) {
|
||||
this.current.openStarts += 1;
|
||||
return;
|
||||
}
|
||||
const run = { id: `run-${++this.runSerial}`, startIdx: idx, settleIdx: null, openStarts: 1, userStarts: 0, lastStop: null, lastError: null, failureBeforeUser: false, slot: false };
|
||||
this.runs.push(run);
|
||||
this.current = run;
|
||||
const s = this.slot;
|
||||
if (!s || s.ackIdx === null || s.runId !== null) {
|
||||
this.overlap("O1", idx, { run: run.id, why: !s ? "no slot held" : s.ackIdx === null ? "before the slot's ack" : "after the slot's run" });
|
||||
return;
|
||||
}
|
||||
s.runId = run.id;
|
||||
run.slot = true;
|
||||
}
|
||||
|
||||
#userStart(idx) {
|
||||
const run = this.current;
|
||||
if (!run) return;
|
||||
run.userStarts += 1;
|
||||
if (run.userStarts === 2) this.overlap("O6", idx, { run: run.id });
|
||||
const s = this.slot;
|
||||
if (s && run.id === s.runId && run.userStarts === 1 && !s.working && !this.overlapped && !this.gap) {
|
||||
s.working = true;
|
||||
this.onWorking(s);
|
||||
}
|
||||
}
|
||||
|
||||
#assistantEnd(message, idx) {
|
||||
const stopReason = typeof message.stopReason === "string" ? message.stopReason : null;
|
||||
// Lead decision 34: only the controller's abort produces `aborted`.
|
||||
if (stopReason === "aborted" && !this.stopLink()) this.overlap("aborted-without-stop", idx, { run: this.current?.id ?? null });
|
||||
const run = this.current;
|
||||
if (!run) {
|
||||
// Pi's run-failure handler can emit a failure message with no
|
||||
// agent_start (pi-agent-core agent.js 349–364).
|
||||
if (this.slot && this.slot.ackIdx !== null && (stopReason === "error" || stopReason === "aborted")) this.slot.failureNoRun = true;
|
||||
return;
|
||||
}
|
||||
run.lastStop = stopReason;
|
||||
run.lastError = typeof message.errorMessage === "string" ? message.errorMessage.slice(0, 2000) : null;
|
||||
if ((stopReason === "error" || stopReason === "aborted") && run.userStarts === 0) run.failureBeforeUser = true;
|
||||
}
|
||||
|
||||
#settled(idx) {
|
||||
const run = this.current;
|
||||
const s = this.slot;
|
||||
if (!run) {
|
||||
this.looseSettles.push(idx);
|
||||
if (!s || s.ackIdx === null) return this.overlap("O2", idx, { why: !s ? "no slot held" : "before the slot's ack" });
|
||||
if (s.settleSeen) return this.overlap("O2", idx, { why: "a second settle for one ack" });
|
||||
s.settleSeen = true;
|
||||
if (s.runId !== null) return this.overlap("O2", idx, { why: "a second settle for one ack" });
|
||||
if (!s.failureNoRun) return this.overlap("O4", idx, { why: "settle after the ack with no agent_start and no failure message" });
|
||||
if (this.overlapped || this.gap) return undefined;
|
||||
return this.onSlotSettled(s, { state: "delivery-unknown", reason: "ack-without-start" });
|
||||
}
|
||||
if (run.openStarts > 0) {
|
||||
// O3: the settle does not close this run; the run stays current.
|
||||
return this.overlap("O3", idx, { run: run.id });
|
||||
}
|
||||
run.settleIdx = idx;
|
||||
this.current = null;
|
||||
for (const w of [...this.waiters]) w();
|
||||
if (!run.slot) {
|
||||
if (!s || s.ackIdx === null) this.overlap("O2", idx, { run: run.id, why: "settle of a run that is not the slot's" });
|
||||
return undefined;
|
||||
}
|
||||
if (!s || s.runId !== run.id) return this.overlap("O2", idx, { run: run.id, why: "the slot's run settled after the slot moved on" });
|
||||
if (s.settleSeen) return this.overlap("O2", idx, { why: "a second settle for one ack" });
|
||||
s.settleSeen = true;
|
||||
if (this.overlapped || this.gap) return undefined;
|
||||
return this.onSlotSettled(s, classify(s, run, this.stopLink));
|
||||
}
|
||||
}
|
||||
|
||||
// §3 rule 4, for a settle that closes the slot's own run with no overlap.
|
||||
export function classify(slot, run, stopLink) {
|
||||
if (!slot.working) return { state: "delivery-unknown", reason: "ack-without-start" };
|
||||
const s = run.lastStop;
|
||||
if (DONE_STOPS.has(s)) return { state: "finished", reason: null };
|
||||
if (s === "error") return { state: "failed", reason: null, nativeError: run.lastError };
|
||||
if (s === "aborted") {
|
||||
const stop = stopLink();
|
||||
// Unreachable without a stop: #assistantEnd raised the overlap already.
|
||||
return stop ? { state: "failed", reason: "interrupted", stop } : { outcomeUnknown: "aborted-without-stop" };
|
||||
}
|
||||
return { outcomeUnknown: s === null ? "no-final-assistant-message" : "unrecognised-stop-reason" };
|
||||
}
|
||||
|
||||
function queueDigestCount(ev) {
|
||||
const all = [...(Array.isArray(ev.steering) ? ev.steering : []), ...(Array.isArray(ev.followUp) ? ev.followUp : [])];
|
||||
return all.length;
|
||||
}
|
||||
Reference in New Issue
Block a user