feat(conversation): CHAT-02 read-only Pi history reader and two board routes (#1507)

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]>
This commit is contained in:
2026-09-26 16:36:20 -05:00
co-authored by Claude Opus 5.5
parent c5db8c8819
commit a5beb6d97d
22 changed files with 3247 additions and 3 deletions
+124
View File
@@ -0,0 +1,124 @@
// CHAT-01 history limits (#1507, CHAT-02): at most 100 parts per page, 8 MiB
// of serialized UTF-8 per page, 64 blocks per part and 262144 characters per
// string. Oversize content splits into fragments and continuation parts; it is
// never clipped. Byte limits are measured on the serialized JSON, not on
// characters.
import { createHash } from "node:crypto";
export const LIMITS = Object.freeze({ parts: 100, pageBytes: 8 * 1024 * 1024, blocks: 64, chars: 262144 });
// Internal budgets that keep any single part well inside one page, so a page
// always holds at least one part.
export const FRAGMENT_BYTES = 1024 * 1024;
export const PART_BYTES = 4 * 1024 * 1024;
export const ID = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/;
// A native value used as a CHAT-01 id. Values that do not fit the id pattern
// (OpenAI tool calls such as "call_x|fc_y") map to a stable hash.
export function safeId(value) {
if (typeof value === "string" && ID.test(value)) return value;
return "h-" + createHash("sha256").update(String(value)).digest("hex").slice(0, 40);
}
// Bytes one code point adds to a JSON string literal.
function jsonCost(cp) {
if (cp === 0x22 || cp === 0x5c) return 2;
if (cp < 0x20) return cp === 8 || cp === 9 || cp === 10 || cp === 12 || cp === 13 ? 2 : 6;
if (cp < 0x80) return 1;
if (cp < 0x800) return 2;
if (cp >= 0xd800 && cp <= 0xdfff) return 6; // lone surrogate, escaped by JSON.stringify
if (cp <= 0xffff) return 3;
return 4;
}
// Splits a string into fragments of at most LIMITS.chars code points and
// FRAGMENT_BYTES of JSON. Never cuts a surrogate pair. "" is one fragment.
export function fragments(str) {
const out = [];
let start = 0, chars = 0, bytes = 0;
for (let i = 0; i < str.length;) {
const cp = str.codePointAt(i);
const width = cp > 0xffff ? 2 : 1;
const cost = jsonCost(cp);
if (chars > 0 && (chars + 1 > LIMITS.chars || bytes + cost > FRAGMENT_BYTES)) {
out.push(str.slice(start, i));
start = i;
chars = 0;
bytes = 0;
}
chars += 1;
bytes += cost;
i += width;
}
out.push(str.slice(start));
return out;
}
const bytesOf = (value) => Buffer.byteLength(JSON.stringify(value), "utf8");
// One native unit (a message, a compaction, a notice) becomes one or more
// CHAT-01 entries. `unit.blocks` holds logical blocks: { fields, key, value },
// where `value` is the string that may split and `key` names its field. A
// block with key null (an attachment) has no string and is one fragment.
export function unitParts(unit, base) {
const items = [];
unit.blocks.forEach((b, ordinal) => {
if (b.key === null) {
items.push({ ...b.fields, block: ordinal, fragment: 0, lastFragment: true });
return;
}
const pieces = fragments(b.value);
pieces.forEach((piece, k) => {
items.push({ ...b.fields, [b.key]: piece, block: ordinal, fragment: k, lastFragment: k === pieces.length - 1 });
});
});
const groups = [];
let current = [], currentBytes = 0;
for (const item of items) {
const size = bytesOf(item) + 1;
if (current.length && (current.length === LIMITS.blocks || currentBytes + size > PART_BYTES)) {
groups.push(current);
current = [];
currentBytes = 0;
}
current.push(item);
currentBytes += size;
}
groups.push(current);
return groups.map((content, part) => ({
version: 2,
kind: "entry",
id: safeId(`${unit.id}:${part}`),
conversation: base.conversation,
branch: base.branch,
parent: unit.parent,
execution: base.execution,
nativeEntry: unit.nativeEntry,
role: unit.role,
request: null,
content,
part,
lastPart: part === groups.length - 1,
createdAt: unit.createdAt,
message: unit.message,
}));
}
// Takes entries from `start` while the page stays within LIMITS. `shell` is
// the page record with an empty `entries` array.
export function takePage(entries, sizes, start, shell) {
let bytes = bytesOf(shell);
let end = start;
while (end < entries.length && end - start < LIMITS.parts) {
const add = sizes[end] + (end > start ? 1 : 0);
if (bytes + add > LIMITS.pageBytes) break;
bytes += add;
end += 1;
}
if (end === start && start < entries.length) throw new Error("internal: a single part exceeds the page byte limit");
return end;
}
export { bytesOf };
+287
View File
@@ -0,0 +1,287 @@
// Pi session parser for read-only histories (#1507, CHAT-02). Pinned against
// @earendil-works/pi-coding-agent 0.85.1 docs/session-format.md and
// dist/core/session-manager.js:
//
// - line 1 is the session header; every other line is an entry with id,
// parentId and timestamp, forming a tree;
// - on load Pi takes the leaf to be the last entry in file order and skips
// malformed lines (_buildIndex, parseSessionEntries). This parser uses the
// same default leaf and shows a malformed line as an unavailable notice;
// - resetLeaf starts a new root entry, so a file can hold several roots;
// - the full branch path is shown, including entries before a compaction;
// the compaction itself is a marker in place. retainedTail copies entries
// that are already on the path, so it is not rendered again.
//
// get_entries order is not a branch transcript and is not used, and
// parentSession is never followed or opened.
import { Refusal } from "./safe-fs.mjs";
import { safeId } from "./parts.mjs";
const toTime = (v) => {
const ms = typeof v === "number" ? v : typeof v === "string" ? Date.parse(v) : NaN;
if (!Number.isFinite(ms)) return null;
const iso = new Date(ms).toISOString();
return /^\d{4}-/.test(iso) ? iso : null;
};
// Branch names. The first root's line of history is "main"; every later root
// (Pi's resetLeaf) and every later child at a fork starts a branch named after
// its first entry. Appending never changes which child came first, so a name
// holds for as long as the file only grows.
export const MAIN = "main";
const branchName = (r) => safeId(`b.${r.entry.id}`);
// `text` holds complete lines only (it ends with "\n" or is empty).
export function parseSnapshot(text) {
const lines = text.split("\n");
lines.pop();
let header;
try {
header = JSON.parse(lines[0] ?? "");
} catch {
header = null;
}
if (!header || typeof header !== "object" || header.type !== "session") throw new Refusal("not-a-pi-session", "the first line is not a Pi session header");
const entries = [], malformed = [], byId = new Map(), children = new Map();
for (let i = 1; i < lines.length; i++) {
const line = lines[i];
if (!line.trim()) continue;
let e;
try {
e = JSON.parse(line);
} catch {
e = null;
}
if (!e || typeof e !== "object" || Array.isArray(e) || typeof e.id !== "string" || typeof e.type !== "string" || e.type === "session") {
malformed.push({ line: i + 1, after: entries.length - 1 });
continue;
}
const record = { entry: e, line: i + 1, index: entries.length };
entries.push(record);
byId.set(e.id, record); // later wins, as in Pi's index
}
for (const r of byId.values()) {
const p = typeof r.entry.parentId === "string" ? r.entry.parentId : null;
if (!children.has(p)) children.set(p, []);
children.get(p).push(r);
}
for (const list of children.values()) list.sort((a, b) => a.index - b.index);
const parentOf = (r) => (typeof r.entry.parentId === "string" ? byId.get(r.entry.parentId) : undefined) ?? null;
const firstRoot = [...byId.values()].filter((r) => !parentOf(r)).sort((a, b) => a.index - b.index)[0] ?? null;
// Walk up while each entry is its parent's earliest child. The walk ends at
// a root or a later child, whose name the branch takes. A walk from a leaf
// always ends there; the loop guard only keeps a hostile file finite.
const branchOf = (leaf) => {
const seen = new Set();
for (let r = leaf; !seen.has(r); ) {
seen.add(r);
const p = parentOf(r);
if (!p) return r === firstRoot ? MAIN : branchName(r);
if (children.get(p.entry.id)[0] !== r) return branchName(r);
r = p;
}
return branchName(leaf);
};
const leaves = [...byId.values()].filter((r) => !children.has(r.entry.id)).sort((a, b) => a.index - b.index);
const defaultLeaf = entries.length ? byId.get(entries[entries.length - 1].entry.id) : null;
const defaultBranch = defaultLeaf ? branchOf(defaultLeaf) : MAIN;
// Branch name to leaf. The default branch ends at Pi's default leaf, and a
// file with no entries yet has an empty "main".
const branches = new Map();
for (const leaf of leaves) {
const name = branchOf(leaf);
if (!branches.has(name)) branches.set(name, leaf);
}
if (defaultLeaf) branches.set(defaultBranch, defaultLeaf);
if (!branches.size) branches.set(MAIN, null);
return { header, entries, malformed, byId, branches, defaultLeaf, defaultBranch };
}
// Root-to-leaf items for one branch: { record } for native entries and
// { notice, lines? } for markers placed where they apply. A null leaf is a
// branch with no entries yet.
export function branchPath(parsed, leaf) {
const path = [];
const seen = new Set();
const named = new Set();
let r = leaf;
while (r) {
if (seen.has(r.entry.id)) {
path.push({ notice: "loop" });
break;
}
seen.add(r.entry.id);
path.push({ record: r });
const parentId = r.entry.parentId;
if (parentId === null || parentId === undefined) break;
const parent = typeof parentId === "string" ? parsed.byId.get(parentId) : null;
if (parent) {
r = parent;
continue;
}
// The parent is missing. Unreadable lines just before this entry may have
// held it; the notice names them. The history does not continue past the
// gap: the entries before it may belong to another branch.
const lost = parsed.malformed.filter((m) => m.after === r.index - 1).map((m) => m.line);
lost.forEach((line) => named.add(line));
path.push({ notice: "missing-parent", lines: lost });
break;
}
path.reverse();
// Every other unreadable line is shown at its file position on every
// branch, after the leaf too: it may belong to any branch. Placing it the
// same way on every branch keeps a branch's earlier parts unchanged while
// the file grows.
const out = [];
const pending = parsed.malformed.filter((m) => !named.has(m.line));
let mi = 0;
const flushBefore = (line) => {
while (mi < pending.length && pending[mi].line < line) out.push({ notice: "malformed", lines: [pending[mi++].line] });
};
for (const item of path) {
const line = item.record?.line ?? item.lines?.[0];
if (line) flushBefore(line);
out.push(item);
}
flushBefore(Infinity);
return out;
}
const NOTICE_TEXT = {
loop: "History before this point is unavailable: the entry chain loops.",
"missing-parent": "History before this point is unavailable: an earlier entry is missing from the file.",
"parent-session": "This session was forked from an earlier session. The earlier session is not opened here.",
};
function text(value) {
return { fields: { type: "text" }, key: "text", value: String(value) };
}
function contentText(content) {
if (typeof content === "string") return [text(content)];
if (!Array.isArray(content)) return [];
return content.flatMap((b, i) => blockFor(b, i));
}
function blockFor(b, i, owner = "") {
if (!b || typeof b !== "object") return [];
if (b.type === "text") return [text(b.text ?? "")];
if (b.type === "image") return [{ fields: { type: "attachment", attachment: safeId(`${owner}image.${i}`) }, key: null, value: null }];
return [text(`[unsupported content block: ${safeId(String(b.type))}]`)];
}
// Unit ids are namespaced so no native id can collide with a notice: "n."
// for native entries, "x." for notices. Entry ids derive from them (parts.mjs).
//
// Converts one native entry to zero or more units. Entries that Pi keeps out
// of the transcript (model and thinking changes, labels, names, extension
// state) produce none.
function unitsFor(record, ctx) {
const e = record.entry;
const nativeEntry = safeId(e.id);
const createdAt = toTime(e.timestamp) ?? ctx.lastTime;
ctx.lastTime = createdAt;
const base = { id: `n.${e.id}`, nativeEntry, parent: typeof e.parentId === "string" ? safeId(e.parentId) : null, createdAt, message: nativeEntry };
const notice = (value, suffix) => ({ ...base, id: `x.${suffix}.${e.id}`, message: safeId(`x.${suffix}.${e.id}`), role: "notice", blocks: [text(value)] });
switch (e.type) {
case "message":
return messageUnits(e.message, base, notice);
case "compaction":
return [{ ...base, role: "compaction", blocks: [{ fields: { type: "compaction", nativeEntry }, key: "summary", value: String(e.summary ?? "") }] }];
case "branch_summary":
return [{ ...base, role: "notice", blocks: [text(`Branch summary\n\n${e.summary ?? ""}`)] }];
case "custom_message":
return e.display ? [{ ...base, role: "notice", blocks: contentText(e.content) }] : [];
case "model_change":
case "thinking_level_change":
case "label":
case "session_info":
case "custom":
return [];
default:
return [{ ...base, role: "notice", blocks: [text(`An entry of type ${safeId(e.type)} is not shown.`)] }];
}
}
function messageUnits(m, base, notice) {
if (!m || typeof m !== "object") return [{ ...base, role: "notice", blocks: [text("This entry has no readable message.")] }];
const owner = `${base.nativeEntry}.`;
switch (m.role) {
case "user":
return [{ ...base, role: "user", blocks: typeof m.content === "string" ? [text(m.content)] : (Array.isArray(m.content) ? m.content : []).flatMap((b, i) => blockFor(b, i, owner)) }];
case "assistant": {
const blocks = (Array.isArray(m.content) ? m.content : []).flatMap((b, i) => {
if (b?.type === "thinking") {
// Redacted reasoning is stored as a placeholder; it is unavailable, and
// thinkingSignature is never read.
const t = b.redacted === true ? "" : typeof b.thinking === "string" ? b.thinking : "";
return [{ fields: { type: "thinking", visibility: t ? "permitted-visible" : "unavailable" }, key: "text", value: t }];
}
if (b?.type === "toolCall") {
return [{ fields: { type: "tool-call", call: safeId(b.id), name: safeId(b.name) }, key: "argumentsText", value: JSON.stringify(b.arguments ?? {}) }];
}
return blockFor(b, i, owner);
});
const units = [{ ...base, role: "assistant", blocks }];
if (m.stopReason === "error" || m.stopReason === "aborted") {
units.push(notice(m.errorMessage ? `The turn ended (${m.stopReason}): ${m.errorMessage}` : `The turn ended (${m.stopReason}).`, "end"));
}
return units;
}
case "toolResult": {
const content = Array.isArray(m.content) ? m.content : [];
const joined = content.filter((b) => b?.type === "text").map((b) => String(b.text ?? "")).join("\n");
const blocks = [{ fields: { type: "tool-result", call: safeId(m.toolCallId), isError: m.isError === true }, key: "text", value: joined }];
content.forEach((b, i) => {
if (b?.type === "image") blocks.push(...blockFor(b, i, owner));
});
return [{ ...base, role: "tool", blocks }];
}
case "bashExecution": {
const call = safeId(`bash:${base.nativeEntry}`);
return [{
...base,
role: "tool",
blocks: [
{ fields: { type: "tool-call", call, name: "bash" }, key: "argumentsText", value: JSON.stringify({ command: String(m.command ?? "") }) },
{ fields: { type: "tool-result", call, isError: m.cancelled === true || (m.exitCode !== 0 && m.exitCode !== undefined) }, key: "text", value: String(m.output ?? "") },
],
}];
}
case "custom":
return m.display ? [{ ...base, role: "notice", blocks: contentText(m.content) }] : [];
case "branchSummary":
return [{ ...base, role: "notice", blocks: [text(`Branch summary\n\n${m.summary ?? ""}`)] }];
case "compactionSummary":
return [{ ...base, role: "compaction", blocks: [{ fields: { type: "compaction", nativeEntry: base.nativeEntry }, key: "summary", value: String(m.summary ?? "") }] }];
default:
return [{ ...base, role: "notice", blocks: [text(`A message with role ${safeId(String(m.role))} is not shown.`)] }];
}
}
// All units for one branch, in order.
export function branchUnits(parsed, path) {
const ctx = { lastTime: toTime(parsed.header.timestamp) ?? new Date(0).toISOString() };
const units = [];
if (parsed.header.parentSession !== undefined && parsed.header.parentSession !== null) {
units.push({ id: "x.parent-session", nativeEntry: "x.parent-session", parent: null, createdAt: ctx.lastTime, message: "x.parent-session", role: "notice", blocks: [text(NOTICE_TEXT["parent-session"])] });
}
for (const item of path) {
if (item.record) {
units.push(...unitsFor(item.record, ctx));
continue;
}
const lines = item.lines ?? [];
const id = item.notice === "malformed" ? `x.line-${lines[0]}` : `x.${item.notice}`;
const value = item.notice === "malformed"
? `Line ${lines[0]} could not be read. It may belong to this branch or another one.`
: item.notice === "missing-parent" && lines.length
? `${NOTICE_TEXT["missing-parent"]} ${lines.length === 1 ? `Line ${lines[0]} could not be read and may have held it.` : `Lines ${lines.join(", ")} could not be read and may have held it.`}`
: NOTICE_TEXT[item.notice];
units.push({ id, nativeEntry: id, parent: null, createdAt: ctx.lastTime, message: id, role: "notice", blocks: [text(value)] });
}
return units;
}
export { toTime };
+401
View File
@@ -0,0 +1,401 @@
// 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 };
}
+123
View File
@@ -0,0 +1,123 @@
// Read-only file access for approved Pi session roots (#1507, CHAT-02).
//
// A root is <projectRoot>/.pi/state/<seat>/sessions. The project root comes
// from the board's own configuration and is trusted as given (it may itself
// be a symlink, like the compatibility path to this checkout). Every
// component below it must be a real directory, never a symlink, and a
// session file must be a regular file directly inside the root.
//
// A file is opened O_RDONLY | O_NOFOLLOW | O_NONBLOCK, and the descriptor's
// (dev, ino) must match the lstat taken before the open. Node has no openat,
// so a directory component swapped between the checks and the open is caught
// by re-checking the components after the open, not prevented outright. The
// seat that owns a root can write there anyway; the checks keep anything
// outside the root from being read through it.
//
// Nothing here writes, renames, creates or migrates a file.
import { lstatSync, openSync, fstatSync, readSync, closeSync, readdirSync, constants } from "node:fs";
import { join, relative, isAbsolute, sep, basename } from "node:path";
export class Refusal extends Error {
constructor(code, message, { reconcile = false } = {}) {
super(message);
this.code = code;
this.reconcile = reconcile;
}
}
const SESSION_NAME = /^[A-Za-z0-9][A-Za-z0-9._:-]*\.jsonl$/;
// Permission errors inside a root are one conversation's problem, not the
// catalogue's: they become a refusal instead of a thrown error.
function denied(err, what) {
if (err.code === "EACCES" || err.code === "EPERM") return new Refusal("unreadable", `${what} is not readable`);
return err;
}
// A directory above the file without search permission is refused the same
// way, so one bad root does not fail the catalogue.
function lstatOrNull(path) {
try {
return lstatSync(path, { bigint: true });
} catch (err) {
if (err.code === "ENOENT" || err.code === "ENOTDIR") return null;
throw denied(err, "a session path component");
}
}
// Every component from the project root down to the sessions directory must
// be a real directory.
export function checkRoot(root) {
const rel = relative(root.projectRoot, root.dir);
if (!rel || rel.startsWith("..") || isAbsolute(rel)) throw new Refusal("unsafe-path", "session root is outside its project");
let path = root.projectRoot;
for (const part of rel.split(sep)) {
path = join(path, part);
const st = lstatOrNull(path);
if (!st) throw new Refusal("unavailable", "session root does not exist");
if (st.isSymbolicLink()) throw new Refusal("unsafe-path", "session root contains a symlink");
if (!st.isDirectory()) throw new Refusal("unsafe-path", "session root is not a directory");
}
}
// Session files directly inside a root, by name. Symlinks and anything that
// is not a regular *.jsonl file are reported, never followed.
export function listSessionFiles(root) {
checkRoot(root);
const files = [], refused = [];
let dirents;
try {
dirents = readdirSync(root.dir, { withFileTypes: true });
} catch (err) {
throw denied(err, "session root");
}
for (const dirent of dirents) {
if (!dirent.name.endsWith(".jsonl")) continue;
if (!SESSION_NAME.test(dirent.name)) refused.push({ name: dirent.name, code: "unsafe-path" });
else if (dirent.isSymbolicLink()) refused.push({ name: dirent.name, code: "unsafe-path" });
else if (dirent.isFile()) files.push(dirent.name);
}
return { files: files.sort(), refused };
}
// Opens one session file read-only. The caller must close the returned fd.
export function openSessionFile(root, name) {
if (typeof name !== "string" || name !== basename(name) || !SESSION_NAME.test(name)) throw new Refusal("unsafe-path", "not a session file name");
checkRoot(root);
const path = join(root.dir, name);
const before = lstatOrNull(path);
if (!before) throw new Refusal("unknown-conversation", "session file no longer exists", { reconcile: true });
if (before.isSymbolicLink()) throw new Refusal("unsafe-path", "session file is a symlink");
if (!before.isFile()) throw new Refusal("unsafe-path", "session file is not a regular file");
let fd;
try {
fd = openSync(path, constants.O_RDONLY | constants.O_NOFOLLOW | constants.O_NONBLOCK);
} catch (err) {
if (err.code === "ELOOP") throw new Refusal("unsafe-path", "session file became a symlink");
if (err.code === "ENOENT") throw new Refusal("unknown-conversation", "session file no longer exists", { reconcile: true });
throw denied(err, "session file");
}
try {
const st = fstatSync(fd, { bigint: true });
if (!st.isFile() || st.dev !== before.dev || st.ino !== before.ino) throw new Refusal("unsafe-path", "session file changed while it was opened");
checkRoot(root);
return { fd, dev: st.dev.toString(), ino: st.ino.toString(), size: Number(st.size) };
} catch (err) {
closeSync(fd);
throw err;
}
}
export function readRange(fd, start, length) {
const buf = Buffer.alloc(length);
let done = 0;
while (done < length) {
const n = readSync(fd, buf, done, length - done, start + done);
if (n === 0) break;
done += n;
}
return done === length ? buf : buf.subarray(0, done);
}
export { closeSync };