Row 47, the S5 follow-up. flows.test gains a test for a hold set inside held input, which kills mutant Mr. claim W5 and W13 reap their own force-stopped scopes, K19 kills its scope if its kill or release fails, and the harness gains liveShims, killShims and shimsGone with a final sweep in reap. claim, cohort and races each end by asserting no shim outlives the file. Author: Dewey. Approved round 1 by Filbert (comment 27041) and Darkwing (comment 27043) on candidate 0587e393. The src side, where a force stop never sends release, is row 50 (#1536). Co-Authored-By: Claude Opus 5.5 <[email protected]>
266 lines
10 KiB
JavaScript
266 lines
10 KiB
JavaScript
// Shared CHAT-03 fixture setup (#1507). Every fixture lives in its own temp
|
||
// directory: fixture root, project, session file, claim root and socket
|
||
// directory. Nothing here touches a live session.
|
||
|
||
import { mkdtempSync, mkdirSync, writeFileSync, readFileSync, rmSync } from "node:fs";
|
||
import { tmpdir } from "node:os";
|
||
import { join, resolve } from "node:path";
|
||
import { spawnSync } from "node:child_process";
|
||
import { Controller } from "../src/controller.mjs";
|
||
import { ConversationClient } from "../src/client.mjs";
|
||
import { AUTHORITY } from "../src/cohort.mjs";
|
||
import { FixtureVerifier } from "../src/records.mjs";
|
||
import { FakeLauncher } from "./fake-pi.mjs";
|
||
|
||
export const REPO = resolve(import.meta.dirname, "..", "..", "..");
|
||
export const SESSION_ID = "0f5e1c2a-1111-4222-8333-944455556666";
|
||
export const FAST = Object.freeze({ ack: 1500, state: 1500, start: 1500, clear: 1500, abort: 3000, settle: 3000, grace: 50, write: 1500, maxRounds: 3 });
|
||
export const time = (i) => new Date(Date.UTC(2026, 9, 4, 12, 0, 0) + i * 1000).toISOString();
|
||
|
||
const tmp = mkdtempSync(join(tmpdir(), "chat03-"));
|
||
let serial = 0;
|
||
|
||
export function cleanupAll() {
|
||
spawnSync("chmod", ["-R", "u+rwx", tmp]);
|
||
rmSync(tmp, { recursive: true, force: true });
|
||
}
|
||
|
||
export const header = (cwd) => ({ type: "session", version: 3, id: SESSION_ID, timestamp: time(0), cwd });
|
||
export const userEntry = (id, parentId, text, i = 1) => ({ type: "message", id, parentId, timestamp: time(i), message: { role: "user", content: [{ type: "text", text }] } });
|
||
// Pi writes this when it creates a session (sdk.js 247–251), so a session Pi
|
||
// created and prompted carries one. Pi appends at startup to a session with
|
||
// messages and no thinking entry, and to any session with no messages.
|
||
export const thinkingEntry = (id = "f0e1d2c3", parentId = null) => ({ type: "thinking_level_change", id, parentId, timestamp: time(0), thinkingLevel: "off" });
|
||
export const assistantEntry = (id, parentId, text, i = 2) => ({ type: "message", id, parentId, timestamp: time(i), message: { role: "assistant", content: [{ type: "text", text }], stopReason: "stop" } });
|
||
|
||
// A fixture tree with one session file. `entries` default to one exchange.
|
||
export function fixture({ seat = "fixture-seat", name = "s1.jsonl", entries } = {}) {
|
||
const base = join(tmp, `f${++serial}`);
|
||
const proj = join(base, "proj");
|
||
const sessions = join(proj, ".pi", "state", seat, "sessions");
|
||
mkdirSync(sessions, { recursive: true });
|
||
const sessionFile = join(sessions, name);
|
||
const list = entries ?? [thinkingEntry(), userEntry("a1b2c3d4", "f0e1d2c3", "hello"), assistantEntry("b2c3d4e5", "a1b2c3d4", "hi")];
|
||
writeFileSync(sessionFile, [header(proj), ...list].map((v) => JSON.stringify(v) + "\n").join(""));
|
||
return { base, proj, sessions, sessionFile, seat, claimRoot: join(base, "claims"), socketDir: join(base, "sock") };
|
||
}
|
||
|
||
export const noUnits = Object.freeze({ lookup: async () => ({ state: "absent" }) });
|
||
|
||
// A controller on the in-process fake engine.
|
||
export function controllerFor(fx, { launcher = new FakeLauncher(), verifier = new FixtureVerifier({ authorities: [AUTHORITY] }), timeouts = FAST, ...opts } = {}) {
|
||
const ctrl = new Controller({
|
||
fixtureRoot: fx.base, claimRoot: fx.claimRoot, socketDir: fx.socketDir, sessionFile: fx.sessionFile, seat: fx.seat,
|
||
launcher, verifier, units: noUnits, timeouts, ...opts,
|
||
});
|
||
return { ctrl, launcher, verifier };
|
||
}
|
||
|
||
// Starts a controller and connects `n` clients; the first takes control.
|
||
export async function started(opts = {}) {
|
||
const fx = opts.fx ?? fixture(opts.fixture);
|
||
const { ctrl, launcher, verifier } = controllerFor(fx, opts);
|
||
const started = await ctrl.start();
|
||
const clients = [];
|
||
for (let i = 0; i < (opts.clients ?? 1); i++) {
|
||
const c = new ConversationClient({ socketPath: ctrl.socketPath });
|
||
await c.connect();
|
||
clients.push(c);
|
||
}
|
||
if (opts.control !== false && clients[0]) {
|
||
const r = await clients[0].takeover();
|
||
if (r.outcome !== "transferred") throw new Error(`takeover: ${JSON.stringify(r)}`);
|
||
}
|
||
const close = async () => {
|
||
for (const c of clients) c.close();
|
||
await ctrl.close({ killEngine: true });
|
||
};
|
||
return { fx, ctrl, launcher, verifier, started, clients, client: clients[0], engine: launcher.last, close };
|
||
}
|
||
|
||
// Waits for a receipt to reach one of `states` (or outcome unknown).
|
||
export async function receiptState(client, receiptId, states, ms = 4000) {
|
||
const want = new Set(Array.isArray(states) ? states : [states]);
|
||
const ok = await client.waitFor(() => {
|
||
const r = client.receipt(receiptId);
|
||
return r && want.has(r.state);
|
||
}, ms);
|
||
if (!ok) throw new Error(`receipt ${receiptId} stayed ${client.receipt(receiptId)?.state}; wanted ${[...want]}`);
|
||
return client.receipt(receiptId);
|
||
}
|
||
|
||
export function outcomeUnknownPush(client, receiptId) {
|
||
return client.pushes.find((p) => p.kind === "receipt" && p.receipt.id === receiptId && p.outcomeUnknown === true) ?? null;
|
||
}
|
||
|
||
export const sessionText = (fx) => readFileSync(fx.sessionFile, "utf8");
|
||
export const tick = (ms = 0) => new Promise((r) => setTimeout(r, ms));
|
||
export { tmp };
|
||
|
||
// ---- controller in a child process -----------------------------------------
|
||
|
||
import { spawn } from "node:child_process";
|
||
|
||
const CHILD = join(import.meta.dirname, "ctrl-child.mjs");
|
||
const childPids = new Set();
|
||
|
||
export function killChildren() {
|
||
for (const pid of childPids) {
|
||
try {
|
||
process.kill(pid, "SIGKILL");
|
||
} catch {
|
||
// gone
|
||
}
|
||
}
|
||
childPids.clear();
|
||
}
|
||
|
||
// Spawns ctrl-child.mjs. `msgs` collects its JSON lines; `next(pred)` waits
|
||
// for one.
|
||
export function spawnController(cfg) {
|
||
const proc = spawn(process.execPath, [CHILD, JSON.stringify({ timeouts: FAST, ...cfg })], { stdio: ["pipe", "pipe", "pipe"] });
|
||
childPids.add(proc.pid);
|
||
const msgs = [];
|
||
const waiters = new Set();
|
||
let buf = "";
|
||
let stderr = "";
|
||
proc.stderr.on("data", (c) => (stderr += c));
|
||
proc.stdout.on("data", (c) => {
|
||
buf += c;
|
||
let i;
|
||
while ((i = buf.indexOf("\n")) >= 0) {
|
||
const line = buf.slice(0, i);
|
||
buf = buf.slice(i + 1);
|
||
try {
|
||
msgs.push(JSON.parse(line));
|
||
} catch {
|
||
msgs.push({ raw: line });
|
||
}
|
||
for (const w of [...waiters]) w();
|
||
}
|
||
});
|
||
const exited = new Promise((r) => proc.on("exit", (code, signal) => {
|
||
childPids.delete(proc.pid);
|
||
r({ code, signal });
|
||
for (const w of [...waiters]) w();
|
||
}));
|
||
let done = false;
|
||
exited.then(() => (done = true));
|
||
const next = (pred, ms = 8000) => new Promise((resolve, reject) => {
|
||
const check = () => {
|
||
const hit = msgs.find(pred);
|
||
if (hit) {
|
||
waiters.delete(check);
|
||
clearTimeout(t);
|
||
resolve(hit);
|
||
} else if (done) {
|
||
waiters.delete(check);
|
||
clearTimeout(t);
|
||
reject(new Error(`controller child exited; lines ${JSON.stringify(msgs)} stderr ${stderr.slice(-2000)}`));
|
||
}
|
||
};
|
||
const t = setTimeout(() => {
|
||
waiters.delete(check);
|
||
reject(new Error(`timeout; lines ${JSON.stringify(msgs)} stderr ${stderr.slice(-2000)}`));
|
||
}, ms);
|
||
waiters.add(check);
|
||
check();
|
||
});
|
||
return { proc, msgs, next, exited, send: (line) => proc.stdin.write(line + "\n"), stderr: () => stderr };
|
||
}
|
||
|
||
// Stops every scope and kills every pgroup engine a fixture's claim records
|
||
// name, so a crashed controller leaves nothing running.
|
||
import { readdirSync as lsdir } from "node:fs";
|
||
import { processStart } from "../../discord/src/journal.mjs";
|
||
|
||
export function claimRecords(fx) {
|
||
const out = [];
|
||
const walk = (dir) => {
|
||
let names = [];
|
||
try {
|
||
names = lsdir(dir, { withFileTypes: true });
|
||
} catch {
|
||
return;
|
||
}
|
||
for (const d of names) {
|
||
const p = join(dir, d.name);
|
||
if (d.isDirectory()) walk(p);
|
||
else if (/^r\d{10}\.json$/.test(d.name)) {
|
||
try {
|
||
out.push({ path: p, record: JSON.parse(readFileSync(p, "utf8")) });
|
||
} catch {
|
||
out.push({ path: p, record: null });
|
||
}
|
||
}
|
||
}
|
||
};
|
||
walk(fx.claimRoot);
|
||
return out;
|
||
}
|
||
|
||
export function reap(fx) {
|
||
const units = new Set(), engines = new Map();
|
||
for (const { record } of claimRecords(fx)) {
|
||
if (!record) continue;
|
||
if (record.unitName && record.spawnMarker) units.add(record.unitName);
|
||
if (record.engine?.kind === "pgroup" && record.engine.pid) engines.set(record.engine.pid, record.engine.start);
|
||
}
|
||
for (const u of units) spawnSync("systemctl", ["--user", "kill", "--signal=SIGKILL", `${u}.scope`], { stdio: "ignore", timeout: 5000 });
|
||
for (const u of units) spawnSync("systemctl", ["--user", "stop", `${u}.scope`], { stdio: "ignore", timeout: 5000 });
|
||
for (const [pid, start] of engines) {
|
||
if (processStart(pid) !== start) continue;
|
||
try {
|
||
process.kill(-pid, "SIGKILL");
|
||
} catch {
|
||
// gone
|
||
}
|
||
}
|
||
// A shim the records don't name (a direct launch), or one the systemctl
|
||
// calls above didn't reach: neither result is checked.
|
||
killShims(fx.base);
|
||
}
|
||
|
||
// The shims under `root` still running, read from /proc. A shim ignores
|
||
// SIGTERM and only exits on `release` (src/shim.mjs), and a force stop
|
||
// never sends `release`, so a stopped scope stays up until it is killed.
|
||
export function liveShims(root = tmp) {
|
||
const out = [];
|
||
for (const pid of lsdir("/proc").filter((n) => /^\d+$/.test(n))) {
|
||
let argv, cgroup;
|
||
try {
|
||
argv = readFileSync(`/proc/${pid}/cmdline`, "utf8").split("\0");
|
||
if (!argv[1]?.endsWith("/shim.mjs") || argv[2] !== "--socket" || !argv[3]?.startsWith(root + "/")) continue;
|
||
cgroup = readFileSync(`/proc/${pid}/cgroup`, "utf8");
|
||
} catch {
|
||
continue; // gone
|
||
}
|
||
const scope = cgroup.trim().split("/").find((s) => s.endsWith(".scope"));
|
||
out.push({ pid: Number(pid), socket: argv[3], unit: scope ? scope.slice(0, -".scope".length) : null });
|
||
}
|
||
return out;
|
||
}
|
||
|
||
// SIGKILLs each shim's scope, which takes the engine with it, and the shim
|
||
// itself in case the scope kill doesn't land.
|
||
export function killShims(root = tmp) {
|
||
for (const s of liveShims(root)) {
|
||
if (s.unit) spawnSync("systemctl", ["--user", "kill", "--signal=SIGKILL", `${s.unit}.scope`], { stdio: "ignore", timeout: 5000 });
|
||
try {
|
||
process.kill(s.pid, "SIGKILL");
|
||
} catch {
|
||
// gone
|
||
}
|
||
}
|
||
}
|
||
|
||
// Waits up to `ms` for every shim under `root` to go; returns those left.
|
||
export async function shimsGone(root = tmp, ms = 5000) {
|
||
const end = Date.now() + ms;
|
||
for (;;) {
|
||
const left = liveShims(root);
|
||
if (left.length === 0 || Date.now() >= end) return left;
|
||
await tick(50);
|
||
}
|
||
}
|