1742 lines
68 KiB
JavaScript
1742 lines
68 KiB
JavaScript
// 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 { 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 { 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>", "</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);
|
|
}
|
|
|
|
// --- 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 "<la domanda dell\'utente nel messaggio sopra>"` 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 \"<domanda originale nel messaggio utente>\" --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 <id>` 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 <t> --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;
|
|
}
|
|
|
|
// 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));
|
|
}
|
|
|
|
// True when a reviewer_decide has no merito options but is allowed to close empty
|
|
// and advance (the empty memory phase F2). The gate then shows an info notice and
|
|
// auto-advances instead of presenting an empty checklist.
|
|
export function shouldSkipEmptyDecide({ meritCount, allowEmpty, advance }) {
|
|
return meritCount === 0 && !!allowEmpty && !!advance;
|
|
}
|
|
|
|
// Phase 2 uses a single positive checkbox meaning: apply this memory now.
|
|
// Candidates not checked by the reviewer remain undecided and may be shown again.
|
|
export function memorySelectionWidgetProps(options) {
|
|
return {
|
|
title: "Seleziona le memory da applicare alla domanda",
|
|
selected: options.filter((option) => option.recommended).map((option) => option.id),
|
|
selectionLabel: "memory da applicare",
|
|
confirmLabel: "Applica le memory selezionate",
|
|
};
|
|
}
|
|
|
|
// --- F8 memory-promotion gate: pure candidate->widget mapping (L1-tested) ------
|
|
// The candidates come from `tht memory promote --preview --json` (deterministic,
|
|
// reviewer-approved decisions only); the model never authors them.
|
|
export function promotionOptions(candidates) {
|
|
return candidates.map((c) => ({
|
|
id: `seq-${c.decision_seq}`,
|
|
label: `${c.type}: ${c.subject}`,
|
|
}));
|
|
}
|
|
|
|
export function promotionContent(candidates) {
|
|
return candidates
|
|
.map(
|
|
(c) =>
|
|
`- **${c.type}: ${c.subject}** (decisione #${c.decision_seq})\n` +
|
|
` ${c.detail || ""}\n` +
|
|
` Motivo: ${c.rationale || "—"}\n` +
|
|
` Domanda di contesto: ${c.question_context || "—"}`,
|
|
)
|
|
.join("\n");
|
|
}
|
|
|
|
export function splitPromotionChoices(candidates, choices) {
|
|
const chosen = new Set(choices ?? []);
|
|
const promote = [];
|
|
const decline = [];
|
|
for (const c of candidates) {
|
|
(chosen.has(`seq-${c.decision_seq}`) ? promote : decline).push(c);
|
|
}
|
|
return { promote, decline };
|
|
}
|
|
|
|
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;
|
|
|
|
// 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 (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` +
|
|
"<tht-sessione-skill>\n" +
|
|
SESSION_SKILL +
|
|
"\n</tht-sessione-skill>\n\n" +
|
|
inject +
|
|
(retrievalPack
|
|
? "\n\n<retrieval-pack>\n" + retrievalPack + "\n</retrieval-pack>"
|
|
: ""),
|
|
};
|
|
});
|
|
|
|
// 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;
|
|
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" },
|
|
);
|
|
}
|
|
});
|
|
|
|
// --- 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. 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: prepareReviewerArguments,
|
|
async execute(_id, params, _signal, _onUpdate, ctx) {
|
|
lockActive = true;
|
|
try {
|
|
const { session, title, options: opts, intro, advance } = params;
|
|
const typeErr = validateDecisionTypes(ctx, opts, session);
|
|
if (typeErr) return textResult(typeErr);
|
|
const phase = phaseId(ctx, currentPhase(ctx, session));
|
|
const recommended = opts.find((o) => o.recommended)?.id ?? null;
|
|
|
|
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);
|
|
const outcome = resolveSelectOutcome(opts, 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") {
|
|
const err = relayIfThtFails(
|
|
ctx,
|
|
decisionAddArgs(session, outcome.decision),
|
|
"",
|
|
);
|
|
if (err) return err;
|
|
if (advance) advanceIfReady(ctx, session);
|
|
return textResult(
|
|
`Decisione registrata (${outcome.decision.type}): ${outcome.option.label}.`,
|
|
);
|
|
}
|
|
return textResult(
|
|
`Scelta del reviewer: ${outcome.option ? outcome.option.label : outcome.choice}`,
|
|
);
|
|
} 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_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()),
|
|
}),
|
|
prepareArguments: prepareReviewerArguments,
|
|
async execute(_id, params, _signal, _onUpdate, ctx) {
|
|
lockActive = true;
|
|
try {
|
|
const { session, title, options: opts, advance } = params;
|
|
const typeErr = validateDecisionTypes(ctx, opts, session);
|
|
if (typeErr) return textResult(typeErr);
|
|
const phase = phaseId(ctx, currentPhase(ctx, session));
|
|
const toAdd = [];
|
|
const meritOptions = opts
|
|
.filter((o) => !isReserved(o.label))
|
|
.map((o) => ({ id: o.id, label: o.label }));
|
|
const joinOnly = meritOptions.length > 0 && opts
|
|
.filter((o) => !isReserved(o.label))
|
|
.every((o) => o.decision.type === "join_modified");
|
|
if (shouldSkipEmptyDecide({ meritCount: meritOptions.length, allowEmpty: params.allow_empty ?? false, advance })) {
|
|
await ctx.ui.notify(
|
|
"Nessuna memory riutilizzabile per questa domanda — passo alla fase successiva.",
|
|
"info",
|
|
);
|
|
advanceIfReady(ctx, session);
|
|
return textResult(
|
|
"Fase memoria vuota: nessuna decisione da registrare, avanzamento automatico alla fase successiva.",
|
|
);
|
|
}
|
|
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: phase === "F2" ? memorySelectionWidgetProps(opts).title : title,
|
|
allowEmpty: params.allow_empty ?? false,
|
|
options: meritOptions,
|
|
...(phase === "F2" ? memorySelectionWidgetProps(opts) : {}),
|
|
});
|
|
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.recommended ?? true,
|
|
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);
|
|
|
|
// 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_*.",
|
|
);
|
|
}
|
|
|
|
pi.registerTool({
|
|
name: "reviewer_memory_promote",
|
|
label: "Promozione memorie riusabili (reviewer)",
|
|
description:
|
|
"F8 (prima della chiusura di fase): propone al reviewer i candidati di promozione " +
|
|
"calcolati dalla CLI (tht memory promote --preview: tipi riusabili concept_clarified/" +
|
|
"table_promoted/table_excluded, max 5, esclusi i gia' promossi/rifiutati). Le selezioni " +
|
|
"vengono salvate nel vectordb (tht memory save-one) e registrate come memory_promoted; " +
|
|
"le deselezioni come memory_promotion_declined (non riproposte). Nessun parametro oltre " +
|
|
"alla sessione: i candidati sono deterministici, NON li scrivi tu. Registrata la " +
|
|
"promozione, in F8 il gate chiude la fase e finalizza la sessione da solo: NON " +
|
|
"presentare un reviewer_confirm dopo.",
|
|
parameters: Type.Object({
|
|
session: Type.String(),
|
|
}),
|
|
async execute(_id, params, _signal, _onUpdate, ctx) {
|
|
lockActive = true;
|
|
try {
|
|
const { session } = params;
|
|
const curNum = currentPhase(ctx, session);
|
|
const phase = phaseId(ctx, curNum);
|
|
let candidates;
|
|
try {
|
|
candidates = JSON.parse(
|
|
tht(ctx, ["memory", "promote", "--session", session, "--preview", "--json"]),
|
|
);
|
|
} catch (e) {
|
|
const msg = (e.stderr || e.message || String(e)).toString().trim();
|
|
return textResult(`Preview di promozione non disponibile: ${msg}`);
|
|
}
|
|
if (!Array.isArray(candidates) || candidates.length === 0) {
|
|
await ctx.ui.notify(
|
|
"Nessuna decisione riusabile da promuovere in memoria per questa sessione.",
|
|
"info",
|
|
);
|
|
const closed = closeAfterPromotion(
|
|
ctx, session, curNum, "Nessun candidato di promozione.",
|
|
);
|
|
if (closed) return closed;
|
|
return textResult(
|
|
"Nessun candidato di promozione: prosegui con la chiusura della sessione.",
|
|
);
|
|
}
|
|
const options = promotionOptions(candidates);
|
|
const widget = buildMultiselectRequest({
|
|
id: `u${Date.now()}`,
|
|
phase,
|
|
title: "Quali decisioni salvare nella memoria riutilizzabile?",
|
|
allowEmpty: true,
|
|
options,
|
|
selected: options.map((o) => o.id),
|
|
content: promotionContent(candidates),
|
|
});
|
|
const resp = await emitAndWait(ctx, widget);
|
|
if (resp.control === "freetext")
|
|
return textResult(`Altro (reviewer): ${resp.text}. Valuta e ripresenta il gate.`);
|
|
if (resp.control === "back")
|
|
return textResult("Il reviewer vuole tornare indietro.");
|
|
if (resp.control === "exit") return textResult("Il reviewer vuole uscire.");
|
|
const { promote, decline } = splitPromotionChoices(candidates, resp.choices);
|
|
let saved = 0;
|
|
for (const c of promote) {
|
|
const err = relayIfThtFails(
|
|
ctx,
|
|
["memory", "save-one", "--session", session,
|
|
"--decision", String(c.decision_seq), "--json"],
|
|
`Recupero manuale (umano): tht memory save-one --session ${session} --decision ${c.decision_seq}. Finora salvate: ${saved}.`,
|
|
);
|
|
if (err) return err;
|
|
const e2 = relayIfThtFails(ctx, decisionAddArgs(session, {
|
|
type: "memory_promoted", subject: c.subject,
|
|
detail: `seq:${c.decision_seq}`, rationale: c.rationale || c.detail || "",
|
|
}), `Memoria salvata nel vectordb ma decisione memory_promoted NON registrata: recupero manuale (umano) con tht decision add --session ${session} --type memory_promoted --subject "${c.subject}" --detail seq:${c.decision_seq}.`);
|
|
if (e2) return e2;
|
|
saved++;
|
|
}
|
|
for (const c of decline) {
|
|
const err = relayIfThtFails(ctx, decisionAddArgs(session, {
|
|
type: "memory_promotion_declined", subject: c.subject,
|
|
detail: `seq:${c.decision_seq}`,
|
|
}), "");
|
|
if (err) return err;
|
|
}
|
|
const summary =
|
|
`Promozione registrata: ${saved} memorie salvate nel vectordb, ` +
|
|
`${decline.length} candidati scartati.`;
|
|
const closed = closeAfterPromotion(ctx, session, curNum, summary);
|
|
if (closed) return closed;
|
|
return textResult(summary);
|
|
} catch (fatal) {
|
|
const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim();
|
|
return textResult(`[reviewer_memory_promote ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`);
|
|
}
|
|
},
|
|
});
|
|
|
|
pi.registerTool({
|
|
name: "rewrite_question",
|
|
label: "Riscrittura domanda e chiusura F3 (deterministica)",
|
|
description:
|
|
"In Fase 3 scrive deterministicamente question.md, registra question_rewritten e " +
|
|
"chiude la fase senza chiedere conferma al reviewer. 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;
|
|
try {
|
|
const { session, question, assumptions } = params;
|
|
const curNum = currentPhase(ctx, session);
|
|
if (curNum !== 3) {
|
|
return textResult("La riscrittura automatica e' disponibile solo in Fase 3.");
|
|
}
|
|
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;
|
|
const decisionErr = relayIfThtFails(
|
|
ctx,
|
|
decisionAddArgs(session, {
|
|
type: "question_rewritten",
|
|
subject: "domanda",
|
|
detail: question,
|
|
}),
|
|
"",
|
|
);
|
|
if (decisionErr) return decisionErr;
|
|
const advanced = advancePhaseAndFinalize(ctx, session, curNum);
|
|
if (advanced.err) return advanced.err;
|
|
return textResult(`Domanda riscritta e Fase 3 completata (sessione ${session}).`);
|
|
} catch (fatal) {
|
|
const msg = (fatal.stderr || fatal.message || String(fatal)).toString().trim();
|
|
return textResult(`[rewrite_question ERRORE INTERNO] ${msg}. Riprova o usa un approccio diverso.`);
|
|
}
|
|
},
|
|
});
|
|
|
|
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 <session_id> [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 <session_id> [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" },
|
|
);
|
|
},
|
|
});
|
|
}
|