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/).
This commit is contained in:
2026-07-05 18:24:26 +02:00
parent 2f68b0d109
commit c3a3cb8da5
6 changed files with 2147 additions and 0 deletions
+480
View File
@@ -0,0 +1,480 @@
// Standalone replay server: serves the built SPA + a tiny SSE/REST shim that
// replays the recorded reviewer gates of a logged session, no VPN/Pi/DWH needed.
//
// Same-origin on :5333. The SPA bundle is built once with
// VITE_BACKEND_URL=http://localhost:5333 (see scripts/replay.sh), so all REST
// + SSE calls land on this server and we control the entire interaction loop.
//
// Run: node tools/replay/server.mjs (after `scripts/replay.sh build`)
import { createServer } from "node:http";
import { readFile, stat } from "node:fs/promises";
import { existsSync, readdirSync, readFileSync } from "node:fs";
import { fileURLToPath } from "node:url";
import { dirname, join, extname, normalize } from "node:path";
import { homedir } from "node:os";
import { createRequire } from "node:module";
const PORT = Number(process.env.PORT ?? 5333);
const HERE = dirname(fileURLToPath(import.meta.url));
const WEB_DIR = join(HERE, "web");
const REPLAY_PATH = join(HERE, "replay.json");
const require = createRequire(import.meta.url);
const replay = JSON.parse(await readFile(REPLAY_PATH, "utf8"));
const REPLAY_SESSION_ID = replay.session.id;
const GATES = replay.gates;
if (!existsSync(WEB_DIR)) {
console.error(`Missing ${WEB_DIR}. Run: scripts/replay.sh build`);
process.exit(1);
}
// --- replay state (single-user, single-session) --------------------------
// One cursor advances across the gate sequence. The SSE connection streams the
// gate at the cursor; a POST /response advances the cursor and the next gate is
// pushed on the same SSE stream. A fresh SSE connection (browser refresh, or
// clicking the session again after completion) resets the cursor to 0, so the
// replay can be walked any number of times.
let cursor = 0;
let sseClient = null; // the active SSE response, if any
function sendSse(res, eventName, data) {
res.write(`event: ${eventName}\ndata: ${JSON.stringify(data)}\n\n`);
}
function currentGate() {
return { index: cursor, gate: GATES[cursor] ?? null };
}
// Replay-completion marker id prefix. The final ui_request carries this id so
// the response handler can recognise "the user wants to restart" and close the
// SSE (browser EventSource reconnects → cursor resets → gate 1 re-emitted).
const RESTART_ID_PREFIX = "replay-restart-";
function emitCurrentGate(res) {
const { index, gate } = currentGate();
if (!gate) {
// End of replay: emit a final ui_request as a one-button select. This
// matters because the reducer sets pendingWidget on ANY ui_request, and
// pendingWidget!=null is exactly what stops the working spinner. Using a
// select (rather than widget=info) gives the user an in-place "Ripeti il
// replay" button: clicking it POSTs a response, which we recognise via the
// RESTART_ID_PREFIX and handle by closing the SSE — the browser EventSource
// then auto-reconnects, the cursor resets to 0, and gate 1 is re-emitted.
// We deliberately do NOT mark the session "finalized": that status makes
// the real backend refuse resume with 409, and we want the replay to be
// re-runnable any number of times.
const id = `${RESTART_ID_PREFIX}${Date.now()}`;
sendSse(res, "ui_request", {
type: "ui_request",
ui_request: {
id,
widget: "select",
title: `✓ Replay completato — ${GATES.length}/${GATES.length} gate riprodotti`,
intro: `Hai attraversato tutti i gate reviewer della sessione. Scegli cosa fare:`,
// Two real options so SelectWidget renders them as equal-weight buttons
// (a reserved "exit" would render as a small muted link instead). The
// response handler tells them apart by ui_response.choices[0].
options: [
{ id: "restart", label: "↻ Ripeti il replay", recommended: true },
{ id: "exit", label: "■ Esci" },
],
},
});
console.error(`[done] replay completato (${GATES.length} gate)`);
return;
}
sendSse(res, "ui_request", {
type: "ui_request",
ui_request: gate.descriptor,
});
console.error(
`[${index + 1}/${GATES.length}] -> ${gate.toolName}: ${gate.descriptor.title}`,
);
}
// --- helpers --------------------------------------------------------------
const MIME = {
".html": "text/html; charset=utf-8",
".js": "text/javascript; charset=utf-8",
".css": "text/css; charset=utf-8",
".json": "application/json; charset=utf-8",
".svg": "image/svg+xml",
".png": "image/png",
".jpg": "image/jpeg",
".ico": "image/x-icon",
".woff2": "font/woff2",
".woff": "font/woff",
".map": "application/json; charset=utf-8",
};
async function serveStatic(req, res, urlPath) {
let p = normalize(join(WEB_DIR, urlPath));
if (!p.startsWith(WEB_DIR)) {
res.writeHead(403);
res.end("forbidden");
return;
}
// Directory requests (incl. "/") and missing files fall back to index.html
// (SPA: client-side routing handles all paths under /).
let isDir = false;
try { isDir = (await stat(p)).isDirectory(); } catch { /* missing */ }
if (isDir || !existsSync(p)) {
p = join(WEB_DIR, "index.html");
}
try {
const data = await readFile(p);
res.writeHead(200, { "Content-Type": MIME[extname(p)] ?? "application/octet-stream" });
res.end(data);
} catch {
res.writeHead(404);
res.end("not found");
}
}
function sendJson(res, code, obj) {
const body = JSON.stringify(obj);
res.writeHead(code, {
"Content-Type": "application/json; charset=utf-8",
"Content-Length": Buffer.byteLength(body),
});
res.end(body);
}
function sendNoContent(res) {
res.writeHead(204);
res.end();
}
// Describe what the reviewer actually picked on 3 Jul, for the info badge.
function describeRealChoice(gate) {
const r = gate.real;
if (r.failed) return { level: "warning", text: `⚠ Gate fallito nella sessione reale: ${r.note}` };
if (r.note && /Fase memoria vuota|avanzamento automatico/i.test(r.note)) {
return { level: "info", text: `ℹ Nella sessione reale questo gate fu saltato (memorie vuote, avanzamento automatico).` };
}
const labels = r.labels ?? [];
const types = r.decisionTypes ?? [];
if (labels.length === 0 && types.length === 0) {
return { level: "info", text: `ℹ Scelta reale non campionata dal log: ${r.note ?? "—"}` };
}
const parts = [];
if (labels.length) parts.push(labels.join(", "));
if (types.length) parts.push(`decisioni: ${types.join(", ")}`);
return { level: "warning", text: `📌 SCELTA REALE (3 lug 2026): ${parts.join(" — ")}` };
}
// Compare the user's click to the real choice. Approve/Reject and select options
// use labels; decide uses decisionTypes (loose match — see compareChoices).
function compareUserChoice(gate, uiResponse) {
const r = gate.real;
const realLabels = new Set((r.labels ?? []).map((s) => s.toLowerCase()));
const realTypes = new Set((r.decisionTypes ?? []));
const choiceIds = new Set(uiResponse.choices ?? []);
const opts = gate.descriptor.options ?? [];
// select: choice id -> option label
if (gate.descriptor.widget === "select") {
const picked = opts.find((o) => choiceIds.has(o.id));
if (!picked || realLabels.size === 0) return null;
const match = [...realLabels].some((rl) => picked.label.toLowerCase().includes(rl.split(" (")[0].toLowerCase()));
return match ? "✓ Hai scelto come il 3 luglio." : null;
}
// artifact-gate: approve/reject
if (gate.descriptor.widget === "artifact-gate") {
if (choiceIds.has("approve") && realLabels.has("salva e procedi")) {
return "✓ Hai scelto come il 3 luglio (Salva e procedi).";
}
if (choiceIds.has("reject") && realLabels.has("rifiuta")) {
return "✓ Hai scelto come il 3 luglio (Rifiuta).";
}
return null;
}
// multiselect: count of picked options vs number of real decision types
if (gate.descriptor.widget === "multiselect") {
if (realTypes.size === 0) return null;
const pickedCount = opts.filter((o) => choiceIds.has(o.id)).length;
return pickedCount === realTypes.size
? `✓ Stesso numero di scelte del 3 luglio (${pickedCount}).`
: null;
}
return null;
}
// --- HTTP routing ---------------------------------------------------------
const server = createServer(async (req, res) => {
const url = new URL(req.url, `http://localhost:${PORT}`);
const path = url.pathname;
const method = req.method;
// CORS preflight (not needed same-origin, but harmless).
res.setHeader("Access-Control-Allow-Origin", req.headers.origin ?? "*");
res.setHeader("Access-Control-Allow-Credentials", "true");
res.setHeader("Access-Control-Allow-Headers", "Content-Type");
res.setHeader("Access-Control-Allow-Methods", "GET,POST,PUT,DELETE,OPTIONS");
if (method === "OPTIONS") {
res.writeHead(204);
return res.end();
}
// --- SSE stream (the heart of the replay) -------------------------------
if (method === "GET" && path === `/sessions/${REPLAY_SESSION_ID}/events`) {
res.writeHead(200, {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
});
sseClient = res;
// A fresh SSE connection restarts the replay from the beginning: this is
// what makes the session re-runnable after completion, and what makes a
// browser refresh re-walk the gates. (Single-user, single-session.)
cursor = 0;
req.on("close", () => {
if (sseClient === res) sseClient = null;
});
emitCurrentGate(res);
return;
}
// --- reviewer response: advance the cursor + push next gate on SSE ------
if (method === "POST" && path === `/sessions/${REPLAY_SESSION_ID}/response`) {
const body = await readBody(req);
const uiResponse = (body?.ui_response ?? {});
const isRestartId = typeof uiResponse.id === "string" && uiResponse.id.startsWith(RESTART_ID_PREFIX);
if (isRestartId) {
// The final "completed" select carries a RESTART_ID_PREFIX id with two
// options: "restart" (re-run the loop) and "exit" (park the session).
// Restart closes the SSE so the browser EventSource reconnects → cursor
// resets → gate 1. Exit keeps the SSE open and emits a parked info widget
// (no spinner, no replay); the user resumes via the sidebar menu.
const choice = (uiResponse.choices ?? [])[0];
if (choice === "exit") {
// Emit a system_event the frontend interprets as "leave the live session
// view and return to the landing". The AppShell has an effect that calls
// stopSession() on system_event { event: "session_exit" }, which resets
// the active session and shows the landing again. The session stays in
// the sidebar list with status "open", so the user can click it to open
// the documents panel and press Resume there to restart the replay.
console.error("[exit] replay fermato, ritorno alla landing");
if (sseClient) {
sendSse(sseClient, "system_event", {
type: "system_event",
event: "session_exit",
});
}
return sendNoContent(res);
}
// Restart (default): close SSE to reconnect.
if (sseClient) {
try { sseClient.end(); } catch { /* already closed */ }
sseClient = null;
}
console.error("[restart] replay richiesto di nuovo dall'utente");
return sendNoContent(res);
}
const { index, gate } = currentGate();
const isLast = index === GATES.length - 1;
// Ack the POST first. Then push info badges + the next gate. The delay lets
// the browser clear its `stepMessages` (reset by WidgetHost's setLastUserEntry
// on POST completion) BEFORE our info badges arrive, so the "scelta reale"
// badge survives. On the LAST gate we skip the delay: the final "completed"
// ui_request (widget=select) must land ASAP to set pendingWidget and stop
// the spinner; setLastUserEntry does NOT reset pendingWidget, so order is safe.
const delay = isLast ? 0 : 150;
setTimeout(() => {
if (!gate) return;
logChoice(gate, uiResponse);
if (sseClient) {
const realDesc = describeRealChoice(gate);
sendSse(sseClient, "info", { type: "info", level: realDesc.level, text: realDesc.text });
const match = compareUserChoice(gate, uiResponse);
if (match) sendSse(sseClient, "info", { type: "info", level: "info", text: match });
}
cursor += 1;
if (sseClient) emitCurrentGate(sseClient);
}, delay);
return sendNoContent(res);
}
// --- REST stubs the SPA needs to boot -----------------------------------
if (method === "GET" && path === "/health") {
return sendJson(res, 200, { ok: true, replay: true, gates: GATES.length });
}
if (method === "GET" && path === "/workspaces") {
return sendJson(res, 200, [{ name: "psd (replay)", file: "psd.yaml" }]);
}
if (method === "GET" && path === "/models") {
return sendJson(res, 200, { models: [{ id: "glm-5.2", provider: "zai", label: "GLM 5.2 (replay)" }] });
}
if (method === "GET" && path === "/settings") {
return sendJson(res, 200, {
workspace: "psd",
provider: "zai",
model: "glm-5.2",
thinking: "medium",
});
}
if (method === "GET" && path === "/sessions") {
// Status stays "open" forever: the replay must remain re-runnable, and the
// real backend refuses resume with 409 when status is "finalized".
return sendJson(res, 200, [
{
id: REPLAY_SESSION_ID,
status: "open",
question: replay.session.question ?? "(replay)",
summary: "Replay 3 lug — 20 gate reviewer",
created_at: replay.session.created_at ?? "2026-07-03T17:07:32Z",
updated_at: null,
author: "replay",
name: "Replay 170732",
group: null,
archived: false,
},
]);
}
if (method === "GET" && path === `/sessions/${REPLAY_SESSION_ID}`) {
return sendJson(res, 200, {
id: REPLAY_SESSION_ID,
status: "open",
phase: 1,
question: replay.session.question ?? "",
});
}
if (method === "GET" && path === `/sessions/${REPLAY_SESSION_ID}/documents`) {
return sendJson(res, 200, buildSessionDocuments(REPLAY_SESSION_ID));
}
// Resume: the frontend's doResume reuses the SAME SSE connection when the
// session is already active (the useSessionStream effect only re-runs on
// sessionId change), so simply resetting the cursor wouldn't reach the
// client. We forcibly CLOSE the current SSE connection: the browser's
// EventSource auto-reconnects, opening a fresh stream that resets the cursor
// to 0 and re-emits gate 1. This makes "resume immediately after the last
// gate" work without a hard refresh.
if (method === "POST" && path === `/sessions/${REPLAY_SESSION_ID}/resume`) {
if (sseClient) {
try { sseClient.end(); } catch { /* already closed */ }
sseClient = null;
}
return sendNoContent(res);
}
if (method === "POST" && /^\/sessions\/[^/]+\/(close|steer|rename|group|archive|unarchive)$/.test(path)) {
return sendNoContent(res);
}
if (method === "POST" && path === "/sessions") {
// New-session form: redirect into the replay instead.
return sendJson(res, 200, { id: REPLAY_SESSION_ID });
}
if (method === "DELETE" && path === `/sessions/${REPLAY_SESSION_ID}`) {
return sendNoContent(res);
}
if (method === "PUT" && path === "/settings") {
return sendNoContent(res);
}
// --- static SPA fallback ------------------------------------------------
if (method === "GET") {
return serveStatic(req, res, path);
}
sendNoContent(res);
});
function logChoice(gate, uiResponse) {
const opts = gate.descriptor.options ?? [];
const choiceIds = uiResponse.choices ?? [];
const picked = opts.filter((o) => choiceIds.includes(o.id)).map((o) => o.label);
const real = gate.real;
const realDesc = real.labels?.length
? real.labels.join(", ")
: real.decisionTypes?.length
? `${real.decisionTypes.length} decisioni: ${real.decisionTypes.join(", ")}`
: real.failed
? "FALLITO"
: "—";
console.error(
` <- user picked: ${picked.join(" | ") || uiResponse.control || "(empty)"} | reale: ${realDesc}`,
);
}
function readBody(req) {
return new Promise((resolve) => {
let raw = "";
req.on("data", (c) => (raw += c));
req.on("end", () => {
if (!raw) return resolve({});
try {
resolve(JSON.parse(raw));
} catch {
resolve({});
}
});
});
}
// --- session documents ----------------------------------------------------
// Mirror of harness/tht/session/store.py build_documents, so the frontend's
// GateArtifactBody can resolve file-reference artifacts (e.g. {file:"question.md"})
// to their full content via GET /sessions/:id/documents.
const DOC_SPEC = [
["question.md", "F3", "revised_question", "Revised question", "markdown"],
["schema_linking.json", "F4", "schema_linking", "Schema linking", "schema-linking"],
["sql_final.sql", "F7", "sql", "Final SQL", "sql"],
["validation_report.md", "finalize", "validation_report", "Validation report", "markdown"],
["review_decisions.jsonl", "—", "decisions", "Decisions", "decisions"],
];
function buildSessionDocuments(sid) {
const dir = findSessionDir(sid);
const docs = [{
phase: "—",
key: "question",
title: "Original question",
format: "text",
content: replay.session.question ?? "",
}];
if (!dir) return docs;
for (const [filename, phase, key, title, fmt] of DOC_SPEC) {
const p = join(dir, filename);
if (existsSync(p)) {
docs.push({ phase, key, title, format: fmt, content: readFileSync(p, "utf8") });
}
}
return docs;
}
// Resolve a session id to its on-disk directory: check every workspace yaml in
// harness/workspaces/ for paths.sessions (absolute or relative), then the legacy
// PSD location. Same logic as extract.mjs findSessionDir.
function findSessionDir(sid) {
const repoRoot = join(HERE, "..", "..");
const wsDir = join(repoRoot, "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 = m[1].replace(/^['"]|['"]$/g, "").replace(/\s+#.*$/, "");
if (!p || p.includes("${")) continue;
candidates.push(p);
}
}
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;
}
server.listen(PORT, "127.0.0.1", () => {
console.error(`\n▶ ThothII replay: http://localhost:${PORT}`);
console.error(` sessione: ${REPLAY_SESSION_ID}`);
console.error(` domanda: ${(replay.session.question ?? "").slice(0, 100)}…`);
console.error(` ${GATES.length} gate reviewer da attraversare.`);
console.error(` Ctrl-C per uscire.\n`);
});