packages/conversation is a library with no server: safe-fs, the Pi session parser, CHAT-01 pages, pinned snapshots, cursors and follow. The control board adds GET /api/conversations and /api/conversation behind the Host and Origin guard. Both are read-only, their queries are validated, and each refusal code maps to a status. Dewey authored it (packet 0cf177b1, revision 2). Filbert reviewed the code: R1 revise (branch ids moving on append, the assumed-link bridge merging branches, one unreadable seat directory turning the catalogue into a 500), then R2 approve (3b14d66c). Darkwing reviewed the routes: R1 approve (07b10ad1), R2 approve (b9d92003). The package lands with the routes, because serve.mjs imports the reader at load. On an index export: the eight suites 24/90/43/17/14/15/63/18, conversation and control-board 153/153, webui 9/9. Co-Authored-By: Claude Opus 5.5 <[email protected]>
402 lines
19 KiB
JavaScript
402 lines
19 KiB
JavaScript
// Read-only Pi conversation histories for the Console (#1507, CHAT-02).
|
|
//
|
|
// A library with no server: the control board owns the two GET routes (D3).
|
|
// Only approved roots are read, <projectRoot>/.pi/state/<seat>/sessions from
|
|
// the board's repository specs. There is no global scan, and no path comes
|
|
// from the caller: a conversation is an opaque id that is resolved by listing
|
|
// the roots again. Opening a conversation never resumes, forks, launches or
|
|
// controls anything, and nothing here writes a file.
|
|
//
|
|
// A seat registration is seat-written, so it is a hint, not authority. It can
|
|
// narrow a root (a non-Pi harness makes the root unsupported, D2) and supply
|
|
// the engine launch time. It never adds a root or names a file.
|
|
|
|
import { createHash, randomBytes } from "node:crypto";
|
|
import { realpathSync } from "node:fs";
|
|
import { join, resolve, basename, sep, isAbsolute } from "node:path";
|
|
import { Refusal, listSessionFiles, openSessionFile, readRange, closeSync } from "./safe-fs.mjs";
|
|
import { parseSnapshot, branchPath, branchUnits, toTime } from "./pi.mjs";
|
|
import { unitParts, takePage, bytesOf, safeId } from "./parts.mjs";
|
|
import { samePath } from "../../seat/src/seat.mjs";
|
|
|
|
export const ACTOR = "local-operator";
|
|
export const HISTORY = "pi";
|
|
export const UNSUPPORTED_HARNESS = "unsupported-harness";
|
|
const NO_STREAM = "no-stream";
|
|
const HEAD_BYTES = 1024 * 1024;
|
|
const TAIL_BYTES = 2 * 1024 * 1024;
|
|
const TITLE_CHARS = 200;
|
|
export const MAX_FILE_BYTES = 256 * 1024 * 1024;
|
|
|
|
const sha256 = (data) => createHash("sha256").update(data).digest("hex");
|
|
|
|
// Approved roots from the board's specs. Only repository specs qualify:
|
|
// sessionsDir must be exactly <projectRoot>/.pi/state/<agent>/sessions and the
|
|
// project must be the root's directory name. Fleet and connector specs do not.
|
|
export function rootsFromSpecs(specs, registrations = []) {
|
|
const roots = [];
|
|
for (const spec of specs) {
|
|
if (!spec || spec.connector || typeof spec.sessionsDir !== "string" || typeof spec.agent !== "string") continue;
|
|
const dir = resolve(spec.sessionsDir);
|
|
const projectRoot = resolve(dir, "..", "..", "..", "..");
|
|
if (join(projectRoot, ".pi", "state", spec.agent, "sessions") !== dir || basename(projectRoot) !== spec.project) continue;
|
|
const reg = registrations.find((r) => r && r.seat === spec.agent && r.layout === "repo" && samePath(r.sessionsDir, dir) && (r.project === null || r.project === spec.project)) ?? null;
|
|
const harness = reg?.harness ?? null;
|
|
roots.push({
|
|
seat: spec.agent,
|
|
project: spec.project,
|
|
projectRoot,
|
|
dir,
|
|
harness: harness ?? HISTORY,
|
|
unsupportedReason: harness !== null && harness !== HISTORY ? UNSUPPORTED_HARNESS : null,
|
|
engineStartedAt: toTime(reg?.startedAt) ?? null,
|
|
});
|
|
}
|
|
// A seat on another harness has no Pi sessions directory, so no spec. Its
|
|
// registration adds a placeholder root, never read, when it names the
|
|
// standard directory under a project root already approved above (D2).
|
|
const projects = new Map(roots.map((r) => [r.projectRoot, r.project]));
|
|
for (const reg of registrations) {
|
|
if (!reg || reg.layout !== "repo" || typeof reg.harness !== "string" || reg.harness === HISTORY) continue;
|
|
if (roots.some((r) => r.seat === reg.seat)) continue;
|
|
for (const [projectRoot, project] of projects) {
|
|
const dir = join(projectRoot, ".pi", "state", reg.seat, "sessions");
|
|
if ((reg.project !== null && reg.project !== project) || !samePath(reg.sessionsDir, dir)) continue;
|
|
roots.push({ seat: reg.seat, project, projectRoot, dir, harness: reg.harness, unsupportedReason: UNSUPPORTED_HARNESS, engineStartedAt: toTime(reg.startedAt) ?? null });
|
|
break;
|
|
}
|
|
}
|
|
return roots;
|
|
}
|
|
|
|
export function conversationId(root, name) {
|
|
return "pi-" + sha256(`${root.projectRoot}\0${root.seat}\0${name}`).slice(0, 32);
|
|
}
|
|
|
|
function unsupportedId(root) {
|
|
return "unsupported-" + sha256(`${root.projectRoot}\0${root.seat}`).slice(0, 32);
|
|
}
|
|
|
|
// The header's cwd must be an absolute path to the project or inside it. Pi
|
|
// writes absolute paths; a relative one would resolve against the board's own
|
|
// directory, so it is refused. Real paths are compared first (the checkout is
|
|
// reachable through a compatibility symlink); a cwd that no longer exists is
|
|
// compared as written.
|
|
export function cwdInProject(cwd, projectRoot) {
|
|
if (typeof cwd !== "string" || !isAbsolute(cwd)) return false;
|
|
const inside = (a, b) => a === b || a.startsWith(b.endsWith(sep) ? b : b + sep);
|
|
if (inside(resolve(cwd), resolve(projectRoot))) return true;
|
|
try {
|
|
return inside(realpathSync(cwd), realpathSync(projectRoot));
|
|
} catch {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
const firstChars = (s, n) => {
|
|
const cps = Array.from(String(s).replace(/\s+/g, " ").trim());
|
|
return cps.length > n ? cps.slice(0, n).join("") + "…" : cps.join("");
|
|
};
|
|
|
|
function userText(m) {
|
|
if (!m || m.role !== "user") return null;
|
|
if (typeof m.content === "string") return m.content;
|
|
if (Array.isArray(m.content)) {
|
|
const t = m.content.find((b) => b?.type === "text" && typeof b.text === "string" && b.text.trim());
|
|
return t ? t.text : null;
|
|
}
|
|
return null;
|
|
}
|
|
|
|
const parseLine = (line) => {
|
|
try {
|
|
const v = JSON.parse(line);
|
|
return v && typeof v === "object" && !Array.isArray(v) ? v : null;
|
|
} catch {
|
|
return null;
|
|
}
|
|
};
|
|
|
|
// A catalogue row. This is a CHAT-02 summary, not a CHAT-01 catalogueItem:
|
|
// that record needs an execution-bound currentView, which D1 moved to CHAT-03.
|
|
function baseRow(root, conversation) {
|
|
return {
|
|
conversation,
|
|
seat: root.seat,
|
|
project: root.project,
|
|
harness: root.harness,
|
|
history: root.unsupportedReason ? null : HISTORY,
|
|
title: null,
|
|
readOnly: true,
|
|
controlMode: "unavailable",
|
|
conversationCreatedAt: null,
|
|
engineStartedAt: null,
|
|
lastActivityAt: null,
|
|
availability: "available",
|
|
unsupportedReason: root.unsupportedReason,
|
|
refusal: null,
|
|
};
|
|
}
|
|
|
|
// One catalogue row from a head and a tail window; the whole file is not read.
|
|
function summarize(root, name) {
|
|
const row = baseRow(root, conversationId(root, name));
|
|
let file;
|
|
try {
|
|
file = openSessionFile(root, name);
|
|
} catch (err) {
|
|
if (!(err instanceof Refusal)) throw err;
|
|
return { ...row, availability: "denied", refusal: err.code };
|
|
}
|
|
try {
|
|
const head = readRange(file.fd, 0, Math.min(file.size, HEAD_BYTES)).toString("utf8");
|
|
const headLines = head.split("\n");
|
|
if (headLines.length < 2) return { ...row, availability: "unavailable", refusal: "incomplete-header" };
|
|
const header = parseLine(headLines[0]);
|
|
if (!header || header.type !== "session") return { ...row, availability: "unavailable", refusal: "not-a-pi-session" };
|
|
if (!cwdInProject(header.cwd, root.projectRoot)) return { ...row, availability: "denied", refusal: "foreign-project" };
|
|
row.conversationCreatedAt = toTime(header.timestamp);
|
|
if (row.conversationCreatedAt && root.engineStartedAt && row.conversationCreatedAt >= root.engineStartedAt) row.engineStartedAt = root.engineStartedAt;
|
|
let name_ = null, firstUser = null;
|
|
for (const line of headLines.slice(1, -1)) {
|
|
const e = parseLine(line);
|
|
if (e?.type === "session_info" && typeof e.name === "string" && e.name.trim()) name_ = e.name;
|
|
if (firstUser === null && e?.type === "message") firstUser = userText(e.message);
|
|
}
|
|
const tailStart = Math.max(0, file.size - TAIL_BYTES);
|
|
const tail = readRange(file.fd, tailStart, file.size - tailStart).toString("utf8");
|
|
const complete = tail.slice(0, tail.lastIndexOf("\n") + 1).split("\n").slice(tailStart > 0 ? 1 : 0, -1);
|
|
let lastName = null;
|
|
for (let i = complete.length - 1; i >= 0; i--) {
|
|
const e = parseLine(complete[i]);
|
|
if (!e || e.type === "session" || typeof e.id !== "string") continue;
|
|
if (row.lastActivityAt === null) row.lastActivityAt = toTime(e.timestamp);
|
|
if (lastName === null && e.type === "session_info" && typeof e.name === "string" && e.name.trim()) lastName = e.name;
|
|
if (row.lastActivityAt !== null && lastName !== null) break;
|
|
}
|
|
const title = lastName ?? name_ ?? firstUser;
|
|
row.title = title ? firstChars(title, TITLE_CHARS) : null;
|
|
return row;
|
|
} finally {
|
|
closeSync(file.fd);
|
|
}
|
|
}
|
|
|
|
// Complete lines only: the snapshot ends at the last "\n" present when the
|
|
// descriptor was read, so growth during the read is cut there.
|
|
function readSnapshot(root, name, pinnedLength = null) {
|
|
const file = openSessionFile(root, name);
|
|
try {
|
|
if (file.size > MAX_FILE_BYTES) throw new Refusal("too-large", "session file is larger than the reader accepts");
|
|
if (pinnedLength !== null) {
|
|
if (file.size < pinnedLength) return { file, shorter: true };
|
|
const buf = readRange(file.fd, 0, pinnedLength);
|
|
return { file, buf, length: pinnedLength };
|
|
}
|
|
const buf = readRange(file.fd, 0, file.size);
|
|
const length = buf.lastIndexOf(0x0a) + 1;
|
|
return { file, buf: buf.subarray(0, length), length, incomplete: buf.length > length };
|
|
} finally {
|
|
closeSync(file.fd);
|
|
}
|
|
}
|
|
|
|
export function createReader({ roots, now = () => Date.now(), ttlMs = 10 * 60 * 1000, maxCursors = 1000, cacheSize = 2 } = {}) {
|
|
const listRoots = typeof roots === "function" ? roots : () => roots ?? [];
|
|
const cursors = new Map();
|
|
const parsedCache = new Map();
|
|
const partsCache = new Map();
|
|
|
|
const remember = (cache, key, make) => {
|
|
if (cache.has(key)) {
|
|
const v = cache.get(key);
|
|
cache.delete(key);
|
|
cache.set(key, v);
|
|
return v;
|
|
}
|
|
const v = make();
|
|
cache.set(key, v);
|
|
while (cache.size > cacheSize) cache.delete(cache.keys().next().value);
|
|
return v;
|
|
};
|
|
|
|
function resolveConversation(conversation) {
|
|
if (typeof conversation !== "string") return null;
|
|
for (const root of listRoots()) {
|
|
if (root.unsupportedReason) {
|
|
if (unsupportedId(root) === conversation) return { root, name: null };
|
|
continue;
|
|
}
|
|
let listing;
|
|
try {
|
|
listing = listSessionFiles(root);
|
|
} catch (err) {
|
|
if (err instanceof Refusal) continue;
|
|
throw err;
|
|
}
|
|
for (const name of [...listing.files, ...listing.refused.map((r) => r.name)]) {
|
|
if (conversationId(root, name) === conversation) return { root, name };
|
|
}
|
|
}
|
|
return null;
|
|
}
|
|
|
|
function catalogue() {
|
|
const conversations = [], refusedRoots = [];
|
|
for (const root of listRoots()) {
|
|
if (root.unsupportedReason) {
|
|
conversations.push({ ...baseRow(root, unsupportedId(root)), engineStartedAt: root.engineStartedAt, availability: "unsupported" });
|
|
continue;
|
|
}
|
|
let listing;
|
|
try {
|
|
listing = listSessionFiles(root);
|
|
} catch (err) {
|
|
if (!(err instanceof Refusal)) throw err;
|
|
refusedRoots.push({ seat: root.seat, project: root.project, refusal: err.code });
|
|
continue;
|
|
}
|
|
for (const name of listing.files) conversations.push(summarize(root, name));
|
|
for (const r of listing.refused) conversations.push({ ...baseRow(root, conversationId(root, r.name)), availability: "denied", refusal: r.code });
|
|
}
|
|
conversations.sort((a, b) => (b.lastActivityAt ?? "").localeCompare(a.lastActivityAt ?? "") || a.conversation.localeCompare(b.conversation));
|
|
return { ok: true, conversations, refusedRoots, generatedAt: new Date(now()).toISOString() };
|
|
}
|
|
|
|
// Parsed snapshot plus one branch's entries, cached by content. A null
|
|
// branch is the default one; an unknown branch gives null.
|
|
function build(snap, root, branchName, conversation) {
|
|
const key = `${snap.file.dev}:${snap.file.ino}:${snap.length}:${snap.digest}`;
|
|
const parsed = remember(parsedCache, key, () => parseSnapshot(snap.buf.toString("utf8")));
|
|
if (!cwdInProject(parsed.header.cwd, root.projectRoot)) throw new Refusal("foreign-project", "the session belongs to another project");
|
|
const branch = branchName ?? parsed.defaultBranch;
|
|
if (!parsed.branches.has(branch)) return null;
|
|
const leaf = parsed.branches.get(branch);
|
|
const { entries, sizes } = remember(partsCache, `${key}:${conversation}:${branch}`, () => {
|
|
const base = { conversation, branch, execution: safeId(parsed.header.id) };
|
|
const entries = branchUnits(parsed, branchPath(parsed, leaf)).flatMap((u) => unitParts(u, base));
|
|
return { entries, sizes: entries.map(bytesOf) };
|
|
});
|
|
return { parsed, leaf, branch, entries, sizes };
|
|
}
|
|
|
|
const idsDigest = (entries, end) => sha256(entries.slice(0, end).map((e) => e.id).join("\n"));
|
|
|
|
function issue(state) {
|
|
while (cursors.size >= maxCursors) cursors.delete(cursors.keys().next().value);
|
|
const id = "c-" + randomBytes(16).toString("hex");
|
|
const expiresAt = new Date(now() + ttlMs).toISOString();
|
|
const record = {
|
|
version: 2, kind: "cursor", id, conversation: state.conversation, branch: state.branch, snapshotDigest: state.digest,
|
|
sourceEpoch: state.epoch, lastEntry: state.lastEntry, expiresAt, actor: state.actor, purpose: state.purpose,
|
|
};
|
|
cursors.set(id, { record, state });
|
|
return record;
|
|
}
|
|
|
|
// One page from `offset`, with a next cursor when more parts remain in this
|
|
// snapshot and a follow cursor when the page reaches its end.
|
|
function pageFrom({ root, name, snap, built, conversation, offset, epoch, actor, purpose, incomplete }) {
|
|
const { entries, sizes } = built;
|
|
const shell = {
|
|
version: 2, kind: "page", conversation, branch: built.branch, snapshotDigest: snap.digest, sourceEpoch: epoch,
|
|
entries: [], nextCursor: "c-" + "0".repeat(32), hasMore: true, readOnly: true, streamEpoch: NO_STREAM, throughSequence: 0,
|
|
};
|
|
const end = takePage(entries, sizes, offset, shell);
|
|
const hasMore = end < entries.length;
|
|
const state = {
|
|
root, name, dev: snap.file.dev, ino: snap.file.ino, length: snap.length, digest: snap.digest, epoch,
|
|
conversation, branch: built.branch,
|
|
offset: end, prefix: idsDigest(entries, end), lastEntry: end > 0 ? entries[end - 1].id : null, actor, purpose, incomplete,
|
|
};
|
|
const cursor = hasMore ? issue({ ...state, follow: false }) : null;
|
|
const follow = hasMore ? null : issue({ ...state, follow: true });
|
|
const page = { ...shell, entries: entries.slice(offset, end), nextCursor: cursor ? cursor.id : null, hasMore };
|
|
const parsed = built.parsed;
|
|
const view = {
|
|
conversation,
|
|
branch: built.branch,
|
|
defaultBranch: parsed.defaultBranch,
|
|
branches: [...parsed.branches]
|
|
.sort(([, a], [, b]) => (a?.index ?? -1) - (b?.index ?? -1))
|
|
.map(([branch, leaf]) => ({ branch, isDefault: branch === parsed.defaultBranch, lastActivityAt: leaf ? toTime(leaf.entry.timestamp) : null })),
|
|
incomplete,
|
|
forkedFromEarlierSession: parsed.header.parentSession !== undefined && parsed.header.parentSession !== null,
|
|
unreadableLines: parsed.malformed.length,
|
|
};
|
|
return { ok: true, page, cursor, follow, view };
|
|
}
|
|
|
|
function refuse(err) {
|
|
if (err instanceof Refusal) return { ok: false, refusal: { code: err.code, reconcile: err.reconcile, message: err.message } };
|
|
throw err;
|
|
}
|
|
|
|
function checkCaller(actor, purpose) {
|
|
if (actor !== ACTOR) throw new Refusal("unknown-actor", "only the local operator reads histories on this route");
|
|
if (purpose !== "history") throw new Refusal("unsupported-purpose", "only history pages are served");
|
|
}
|
|
|
|
function snapshot(root, name, pinnedLength = null) {
|
|
const snap = readSnapshot(root, name, pinnedLength);
|
|
if (snap.shorter) return snap;
|
|
if (snap.length === 0) throw new Refusal("incomplete-header", "the session header is not complete yet", { reconcile: true });
|
|
return { ...snap, digest: sha256(snap.buf) };
|
|
}
|
|
|
|
function open({ conversation, branch = null, actor = ACTOR, purpose = "history" } = {}) {
|
|
try {
|
|
checkCaller(actor, purpose);
|
|
const found = resolveConversation(conversation);
|
|
if (!found) throw new Refusal("unknown-conversation", "no such conversation in the approved roots", { reconcile: true });
|
|
if (found.root.unsupportedReason) throw new Refusal(found.root.unsupportedReason, "this harness has no history reader yet");
|
|
const snap = snapshot(found.root, found.name);
|
|
const built = build(snap, found.root, branch ?? null, conversation);
|
|
if (!built) throw new Refusal("unknown-branch", "no such branch in this conversation", { reconcile: true });
|
|
const epoch = "e-" + sha256(`${snap.file.dev}:${snap.file.ino}:${snap.digest}`).slice(0, 40);
|
|
return pageFrom({ root: found.root, name: found.name, snap, built, conversation, offset: 0, epoch, actor, purpose, incomplete: snap.incomplete });
|
|
} catch (err) {
|
|
return refuse(err);
|
|
}
|
|
}
|
|
|
|
function next({ cursor, conversation, branch, actor = ACTOR, purpose = "history" } = {}) {
|
|
try {
|
|
const held = typeof cursor === "string" ? cursors.get(cursor) : undefined;
|
|
if (!held) throw new Refusal("cursor-unknown", "the cursor is unknown", { reconcile: true });
|
|
const { record, state } = held;
|
|
if (Date.parse(record.expiresAt) <= now()) {
|
|
cursors.delete(cursor);
|
|
throw new Refusal("cursor-expired", "the cursor has expired", { reconcile: true });
|
|
}
|
|
if (actor !== record.actor || purpose !== record.purpose || conversation !== record.conversation || branch !== record.branch) {
|
|
throw new Refusal("cursor-foreign", "the cursor belongs to another view", { reconcile: true });
|
|
}
|
|
// The source must still be the snapshot's file with the same prefix.
|
|
const pinned = snapshot(state.root, state.name, state.length);
|
|
if (pinned.shorter) throw new Refusal("source-replaced", "the session file is shorter than the snapshot", { reconcile: true });
|
|
if (pinned.file.dev !== state.dev || pinned.file.ino !== state.ino) throw new Refusal("source-replaced", "the session file was replaced", { reconcile: true });
|
|
if (pinned.digest !== state.digest) throw new Refusal("source-replaced", "the session file was rewritten", { reconcile: true });
|
|
const args = { root: state.root, name: state.name, conversation, epoch: state.epoch, actor, purpose };
|
|
if (!state.follow) {
|
|
const built = build(pinned, state.root, state.branch, conversation);
|
|
return pageFrom({ ...args, snap: pinned, built, offset: state.offset, incomplete: state.incomplete });
|
|
}
|
|
// Follow: a fresh snapshot that extends the verified prefix. The view
|
|
// stays on its branch; view.defaultBranch shows where Pi's default is.
|
|
const fresh = snapshot(state.root, state.name);
|
|
if (fresh.file.dev !== state.dev || fresh.file.ino !== state.ino || fresh.length < state.length) throw new Refusal("source-replaced", "the session file was replaced", { reconcile: true });
|
|
if (sha256(fresh.buf.subarray(0, state.length)) !== state.digest) throw new Refusal("source-replaced", "the session file was rewritten", { reconcile: true });
|
|
const built = build(fresh, state.root, state.branch, conversation);
|
|
if (!built || built.entries.length < state.offset || idsDigest(built.entries, state.offset) !== state.prefix) {
|
|
throw new Refusal("source-replaced", "the history before this point changed", { reconcile: true });
|
|
}
|
|
return pageFrom({ ...args, snap: fresh, built, offset: state.offset, incomplete: fresh.incomplete });
|
|
} catch (err) {
|
|
return refuse(err);
|
|
}
|
|
}
|
|
|
|
return { catalogue, open, next, cursorCount: () => cursors.size };
|
|
}
|