Files
ThothII/tools/replay/server.mjs
T
marcopanandClaude Fable 5 c4951e2aa5 chore: hygiene pass — ruff clean, docs storage-model truth, replay /me, failSession log
Audit findings 6.1-6.4 + the audit's remediation plan itself
(docs/superpowers/plans/2026-07-20-full-audit-remediation-plan.md).

- ruff: 34 → 0 (unused imports/f-strings auto-fixed; E702 semicolon lines
  split in test files; one unused local dropped). Suite still 819 green.
- CLAUDE.md + PROJECT_STATE.md no longer claim "no database / settings in
  settings.json": the harness selects filesystem OR PostgreSQL session
  storage (repository.py, server mode), and settings flow through harness
  preferences with the JSON file as fallback only.
- tools/replay: stub /me (SPA boot was parsing the SPA's own HTML as JSON)
  and /runtime/prewarm.
- failSession best-effort persistence now logs its failure server-side
  instead of vanishing.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-20 01:55:41 +02:00

534 lines
22 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 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;
console.error(`[sse] nuova connessione SSE, reset cursore a 0`);
// 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", () => {
console.error(`[sse] connessione chiusa`);
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 ?? {});
// Exit (the reserved "Exit" control on any gate, OR the "Esci" option on the
// final completed widget): stop the replay and ask the frontend to return to
// the landing. The AppShell has an effect that calls stopSession() on
// system_event { event: "session_exit" }. We do NOT advance the cursor.
const isExit =
uiResponse.control === "exit" ||
(typeof uiResponse.id === "string" && uiResponse.id.startsWith(RESTART_ID_PREFIX) && (uiResponse.choices ?? [])[0] === "exit");
if (isExit) {
console.error("[exit] replay fermato, ritorno alla landing");
if (sseClient) {
sendSse(sseClient, "system_event", { type: "system_event", event: "session_exit" });
}
return sendNoContent(res);
}
// Back: the reserved "Go back" control steps ONE gate backwards and
// re-emits it (so the reviewer can revise the previous answer). The cursor
// is clamped at 0, so "back" on the first gate re-presenta the same gate.
// The re-emit is delayed for the same reason the normal advance is: the
// browser's WidgetHost calls setLastUserEntry (which clears stepMessages)
// and clearPending on POST completion, and our ui_request must land AFTER
// those local state updates so the gate actually renders (otherwise the
// spinner stays on: pendingWidget was cleared and the new ui_request races
// the clear, sometimes landing before it).
if (uiResponse.control === "back") {
cursor = Math.max(0, cursor - 1);
const backIdx = cursor;
console.error(`[back] torno al gate ${backIdx + 1}/${GATES.length}`);
setTimeout(() => {
if (sseClient) emitCurrentGate(sseClient);
}, 150);
return sendNoContent(res);
}
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): reset the cursor and re-emit gate 1 in-band. We
// deliberately do NOT close the SSE — see the resume handler above for
// why closing caused a reconnect loop that hung the spinner.
cursor = 0;
if (sseClient) {
console.error("[restart] replay richiesto di nuovo, ri-emetto gate 1 in-band");
emitCurrentGate(sseClient);
}
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 -----------------------------------
// The SPA queries /me at boot (principal for session scoping/admin UI). Falling
// through to the static handler served HTML that apiFetch tried to parse as JSON.
if (method === "GET" && path === "/me") {
return sendJson(res, 200, {
issuer: "replay", subject: "reviewer@replay", displayName: "Replay reviewer", isAdmin: true,
});
}
// Boot-time best-effort warmup: a no-op in replay (no Pi runtime exists).
if (method === "POST" && path === "/runtime/prewarm") {
return sendNoContent(res);
}
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: reset the cursor to 0 and re-emit gate 1 on the current SSE
// connection (if any). We deliberately do NOT close the SSE: the frontend's
// doResume reuses the SAME SSE connection when the session is already active
// (useSessionStream only re-runs on sessionId change), and forcing a close
// here triggered an EventSource auto-reconnect loop that raced with the
// POST /response flow, leaving the spinner stuck after the first gate.
// Re-emitting gate 1 in-band is enough: applyEvent sets pendingWidget on
// any ui_request, so the widget renders and the spinner stops.
if (method === "POST" && path === `/sessions/${REPLAY_SESSION_ID}/resume`) {
cursor = 0;
if (sseClient) {
console.error("[resume] reset cursore a 0, ri-emetto gate 1 in-band");
emitCurrentGate(sseClient);
} else {
console.error("[resume] reset cursore a 0 (nessun client SSE attivo)");
}
// The frontend's runResume reads result.alreadyActive (post-f979ada contract):
// a bare 204 makes it throw and toast "Failed to resume session".
return sendJson(res, 200, { id: REPLAY_SESSION_ID, alreadyActive: Boolean(sseClient) });
}
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`);
});