Files
ThothII/tools/replay/extract.mjs
T
marcopan c3a3cb8da5 feat(replay): standalone reviewer-gate replay server on :5333
Reproduces the real reviewer UI for any recorded session, with no VPN/Pi/
Python/DWH. The server (node:http, zero deps) serves the built SPA and a
tiny SSE/REST shim that re-emits the reviewer gates captured in a Pi
transcript, in their original order, with the reviewer's real 3-Jul
choices shown as comparison badges.

- tools/replay/extract.mjs: extracts gates from one or many transcripts
  (session-id, file, or directory). Handles sessions split across resume
  re-entries by sorting on message timestamp and dropping unanswered
  gates. Reads the question from session_manifest.yaml, resolving the
  sessions dir from any workspace yaml (no hardcoded paths).
- tools/replay/server.mjs: same-origin :5333. SSE streams gates; POST
  /response advances the cursor and pushes info badges (scelta reale).
  POST /resume and the final "Ripeti/Esci" widget close the SSE so the
  browser EventSource reconnects (cursor resets, gate 1 re-emitted) — the
  replay is re-runnable any number of times. GET /sessions/:id/documents
  reads the real session files so GateArtifactBody resolves file-reference
  artifacts. Exit emits system_event {event:"session_exit"} to return to
  the landing.
- scripts/replay.sh: launcher (extract / build / run / all).
- tools/replay/README.md: data flow, commands, fidelity notes.
- .gitignore: ignore tools/replay/web/ (built artifact, like dist/).
2026-07-05 18:24:26 +02:00

363 lines
14 KiB
JavaScript

