// tht-gate.js -- Pi extension: the HITL gate for the ThothII NL->SQL workflow. // // REWRITE of the reference implementation's tht-gate.js. Two changes vs the source: // (1) F2 single source: workflow facts (max_phase, phase names, the schema-linking // phase) come from `tht phase meta --json`, NOT from JS-mirrored constants. // the reference implementation PHASE_NAMES array (truncated to 7) is gone; F8/datamart can no // longer drift out of sync. // (2) D2/D4 widget-descriptor: the reviewer interaction is emitted as a // widget-descriptor JSON (built by ./gate/builders.js) and awaited by id, instead // of rendered by blocking native TUI primitives (ctx.ui.select/custom). // // PRESERVED VERBATIM from the source (load-bearing runtime glue, spec D4): // - the anti-bypass tool_call hook (FORBIDDEN + PROTECTED_FILES) // - the input lock + the `input` hook (entry detection + free-input block + `!` steer) // - the kickoff injection (before_agent_start) + the two kickoff payloads // - the agent_end prose safety net // - the exit-code contracts with the CLI (5 = gate refusal, 6 = needs human / // auto-advance not ready -> silent no-op) // - textResult / tht() / relayIfThtFails / advanceIfReady helpers // // TESTING: the pure builders are L1-tested (./gate/__tests__/). This file is the // GLUE -- it depends on the Pi runtime (pi.on, pi.registerTool, ctx.sendRaw) and is // verified end-to-end at L2 (Task D4). It is intentionally NOT unit-tested here; a // fake-Pi runtime mock (cross-cutting follow-up) would let it run in CI. import { execFileSync } from "node:child_process"; import { Type } from "typebox"; import { buildSelectRequest, buildMultiselectRequest, buildArtifactGate, } from "./gate/builders.js"; import { isReserved, stripReserved } from "./reserved-labels.mjs"; // --- anti-bypass block lists (spec D4, verbatim from source L169-177) ----------- const FORBIDDEN = [ /\btht\s+phase\s+(advance|reopen)\b/, /\btht\s+decision\s+add\b/, /\btht\s+cte\s+plan\b/, ]; const PROTECTED_FILES = /(review_decisions\.jsonl|session_manifest\.yaml|cte_plan\.json)/; // --- kickoff payloads (verbatim from source L184-212, load-bearing model prose) - const NUOVA_DOMANDA_KICKOFF = "Istruzioni operative — nuova sessione ThothII (workflow human-in-the-middle: tu " + "orchestri, il reviewer decide, la CLI `tht` persiste; NON sei in modalita' autonoma).\n" + '1. Esegui `tht session new ""` e annota ' + "l'id stampato nell'ultima riga.\n" + "2. Carica la skill leggendo il file con il tool `read`: `.pi/skills/tht-sessione/SKILL.md` " + "(NON come comando di shell ne' come `/skill:...`). Poi segui il suo workflow dalla Fase 1 " + "(Chiarimento), usando l'id di sessione in ogni comando `tht`.\n" + "Se la skill non si carica (per qualsiasi motivo): FERMATI. Non proseguire da solo, non " + "improvvisare analisi o query. Comunica al reviewer che la skill tht-sessione non e' " + "disponibile e attendi istruzioni.\n" + "Regole non negoziabili (valgono SEMPRE, anche senza la skill):\n" + "- Una domanda al reviewer per volta; attendi la sua risposta prima di proseguire.\n" + "- MAI promuovere, escludere, correggere o applicare alcunche' senza conferma esplicita.\n" + "- Le interazioni col reviewer passano dai tool reviewer_select/reviewer_decide/" + "reviewer_confirm (widget-descriptor); il reviewer invia testo libero prefissando '!'. " + "NON eseguire mai `tht phase advance|reopen` ne' `tht decision add` da shell.\n" + "- Procedi una fase alla volta: a fine fase cedi il turno al reviewer, non incatenare le fasi."; const RIPRENDI_KICKOFF = "Istruzioni operative — ripresa di una sessione ThothII esistente (id nel messaggio sopra).\n" + "1. Esegui `tht session show ` e leggi: stato, domanda, decisioni registrate, presenza " + "di schema_linking.json.\n" + "2. Carica la skill leggendo `.pi/skills/tht-sessione/SKILL.md` con il tool `read` (non come " + "comando di shell ne' `/skill:...`).\n" + "3. Determina l'ultima fase completata dai fatti persistiti (le decisioni sono la verita': " + "cio' che non e' registrato non e' avvenuto) e riprendi da li'.\n" + "Valgono le stesse regole non negoziabili: una domanda per volta, conferma esplicita, tool " + "reviewer_*, niente phase advance/reopen o decision add da shell, una fase alla volta."; // Recovery hint shown when `tht phase advance` refuses with exit 5 (gate not satisfied). const PHASE_RECOVERY = "Completa i prerequisiti della fase (decisioni/artefatti) e riprova."; // --- helpers (verbatim from source) ------------------------------------------- function textResult(text) { return { content: [{ type: "text", text }], details: {} }; } // Single chokepoint for all CLI calls. cwd is the Pi project root (harness/). function tht(ctx, args) { return execFileSync("tht", args, { cwd: ctx.cwd, encoding: "utf8" }); } // Runs a privileged tht call; converts any failure into an actionable textResult // (never propagates a raw "Command failed" to the model). function relayIfThtFails(ctx, args, recovery) { try { tht(ctx, args); return null; } catch (e) { const cliMsg = (e.stderr || e.message || String(e)).toString().trim(); return textResult(`${cliMsg} ${recovery}`); } } // F2 single-source: workflow facts from `tht phase meta --json`. Cached per session. // Replaces the source's mirrored PHASE_NAMES constant (which drifted to 7 entries) // and the regex-parse of `tht phase show` text. let _phaseMetaCache = null; function phaseMeta(ctx) { if (_phaseMetaCache) return _phaseMetaCache; const raw = tht(ctx, ["phase", "meta", "--json"]); const meta = JSON.parse(raw); _phaseMetaCache = meta; return meta; } function phaseName(ctx, num) { const meta = phaseMeta(ctx); const p = meta.phases.find((x) => x.num === num); return p ? p.name : "?"; } function currentPhase(ctx, session) { const out = tht(ctx, ["phase", "show", "--session", session]); const m = out.match(/Fase corrente:\s*(\d+)/); return m ? parseInt(m[1], 10) : 1; } // tht phase advance --auto: exit 6 = not ready / needs human (silent no-op), others propagated. // Used for the fire-and-forget auto-advance of the auto phases (F2 memory, F6 cte) after a // reviewer_decide: it only advances when the phase is auto-eligible (zero substantive // decisions + prerequisites met), otherwise the CLI exits 6 and we no-op. function advanceIfReady(ctx, session) { try { tht(ctx, ["phase", "advance", "--auto", "--session", session]); return { advanced: true }; } catch (e) { if (e.status === 6) return { advanced: false }; return { advanced: false, error: (e.stderr || e.message || String(e)).toString().trim() }; } } // --- widget emission + wait (NEW: replaces ctx.ui.* blocking primitives) ------- // // emitAndWait(ctx, descriptor) sends a ui_request widget and awaits the correlated // ui_response by id. This is the correlation-by-id layer the reference implementation never had (its // native TUI primitives handled it implicitly). The descriptor is built by the pure // ./gate/builders.js (L1-tested). // // The no-limbo invariant: a control response of "cancel"/undefined is NEVER accepted // as a final answer -- the widget is re-presented. Real escapes (Back/Exit/Other) are // always present as selectable options in the descriptor (reserved field), so cancel // has no legitimate meaning. let _pending = new Map(); // id -> {resolve, reject} async function emitAndWait(ctx, descriptor) { for (;;) { const promise = new Promise((resolve, reject) => { _pending.set(descriptor.id, { resolve, reject }); }); ctx.sendRaw({ type: "extension_ui_request", ui_request: descriptor }); const resp = await promise; // no-limbo: a cancel/undefined response re-presents the same widget if (resp && resp.control !== "cancel" && resp.id === descriptor.id) { return resp; } if (ctx.hasUI) { await ctx.ui.notify( "Esc non chiude il gate: usa Torna indietro / Esci / Altro dalle opzioni.", "warning", ); } } } // The runtime invokes this when a ui_response arrives (RPC reply or SSE from FE). function handleUiResponse(event) { const slot = _pending.get(event.id); if (slot) { _pending.delete(event.id); slot.resolve(event); } } // --- the extension ------------------------------------------------------------ export default function (pi) { let lockActive = false; let lastSteered = false; let pendingKickoff = null; let activeSessionId = null; // 1) ANTI-BYPASS tool_call hook (spec D4, verbatim). Blocks direct phase/decision // calls and writes to protected files so the reviewer cannot bypass the gate. pi.on("tool_call", (event) => { if (event.toolName === "bash") { const cmd = event.input?.command ?? ""; // tht session finalize: legit session-close path -- unlock and let through. if (/\btht\s+session\s+finalize\b/.test(cmd)) { lockActive = false; lastSteered = false; return; } if (FORBIDDEN.some((re) => re.test(cmd))) { return { block: true, reason: "Avanzamento/ritorno di fase e registrazione decisioni passano dal " + "reviewer: usa i tool reviewer_confirm / reviewer_select.", }; } } if (event.toolName === "write" || event.toolName === "edit") { const path = event.input?.path ?? event.input?.file_path ?? ""; if (PROTECTED_FILES.test(path)) { return { block: true, reason: "Questo file di stato e' gestito dal gate: non scriverlo direttamente.", }; } } }); // 2) INPUT hook: workflow entry detection + free-input block + `!` steer channel. pi.on("input", async (event, ctx) => { const raw = (event.text ?? "").trimStart(); if ( event.source === "interactive" && /^\/(nuova-domanda|riprendi-sessione)\b/.test(raw) ) { // NOTE: the source runs an ollama preflight here; ThothII defers embeddings // readiness to the session's first vector op. Entry detection only: lockActive = true; lastSteered = false; pendingKickoff = /^\/nuova-domanda\b/.test(raw) ? NUOVA_DOMANDA_KICKOFF : RIPRENDI_KICKOFF; } if (!lockActive || event.source !== "interactive") return { action: "continue" }; const trimmed = (event.text ?? "").trimStart(); if (trimmed.length === 0) return { action: "continue" }; if (trimmed.startsWith("/")) return { action: "continue" }; // `!`-prefixed free text -> forward to the model (the one sanctioned steer channel). if (trimmed.startsWith("!")) return { action: "transform", text: trimmed.slice(1).trimStart() }; if (ctx.hasUI) { await ctx.ui.notify( "Durante la sessione rispondi con i widget del gate. " + "Per inviare testo libero al modello inizia la riga con '!'.", "warning", ); } return { action: "handled" }; }); // 3) BEFORE_AGENT_START: one-shot kickoff injection (appends the operational // instructions to the system prompt on the turn that starts/resumes a session). pi.on("before_agent_start", (event) => { if (!pendingKickoff) return undefined; const inject = pendingKickoff; pendingKickoff = null; return { systemPrompt: `${event.systemPrompt}\n\n${inject}` }; }); // 4) SESSION_START: reset all session-scoped state. pi.on("session_start", (_event, _ctx) => { lockActive = false; lastSteered = false; pendingKickoff = null; activeSessionId = null; _phaseMetaCache = null; }); // 5) AGENT_END prose safety net: if the model emits prose instead of a reviewer_* // tool call (and the lock is active), nudge it back to the gate tools. pi.on("agent_end", async (event) => { if (!lockActive) return; const msgs = event.messages ?? []; const last = msgs[msgs.length - 1]; const isProse = last && last.role === "assistant" && Array.isArray(last.content) && last.content.some((b) => b.type === "text" && (b.text ?? "").trim().length > 0) && !last.content.some((b) => b.type === "toolCall"); if (isProse && !lastSteered) { lastSteered = true; await pi.sendUserMessage( "Le risposte del reviewer arrivano solo dai widget del gate. Riproponi la " + "richiesta come reviewer_select (opzioni + 'Altro'), non in chat.", { deliverAs: "followUp" }, ); } }); // 6) UI_RESPONSE dispatch: correlate an incoming response to its awaiting widget. pi.on("extension_ui_response", handleUiResponse); // --- the four reviewer tools (widget-descriptor emit + await) --------------- pi.registerTool({ name: "reviewer_select", label: "Domanda a scelta (reviewer)", description: "Pone una domanda a scelta singola al reviewer via un widget select. Le opzioni di " + "controllo (Altro/Torna indietro/Esci) sono sempre presenti. NON persiste: serve a " + "chiedere, non a decidere.", parameters: Type.Object({ session: Type.String({ description: "Id sessione (per determinare la fase)." }), title: Type.String(), options: Type.Array( Type.Object({ id: Type.String(), label: Type.String(), recommended: Type.Optional(Type.Boolean()), }), ), intro: Type.Optional(Type.String()), }), async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; const { session, title, options: opts, intro } = params; const phase = phaseName(ctx, currentPhase(ctx, session)); const recommended = opts.find((o) => o.recommended)?.id ?? null; const clean = stripReserved(opts.map((o) => o.label)); const widget = buildSelectRequest({ id: `u${Date.now()}`, phase, title, intro: intro ?? null, recommended, options: opts.filter((o) => !isReserved(o.label)).map((o) => ({ id: o.id, label: o.label })), }); const resp = await emitAndWait(ctx, widget); // control responses (back/exit/other) are surfaced as text for the model to act on. if (resp.control === "freetext") return textResult(`Altro (reviewer): ${resp.text}`); if (resp.control === "back") return textResult("Il reviewer vuole tornare indietro."); if (resp.control === "exit") return textResult("Il reviewer vuole uscire."); const chosen = opts.find((o) => o.id === resp.choice); return textResult(`Scelta del reviewer: ${chosen ? chosen.label : resp.choice}`); }, }); pi.registerTool({ name: "reviewer_decide", label: "Decisione di merito (reviewer)", description: "Pone una decisione di merito al reviewer via widget multiselect e PERSISTE le " + "scelte (tht decision add). Ogni opzione porta un payload decision {type, subject, " + "detail, rationale}. Dopo la conferma, se p.advance e' vero tenta tht phase advance --if-ready.", parameters: Type.Object({ session: Type.String(), title: Type.String(), options: Type.Array( Type.Object({ id: Type.String(), label: Type.String(), decision: Type.Object({ type: Type.String(), subject: Type.String(), detail: Type.Optional(Type.String()), rationale: Type.Optional(Type.String()), }), recommended: Type.Optional(Type.Boolean()), }), ), allow_empty: Type.Optional(Type.Boolean()), advance: Type.Optional(Type.Boolean()), }), async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; const { session, title, options: opts, advance } = params; const phase = phaseName(ctx, currentPhase(ctx, session)); const toAdd = []; const widget = buildMultiselectRequest({ id: `u${Date.now()}`, phase, title, allowEmpty: params.allow_empty ?? false, options: opts.filter((o) => !isReserved(o.label)).map((o) => ({ id: o.id, label: o.label })), recommended: opts.find((o) => o.recommended)?.id ?? null, }); const resp = await emitAndWait(ctx, widget); if (resp.control === "freetext") { return textResult(`Altro (reviewer): ${resp.text}. Riformula la proposta tenendo conto.`); } if (resp.control === "back") return textResult("Il reviewer vuole tornare indietro."); if (resp.control === "exit") return textResult("Il reviewer vuole uscire."); const chosen = opts.filter((o) => (resp.choices ?? []).includes(o.id)); for (const c of chosen) { const d = c.decision; const args = ["decision", "add", "--session", session, "--type", d.type, "--subject", d.subject]; if (d.detail) args.push("--detail", d.detail); if (d.rationale) args.push("--rationale", d.rationale); const err = relayIfThtFails(ctx, args, ""); if (err) return err; toAdd.push(d); } if (advance) advanceIfReady(ctx, session); return textResult( toAdd.length ? `Registrate ${toAdd.length} decisioni: ${toAdd.map((d) => d.type).join(", ")}.` : "Nessuna decisione registrata (il reviewer non ha selezionato opzioni di merito).", ); }, }); pi.registerTool({ name: "reviewer_confirm", label: "Gate di avanzamento (reviewer)", description: "Checkpoint di fase/CTE/SQL: presenta un widget artifact-gate (artefatto + " + "Approva/Rifiuta/Altro) ed esegue l'azione privilegiata (tht phase advance, cte plan, " + "decision add sql_approved) solo su approvazione. kind: phase | cte_plan | cte_result | sql.", parameters: Type.Object({ session: Type.String(), kind: Type.Union([ Type.Literal("phase"), Type.Literal("cte_plan"), Type.Literal("cte_result"), Type.Literal("sql"), ]), title: Type.String(), artifact: Type.Object({ kind: Type.String(), data: Type.Any(), version: Type.Optional(Type.Number()), }), // kind:"cte_plan" only -- ordered list of CTE names to persist (tht cte plan --name). names: Type.Optional(Type.Array(Type.String())), }), async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; const { session, kind, title, artifact } = params; const phase = phaseName(ctx, currentPhase(ctx, session)); const widget = buildArtifactGate({ id: `u${Date.now()}`, phase, title, artifact, action: { kind: "approve_reject", prompt: "Approvi o rifiuti?" }, }); const resp = await emitAndWait(ctx, widget); if (resp.control === "freetext" || resp.choice === "reject") { return textResult( `Rifiutato${resp.text ? ` (motivo: ${resp.text})` : ""}: rivedi e riprova.`, ); } if (resp.control === "back") return textResult("Il reviewer vuole tornare indietro."); if (resp.control === "exit") return textResult("Il reviewer vuole uscire."); // approved -> execute the privileged action via the CLI. if (kind === "phase") { // Explicit human approval: advance unconditionally except for unmet // prerequisites. Plain `phase advance` (no --auto) enforces advance_problems // and exits 6 with the missing items, which relayIfThtFails surfaces. const err = relayIfThtFails(ctx, ["phase", "advance", "--session", session], PHASE_RECOVERY); if (err) return err; return textResult(`Fase approvata (sessione ${session}).`); } if (kind === "cte_plan") { // The CTE plan is the ordered list of CTE names; the model passes them in // params.names (tht cte plan requires at least one --name). const names = Array.isArray(params.names) ? params.names : []; if (names.length === 0) { return textResult( "Nessun nome CTE fornito: il piano CTE richiede l'elenco ordinato dei CTE " + "(parametro names di reviewer_confirm).", ); } const planArgs = ["cte", "plan", "--session", session]; for (const n of names) planArgs.push("--name", n); const err = relayIfThtFails(ctx, planArgs, ""); if (err) return err; return textResult(`CTE plan approvato (${names.length} CTE, sessione ${session}).`); } if (kind === "cte_result" || kind === "sql") { const dt = kind === "sql" ? "sql_approved" : "cte_approved"; const err = relayIfThtFails( ctx, ["decision", "add", "--session", session, "--type", dt, "--subject", `phase:${currentPhase(ctx, session)}`], "", ); if (err) return err; return textResult(`${kind} approvato (sessione ${session}).`); } return textResult(`Approvato (kind=${kind}, sessione ${session}).`); }, }); pi.registerTool({ name: "rewrite_question", label: "Riscrittura domanda (deterministica)", description: "Scrive deterministicamente question.md via tht session set-question (evita il tool " + "di edit unreliable). assumptions puo' essere array o stringa JSON.", parameters: Type.Object({ session: Type.String(), question: Type.String(), assumptions: Type.Optional(Type.Any()), }), async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; const { session, question, assumptions } = params; let assumps = assumptions; if (typeof assumps === "string") { try { assumps = JSON.parse(assumps); } catch { assumps = [assumps]; } } // `tht session set-question` takes the session id as a positional argument // (the `session` command group uses positional ids, unlike phase/cte/decision // which use --session). const args = ["session", "set-question", session, "--question", question]; if (Array.isArray(assumps)) { for (const a of assumps) args.push("--assumption", String(a)); } const err = relayIfThtFails(ctx, args, ""); if (err) return err; return textResult(`Domanda riscritta per la sessione ${session}.`); }, }); // --- slash command: /torna [session_id] [N] (rollback to a previous phase) -- pi.registerCommand("torna", { description: "Torna a una fase precedente: /torna [N] (default: un passo).", handler: async (args, ctx) => { const parts = args.trim().split(/\s+/).filter(Boolean); const sessionId = parts.length >= 2 ? parts[0] : parts.length === 1 && /^\d/.test(parts[0]) ? undefined : parts[0]; const sid = sessionId ?? activeSessionId ?? process.env.THT_SESSION; if (!sid) { await ctx.ui.notify("Uso: /torna [N]", "warning"); return; } const cur = currentPhase(ctx, sid); const targetArg = parts.find((x) => /^\d+$/.test(x)); const target = targetArg ? parseInt(targetArg, 10) : cur - 1; if (target < 1 || target >= cur) { await ctx.ui.notify(`Target non valido (fase corrente ${cur}).`, "warning"); return; } tht(ctx, ["phase", "reopen", "--session", sid, "--phase", String(target)]); await pi.sendUserMessage( `Ho riaperto la Fase ${target} della sessione ${sid}. Riprendi il protocollo da quella fase.`, { deliverAs: "followUp" }, ); }, }); }