425 lines
16 KiB
JavaScript
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/core/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}`);
|