// 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: pure builders and domain modules are L1-tested under ./gate/__tests__/. // The composition root is exercised through the fake Pi runtime and live at L2. import { execFileSync } from "node:child_process"; import { readFileSync } from "node:fs"; import { Type } from "typebox"; import { buildSelectRequest, buildMultiselectRequest, buildArtifactGate, buildSchemaLinkingRequest, buildJoinReviewRequest, } from "./gate/builders.js"; import { validateCtePlanV2, validateCteResultThin, validatePhaseSummaryV2, } from "./gate/artifact-contracts.js"; import { memoizeGetColumns, enrichCtePlanV2, buildCteResultV2, enrichCteResultColumns, enrichPhaseSummaryV2, appendLedgerSection, } from "./gate/enrich.js"; import { createMemoryGate } from "./gate/memory/index.js"; import { createDisambiguationGate } from "./gate/disambiguation/index.js"; import { isReserved } from "./reserved-labels.mjs"; // Load the workflow contract once when Pi loads the extension. Asking the model to // discover/read the skill as its first action proved unreliable with remote models: // they can spend the whole RPC turn exploring the repository before opening the exact // file named in the kickoff. Embedding the canonical file keeps SKILL.md as the single // source of truth while making session bootstrap deterministic. const SESSION_SKILL = readFileSync( new URL("../skills/tht-sessione/SKILL.md", import.meta.url), "utf8", ); // Read through the CLI so workspace/config resolution stays canonical (sessions may // live outside this repository). A missing pack is a soft miss: standalone/TUI flows // retain the existing `tht search pack` bootstrap path. export function readRetrievalPack(ctx, sessionId) { if (!sessionId) return null; try { const content = execFileSync( "tht", ["session", "retrieval-pack", sessionId], { cwd: ctx.cwd, encoding: "utf8" }, ).trim(); return content ? content.replaceAll("", "</retrieval-pack>") : null; } catch { return null; } } // --- prepareArguments: parse stringified arrays (workaround for models that send // arrays as JSON strings -- same pattern as pi-core's edit tool prepareEditArguments) // Coerce a value that may arrive as a JSON-encoded string back into an object/array. // Models sometimes stringify object params (Type.Object / Type.Any); the parsed value // is returned only when it is genuinely an object, so plain strings pass through intact. export function jsonObjectOrSelf(value) { if (typeof value !== "string") return value; try { const parsed = JSON.parse(value); return parsed && typeof parsed === "object" ? parsed : value; } catch { return value; } } export function prepareReviewerArguments(input) { if (!input || typeof input !== "object") return input; const args = { ...input }; // Parse stringified JSON arrays (workaround for models that send arrays as strings) if (typeof args.options === "string") { try { const parsed = JSON.parse(args.options); if (Array.isArray(parsed)) args.options = parsed; } catch { /* not JSON */ } } if (typeof args.names === "string") { try { const parsed = JSON.parse(args.names); if (Array.isArray(parsed)) args.names = parsed; } catch { /* not JSON */ } } if (typeof args.tables === "string") { try { const parsed = JSON.parse(args.tables); if (Array.isArray(parsed)) args.tables = parsed; } catch { /* not JSON */ } } // reviewer_confirm's `artifact` is a Type.Object; some models send it as a // stringified JSON object, which fails validation BEFORE execute() and makes the // model loop. Coerce it back to an object so the gate proceeds. if (args.artifact !== undefined) args.artifact = jsonObjectOrSelf(args.artifact); // GLM sometimes double-stringifies: artifact.data itself arrives as a JSON // string (v2 payloads) even once artifact is an object. Coerce it too; a // legacy markdown string in artifact.data is not valid JSON and jsonObjectOrSelf // already leaves non-JSON strings unchanged. if (args.artifact && typeof args.artifact === "object" && args.artifact.data !== undefined) { args.artifact = { ...args.artifact, data: jsonObjectOrSelf(args.artifact.data) }; } // Models frequently forget artifact.kind while providing it at top level — // copy it over so validation doesn't loop. if (args.artifact && typeof args.artifact === "object" && !args.artifact.kind && args.kind) { args.artifact = { ...args.artifact, kind: args.kind }; } return args; } // --- 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/, // La promozione in memoria passa dal reviewer (reviewer_memory_promote), mai da shell. /\btht\s+memory\s+(promote|save-one)\b/, // La re-introspezione del DWH (~3 min) è manutenzione fuori sessione, mai in workflow. /\btht\s+schema\s+introspect\b[^\n]*--refresh/, ]; const PROTECTED_FILES = /(review_decisions\.jsonl|session_manifest\.yaml|cte_plan\.json)/; // The model orchestrates the workflow; it must never rewrite the gate/extension code // itself. pi's write/edit tools are otherwise unrestricted, so a confused model can // (and did) patch tht-gate.js mid-loop. Block any write under the extensions dir. export const GATE_CODE_FILES = /\.pi[\\/]extensions[\\/]/; // Bash can mutate protected state around the write/edit hook (`echo >> ledger`, // `sed -i` on the manifest, `cat > tht-gate.js`). Block any bash command that names // a protected path in a mutating context; read-only mentions (cat/grep/ls) stay // allowed. Defense against a CONFUSED model, not a sandbox. const PROTECTED_PATH_TOKEN = /(review_decisions\.jsonl|session_manifest\.yaml|cte_plan\.json|\.pi[\\/]extensions)/; const BASH_MUTATION = /(>>?|\btee\b|\bsed\b[^|;&]*\s-i\b|\bmv\b|\bcp\b|\brm\b|\btruncate\b|\bdd\b|open\([^)]*['"][wa]\+?b?['"]|\bchmod\b|\bln\b)/; export function isProtectedBashMutation(cmd) { return PROTECTED_PATH_TOKEN.test(cmd) && BASH_MUTATION.test(cmd); } // Session artifacts are exposed by deterministic `tht session documents`; shell // discovery is both unnecessary and dangerous (a model once escalated to `find /` // and wedged the whole RPC turn). Match shell command boundaries so the supported // `tht search find` subcommand remains available. export function isFilesystemFind(cmd) { return /(?:^|[;&|\n])\s*(?:(?:sudo|command)\s+)?find(?:\s|$)/.test(cmd); } // --- kickoff payloads (verbatim from source L184-212, load-bearing model prose) - const NUOVA_DOMANDA_KICKOFF = "AZIONE IMMEDIATA OBBLIGATORIA: la skill canonica è già inclusa integralmente nel " + "system prompt. Non esplorare il repository e non leggere altri file di istruzioni.\n" + "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. Segui il workflow dalla Fase 1 (Chiarimento), usando l'id di sessione in ogni " + "comando `tht`.\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 NUOVA_DOMANDA_KICKOFF_PROVIDED = (sessionId, hasRetrievalPack = false) => "AZIONE IMMEDIATA OBBLIGATORIA: non esplorare il repository, non usare `ls`, `find`, " + "`--help` e non leggere README, sorgenti, test o altri file di istruzioni. " + "La skill canonica è già inclusa integralmente nel system prompt qui sotto.\n" + "Istruzioni operative — sessione ThothII (workflow human-in-the-middle). " + `La sessione è GIÀ creata con id \`${sessionId}\`: usalo in OGNI comando \`tht\`. ` + "NON creare una nuova sessione.\n" + (hasRetrievalPack ? "1. Il retrieval pack persistito è già incluso nel system prompt: usalo direttamente. " + "NON eseguire `tht search pack` e NON chiamare un tool per leggere `retrieval_pack.md`.\n" : "1. Il retrieval pack non era ancora disponibile: come PRIMA chiamata tool esegui " + "subito `tht search pack \"\" --session " + sessionId + "` e usa il risultato.\n") + "2. Nel primo turno identifica la SOLA ambiguità con maggiore impatto sulla query e " + "chiama subito `reviewer_select` con opzioni concrete. Niente lunga narrazione, elenco " + "di tutte le ambiguità o ricapitolazione preliminare.\n" + "Regole non negoziabili: una domanda al reviewer per volta; mai promuovere/escludere/" + "correggere senza conferma; le interazioni passano dai tool reviewer_*; testo libero col prefisso '!'. " + "MAI `tht phase advance|reopen` né `tht decision add` da shell."; const RIPRENDI_KICKOFF = (hasRetrievalPack = false) => "AZIONE IMMEDIATA OBBLIGATORIA: non esplorare il repository, non usare `ls`, `find`, " + "`--help` e non leggere README, sorgenti, test o altri file di istruzioni. " + "La skill canonica è già inclusa integralmente nel system prompt qui sotto.\n" + "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" + (hasRetrievalPack ? "2. Il retrieval pack persistito e' gia' incluso nel system prompt: usalo direttamente. NON eseguire `tht search pack` e NON chiamare un tool per leggere `retrieval_pack.md`.\n" : "2. Se serve il retrieval pack, esegui `tht search pack` una sola volta, senza esplorare altri file.\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" + "Esegui il passo 1 ORA, in QUESTO stesso turno, chiamando subito il tool `bash` per " + "`tht session show`: NON limitarti a dichiarare l'intenzione e NON " + "terminare il turno prima di aver chiamato i tool. Se la sessione e' in fase 1 e non ha decisioni, dopo `session show` chiama `reviewer_select` subito: non produrre testo libero.\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: {} }; } // Load THT_* env vars from .env (project root = ctx.cwd). function loadEnvFromDotenv(ctx) { const fs = require("node:fs"); const path = require("node:path"); const envPath = path.join(ctx.cwd, ".env"); try { const raw = fs.readFileSync(envPath, "utf8"); for (const line of raw.split("\n")) { const trimmed = line.trim(); if (!trimmed || trimmed.startsWith("#")) continue; const idx = trimmed.indexOf("="); if (idx === -1) continue; const key = trimmed.slice(0, idx).trim(); const value = trimmed.slice(idx + 1).trim(); if (key.startsWith("THT_") || key.startsWith("PSD_")) { process.env[key] = value; } } } catch { // .env not found or unreadable — proceed with existing env. } } // Single chokepoint for all CLI calls. cwd is the Pi project root (harness/). function tht(ctx, args, input) { loadEnvFromDotenv(ctx); return execFileSync("tht", args, { cwd: ctx.cwd, encoding: "utf8", input }); } // 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, input) { try { tht(ctx, args, input); 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; } // Returns the workflow.yaml phase id (e.g. "F1"), not the descriptive name — this is // what feeds the widget's `phase` field, which the frontend WorkflowBar matches against // its own "F1".."F8" ids to light up the stepper pills. function phaseId(ctx, num) { const meta = phaseMeta(ctx); const p = meta.phases.find((x) => x.num === num); return p ? p.id : "?"; } // Known decision types from workflow.yaml emits + meta types. Pre-validates decision // payloads BEFORE showing the reviewer widget so a typo in the type name doesn't waste // a human interaction (the CLI would reject it only after the reviewer has already picked). let _knownTypes = null; function knownDecisionTypes(ctx) { if (_knownTypes) return _knownTypes; const meta = phaseMeta(ctx); const emitted = new Set(); for (const p of meta.phases) { for (const t of (p.emits || [])) emitted.add(t); } // A workflow meta with NO emits (older tht, minimal test stubs) gives the gate no // vocabulary to validate against: skip validation instead of rejecting every // substantive type and bricking the gates. if (emitted.size === 0) return null; for (const t of [ "phase_approved", "phase_auto_approved", "phase_reopened", "phase_skipped", "decision_retracted", ]) emitted.add(t); _knownTypes = emitted; return emitted; } function decisionMinPhaseMap(ctx) { const meta = phaseMeta(ctx); const mins = {}; for (const p of meta.phases) { for (const t of (p.emits || [])) { if (!(t in mins) || p.num < mins[t]) mins[t] = p.num; } } return mins; } function validateDecisionTypes(ctx, options, session) { const known = knownDecisionTypes(ctx); if (!known) return null; for (const o of options) { if (o.decision && !known.has(o.decision.type)) { return `Tipo di decisione '${o.decision.type}' non valido. Tipi ammessi: ${[...known].join(", ")}. Correggi e riprova.`; } } if (session) { const cur = currentPhase(ctx, session); const mins = decisionMinPhaseMap(ctx); for (const o of options) { if (o.decision && mins[o.decision.type] && mins[o.decision.type] > cur) { return `Tipo '${o.decision.type}' ammesso dalla Fase ${mins[o.decision.type]}, sessione alla Fase ${cur}. Chiudi prima la fase corrente.`; } } } return null; } // Full phase descriptor {id, num, name} for the v2 phase-summary payload's `phase`. function phaseMetaForNum(ctx, num) { const meta = phaseMeta(ctx); const p = meta.phases.find((x) => x.num === num); return p ? { id: p.id, num: p.num, name: p.name } : { id: "?", num, name: "" }; } // Memoized, never-throwing catalog lookup for enrich: `tht schema columns --json` // -> {description, columns:[...]} or null on any miss (same hardened pattern as // reviewer_schema_linking, but a miss is a soft null here, not a hard fail). function makeGetColumns(ctx) { return memoizeGetColumns((tableName) => { try { return JSON.parse(tht(ctx, ["schema", "columns", tableName, "--json"])); } catch { return null; } }); } 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; } // --- fuzzy table-name matching (prevents multi-turn debug loops on typos) --- function levenshtein(a, b) { if (a === b) return 0; const la = a.length, lb = b.length; if (!la) return lb; if (!lb) return la; let prev = Array.from({ length: lb + 1 }, (_, i) => i); for (let i = 1; i <= la; i++) { const cur = [i]; for (let j = 1; j <= lb; j++) { cur[j] = Math.min(prev[j] + 1, cur[j - 1] + 1, prev[j - 1] + (a[i - 1] !== b[j - 1] ? 1 : 0)); } prev = cur; } return prev[lb]; } const _catalogNamesCache = new Map(); function catalogTableNames(ctx) { const key = ctx.cwd ?? ""; if (_catalogNamesCache.has(key)) return _catalogNamesCache.get(key); let names; try { const raw = tht(ctx, ["schema", "render", "--format", "mschema-text"]); names = [...raw.matchAll(/^CREATE TABLE (\S+)/gm)].map((m) => m[1]); } catch { names = []; } _catalogNamesCache.set(key, names); return names; } function closestTableName(ctx, name, maxDist = 3) { const names = catalogTableNames(ctx); let best = null, bestDist = maxDist + 1; for (const n of names) { if (Math.abs(n.length - name.length) > maxDist) continue; const d = levenshtein(n, name); if (d < bestDist) { bestDist = d; best = n; } } return best; } // 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(), }; } } // Unconditional phase advance — the reviewer already approved via the widget interaction. // ONLY for gates whose selection IS the phase approval by design: reviewer_schema_linking // (F4 curation) and the F8 close path. Generic reviewer_decide must use advanceIfReady — // forcing there bypasses the reviewer_confirm phase gate and contradicts SKILL.md. function forceAdvance(ctx, session) { try { tht(ctx, ["phase", "advance", "--session", session]); return { advanced: true }; } catch (e) { return { advanced: false, error: (e.stderr || e.message || String(e)).toString().trim(), }; } } // --- widget emission + wait — usa l'API UI NATIVA di Pi (ctx.ui.input) -------- // // emitAndWait(ctx, descriptor) sends the widget-descriptor as JSON in the `title` // of ctx.ui.input and awaits the response as a parsed string value. No ctx.sendRaw, // no pi.on("extension_ui_response"): Pi routes the response via pendingExtensionRequests. // // The no-limbo invariant: undefined/null/invalid-JSON/control:"cancel" are 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, so cancel has no // legitimate meaning. // The frontend's reviewer widgets always carry the picked option(s) in a `choices` // ARRAY (SelectWidget/ArtifactGateWidget send `choices: [optionId]`). Single-pick // handlers read the first element; `choice` (singular) is accepted as a legacy form. export function selectedChoice(resp) { if (Array.isArray(resp?.choices)) return resp.choices[0]; return resp?.choice; } // Builds the `tht decision add` argv for one decision payload {type, subject, detail?, // rationale?}. Shared by reviewer_decide (multiselect) and reviewer_select (auto-confirm). export function decisionAddArgs(session, d) { 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); return args; } // Keep the model inside the gate workflow immediately after a reviewer selection. // The agent_end prose safety net is too late for providers that keep streaming a // single, very long assistant turn: put the continuation contract in the tool result // that unlocks the model after the reviewer response. export function decisionRecordedResultText(decision, option) { return ( `Decisione registrata (${decision.type}): ${option.label}. ` + "Ora rileggi lo stato persistito con `tht session show` e invoca immediatamente " + "il prossimo tool reviewer_ richiesto dal workflow. " + "Non scrivere analisi o spiegazioni visibili." ); } export function hasUngatedAssistantProse(messages) { const last = messages?.[messages.length - 1]; if (!last || last.role !== "assistant" || !Array.isArray(last.content)) return false; const hasText = last.content.some( (block) => block.type === "text" && (block.text ?? "").trim().length > 0, ); const hasReviewerTool = last.content.some( (block) => block.type === "toolCall" && typeof block.name === "string" && block.name.startsWith("reviewer_"), ); return hasText && !hasReviewerTool; } // Workstream F: classifies a reviewer_select response into the action the gate takes. // A concrete choice that carries a `decision` payload auto-confirms (persist directly, // no second gate); a bare choice stays ask-only; back/exit/Other never persist. export function resolveSelectOutcome(opts, resp) { if (resp?.control === "freetext") return { kind: "freetext", text: resp.text }; if (resp?.control === "back") return { kind: "back" }; if (resp?.control === "exit") return { kind: "exit" }; const choice = selectedChoice(resp); const option = (opts || []).find((o) => o.id === choice) || null; if (option && option.decision) return { kind: "decision", option, decision: option.decision }; return { kind: "choice", option, choice }; } // Classifies a reviewer_confirm response. Advance happens ONLY on an explicit // "approve" choice; a bare/unknown response resolves to {kind:"unknown"} and is // re-presented (never an accidental approve). freetext is actionable feedback, // distinct from an explicit "reject". export function resolveConfirmOutcome(resp) { if (resp?.control === "freetext") return { kind: "freetext", text: resp.text }; if (resp?.control === "back") return { kind: "back" }; if (resp?.control === "exit") return { kind: "exit" }; const choice = selectedChoice(resp); if (choice === "approve") return { kind: "approve" }; if (choice === "reject") return { kind: "reject" }; return { kind: "unknown", choice }; } export function isJoinReviewApproval(resp, optionIds) { if (resp?.kind !== "join-review" || !Array.isArray(resp.choices)) return false; const expected = new Set(optionIds); const chosen = new Set(resp.choices); return expected.size === optionIds.length && chosen.size === resp.choices.length && chosen.size === expected.size && [...chosen].every((id) => expected.has(id)); } export async function emitAndWait(ctx, descriptor) { for (;;) { const value = await ctx.ui.input(JSON.stringify(descriptor), ""); if (value === undefined || value === null) { await reLoop(ctx); continue; } let resp; try { resp = JSON.parse(value); } catch { await reLoop(ctx); continue; } if (resp && resp.control !== "cancel" && resp.id === descriptor.id) return resp; await reLoop(ctx); } } async function reLoop(ctx) { // TUI-only: "Esc"/terminal navigation has no meaning in RPC (the browser // renders widget buttons). Guard on ctx.mode, NOT ctx.hasUI — on pi >=0.80 // hasUI is true in RPC too, so this warning would leak to the frontend. if (ctx.mode === "tui") await ctx.ui.notify( "Esc non chiude il gate: usa Torna indietro / Esci / Altro dalle opzioni.", "warning", ); } // --- the extension ------------------------------------------------------------ export default function (pi) { let lockActive = false; let lastSteered = false; let pendingKickoff = null; let activeSessionId = null; const memoryGate = createMemoryGate({ workflow: { activate: () => { lockActive = true; }, phase: (ctx, session) => { const number = currentPhase(ctx, session); return { number, id: phaseId(ctx, number) }; }, advance: advanceIfReady, close: (...args) => closeAfterPromotion(...args), }, memory: { execute: (ctx, args) => tht(ctx, ["memory", ...args]), mutate: (ctx, args, recovery) => relayIfThtFails( ctx, ["memory", ...args], recovery, ), }, ledger: { validate: validateDecisionTypes, record: (ctx, session, decision, recovery) => relayIfThtFails( ctx, decisionAddArgs(session, decision), recovery, ), }, waitForReviewer: emitAndWait, toTextResult: textResult, }); const disambiguationGate = createDisambiguationGate({ workflow: { activate: () => { lockActive = true; }, phase: (ctx, session) => ({ number: currentPhase(ctx, session) }), describe: (ctx, session) => phaseId(ctx, currentPhase(ctx, session)), decisionTypes: knownDecisionTypes, decisionMinimumPhases: decisionMinPhaseMap, advance: advancePhaseAndFinalize, }, session: { mutate: (ctx, args, recovery) => relayIfThtFails( ctx, ["session", ...args], recovery, ), }, ledger: { record: (ctx, session, decision, recovery) => relayIfThtFails( ctx, decisionAddArgs(session, decision), recovery, ), }, reviewer: { buildSelect: buildSelectRequest, isReserved, }, toTextResult: textResult, }); // 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 (isFilesystemFind(cmd)) { return { block: true, reason: "Non cercare file con `find`: gli artefatti persistiti della sessione " + "si leggono con `tht session documents --json`; per catalogo ed " + "evidence usa `tht schema render` / `tht search find`.", }; } 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 (isProtectedBashMutation(cmd)) { return { block: true, reason: "I file di stato della sessione e il codice del gate non si modificano " + "da shell: lo stato passa SOLO dai tool del gate (reviewer_*/write_*). " + "La lettura (cat/grep) resta permessa.", }; } } if (event.toolName === "write" || event.toolName === "edit") { const path = event.input?.path ?? event.input?.file_path ?? ""; if (GATE_CODE_FILES.test(path)) { return { block: true, reason: "Il codice del gate (.pi/extensions) non va modificato: sei l'orchestratore " + "del workflow, non uno sviluppatore del gate. Usa i tool del gate.", }; } 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(); // entry detection: workflow-start funziona sia da TUI sia da comando RPC `prompt`. if (/^\/(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) ? process.env.THT_SESSION ? NUOVA_DOMANDA_KICKOFF_PROVIDED(process.env.THT_SESSION) : NUOVA_DOMANDA_KICKOFF : RIPRENDI_KICKOFF(); } // free-input block: attivo quando il lock è su, per qualsiasi input utente (non solo interattivo). if (!lockActive) 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.mode === "tui") { 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, ctx) => { if (!pendingKickoff) return undefined; const sessionId = process.env.THT_SESSION; const isProvidedNewSession = Boolean(sessionId) && pendingKickoff === NUOVA_DOMANDA_KICKOFF_PROVIDED(sessionId); const isResumeSession = pendingKickoff === RIPRENDI_KICKOFF(); const retrievalPack = isProvidedNewSession || isResumeSession ? readRetrievalPack(ctx, sessionId) : null; const inject = isProvidedNewSession ? NUOVA_DOMANDA_KICKOFF_PROVIDED(sessionId, Boolean(retrievalPack)) : isResumeSession ? RIPRENDI_KICKOFF(Boolean(retrievalPack)) : pendingKickoff; pendingKickoff = null; return { message: { role: "user", content: [{ type: "text", text: "Esegui ora il workflow richiesto. Non scrivere analisi, spiegazioni o un elenco: usa il tool bash per `tht session show` e, se la sessione e' in Fase 1 senza decisioni, invoca immediatamente reviewer_select. La tua prossima risposta visibile deve essere una tool call." }], }, systemPrompt: `${event.systemPrompt}\n\n` + "\n" + SESSION_SKILL + "\n\n\n" + inject + (retrievalPack ? "\n\n\n" + retrievalPack + "\n" : ""), }; }); // 4) SESSION_START: Pi 0.80 emits this AFTER the RPC input hook. Preserve a // just-armed kickoff until before_agent_start injects it; clearing it here // silently drops the F1/resume contract and lets the model meander. pi.on("session_start", (_event, _ctx) => { if (!pendingKickoff) { lockActive = false; lastSteered = false; } 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; if (hasUngatedAssistantProse(event.messages ?? []) && !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" }, ); } }); // --- 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. La scelta su un'opzione " + "concreta È la conferma: se quell'opzione porta un payload `decision` {type, subject, " + "detail?, rationale?}, la decisione viene PERSISTITA direttamente (tht decision add) " + "senza un secondo gate di conferma; se p.advance è vero, tenta tht phase advance --if-ready. " + "Un'opzione SENZA `decision` resta solo-richiesta (non persiste). Altro/Torna indietro/Esci " + "non persistono mai e tornano come testo da gestire.", 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(), decision: Type.Optional( Type.Object({ type: Type.String(), subject: Type.String(), detail: Type.Optional(Type.String()), rationale: Type.Optional(Type.String()), }), ), recommended: Type.Optional(Type.Boolean()), }), ), intro: Type.Optional(Type.String()), advance: Type.Optional(Type.Boolean()), }), prepareArguments: disambiguationGate.prepareClarificationArguments, async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; try { const prepared = disambiguationGate.prepareClarification( ctx, params, `u${Date.now()}`, ); if (prepared.result) return prepared.result; const response = await emitAndWait(ctx, prepared.widget); const outcome = disambiguationGate.resolveClarification(prepared, response); if (outcome.decision) { const err = relayIfThtFails( ctx, decisionAddArgs(prepared.session, outcome.decision), "", ); if (err) return err; if (prepared.advance) advanceIfReady(ctx, prepared.session); } return textResult(outcome.text); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[reviewer_select ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); pi.registerTool({ name: "reviewer_datamart", label: "Scelta datamart (reviewer)", description: "F8: gestisce deterministicamente la scelta di generazione del datamart. " + "Con THT_PROFILE=workstation non mostra alcun widget e registra automaticamente " + "datamart_declined; con profile=server mostra sempre entrambe le opzioni Si/No e " + "registra la decisione scelta. Dopo questo tool prosegui con reviewer_memory_promote.", parameters: Type.Object({ session: Type.String(), }), async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; try { const { session } = params; const curNum = currentPhase(ctx, session); if (curNum !== 8) { return textResult( `La scelta datamart e' disponibile solo in Fase 8 (fase corrente: ${curNum}).`, ); } const subject = "phase:8"; if (process.env.THT_PROFILE === "workstation") { const decision = { type: "datamart_declined", subject }; const err = relayIfThtFails(ctx, decisionAddArgs(session, decision), ""); if (err) return err; return textResult( "Profilo workstation: datamart saltato automaticamente " + "e decisione datamart_declined registrata. Invoca ora reviewer_memory_promote.", ); } const options = [ { id: "generate", label: "Sì, genera il datamart", decision: { type: "datamart_requested", subject }, }, { id: "skip", label: "No, salta il datamart", decision: { type: "datamart_declined", subject }, recommended: true, }, ]; const widget = buildSelectRequest({ id: `u${Date.now()}`, phase: phaseId(ctx, curNum), title: "Vuoi generare un datamart?", intro: null, recommended: "skip", options: options.map((option) => ({ id: option.id, label: option.label })), }); const resp = await emitAndWait(ctx, widget); const outcome = resolveSelectOutcome(options, resp); if (outcome.kind === "freetext") return textResult(`Altro (reviewer): ${outcome.text}`); if (outcome.kind === "back") return textResult("Il reviewer vuole tornare indietro."); if (outcome.kind === "exit") return textResult("Il reviewer vuole uscire."); if (outcome.kind !== "decision") return textResult("Scelta datamart non valida: ripresenta reviewer_datamart."); const err = relayIfThtFails( ctx, decisionAddArgs(session, outcome.decision), "", ); if (err) return err; return textResult( decisionRecordedResultText(outcome.decision, outcome.option), ); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult( `[reviewer_datamart ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`, ); } }, }); 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(), description: Type.Optional(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()), }), prepareArguments: prepareReviewerArguments, async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; try { const { session, title, advance } = params; const phase = phaseId(ctx, currentPhase(ctx, session)); if (phase === "F2") return await memoryGate.reviewRecall(ctx, params, phase); const opts = params.options; const typeErr = validateDecisionTypes(ctx, opts, session); if (typeErr) return textResult(typeErr); const toAdd = []; const meritOptions = opts .filter((o) => !isReserved(o.label)) .map((o) => ({ id: o.id, label: o.label, ...(o.detail ? { detail: o.detail } : {}), ...(o.rationale ? { rationale: o.rationale } : {}), ...(o.meta ? { meta: o.meta } : {}), })); const joinOnly = meritOptions.length > 0 && opts .filter((o) => !isReserved(o.label)) .every((o) => o.decision.type === "join_modified"); const widget = joinOnly ? buildJoinReviewRequest({ id: `u${Date.now()}`, phase, title, options: opts .filter((o) => !isReserved(o.label)) .map((o) => ({ id: o.id, label: o.label, detail: o.decision.detail ?? "", rationale: o.decision.rationale ?? "", })), }) : buildMultiselectRequest({ id: `u${Date.now()}`, phase, title, allowEmpty: params.allow_empty ?? false, options: meritOptions, }); let resp; for (;;) { resp = await emitAndWait(ctx, widget); if (resp.control === "freetext" || resp.control === "back" || resp.control === "exit") break; if (!joinOnly || isJoinReviewApproval(resp, meritOptions.map((option) => option.id))) break; await reLoop(ctx); } 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 = joinOnly ? opts.filter((o) => !isReserved(o.label)) : opts.filter((o) => (resp.choices ?? []).includes(o.id)); if (joinOnly) { const decisions = chosen.map((choice) => choice.decision); const err = relayIfThtFails( ctx, ["decision", "add-join-set", "--session", session, "--doc", "-"], "Il set di join non è stato registrato; correggi l'errore e riprova.", JSON.stringify(decisions), ); if (err) return err; toAdd.push(...decisions); } else { for (const c of chosen) { const d = c.decision; const err = relayIfThtFails(ctx, decisionAddArgs(session, d), ""); if (err) return err; toAdd.push(d); } } if (!toAdd.length) { const adv = advance ? advanceIfReady(ctx, session) : { advanced: false }; const msg = "Nessuna decisione registrata (il reviewer non ha selezionato opzioni di merito)."; return textResult(adv.advanced ? msg + " Fase avanzata automaticamente." : msg); } // SKILL.md contract: advance:true means auto-advance-IF-ELIGIBLE (empty F2 / // skipped F6), never a bypass of the phase gate once substantive decisions // exist. The only gates that close their phase on selection are // reviewer_schema_linking (F4) and reviewer_memory_promote (F8). const adv = advance ? advanceIfReady(ctx, session) : { advanced: false }; const parts = [`Registrate ${toAdd.length} decisioni: ${toAdd.map((d) => d.type).join(", ")}.`]; parts.push(adv.advanced ? "Fase avanzata automaticamente." : 'La fase resta aperta: chiudila con reviewer_confirm kind:"phase" quando pronta.'); return textResult(parts.join(" ")); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[reviewer_decide ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); pi.registerTool({ name: "reviewer_schema_linking", label: "Schema linking: tabelle/colonne/esclusioni (reviewer)", description: "F4: presenta al reviewer le tabelle da promuovere/escludere con la loro descrizione " + "e, per ogni tabella, le colonne dal catalogo (le suggerite pre-selezionate). Il reviewer " + "cura le colonne di ogni tabella promossa. PERSISTE table_promoted/table_excluded e " + "column_promoted/column_excluded, poi riproietta schema_linking.json in modo deterministico " + "(tht session sync-schema-linking). `tables[]`: {id, name, kind: 'promote'|'exclude', " + "rationale?, suggested_columns?: string[]}. Le colonne complete arrivano dal catalogo, non dal modello.", parameters: Type.Object({ session: Type.String(), title: Type.String(), tables: Type.Array( Type.Object({ id: Type.String(), name: Type.String(), kind: Type.Union([Type.Literal("promote"), Type.Literal("exclude")]), rationale: Type.Optional(Type.String()), suggested_columns: Type.Optional(Type.Array(Type.String())), recommended: Type.Optional(Type.Boolean()), }), ), advance: Type.Optional(Type.Boolean()), }), prepareArguments: prepareReviewerArguments, async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; try { const { session, title, tables: rawTables, advance } = params; const tables = Array.isArray(rawTables) ? rawTables : (() => { try { const p = JSON.parse(rawTables); return Array.isArray(p) ? p : []; } catch { return []; } })(); if (!tables.length) return textResult("Parametro `tables` vuoto o non parsabile. Riprova con un array JSON di tabelle."); const phase = phaseId(ctx, currentPhase(ctx, session)); // Enrich each table with its full catalog columns (deterministic source). const enriched = []; for (const t of tables) { let cat; try { cat = JSON.parse(tht(ctx, ["schema", "columns", t.name, "--json"])); } catch (e) { // Fuzzy recovery: the model often misspells Italian table names // (e.g. "abellazione" vs "ablazione"). Auto-correct if a close // catalog match exists, preventing a multi-turn debug spiral. const fix = closestTableName(ctx, t.name); if (fix) { try { cat = JSON.parse(tht(ctx, ["schema", "columns", fix, "--json"])); t.id = fix; t.name = fix; } catch { /* fall through to error */ } } if (!cat) { const msg = (e.stderr || e.message || String(e)).toString().trim(); return textResult( `Tabella '${t.name}' non caricabile dal catalogo (${msg}). Proponi solo tabelle presenti nel catalogo (usa 'tht schema render' / 'tht search' per verificarne i nomi).`, ); } } if (!Array.isArray(cat.columns)) { return textResult(`Catalogo per '${t.name}' non contiene colonne valide. Verifica con 'tht schema columns ${t.name} --json'.`); } const suggested = new Set(t.suggested_columns ?? []); enriched.push({ id: t.id, name: t.name, kind: t.kind, recommended: t.kind === "promote" ? (t.recommended ?? true) : false, description: cat.description ?? "", rationale: t.rationale ?? "", columns: cat.columns.map((c) => ({ ...c, suggested: suggested.has(c.name) })), }); } const widget = buildSchemaLinkingRequest({ id: `u${Date.now()}`, phase, title, tables: enriched, }); const resp = await emitAndWait(ctx, widget); if (resp.control === "back") return textResult("Il reviewer vuole tornare indietro."); if (resp.control === "exit") return textResult("Il reviewer vuole uscire."); if (resp.control === "freetext") return textResult(`Altro (reviewer): ${resp.text}. Riformula tenendone conto.`); // Build the COMPLETE decision set first, persist it with ONE atomic ledger // write (decision add-batch): a per-item loop could fail halfway and leave // the audit ledger half-written, with duplicates on retry. const byId = new Map(enriched.map((t) => [t.id, t])); const toPersist = []; let n = 0; for (const rt of resp.tables ?? []) { const t = byId.get(rt.id); if (!t || !rt.enacted) continue; if (t.kind === "promote") { toPersist.push({ type: "table_promoted", subject: t.name, detail: t.description, rationale: t.rationale, }); n++; const sel = new Set(rt.columns ?? []); for (const c of t.columns) { const type = sel.has(c.name) ? "column_promoted" : (c.suggested ? "column_excluded" : null); if (!type) continue; toPersist.push({ type, subject: `${t.name}.${c.name}`, detail: c.description ?? "", }); } } else { toPersist.push({ type: "table_excluded", subject: t.name, detail: t.description, rationale: t.rationale, }); n++; } } if (toPersist.length) { const eBatch = relayIfThtFails( ctx, ["decision", "add-batch", "--session", session, "--doc", "-"], "Nessuna decisione registrata (batch atomico fallito): correggi e ripresenta il gate.", JSON.stringify(toPersist), ); if (eBatch) return eBatch; } // Deterministic projection of the ledger into schema_linking.json. const eSync = relayIfThtFails(ctx, ["session", "sync-schema-linking", session], ""); if (eSync) return eSync; const adv = advance ? forceAdvance(ctx, session) : { advanced: false }; const parts = [`Schema linking registrato dal reviewer (${n} tabelle + colonne curate). schema_linking.json scritto.`]; if (adv.advanced) parts.push("Fase avanzata automaticamente — nessun gate aggiuntivo necessario."); else if (advance && adv.error) parts.push(`Fase resta aperta: ${adv.error}`); return textResult(parts.join(" ")); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[reviewer_schema_linking ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); 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())), }), prepareArguments: prepareReviewerArguments, async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; try { const { session, kind, title } = params; let artifact = params.artifact; const curNum = currentPhase(ctx, session); const phase = phaseId(ctx, curNum); // F7's review target is the persisted artifact, never the model-authored // descriptor payload. A model can (and did) call reviewer_confirm with // `artifact.data: {}` even though write_final_sql had already persisted a // non-empty sql_final.sql; trusting that payload rendered a blank approval // dialog. Hydrate from the repository boundary and fail closed before any // widget is shown if the artifact is unavailable or empty. if (kind === "sql") { let docs; try { docs = JSON.parse(tht(ctx, ["session", "documents", session, "--json"])); } catch (e) { const msg = (e.stderr || e.message || String(e)).toString().trim(); return textResult( `Impossibile leggere sql_final.sql dalla sessione ${session}: ${msg}. ` + "Verifica l'artefatto persistito e riprova.", ); } const sqlDoc = Array.isArray(docs) ? docs.find((doc) => doc && doc.key === "sql" && doc.format === "sql") : null; const sql = typeof sqlDoc?.content === "string" ? sqlDoc.content : ""; if (!sql.trim()) { return textResult( `sql_final.sql assente o vuoto per la sessione ${session}: ` + "scrivi un SQL finale valido prima di ripresentare il gate F7.", ); } artifact = { ...artifact, kind: "sql", data: sql }; } // Sessione gia' finalizzata: non c'e' piu' nulla da approvare — mai // ripresentare il phase-gate (la form "approve") a workflow chiuso. if (kind === "phase" && curNum >= phaseMeta(ctx).max_phase) { try { const show = JSON.parse(tht(ctx, ["session", "show", session, "--json"])); if (show?.status === "finalized") { lockActive = false; return textResult( `La sessione ${session} e' gia' finalizzata: il workflow e' completo. ` + "Comunica al reviewer il riepilogo finale e termina il turno; NON " + "presentare altri gate reviewer_*.", ); } } catch { // show non leggibile: prosegui col gate normale } } // --- v2 payload construction (WS2) ------------------------------------- // The gate builds/validates the structured payload BEFORE showing the // widget, so the reviewer approves a deterministic artifact (not the raw // model data). Legacy (non-v2) payloads fall through unchanged. const data = artifact && typeof artifact === "object" ? artifact.data : undefined; const isV2 = data && typeof data === "object" && data.schema_version === 2; // derived from the (possibly rebuilt) cte_plan doc; the names list to persist. let ctePlanNames = null; if (kind === "cte_plan" && isV2) { const v = validateCtePlanV2(data); if (!v.legacy && !v.ok) { return textResult( "Piano CTE v2 non valido:\n- " + v.errors.join("\n- ") + "\nCorreggi il payload e ripresenta il gate.", ); } const getColumns = makeGetColumns(ctx); const enriched = enrichCtePlanV2(data, getColumns); ctePlanNames = enriched.ctes.map((c) => c.name); artifact = { ...artifact, data: enriched }; } else if (kind === "cte_result" && isV2) { const thinCheck = validateCteResultThin(data); if (!thinCheck.legacy && !thinCheck.ok) { return textResult( "Dati cte_result non validi:\n- " + thinCheck.errors.join("\n- "), ); } const cteName = tht(ctx, ["cte", "next", "--session", session]).trim(); if (!cteName) { return textResult( `Nessun CTE in attesa di approvazione (sessione ${session}).`, ); } let cteInfo; try { cteInfo = JSON.parse( tht(ctx, ["cte", "info", cteName, "--session", session, "--json"]), ); } catch (e) { const msg = (e.stderr || e.message || String(e)).toString().trim(); return textResult( `Impossibile leggere le info del CTE '${cteName}': ${msg}`, ); } if (!cteInfo.last_test || cteInfo.last_test.status === "error") { return textResult( `Il CTE '${cteName}' non ha un test con esito ok: devi eseguire ` + "`tht cte test` con esito ok prima di presentare il gate cte_result.", ); } let built = buildCteResultV2(data, cteInfo); // Enrich column descriptions against the plan's tables (best-effort). const planTables = []; if (cteInfo.doc && Array.isArray(cteInfo.doc.tables)) { for (const t of cteInfo.doc.tables) if (t && t.name) planTables.push(t.name); } built = enrichCteResultColumns(built, planTables, makeGetColumns(ctx)); artifact = { ...artifact, data: built }; } else if (kind === "phase" && isV2) { const v = validatePhaseSummaryV2(data); if (!v.legacy && !v.ok) { return textResult( "Riepilogo di fase v2 non valido:\n- " + v.errors.join("\n- "), ); } let enriched = enrichPhaseSummaryV2( data, phaseMetaForNum(ctx, curNum), makeGetColumns(ctx), ); // Recap deterministico: appende la sezione "decisioni di questa fase" // dal ledger (best-effort: un errore di lettura non blocca il gate). try { const show = JSON.parse( tht(ctx, ["session", "show", session, "--json"]), ); const meta = phaseMeta(ctx).phases.find((p) => p.num === curNum); enriched = appendLedgerSection( enriched, show.decisions, meta ? meta.emits : [], ); } catch { // ledger non leggibile: il riepilogo resta quello del modello } artifact = { ...artifact, data: enriched }; } const widget = buildArtifactGate({ id: `u${Date.now()}`, phase, title, artifact, action: { kind: "approve_reject", prompt: "Approvi o rifiuti?" }, }); let outcome; for (;;) { const resp = await emitAndWait(ctx, widget); outcome = resolveConfirmOutcome(resp); if (outcome.kind !== "unknown") break; await ctx.ui.notify("Scegli «Salva e procedi» o «Rifiuta».", "warning"); } if (outcome.kind === "freetext") return textResult( `Altro (reviewer): ${outcome.text}. Valuta e agisci, poi ri-presenta il gate.`, ); if (outcome.kind === "reject") return textResult("Rifiutato: rivedi e riprova."); if (outcome.kind === "back") return textResult("Il reviewer vuole tornare indietro."); if (outcome.kind === "exit") return textResult("Il reviewer vuole uscire."); // outcome.kind === "approve" -> execute the privileged action via the CLI. if (kind === "phase") { const r = advancePhaseAndFinalize(ctx, session, curNum); if (r.err) return r.err; if (r.finalized) { lockActive = false; lastSteered = false; return textResult( `Fase ${curNum} approvata e sessione finalizzata (${session}).`, ); } return textResult(`Fase approvata (sessione ${session}).`); } if (kind === "cte_plan") { // v2: names derive from data.ctes[]; persist the plan AND the chain doc // (payload A, enriched) via `tht cte plan --name ... --doc -`. // Legacy: the ordered names come from params.names. const names = ctePlanNames ?? (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, o data.ctes[].name nel payload v2).", ); } const planArgs = ["cte", "plan", "--session", session]; for (const n of names) planArgs.push("--name", n); if (ctePlanNames) { // v2: pass the enriched doc on stdin; the CLI validates the name list // matches (same order) and writes cte_plan_doc.json. planArgs.push("--doc", "-"); try { tht(ctx, planArgs, JSON.stringify(artifact.data)); } catch (e) { const cliMsg = (e.stderr || e.message || String(e)).toString().trim(); return textResult(`${cliMsg} Correggi il piano CTE e riprova.`); } } else { const err = relayIfThtFails(ctx, planArgs, ""); if (err) return err; } return textResult( `CTE plan approvato (${names.length} CTE, sessione ${session}).`, ); } if (kind === "cte_result") { // The CTE under review is always next_cte (plan order is enforced by // `tht cte test`). Approve it BY NAME: `cte_approved` keys on the CTE // name (phase.approved_ctes / next_cte); a `phase:N` subject is rejected // by `decision add` (exit 5) and never satisfies F6's advance prereq. const cteName = tht(ctx, [ "cte", "next", "--session", session, ]).trim(); if (!cteName) { return textResult( `Nessun CTE in attesa di approvazione (sessione ${session}).`, ); } const err = relayIfThtFails( ctx, ["decision", "add", "--session", session, "--type", "cte_approved", "--subject", cteName], "", ); if (err) return err; // Deterministic completeness: when no CTE remains in the plan the reviewer // has approved everything F6 can ask — close the phase here instead of // echoing a reviewer_confirm kind:"phase" that decides nothing new. const remaining = tht(ctx, ["cte", "next", "--session", session]).trim(); if (remaining) { return textResult( `CTE '${cteName}' approvato (sessione ${session}). Prossimo CTE del piano: '${remaining}'.`, ); } const closed = advancePhaseAndFinalize(ctx, session, curNum); if (closed.err) return closed.err; if (closed.finalized) { lockActive = false; lastSteered = false; return textResult( `CTE '${cteName}' approvato — era l'ultimo del piano: fase chiusa e sessione finalizzata (${session}).`, ); } return textResult( `CTE '${cteName}' approvato — era l'ultimo del piano: Fase ${curNum} chiusa automaticamente. ` + 'NON presentare reviewer_confirm kind:"phase" per questa fase; prosegui con la successiva.', ); } if (kind === "sql") { const err = relayIfThtFails( ctx, [ "decision", "add", "--session", session, "--type", "sql_approved", "--subject", `phase:${currentPhase(ctx, session)}`, ], "", ); if (err) return err; // F7's only prerequisite IS sql_approved (workflow.yaml): the phase gate // after it decides nothing new — close the phase here. const closed = advancePhaseAndFinalize(ctx, session, curNum); if (closed.err) return closed.err; if (closed.finalized) { lockActive = false; lastSteered = false; return textResult(`SQL approvato — fase chiusa e sessione finalizzata (${session}).`); } return textResult( `SQL approvato — Fase ${curNum} chiusa automaticamente. ` + 'NON presentare reviewer_confirm kind:"phase" per questa fase; prosegui con la successiva (datamart/memorie).', ); } return textResult(`Approvato (kind=${kind}, sessione ${session}).`); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[reviewer_confirm ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); // 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. // Auto-finalize after advancing the LAST phase: the model may not follow // through (token budget, turn end) leaving the session open. function advancePhaseAndFinalize(ctx, session, curNum) { const err = relayIfThtFails( ctx, ["phase", "advance", "--session", session], PHASE_RECOVERY, ); if (err) return { err }; if (curNum < phaseMeta(ctx).max_phase) return { finalized: false }; const fErr = relayIfThtFails( ctx, ["session", "finalize", session], "Fase approvata ma finalizzazione fallita: esegui manualmente " + `\`tht session finalize ${session}\`.`, ); if (fErr) return { err: fErr }; return { finalized: true }; } // La promozione memorie e' l'ultima interazione umana di F8: chiudere qui // (advance + finalize) evita il reviewer_confirm finale a workflow gia' // deciso (la form "approve" ridondante a fine sessione). Fuori dall'ultima // fase ritorna null e il chiamante mantiene il comportamento precedente. function closeAfterPromotion(ctx, session, curNum, summary) { if (curNum < phaseMeta(ctx).max_phase) return null; const r = advancePhaseAndFinalize(ctx, session, curNum); if (r.err) return r.err; lockActive = false; lastSteered = false; return textResult( `${summary} Fase ${curNum} approvata e sessione finalizzata (${session}): ` + "il workflow e' completo. Comunica al reviewer il riepilogo finale " + "della sessione e termina il turno; NON presentare altri gate reviewer_*.", ); } memoryGate.install(pi); disambiguationGate.install(pi); pi.registerTool({ name: "write_schema_linking", label: "Scrittura schema_linking.json (validata)", description: "Scrive deterministicamente schema_linking.json validandolo contro il modello " + "SchemaLinking via tht session set-schema-linking (evita edit a mano e la " + "validazione manuale). schema_linking e' l'oggetto JSON completo: {question, " + "candidates:[{kind:'table'|'column', name, evidence?, decision?}], joins:[{from, " + "to, source?}], excluded:[{kind, name}], open_questions:[], concept_formulas:[]}.", parameters: Type.Object({ session: Type.String(), schema_linking: Type.Any(), }), async execute(_id, params, _signal, _onUpdate, ctx) { lockActive = true; try { const { session } = params; // schema_linking is Type.Any(): a stringified JSON object passes validation // but would be double-encoded here and rejected by the CLI. Normalize first. const schema_linking = jsonObjectOrSelf(params.schema_linking); try { tht( ctx, ["session", "set-schema-linking", session, "--file", "-"], JSON.stringify(schema_linking), ); const message = `schema_linking.json scritto e validato per la sessione ${session}.`; if (currentPhase(ctx, session) !== 4) return textResult(message); const advanced = forceAdvance(ctx, session); if (advanced.error) { return textResult(`${message} Fase non avanzata: ${advanced.error}`); } return textResult(`${message} Fase avanzata automaticamente.`); } catch (e) { const cliMsg = (e.stderr || e.message || String(e)).toString().trim(); return textResult(`${cliMsg} Correggi schema_linking e riprova.`); } } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[write_schema_linking ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); pi.registerTool({ name: "write_cte_sql", label: "Scrittura CTE SQL", description: "Persiste un blocco CTE tramite tht cte save; non scrivere mai file di sessione direttamente.", parameters: Type.Object({ session: Type.String(), name: Type.String(), sql: Type.String() }), async execute(_id, params, _signal, _onUpdate, ctx) { try { const err = relayIfThtFails( ctx, ["cte", "save", "--session", params.session, "--name", params.name, "--file", "-"], "Correggi il blocco CTE e riprova.", params.sql, ); return err || textResult(`CTE ${params.name} salvato per la sessione ${params.session}.`); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[write_cte_sql ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); pi.registerTool({ name: "write_final_sql", label: "Scrittura SQL finale", description: "Persiste il SQL finale tramite tht sql set-final; non scrivere mai file di sessione direttamente.", parameters: Type.Object({ session: Type.String(), sql: Type.String() }), async execute(_id, params, _signal, _onUpdate, ctx) { try { const err = relayIfThtFails( ctx, ["sql", "set-final", "--session", params.session, "--file", "-"], "Correggi il SQL finale e riprova.", params.sql, ); return err || textResult(`SQL finale salvato per la sessione ${params.session}.`); } catch (fatal) { const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim(); return textResult(`[write_final_sql ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`); } }, }); // --- 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" }, ); }, }); }