Files
marcopanandClaude Fable 5 87e875bc81 feat(replay): extract reviewer_schema_linking gates + de-slug question fallback
extract.mjs now emits a schema-linking ui_request descriptor (columns rendered
from the model's suggested_columns, since the live catalog enrichment needs the
DWH) and parses the reviewer's "N tabelle" outcome; when neither the manifest
nor question.md is reachable, the NL question is approximated by de-slugging
the session id. replay.json regenerated from the 2026-07-06-175012 session
(18 answered gates, includes the F4 schema-linking gate).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-07 16:04:00 +02:00

425 lines
16 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) ----
// Infer the workflow phase (F1..F8) of a gate from its title, so the frontend's
// WorkflowBar lights up the correct dot during replay. Patterns (from real
// session titles):
// "F1 — Chiarimento 1/3: …" → F1
// "Fase 3 completata — passo alla → F3 (the gate that CLOSES phase 3)
// Fase 4 (Schema linking)?"
// "F6 — CTE 2 approvato: …" → F6
// Returns null when no phase can be inferred (the frontend keeps the previous
// phase in that case, which is fine for transitional gates).
function inferPhase(title) {
if (typeof title !== "string") return null;
let m = title.match(/Fase\s+(\d+)\s+completata/i);
if (m) return `F${m[1]}`;
m = title.match(/\bF([1-8])\b/);
if (m) return `F${m[1]}`;
return null;
}
function buildDescriptor(name, args, idx) {
const id = `u${idx}`;
const phase = inferPhase(args.title) ?? "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,
};
}
// F4 schema-linking gate (column curation, v2). The model supplies tables with
// kind/rationale/suggested_columns; the live gate enriches columns from the
// catalog via `tht schema columns` (needs DWH), so here we render only what the
// model provided — columns come from suggested_columns as suggested=true.
if (name === "reviewer_schema_linking") {
const tables = parseOptions(args.tables).map((t) => ({
id: t.id,
name: t.name,
kind: t.kind ?? "promote",
recommended: !!t.recommended,
description: t.description ?? null,
rationale: t.rationale ?? null,
columns: (parseOptions(t.suggested_columns).map((c) =>
typeof c === "string" ? { name: c, description: null, suggested: true } : c,
)),
}));
return {
type: "ui_request",
id,
phase,
schema_version: SCHEMA_VERSION,
widget: "schema-linking",
title: args.title,
tables,
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 };
}
if (name === "reviewer_schema_linking") {
// "Schema linking registrato dal reviewer (6 tabelle + colonne curate)."
const m = t.match(/registrato dal reviewer \((\d+)\s+tabelle/);
return { labels: m ? [`${m[1]} tabelle`] : [], 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;
}
// Last-resort fallback when neither the manifest nor question.md can be found
// (e.g. the workspace repo is not checked out here). `tht session new` derives
// the session id slug from the question text, so de-slugging it recovers a
// readable approximation of the original NL question.
function questionFromSessionId(sid) {
if (typeof sid !== "string") return null;
const slug = sid.replace(/^\d{4}-\d{2}-\d{2}-\d{6}-/, "");
if (!slug) return null;
return slug.replace(/-/g, " ");
}
// 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) ?? questionFromSessionId(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}`);