From ad788278ddb2210a8688d649c75953e4f87b90c7 Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 6 Jul 2026 19:13:50 +0200 Subject: [PATCH] feat(gate): reviewer_schema_linking tool with per-column curation + deterministic sync Co-Authored-By: Claude Opus 4.8 --- .../__tests__/gate_schema_linking.test.js | 122 ++++++++++++++++++ .../__tests__/gate_stringified_params.test.js | 6 + harness/.pi/extensions/tht-gate.js | 107 +++++++++++++++ 3 files changed, 235 insertions(+) create mode 100644 harness/.pi/extensions/gate/__tests__/gate_schema_linking.test.js diff --git a/harness/.pi/extensions/gate/__tests__/gate_schema_linking.test.js b/harness/.pi/extensions/gate/__tests__/gate_schema_linking.test.js new file mode 100644 index 00000000..f234dcb6 --- /dev/null +++ b/harness/.pi/extensions/gate/__tests__/gate_schema_linking.test.js @@ -0,0 +1,122 @@ +// L1.5 roundtrip for the reviewer_schema_linking tool: drives the registered tool's +// execute() end-to-end against a fake Pi runtime, stubbing the `tht` shell (execFileSync) +// and ctx.ui.input, then asserts the recorded decision-add / sync-schema-linking argv. +// +// HARNESS ADAPTATIONS (no existing test invokes a tool's execute() end-to-end): +// 1. tht-gate.js uses ESM `import { execFileSync } from "node:child_process"`. Under +// CommonJS `require(esm)` the named binding tracks the module namespace, so patching +// `child_process.execFileSync` BEFORE requiring tht-gate.js is observed by tht(). +// 2. tht()'s loadEnvFromDotenv calls `require("node:fs")`; that `require` is undefined +// in the ESM module scope when loaded via require(esm) (in production Pi loads the +// extension as CommonJS, where require is native). We install a globalThis.require +// shim so loadEnvFromDotenv works. cwd points at a missing dir so the .env read +// fails gracefully inside its own try/catch. +const test = require("node:test"); +const assert = require("node:assert"); +const cp = require("node:child_process"); +const { createRequire } = require("node:module"); +const path = require("node:path"); + +const GATE = path.join(__dirname, "..", "..", "tht-gate.js"); + +// Adaptation #2: shim require() for the ESM-loaded-via-require module scope. +if (typeof globalThis.require === "undefined") { + globalThis.require = createRequire(GATE); +} + +const CATALOG = { + table: "dim_patient", + description: "Anagrafica", + columns: [ + { name: "cod_paz", description: "Codice", type: "bigint", pk: true }, + { name: "nome", description: "Nome", type: "text", pk: false }, + ], +}; + +test("reviewer_schema_linking records table/column decisions and syncs schema_linking", async () => { + const calls = []; + // Adaptation #1: stub the shell BEFORE requiring tht-gate.js. + const origExecFileSync = cp.execFileSync; + cp.execFileSync = (file, args) => { + calls.push(args.join(" ")); + if (args[0] === "phase" && args[1] === "meta") + return JSON.stringify({ phases: [{ num: 4, id: "F4" }] }); + if (args[0] === "phase" && args[1] === "show") return "Fase corrente: 4\n"; + if (args[0] === "schema" && args[1] === "columns" && args[2] === "dim_patient") + return JSON.stringify(CATALOG); + return ""; // decision add, session sync-schema-linking, ... + }; + + try { + const gate = require(GATE); + const { createFakePi } = require("./fake_pi_runtime.js"); + const { pi, ctx, tools } = createFakePi(); + ctx.cwd = "/nonexistent-thothii-test-cwd"; // .env read fails gracefully + gate.default(pi); + + // Reflect the descriptor id back and return the reviewer response so emitAndWait + // accepts it (its resp.id === descriptor.id invariant). The descriptor id is a + // runtime-generated `u${Date.now()}` — capture it, do not hardcode. + let capturedDescriptor = null; + ctx.ui.input = async (title) => { + capturedDescriptor = JSON.parse(title); + return JSON.stringify({ + id: capturedDescriptor.id, + kind: "schema-linking", + tables: [{ id: "t-pat", enacted: true, columns: ["cod_paz"] }], + }); + }; + + const tool = tools.get("reviewer_schema_linking"); + assert.ok(tool, "reviewer_schema_linking must be registered"); + + const result = await tool.def.execute( + "call-1", + { + session: "s1", + title: "Schema linking", + tables: [ + { + id: "t-pat", + name: "dim_patient", + kind: "promote", + suggested_columns: ["cod_paz", "nome"], + }, + ], + }, + null, + null, + ctx, + ); + + // The emitted descriptor carried the enriched catalog columns. + assert.equal(capturedDescriptor.widget, "schema-linking"); + assert.equal(capturedDescriptor.tables[0].columns.length, 2); + + // cod_paz was selected -> promoted; nome was suggested but deselected -> excluded. + assert.ok( + calls.some((c) => /decision add .*--type table_promoted --subject dim_patient\b/.test(c)), + `expected table_promoted dim_patient in: ${JSON.stringify(calls)}`, + ); + assert.ok( + calls.some((c) => /decision add .*--type column_promoted --subject dim_patient\.cod_paz\b/.test(c)), + `expected column_promoted dim_patient.cod_paz in: ${JSON.stringify(calls)}`, + ); + assert.ok( + calls.some((c) => /decision add .*--type column_excluded --subject dim_patient\.nome\b/.test(c)), + `expected column_excluded dim_patient.nome in: ${JSON.stringify(calls)}`, + ); + assert.ok( + calls.some((c) => /^session sync-schema-linking s1$/.test(c)), + `expected session sync-schema-linking s1 in: ${JSON.stringify(calls)}`, + ); + // sync runs AFTER the decisions are recorded. + const syncIdx = calls.findIndex((c) => c.startsWith("session sync-schema-linking")); + const lastDecisionIdx = calls.map((c) => c.startsWith("decision add")).lastIndexOf(true); + assert.ok(syncIdx > lastDecisionIdx, "sync must run after all decision adds"); + + assert.match(result.content[0].text, /Schema linking registrato/); + } finally { + cp.execFileSync = origExecFileSync; + } +}); diff --git a/harness/.pi/extensions/gate/__tests__/gate_stringified_params.test.js b/harness/.pi/extensions/gate/__tests__/gate_stringified_params.test.js index 8acc417a..3e87e645 100644 --- a/harness/.pi/extensions/gate/__tests__/gate_stringified_params.test.js +++ b/harness/.pi/extensions/gate/__tests__/gate_stringified_params.test.js @@ -40,6 +40,12 @@ test("prepareReviewerArguments leaves an object artifact and parses stringified assert.equal(out.options[0].id, "a"); }); +test("prepareReviewerArguments parses stringified tables array", () => { + const out = prepareReviewerArguments({ tables: JSON.stringify([{ id: "t1", name: "x" }]) }); + assert.ok(Array.isArray(out.tables)); + assert.equal(out.tables[0].id, "t1"); +}); + test("GATE_CODE_FILES blocks writes to gate extensions, not session artifacts", () => { assert.ok(GATE_CODE_FILES.test("harness/.pi/extensions/tht-gate.js")); assert.ok(GATE_CODE_FILES.test(".pi/extensions/reserved-labels.mjs")); diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index 07708fd4..7d781fc8 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -29,6 +29,7 @@ import { buildSelectRequest, buildMultiselectRequest, buildArtifactGate, + buildSchemaLinkingRequest, } from "./gate/builders.js"; import { isReserved } from "./reserved-labels.mjs"; @@ -67,6 +68,14 @@ export function prepareReviewerArguments(input) { /* 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. @@ -610,6 +619,104 @@ export default function (pi) { }, }); + 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; + const { session, title, tables, advance } = params; + const phase = phaseId(ctx, currentPhase(ctx, session)); + + // Enrich each table with its full catalog columns (deterministic source). + const enriched = tables.map((t) => { + const cat = JSON.parse(tht(ctx, ["schema", "columns", t.name, "--json"])); + const suggested = new Set(t.suggested_columns ?? []); + return { + 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.`); + + const byId = new Map(enriched.map((t) => [t.id, t])); + let n = 0; + for (const rt of resp.tables ?? []) { + const t = byId.get(rt.id); + if (!t || !rt.enacted) continue; + if (t.kind === "promote") { + const e1 = relayIfThtFails(ctx, decisionAddArgs(session, { + type: "table_promoted", subject: t.name, detail: t.description, rationale: t.rationale, + }), ""); + if (e1) return e1; + 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; + const e2 = relayIfThtFails(ctx, decisionAddArgs(session, { + type, subject: `${t.name}.${c.name}`, detail: c.description ?? "", + }), ""); + if (e2) return e2; + } + } else { + const e3 = relayIfThtFails(ctx, decisionAddArgs(session, { + type: "table_excluded", subject: t.name, detail: t.description, rationale: t.rationale, + }), ""); + if (e3) return e3; + n++; + } + } + + // Deterministic projection of the ledger into schema_linking.json. + const eSync = relayIfThtFails(ctx, ["session", "sync-schema-linking", session], ""); + if (eSync) return eSync; + + if (advance) advanceIfReady(ctx, session); + return textResult(`Schema linking registrato dal reviewer (${n} tabelle + colonne curate).`); + }, + }); + pi.registerTool({ name: "reviewer_confirm", label: "Gate di avanzamento (reviewer)",