// 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`); });