// Extract a replay fixture from one or more Pi agent transcripts.
//
// A single ThothII session can span multiple transcript files (Pi re-enters at
// resume, re-emitting any gate the reviewer left unanswered). To reconstruct
// the "what the reviewer actually saw and decided" sequence we:
// 1. Collect reviewer_* toolCalls across ALL transcripts for the session,
// tagged with the message timestamp (NOT the filename — a transcript's
// own timestamp is the resume moment, not the gate moment).
// 2. Sort chronologically by message timestamp.
// 3. Skip any gate without a toolResult (the reviewer closed Pi without
// answering; the gate was re-presented later in a fresher form). What
// remains is exactly the gates the reviewer acted on, in order.
//
// Usage:
// node tools/replay/extract.mjs [<session-id-or-path>]
//
// <session-id> e.g. 2026-07-03-170732-fammi-...
// Scans every *.jsonl in the Pi agent session dir for that
// session id. Default = the 17:07 full session.
// <path> a transcript file OR a directory of transcripts.
//
// Output: tools/replay/replay.json (descriptor per gate + the reviewer's real
// choice recovered from the toolResult text).
import { readFileSync, writeFileSync, readdirSync, existsSync, statSync } from "node:fs";
import { join, dirname } from "node:path";
import { homedir } from "node:os";
const RESERVED = ["back", "exit", "other"];
const SCHEMA_VERSION = 1;
const PI_SESSIONS_DIR = join(
homedir(),
".pi/agent/sessions/--Users-mp-projects-ThothII-harness--",
);
// --- resolve which transcript files to read --------------------------------
function resolveTranscriptFiles(arg) {
if (!arg) {
// default: the 17:07 full session
return readdirSync(PI_SESSIONS_DIR)
.filter((f) => f.endsWith(".jsonl") && f.startsWith("2026-07-03T17-07-32"))
.map((f) => join(PI_SESSIONS_DIR, f))
.sort();
}
if (existsSync(arg)) {
const s = statSync(arg);
if (s.isDirectory()) {
return readdirSync(arg)
.filter((f) => f.endsWith(".jsonl"))
.map((f) => join(arg, f))
.sort();
}
return [arg];
}
// Treat as a session id: scan the Pi session dir for files mentioning it.
// (Cheap prefilter by substring to avoid parsing every historical transcript.)
const hits = [];
for (const f of readdirSync(PI_SESSIONS_DIR)) {
if (!f.endsWith(".jsonl")) continue;
const p = join(PI_SESSIONS_DIR, f);
// Quick text scan; the session id appears in toolCall arguments.
const head = readFileSync(p, "utf8");
if (head.includes(arg)) hits.push(p);
}
if (hits.length === 0) {
throw new Error(
`No transcripts found for "${arg}". Pass a session id, a transcript file, or a directory.`,
);
}
return hits.sort();
}
const files = resolveTranscriptFiles(process.argv[2]);
// --- first pass: gather reviewer calls (with msg timestamp) + toolResults --
let sessionId = null;
let createdAt = null;
const calls = []; // {callId, name, args, ts, file}
const resultsByCallId = new Map(); // callId -> toolResult text
for (const fn of files) {
const raw = readFileSync(fn, "utf8").split("\n").filter(Boolean);
for (const line of raw) {
let m;
try {
m = JSON.parse(line);
} catch {
continue;
}
if (m.type === "session" && !createdAt) createdAt = m.timestamp;
if (m.type !== "message") continue;
const msg = m.message;
if (!msg || !Array.isArray(msg.content)) continue;
const ts = m.timestamp ?? null;
for (const p of msg.content) {
if (p.type !== "toolCall") continue;
if (typeof p.name !== "string" || !p.name.startsWith("reviewer_")) continue;
calls.push({ callId: p.id, name: p.name, args: p.arguments ?? {}, ts, file: fn });
}
if (msg.role === "toolResult" && typeof msg.toolCallId === "string") {
if (typeof msg.toolName === "string" && msg.toolName.startsWith("reviewer_")) {
const c = Array.isArray(msg.content) ? msg.content[0] : null;
const text = c && typeof c.text === "string" ? c.text : "";
// First writer wins: a toolResult is unique per callId across files.
if (!resultsByCallId.has(msg.toolCallId)) {
resultsByCallId.set(msg.toolCallId, text);
}
}
}
}
}
// --- chronological order, then drop gates the reviewer never answered -------
// A gate without a toolResult = the reviewer closed Pi at that gate; it was
// re-presented (possibly reworded) at the next resume. Keeping only answered
// gates yields the sequence the reviewer actually experienced end-to-end.
calls.sort((a, b) => (a.ts ?? "").localeCompare(b.ts ?? ""));
// When reading a whole directory, calls may span multiple sessions. Pick the
// session with the most reviewer calls (the "subject" of that directory) and
// keep only its gates — otherwise replay.json would interleave two sessions.
const sessionCounts = new Map();
for (const c of calls) {
const s = c.args.session ?? "(none)";
sessionCounts.set(s, (sessionCounts.get(s) ?? 0) + 1);
}
if (sessionCounts.size > 1) {
let best = null;
let bestN = -1;
for (const [s, n] of sessionCounts) if (n > bestN) { best = s; bestN = n; }
if (best) {
const filtered = calls.filter((c) => (c.args.session ?? "(none)") === best);
console.error(` directory mode: ${sessionCounts.size} sessions found, keeping "${best}" (${filtered.length}/${calls.length} calls)`);
calls.splice(0, calls.length, ...filtered);
if (sessionId && sessionId !== best) sessionId = best;
}
}
const answered = calls.filter((c) => resultsByCallId.has(c.callId));
// --- build descriptors (mirrors harness/.pi/extensions/gate/builders.js) ----
function buildDescriptor(name, args, idx) {
const id = `u${idx}`;
const phase = "replay";
if (name === "reviewer_select") {
const opts = parseOptions(args.options).filter((o) => !isReserved(o.label));
const recommendedId = parseOptions(args.options).find((o) => o.recommended)?.id ?? null;
return {
type: "ui_request",
id,
phase,
schema_version: SCHEMA_VERSION,
widget: "select",
title: args.title,
intro: args.intro ?? null,
recommended: recommendedId,
options: opts.map((o) => ({ id: o.id, label: o.label, recommended: !!o.recommended })),
reserved: RESERVED,
};
}
if (name === "reviewer_decide") {
const opts = parseOptions(args.options).filter((o) => !isReserved(o.label));
return {
type: "ui_request",
id,
phase,
schema_version: SCHEMA_VERSION,
widget: "multiselect",
title: args.title,
allow_empty: args.allow_empty ?? false,
options: opts.map((o) => ({
id: o.id,
label: o.label,
recommended: !!o.recommended,
decision: o.decision,
})),
selected: [],
content: null,
reserved: RESERVED,
};
}
if (name === "reviewer_confirm") {
return {
type: "ui_request",
id,
phase,
schema_version: SCHEMA_VERSION,
widget: "artifact-gate",
title: args.title,
artifact: args.artifact,
action: { kind: "approve_reject", prompt: "Approvi o rifiuti?" },
options: [
{ id: "approve", label: "Salva e procedi", recommended: true },
{ id: "reject", label: "Rifiuta" },
],
reserved: RESERVED,
};
}
return null;
}
function isReserved(label) {
if (typeof label !== "string") return false;
const l = label.toLowerCase();
return RESERVED.some((r) => l === r || l.startsWith(r + ":") || l.startsWith(r + " —"));
}
// Some models serialize arrays as JSON strings; the gate's prepareReviewerArguments
// coerces them back to arrays pre-validation. Mirror that here so a stringified
// `options` (or `names`) doesn't blow up the descriptor builder.
function parseOptions(v) {
if (Array.isArray(v)) return v;
if (typeof v === "string") {
try {
const p = JSON.parse(v);
if (Array.isArray(p)) return p;
} catch { /* not JSON */ }
}
return [];
}
// --- recover the real reviewer choice from the toolResult text -------------
function parseRealChoice(name, descriptor, resultText) {
if (!resultText) return { labels: [], note: "no toolResult captured" };
const t = resultText.trim();
const askOnly = t.match(/^Scelta del reviewer:\s*(.+?)\s*\.\s*$/);
if (name === "reviewer_select") {
let label = null;
const m = t.match(/\):\s*(.+?)\s*\.\s*$/);
if (m) label = m[1];
else if (askOnly) label = askOnly[1];
return { labels: label ? [label] : [], note: t };
}
if (name === "reviewer_decide") {
const m = t.match(/Registrate\s+(\d+)\s+decisioni:\s*(.+?)\.\s*$/);
if (m) return { labels: [], decisionTypes: m[2].split(/,\s*/), note: t };
if (/Nessuna decisione registrata/.test(t)) return { labels: [], decisionTypes: [], note: t };
return { labels: [], note: t };
}
if (name === "reviewer_confirm") {
if (/ERRORE/i.test(t)) return { labels: [], note: t, failed: true };
if (/Nessun CTE in attesa/i.test(t)) return { labels: [], note: t, failed: true };
if (/approvato|approvata/i.test(t)) return { labels: ["Salva e procedi"], note: t };
if (/Rifiutato/i.test(t)) return { labels: ["Rifiuta"], note: t };
return { labels: [], note: t };
}
return { labels: [], note: t };
}
// --- assemble the gates array ----------------------------------------------
// sessionId = the dominant session across the surviving calls (post-filter).
sessionId = sessionId ?? null;
{
const sc = new Map();
for (const c of answered) sc.set(c.args.session, (sc.get(c.args.session) ?? 0) + 1);
let best = null, bn = -1;
for (const [s, n] of sc) if (n > bn) { best = s; bn = n; }
if (best) sessionId = best;
}
const gates = answered.map((c, idx) => {
const descriptor = buildDescriptor(c.name, c.args, idx);
const real = parseRealChoice(c.name, descriptor, resultsByCallId.get(c.callId));
return { callId: c.callId, toolName: c.name, ts: c.ts, descriptor, real };
});
// --- session question: prefer the manifest, fall back to question.md -------
function readSessionQuestion(sid) {
if (!sid) return null;
const sessionDir = findSessionDir(sid);
if (!sessionDir) return null;
// session_manifest.yaml has a top-level `question:` field with the original
// NL question. Cheaper and more authoritative than parsing question.md.
const manifestPath = join(sessionDir, "session_manifest.yaml");
if (existsSync(manifestPath)) {
const m = readFileSync(manifestPath, "utf8");
const q = m.match(/^question:\s*(.+?)\s*$/m);
if (q) return stripYaml(q[1]);
}
const qmd = join(sessionDir, "question.md");
if (existsSync(qmd)) {
const md = readFileSync(qmd, "utf8").replace(/^#.*\n+/, "").trim();
return md.split(/\n\s*\n/)[0].trim();
}
return null;
}
// Resolve a session id to its on-disk directory. We check every workspace yaml
// in harness/workspaces/ for a `paths.sessions` entry (relative OR absolute),
// then look for <sessions>/<session-id>/ in each. Workspaces are symlinks to
// per-client repos (uncommitted), so this works for any client without env.
function findSessionDir(sid) {
const wsDir = join(process.cwd(), "harness/workspaces");
const candidates = [];
if (existsSync(wsDir)) {
for (const f of readdirSync(wsDir)) {
if (!/\.(ya?ml)$/i.test(f)) continue;
const ypath = join(wsDir, f);
let y;
try { y = readFileSync(ypath, "utf8"); } catch { continue; }
const m = y.match(/^paths:\s*\n(?:[ \t]+.*\n)*?[ \t]+sessions:\s*(\S+)/m);
if (!m) continue;
let p = stripYaml(m[1]);
if (!p) continue;
// Strip ${VAR} placeholders (unresolvable here) — those workspaces can't
// be located offline and are skipped.
if (p.includes("${")) continue;
candidates.push(p);
}
}
// Plus the legacy hardcoded PSD location as a last resort.
candidates.push(join(homedir(), "projects", "tht-workspace-psd", "sessions"));
for (const base of candidates) {
const abs = join(base, sid);
if (existsSync(abs)) return abs;
}
return null;
}
function stripYaml(s) {
// drop surrounding quotes + trailing comment
let v = s.trim().replace(/^['"]|['"]$/g, "");
v = v.replace(/\s+#.*$/, "");
return v;
}
const replay = {
source: {
transcripts: files,
extractedAt: new Date().toISOString(),
note: "gates without a reviewer answer (resume re-entries) are dropped",
},
session: {
id: sessionId,
question: readSessionQuestion(sessionId),
created_at: createdAt,
},
stats: {
transcriptFiles: files.length,
totalReviewerCalls: calls.length,
answeredGates: answered.length,
droppedUnanswered: calls.length - answered.length,
},
gates,
};
const outPath = new URL("./replay.json", import.meta.url);
writeFileSync(outPath, JSON.stringify(replay, null, 2) + "\n", "utf8");
// --- coverage summary ------------------------------------------------------
const byTool = {};
let withChoice = 0;
for (const g of gates) {
byTool[g.toolName] = (byTool[g.toolName] ?? 0) + 1;
if (g.real.labels.length || (g.real.decisionTypes && g.real.decisionTypes.length)) withChoice++;
}
console.error(`scanned ${files.length} transcript file(s)`);
console.error(` reviewer calls: ${calls.length} | answered: ${answered.length} | dropped (no answer): ${calls.length - answered.length}`);
console.error(` by tool (answered):`, byTool);
console.error(` with a real choice recovered: ${withChoice}/${gates.length}`);
console.error(` session: ${sessionId ?? "(unknown)"}`);
console.error(` question: ${(replay.session.question ?? "(not found)").slice(0, 90)}`);
console.error(`wrote ${outPath.pathname}`);