fix(harness): remediation difetti review — gate↔CLI, D15, D7/D6, D14, robustezza

Implementazione del piano di remediation progressiva sui difetti emersi
dall'analisi dell'harness. Tutto verificato: 214 test Python (incl. L0 su
Postgres reale), 14 test JS del gate, ruff pulito.

Blocco 1 (CRITICA, integrazione gate↔CLI):
- phase advance: gate usa --auto + exit 6; reviewer_confirm kind:phase fa
  advance esplicito che applica i prerequisiti (prima non avanzava per le
  fasi a conferma umana).
- cte plan riceve i --name dal gate (param names); set-question con id
  posizionale; skill `tht search find`; nuovo comando `tht memory save-one`
  con dedup hash client-side in save_one_memory.

Blocco 2 (D15, stato post-rollback):
- campo `phase` su DecisionRecord + effective_decisions phase-aware per i
  subject "a nome" (cte_approved ecc.); _compute_promotions e finalize sulla
  vista effective; finalize confronta col piano CTE effettivo, non glob;
  `decision add --retracts` + comando `decision retract`.

Blocco 3 (D7 read-only + D6 manifest):
- assert_read_only su tutti e quattro i codepath (direct + REST);
- manifest author/summary/updated_at/updated_by/schema_version popolati +
  helper touch_manifest sulle mutazioni.

Blocco 4-5 (D14a/D14b):
- decision_min_phase data-driven via `emits:` in workflow.yaml;
- formula evidence: status auto, search_formulas, gruppo CLI `tht formula`,
  `search find --kind formula`, load_evidence_dir salta i .sql.md.

Blocco 6 (robustezza):
- taskdoc slice promoted_tables + bound enforced; report escaping/bound +
  rsplit note; filtro kind reader REST/direct; conteggio upserted robusto;
  guard REST run_query non-list; LSH disallineato -> LshIndexError.

Blocco 7 (pulizia):
- dead code gate e KIND_TO_TABLE morto rimossi; doc Postgres-only
  (README + connection.py).

Blocco 0 (parziale): test di compatibilità firma gate↔CLI
(tests/integration). Rinviati: fake-Pi runtime completo, artifact-gate da
disco (#23), parità eligibility REST/direct (#28), unificazione
reserved-labels (#30), memory_rejected da deselezione (#33).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-06-27 17:16:51 +02:00
co-authored by Claude Opus 4.8
parent daf33f77fe
commit c4d130828f
36 changed files with 912 additions and 83 deletions
+33 -27
View File
@@ -14,8 +14,8 @@
// - 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,
// 7 = not-ready silent no-op)
// - 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
@@ -29,11 +29,8 @@ import {
buildSelectRequest,
buildMultiselectRequest,
buildArtifactGate,
buildInfoRequest,
buildFreetextRequest,
withChildLinkage,
} from "./gate/builders.js";
import { ALTRO, BACK_LABEL, QUIT_LABEL, CONTROL_LABELS, isReserved, stripReserved } from "./reserved-labels.mjs";
import { isReserved, stripReserved } from "./reserved-labels.mjs";
// --- anti-bypass block lists (spec D4, verbatim from source L169-177) -----------
const FORBIDDEN = [
@@ -116,29 +113,22 @@ function phaseName(ctx, num) {
const p = meta.phases.find((x) => x.num === num);
return p ? p.name : "?";
}
function maxPhase(ctx) {
return phaseMeta(ctx).max_phase;
}
// The schema-linking phase = the phase whose artifacts_out contains schema_linking.json.
function schemaLinkingPhase(ctx) {
const meta = phaseMeta(ctx);
const p = meta.phases.find((x) => x.artifacts_out.includes("schema_linking.json"));
return p ? p.num : 5;
}
function currentPhase(ctx, session) {
const out = tht(ctx, ["phase", "show", "--session", session]);
const m = out.match(/Fase corrente:\s*(\d+)/);
return m ? parseInt(m[1], 10) : 1;
}
// tht phase advance --if-ready: exit 7 = not ready (silent no-op), others propagated.
// 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", "--if-ready", "--session", session]);
tht(ctx, ["phase", "advance", "--auto", "--session", session]);
return { advanced: true };
} catch (e) {
if (e.status === 7) return { advanced: false };
if (e.status === 6) return { advanced: false };
return { advanced: false, error: (e.stderr || e.message || String(e)).toString().trim() };
}
}
@@ -191,7 +181,6 @@ export default function (pi) {
let lockActive = false;
let lastSteered = false;
let pendingKickoff = null;
let ollamaReady = false;
let activeSessionId = null;
// 1) ANTI-BYPASS tool_call hook (spec D4, verbatim). Blocks direct phase/decision
@@ -270,7 +259,6 @@ export default function (pi) {
lockActive = false;
lastSteered = false;
pendingKickoff = null;
ollamaReady = false;
activeSessionId = null;
_phaseMetaCache = null;
});
@@ -431,6 +419,8 @@ export default function (pi) {
data: Type.Any(),
version: Type.Optional(Type.Number()),
}),
// kind:"cte_plan" only -- ordered list of CTE names to persist (tht cte plan --name).
names: Type.Optional(Type.Array(Type.String())),
}),
async execute(_id, params, _signal, _onUpdate, ctx) {
lockActive = true;
@@ -453,15 +443,28 @@ export default function (pi) {
if (resp.control === "exit") return textResult("Il reviewer vuole uscire.");
// approved -> execute the privileged action via the CLI.
if (kind === "phase") {
// Try auto-advance first; exit 6 = needs human (already handled by this dialog).
const r = advanceIfReady(ctx, session);
if (!r.advanced && r.error) return textResult(`${r.error} ${PHASE_RECOVERY}`);
// Explicit human approval: advance unconditionally except for unmet
// prerequisites. Plain `phase advance` (no --auto) enforces advance_problems
// and exits 6 with the missing items, which relayIfThtFails surfaces.
const err = relayIfThtFails(ctx, ["phase", "advance", "--session", session], PHASE_RECOVERY);
if (err) return err;
return textResult(`Fase approvata (sessione ${session}).`);
}
if (kind === "cte_plan") {
const err = relayIfThtFails(ctx, ["cte", "plan", "--session", session], "");
// The CTE plan is the ordered list of CTE names; the model passes them in
// params.names (tht cte plan requires at least one --name).
const names = Array.isArray(params.names) ? params.names : [];
if (names.length === 0) {
return textResult(
"Nessun nome CTE fornito: il piano CTE richiede l'elenco ordinato dei CTE " +
"(parametro names di reviewer_confirm).",
);
}
const planArgs = ["cte", "plan", "--session", session];
for (const n of names) planArgs.push("--name", n);
const err = relayIfThtFails(ctx, planArgs, "");
if (err) return err;
return textResult(`CTE plan approvato (sessione ${session}).`);
return textResult(`CTE plan approvato (${names.length} CTE, sessione ${session}).`);
}
if (kind === "cte_result" || kind === "sql") {
const dt = kind === "sql" ? "sql_approved" : "cte_approved";
@@ -499,7 +502,10 @@ export default function (pi) {
assumps = [assumps];
}
}
const args = ["session", "set-question", "--session", session, "--question", question];
// `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));
}
+7 -6
View File
@@ -60,7 +60,7 @@ about a domain term, ask the reviewer.
`schema_linking.json`. Write the artifacts with care; they are the decision surface.
8. **Candidates are candidates, not truth.** Present LSH/vector/evidence matches with
their **provenance** (LSH / vector / evidence) and their scores, never as absolute
truth. The reviewer may reject them. Verify filter values with `tht search
truth. The reviewer may reject them. Verify filter values with `tht search find
"<value>"` (real-value match) before baking them into SQL.
9. **Open ambiguities are explicit.** If an ambiguity can't be resolved, offer a
`reviewer_decide` option "Leave ambiguity open" with a rationale, so the reviewer
@@ -77,8 +77,8 @@ about a domain term, ask the reviewer.
Prerequisite: you must already be in Phase 1.
1. Explore the DWH and knowledge base: `tht search "<term>"` (evidence + schema, LSH
over real values) and `tht search --kind evidence "<term>"`. The LSH exposes
1. Explore the DWH and knowledge base: `tht search find "<term>"` (evidence + schema, LSH
over real values) and `tht search find --kind evidence "<term>"`. The LSH exposes
EVERY column where a value appears — it does not collapse to a single best match,
so a value like "ablazione" may anchor on multiple columns.
2. For each ambiguity (clinical term, population, time window, outcome), present a
@@ -154,9 +154,10 @@ Prerequisite: Phase 3 closed.
with a `value_grounded` option for each candidate column (the LSH exposes all of
them, not collapsed to the best match). The reviewer chooses the anchor(s).
4. **Concept formula (D14b).** If a concept (e.g. "fascia pediatrica", "stesso anno")
has a candidate SQL formula (found in the evidence or derived from context),
present it and let the reviewer approve/reject
(`concept_formula_approved`/`concept_formula_rejected`).
has a candidate SQL formula, retrieve it with `tht search find --kind formula
"<concept>"` (or derive it from the evidence/context), present it, and let the
reviewer approve/reject (`concept_formula_approved`/`concept_formula_rejected`).
Reflect the approved formula in `schema_linking.json` (`concept_formulas`).
5. Write `schema_linking.json` (the Phase 4 artifact) and close with
`reviewer_confirm kind:"phase"`. Do NOT run `tht session check` (that's Phase 5).
@@ -42,7 +42,7 @@ by hand to the join.
- Do the filters (WHERE/HAVING) reflect ALL the conditions of the rewritten question?
- Are aggregations, groupings and orderings the required ones?
- Empty or zero result: almost always indicates a problem in conditions or joins.
Verify the filter values with `tht search "<value>"` (match on real values).
Verify the filter values with `tht search find "<value>"` (match on real values).
- Do the joins follow those promoted in `schema_linking.json`? (exception: the FK
`data_time_key → dim_time.day_key` is not declared, see the time section.)
- Time analyses: are you using `JOIN dim_time` and not key arithmetic?
+5
View File
@@ -33,6 +33,11 @@ Keys are never logged; URLs are fine. Rotate any key that appeared in chat.
A workspace wires the relational DWH + the pgvector (dual-key) + embeddings + evidence.
See `workspaces/tht.example.yaml`. `${THT_*}}` tokens expand from `.env`.
> **DB support (MVP):** the `direct` transport supports **PostgreSQL only** (psycopg2
> driver, `pg_*` catalog introspection, postgres-dialect sqlcheck/EXPLAIN). The central
> production DWH is reached via the `rest` transport. Multi-dialect `direct` support
> (sqlserver/mariadb/informix, cf. `Thoth/thoth_sqldb2`) is separate future work.
### Per-customer workspace repo (deployment shape)
ThothII is generic; **evidence content and LSH indexes are per-customer** and live in a
@@ -0,0 +1,74 @@
"""Integration: every `tht ...` command the gate invokes must exist in the Typer CLI.
Root cause of the Blocco 1 critical bugs: the gate (tht-gate.js) and the Python CLI
were ported separately and never run together, so the gate called commands/flags that
did not exist (`phase advance --if-ready`, `cte plan` without `--name`, `set-question
--session` on a positional arg). This test extracts every `["group","sub",...,"--flag"]`
array literal from tht-gate.js and asserts, via `tht <group> <sub> --help`, that the
subcommand exists (exit 0) and that each long flag used is offered. It makes that whole
class of drift impossible to reintroduce silently.
"""
from __future__ import annotations
import re
import subprocess
import sys
from pathlib import Path
import pytest
_ROOT = Path(__file__).resolve().parent.parent.parent
_GATE = _ROOT / ".pi" / "extensions" / "tht-gate.js"
_THT = Path(sys.executable).parent / "tht"
# The command groups the gate drives. Anything else in an array literal is data, not a CLI call.
_GROUPS = {"phase", "session", "cte", "decision", "memory", "search", "schema", "vector"}
# Flags that are framework/JS artifacts, never real CLI options (skip from the check).
_SKIP_FLAGS: set[str] = set()
def _extract_invocations() -> list[tuple[str, str, list[str], str]]:
"""Returns (group, subcommand, long_flags, raw) for each gate CLI call site."""
src = _GATE.read_text()
# array literal opening with two string literals: ["group", "sub" ...
pattern = re.compile(r'\[\s*"([a-z-]+)"\s*,\s*"([a-z][a-z-]*)"((?:\s*,\s*[^\]\[]+?)?)\]')
out: list[tuple[str, str, list[str], str]] = []
for m in pattern.finditer(src):
group, sub, tail = m.group(1), m.group(2), m.group(3)
if group not in _GROUPS:
continue
flags = [f for f in re.findall(r'"(--[a-z][a-z-]*)"', tail) if f not in _SKIP_FLAGS]
out.append((group, sub, flags, m.group(0)))
return out
_INVOCATIONS = _extract_invocations()
def test_gate_invokes_at_least_the_known_commands():
"""Guards the extractor itself: if it silently matches nothing, the test is useless."""
pairs = {(g, s) for g, s, _, _ in _INVOCATIONS}
assert ("phase", "advance") in pairs
assert ("cte", "plan") in pairs
assert ("session", "set-question") in pairs
assert ("decision", "add") in pairs
@pytest.mark.skipif(not _THT.exists(), reason="tht CLI not installed in this venv")
@pytest.mark.parametrize("group,sub,flags,raw", _INVOCATIONS, ids=lambda v: v if isinstance(v, str) else None)
def test_gate_call_site_matches_cli(group: str, sub: str, flags: list[str], raw: str):
res = subprocess.run(
[str(_THT), group, sub, "--help"],
capture_output=True, text=True, cwd=str(_ROOT),
)
assert res.returncode == 0, (
f"gate calls `tht {group} {sub}` but it does not exist in the CLI.\n"
f"call site: {raw}\nstderr: {res.stderr}"
)
help_text = res.stdout + res.stderr
for flag in flags:
assert flag in help_text, (
f"gate passes `{flag}` to `tht {group} {sub}` but the CLI does not offer it.\n"
f"call site: {raw}"
)
+54
View File
@@ -0,0 +1,54 @@
"""Blocco 6: robustness fixes -- taskdoc slice/bound, report escaping, upsert count."""
from tht.report import _markdown_table, extract_reviewer_notes
from tht.taskdoc import generate_task_doc
def test_taskdoc_slices_to_promoted_tables(tmp_path):
s = tmp_path / "sess"
s.mkdir()
(s / "question.md").write_text("q")
(s / "schema_linking.json").write_text(
'{"question":"q","candidates":['
'{"kind":"table","name":"pazienti","decision":"promoted"},'
'{"kind":"table","name":"ricoveri","decision":"promoted"}],'
'"joins":[],"excluded":[],"open_questions":[]}'
)
doc = generate_task_doc(session_dir=s, phase=4, promoted_tables=["pazienti"])
assert "pazienti" in doc.body
assert "ricoveri" not in doc.body # sliced out
def test_taskdoc_truncates_over_budget(tmp_path):
s = tmp_path / "sess"
s.mkdir()
(s / "question.md").write_text("# Domanda\n" + "x" * 200_000)
doc = generate_task_doc(session_dir=s, phase=1)
assert doc.byte_budget_ok is False
assert len(doc.body.encode()) <= 80_000
assert "troncato" in doc.body
def test_markdown_table_escapes_pipes_and_newlines():
table = _markdown_table(["c"], [("a|b\nc",)])
# the cell must not introduce a raw pipe or newline that breaks the row
body_line = table.splitlines()[2]
assert "\\|" in body_line
assert "\n" not in body_line
def test_extract_reviewer_notes_uses_last_heading():
report = (
"## Note del reviewer\nnella cella di dati appariva questo testo\n"
"## Note del reviewer\nnota vera del reviewer"
)
assert extract_reviewer_notes(report) == "nota vera del reviewer"
def test_upsert_count_handles_postgrest_list_wrapping():
from unittest.mock import MagicMock
from tht.vectorstore.rest_client import VectorRestClient
client = VectorRestClient.__new__(VectorRestClient)
client._call = MagicMock(return_value=[{"upserted": 7}]) # list-wrapped scalar
assert client.upsert_records("memory", [{}, {}]) == 7
+37
View File
@@ -0,0 +1,37 @@
"""decision_min_phase must cover every substantive decision type (data-driven via emits).
The reference bug: only the ~5 types referenced in prerequisites had a min phase; all
others (table_promoted, cte_approved, value_grounded, concept_formula_*, ...) defaulted
to phase 1, so require_phase_or_exit would accept them far too early. workflow.yaml now
declares `emits` per phase as the source of truth.
"""
import pytest
from tht.workflow import load_workflow
@pytest.mark.parametrize(
"dtype,expected",
[
("concept_clarified", 1),
("memory_rejected", 2),
("question_rewritten", 3),
("table_promoted", 4),
("column_corrected", 4),
("evidence_accepted", 4),
("value_grounded", 4),
("concept_formula_approved", 4),
("concept_formula_rejected", 4),
("cte_approved", 6),
("sql_approved", 7),
("sql_revised", 7),
("datamart_requested", 8),
],
)
def test_decision_min_phase_from_emits(dtype, expected):
assert load_workflow().decision_min_phase(dtype) == expected
@pytest.mark.parametrize("meta", ["phase_approved", "phase_reopened", "decision_retracted"])
def test_meta_types_not_phase_gated(meta):
assert load_workflow().decision_min_phase(meta) == 1
@@ -0,0 +1,46 @@
"""D15 granularity (a): `tht decision retract` tombstones the last substantive decision.
Wires the step-granularity rollback ("re-ask current widget, discard last answer") to a
reachable CLI command. The data model (decision_retracted + retracts, honored by
effective_decisions) already existed; this pins the command that emits it.
"""
from pathlib import Path
from tht.cli.decision_cmd import retract_cmd
from tht.decisions import append_decision, list_decisions
from tht.phase import current_phase, effective_decisions
def _walk_to_phase(session: Path, target: int) -> None:
while current_phase(session) < target:
append_decision(session, type="phase_approved", subject=f"phase:{current_phase(session)}")
def test_retract_drops_last_substantive_decision(tmp_path, monkeypatch):
s = tmp_path / "2026-01-01-000000-x"
s.mkdir(parents=True)
_walk_to_phase(s, 4)
append_decision(s, type="table_promoted", subject="phase:4", detail="dim_pazienti")
append_decision(s, type="table_promoted", subject="phase:4", detail="fact_ricoveri")
# stub config + session loading (the command only needs a session dir)
import tht.cli.decision_cmd as mod
class _Cfg:
class paths:
sessions = tmp_path
monkeypatch.setattr(mod, "_load_config_or_exit", lambda _c: _Cfg())
monkeypatch.setattr(mod, "load_session_or_exit", lambda _cfg, _s: None)
monkeypatch.setattr(mod, "session_dir", lambda _cfg, sid: tmp_path / sid)
retract_cmd(session="2026-01-01-000000-x", config=Path("x"))
eff = effective_decisions(s)
promoted = [d for d in eff if d.type == "table_promoted"]
assert len(promoted) == 1 # the last one was retracted
assert promoted[0].detail == "dim_pazienti"
# the audit log keeps everything (append-only): 2 promotions + the retract marker
raw = [d.type for d in list_decisions(s)]
assert raw.count("table_promoted") == 2
assert "decision_retracted" in raw
@@ -0,0 +1,47 @@
"""D15 (§4.8): effective_decisions must stale-filter name-subject decisions too.
The reference bug (CRITICA #1): effective_decisions only filtered decisions whose
subject was "phase:N". cte_approved uses the CTE *name* as subject, so a cte_approved
from F6 stayed "effective" after a rollback to F4 -- a stale decision that still counted
(the exact invariant §4.8 forbids). The fix records the emitting phase on each
DecisionRecord (high-water-mark) and effective_decisions falls back to it when the
subject is not "phase:N".
"""
from pathlib import Path
from tht.decisions import append_decision
from tht.phase import approved_ctes, current_phase, effective_decisions
def _walk_to_phase(session: Path, target: int) -> None:
"""Approva in ordine fino a raggiungere `target` (il fold avanza solo su n==cur)."""
while current_phase(session) < target:
append_decision(session, type="phase_approved", subject=f"phase:{current_phase(session)}")
def test_cte_approved_name_subject_is_stale_after_reopen(tmp_path):
s = tmp_path / "sess"
s.mkdir()
_walk_to_phase(s, 6) # ora in F6
# piano CTE + approvazione di un CTE (subject = NOME del CTE, non "phase:6")
(s / "cte_plan.json").write_text('["pazienti_base"]')
append_decision(s, type="cte_approved", subject="pazienti_base")
assert "pazienti_base" in approved_ctes(s)
assert current_phase(s) == 6
# rollback a F4: la cte_approved di F6 deve diventare stale (esclusa dalla vista)
append_decision(s, type="phase_reopened", subject="phase:4")
assert current_phase(s) == 4
assert "pazienti_base" not in approved_ctes(s), (
"cte_approved (subject a nome) di F6 NON deve contare dopo un reopen a F4"
)
types = [d.type for d in effective_decisions(s)]
assert "cte_approved" not in types
def test_phase_field_recorded_on_append(tmp_path):
s = tmp_path / "sess"
s.mkdir()
_walk_to_phase(s, 4)
rec = append_decision(s, type="value_grounded", subject="ablazione")
assert rec.phase == 4 # emitting phase recorded for the high-water-mark filter
+37
View File
@@ -0,0 +1,37 @@
"""D14b wiring: status=auto, search_formulas, and evidence loader skips formula files.
Completes the formula layer beyond the store: the `auto` status the spec requires
(§4.7.2), the concept-substring retrieval that `tht search find --kind formula` uses,
and the guarantee that load_evidence_dir does NOT choke on *.sql.md formula files when
they live under the evidence root.
"""
from tht.evidence.formula_store import ConceptFormula, save_formula, search_formulas
from tht.evidence.model import EvidenceDoc, load_evidence_dir
def test_status_auto_is_valid():
f = ConceptFormula(concept="x", sql="SELECT 1", status="auto")
assert f.status == "auto"
# round-trips through parse/dump
assert ConceptFormula.parse(f.dump()).status == "auto"
def test_search_formulas_substring_case_insensitive(tmp_path):
save_formula(tmp_path, ConceptFormula(concept="fascia pediatrica", sql="SELECT 1"))
save_formula(tmp_path, ConceptFormula(concept="indice di Charlson", sql="SELECT 2"))
hits = search_formulas(tmp_path, "PEDIATRICA")
assert len(hits) == 1
assert hits[0].concept == "fascia pediatrica"
assert search_formulas(tmp_path, "charlson")[0].concept == "indice di Charlson"
def test_load_evidence_dir_skips_formula_files(tmp_path):
# a real evidence doc + a formula file under the same root
(tmp_path / "ev1.md").write_text(
"---\nid: ev1\ntitle: T\n---\nbody text\n"
)
save_formula(tmp_path, ConceptFormula(concept="ablazione", sql="SELECT 1"))
docs = load_evidence_dir(tmp_path)
ids = [d.id for d in docs]
assert ids == ["ev1"] # the .sql.md formula file is skipped, no crash
assert all(isinstance(d, EvidenceDoc) for d in docs)
+43
View File
@@ -0,0 +1,43 @@
"""D6 §5.2: the manifest's new fields are actually written (not just declared).
The reference bug: author/summary/updated_at/updated_by/schema_version existed on the
model but no code ever populated them. These tests pin that create_session fills them
and that mutations (set_question, touch_manifest) bump updated_at/updated_by.
"""
from tht.config import DatabaseConfig
from tht.session.store import create_session, load_session, set_question, touch_manifest
def _db() -> DatabaseConfig:
return DatabaseConfig(host="h", port=5432, database="dw", user="u", password="p", schema="dwh")
def test_create_session_populates_new_fields(tmp_path, monkeypatch):
monkeypatch.setenv("THT_AUTHOR", "alice@psd")
m = create_session("Quanti ricoveri per ablazione nel 2024?", _db(), tmp_path)
assert m.author == "alice@psd"
assert m.updated_by == "alice@psd"
assert m.summary and m.summary.startswith("Quanti ricoveri")
assert m.created_at is not None
assert m.updated_at == m.created_at
assert m.schema_version is not None # from workflow.yaml
def test_default_author_when_env_absent(tmp_path, monkeypatch):
monkeypatch.delenv("THT_AUTHOR", raising=False)
m = create_session("x", _db(), tmp_path)
assert m.author == "dev@local"
def test_touch_manifest_bumps_updated(tmp_path, monkeypatch):
monkeypatch.setenv("THT_AUTHOR", "alice@psd")
m = create_session("domanda", _db(), tmp_path)
created = m.updated_at
monkeypatch.setenv("THT_AUTHOR", "bob@psd")
set_question(m.id, "domanda riscritta", [], tmp_path)
reloaded = load_session(m.id, tmp_path)
assert reloaded.updated_by == "bob@psd"
assert reloaded.updated_at >= created
# touch_manifest helper is reachable and updates the field
touch_manifest(m.id, tmp_path, updated_by="carol@psd")
assert load_session(m.id, tmp_path).updated_by == "carol@psd"
+37
View File
@@ -0,0 +1,37 @@
"""D7: read-only enforcement lives in the execute layer, not only in CLI callers.
assert_read_only is the shared structural guard both codepaths (direct + REST) call,
so reusing the execute API (e.g. from the backend) cannot bypass single-statement +
SELECT-only. Pins that writes/multi-statement are rejected with ExecutionError.
"""
import pytest
from tht.execute import ExecutionError, assert_read_only
@pytest.mark.parametrize(
"sql",
[
"DELETE FROM pazienti",
"UPDATE pazienti SET x = 1",
"INSERT INTO pazienti VALUES (1)",
"DROP TABLE pazienti",
"SELECT 1; DROP TABLE pazienti", # multi-statement
"TRUNCATE pazienti",
],
)
def test_write_or_multistatement_rejected(sql):
with pytest.raises(ExecutionError):
assert_read_only(sql)
@pytest.mark.parametrize(
"sql",
[
"SELECT 1",
"WITH x AS (SELECT 1) SELECT * FROM x",
"SELECT a FROM t UNION SELECT b FROM u",
],
)
def test_select_allowed(sql):
assert_read_only(sql) # no raise
+2
View File
@@ -42,6 +42,7 @@ from tht.cli.datamart_cmd import datamart_app # noqa: E402
from tht.cli.db_cmd import db_app # noqa: E402
from tht.cli.decision_cmd import decision_app # noqa: E402
from tht.cli.evidence_cmd import evidence_app # noqa: E402
from tht.cli.formula_cmd import formula_app # noqa: E402
from tht.cli.lsh_cmd import lsh_app # noqa: E402
from tht.cli.memory_cmd import memory_app # noqa: E402
from tht.cli.phase_cmd import phase_app # noqa: E402
@@ -59,6 +60,7 @@ app.add_typer(vector_app, name="vector")
app.add_typer(memory_app, name="memory")
app.add_typer(search_app, name="search")
app.add_typer(evidence_app, name="evidence")
app.add_typer(formula_app, name="formula")
app.add_typer(db_app, name="db")
app.add_typer(decision_app, name="decision")
app.add_typer(sql_app, name="sql")
+52 -1
View File
@@ -17,6 +17,10 @@ def add_cmd(
subject: str = typer.Option(..., "--subject", help="Oggetto (tabella, colonna, concetto)."),
detail: str = typer.Option("", "--detail"),
rationale: str = typer.Option("", "--rationale"),
retracts: int = typer.Option(
None, "--retracts",
help="Solo per type=decision_retracted: seq della decisione da ritirare (D15).",
),
config: Path = CONFIG_OPT,
) -> None:
"""Registra una decisione del reviewer nella sessione."""
@@ -31,6 +35,12 @@ def add_cmd(
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=1)
if type == "decision_retracted" and retracts is None:
typer.secho(
"ERRORE: decision_retracted richiede --retracts <seq> (la decisione da ritirare).",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=1)
from tht.cli.phase_cmd import require_phase_or_exit
from tht.workflow import load_workflow
@@ -60,12 +70,53 @@ def add_cmd(
raise typer.Exit(code=5)
record = append_decision(
sdir, type=type, subject=subject,
detail=detail, rationale=rationale,
detail=detail, rationale=rationale, retracts=retracts,
)
typer.secho(f"OK: decisione [{record.seq}] {record.type}: {record.subject}",
fg=typer.colors.GREEN)
@decision_app.command("retract")
def retract_cmd(
session: str = typer.Option(..., "--session", help="Id della sessione."),
config: Path = CONFIG_OPT,
) -> None:
"""Ritira l'ultima decisione sostanziale della fase corrente (D15 granularita' step).
Granularita' (a) del rollback §4.8: 'rispondi di nuovo a questa domanda'. Scrive un
marker decision_retracted (append-only, l'audit resta) che effective_decisions onora;
il widget corrente puo' essere riproposto. Non cambia la fase."""
from tht.decisions import append_decision
from tht.phase import effective_decisions
cfg = _load_config_or_exit(config)
load_session_or_exit(cfg, session)
sdir = session_dir(cfg, session)
# Ultima decisione NON-meta della vista effective = quella associata al widget corrente.
meta = {
"phase_approved", "phase_auto_approved", "phase_reopened",
"phase_skipped", "decision_retracted",
}
substantive = [d for d in effective_decisions(sdir) if d.type not in meta]
if not substantive:
typer.secho(
"Nessuna decisione sostanziale da ritirare nella fase corrente.",
fg=typer.colors.YELLOW, err=True,
)
raise typer.Exit(code=6)
target = substantive[-1]
record = append_decision(
sdir, type="decision_retracted", subject=target.subject,
rationale=f"ritira [{target.seq}] {target.type}", retracts=target.seq,
)
typer.secho(
f"OK: ritirata decisione [{target.seq}] {target.type}: {target.subject} "
f"(marker #{record.seq}).",
fg=typer.colors.GREEN,
)
@decision_app.command("list")
def list_cmd(
session: str = typer.Option(..., "--session"),
+81
View File
@@ -0,0 +1,81 @@
"""tht formula -- concept->SQL formula evidence store (spec D14b, §4.7.2).
Authoring + listing of reusable concept formulas. Retrieval for the workflow is via
`tht search find --kind formula "<concept>"`. The reviewer approves/rejects a candidate
formula in F4 (decisions concept_formula_approved / concept_formula_rejected), and the
approved one is reflected in schema_linking.json (SchemaLinking.concept_formulas).
"""
from pathlib import Path
import typer
from tht.cli.config_cmd import CONFIG_OPT
from tht.cli.evidence_cmd import evidence_root
from tht.cli.schema_cmd import _load_config_or_exit
formula_app = typer.Typer(help="Formule di concetto (concept -> SQL) riusabili (D14b)")
@formula_app.command("save")
def save_cmd(
concept: str = typer.Option(..., "--concept", help="Il concetto, es. 'fascia pediatrica'."),
sql_file: Path = typer.Option(
None, "--sql-file", help="File con l'espressione SQL della formula."
),
sql: str = typer.Option(None, "--sql", help="Espressione SQL inline (alternativa a --sql-file)."),
column: list[str] = typer.Option(None, "--column", help="Colonna usata (ripetibile)."),
status: str = typer.Option("draft", "--status", help="auto | draft | reviewed."),
source: list[str] = typer.Option(None, "--source", help="Fonte/evidence (ripetibile)."),
config: Path = CONFIG_OPT,
) -> None:
"""Salva una formula di concetto sotto artifacts/evidence/formulas/."""
from tht.evidence.formula_store import ConceptFormula, save_formula
cfg = _load_config_or_exit(config)
if sql_file is None and not sql:
typer.secho("ERRORE: indica --sql-file <path> oppure --sql \"<espr>\".",
fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
body = sql_file.read_text() if sql_file else sql
try:
formula = ConceptFormula(
concept=concept, columns=column or [], sql=body.strip(),
status=status, sources=source or [],
)
except ValueError as e:
typer.secho(f"ERRORE: formula non valida: {e}", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
path = save_formula(evidence_root(cfg), formula)
typer.secho(f"OK: formula '{concept}' salvata in {path}.", fg=typer.colors.GREEN)
@formula_app.command("list")
def list_cmd(
concept: str = typer.Option(None, "--concept", help="Filtra per concetto (substring)."),
json_out: bool = typer.Option(False, "--json"),
config: Path = CONFIG_OPT,
) -> None:
"""Elenca le formule salvate (tutte, o filtrate per concetto)."""
import json as _json
from rich.console import Console
from rich.table import Table
from tht.evidence.formula_store import _load_all, search_formulas
cfg = _load_config_or_exit(config)
root = evidence_root(cfg)
formulas = search_formulas(root, concept) if concept else _load_all(root)
if json_out:
typer.echo(_json.dumps([f.model_dump(mode="json") for f in formulas],
ensure_ascii=False, indent=2))
return
if not formulas:
typer.secho("Nessuna formula salvata.", fg=typer.colors.YELLOW)
return
t = Table(title="Formule di concetto")
for col in ("Concetto", "Status", "Colonne"):
t.add_column(col)
for f in formulas:
t.add_row(f.concept, f.status, ", ".join(f.columns))
Console().print(t)
+59 -1
View File
@@ -12,7 +12,11 @@ from sqlalchemy.exc import OperationalError, ProgrammingError
from tht.cli.config_cmd import CONFIG_OPT
from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli._guards import require_server_profile, require_vector_write_allowed
from tht.cli._guards import (
has_vector_write_rest,
require_server_profile,
require_vector_write_allowed,
)
from tht.cli.session_cmd import load_session_or_exit, session_dir
from tht.cli.vector_cmd import require_vector_cfg
@@ -132,6 +136,60 @@ def promote_cmd(
typer.secho(f"ATTENZIONE: {warning}", fg=typer.colors.YELLOW)
@memory_app.command("save-one")
def save_one_cmd(
session: str = typer.Option(..., "--session"),
decision: int = typer.Option(
..., "--decision", help="decision_seq della decisione da salvare come memoria."
),
json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."),
config: Path = CONFIG_OPT,
) -> None:
"""Upsert mirato (una riga) della memoria di una decisione su pgvector via writer key (D11).
Promuove la decisione nel registro locale (idempotente) e fa un singolo upsert
remoto con dedup hash client-side -- niente full-resync. Abilita il salvataggio
di una memoria da postazione remota (workstation) con la sola writer key.
"""
import json as _json
from tht.cli.vector_cmd import make_embedder
from tht.memory import load_registry, promote, save_one_memory
from tht.vectorstore.rest_client import VectorRestClient
cfg = _load_config_or_exit(config)
manifest = load_session_or_exit(cfg, session)
require_vector_write_allowed(cfg, "memory save-one")
if not has_vector_write_rest(cfg):
typer.secho(
"ERRORE: `memory save-one` richiede la sezione `vector_write_rest` con una "
"API key di upsert nel workspace yaml (upsert remoto via writer key).",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=4)
sdir = session_dir(cfg, session)
# Promuove la decisione scelta nel registro locale (idempotente: salta se gia' presente
# o se stale post-rollback, perche' _compute_promotions usa la vista effective).
promote(sdir, manifest, seqs=[decision], registry_path=registry_path(cfg))
records = [r for r in load_registry(registry_path(cfg)) if r.session_id == manifest.id]
writer = VectorRestClient(cfg.vector_write_rest)
embedder = make_embedder(cfg.embeddings)
count = save_one_memory(records, decision, writer=writer, embedder=embedder)
msg = (
f"{count} memoria salvata su pgvector (decision_seq {decision})."
if count
else f"Nessun upsert (decisione {decision} assente/stale o memoria gia' aggiornata)."
)
if json_out:
typer.echo(_json.dumps(
{"upserted": count, "decision_seq": decision, "message": msg}, ensure_ascii=False))
return
typer.secho(f"OK: {msg}", fg=typer.colors.GREEN if count else typer.colors.YELLOW)
@memory_app.command("clear")
def clear_cmd(
yes: bool = typer.Option(False, "--yes", "-y", help="Salta la richiesta di conferma."),
+26 -1
View File
@@ -10,6 +10,7 @@ KIND_MAP = {
"evidence": ["evidence"],
"schema": ["schema_table", "schema_column"],
"values": [], # solo LSH
"formula": [], # solo formula store (D14b), niente LSH/vector
}
# Default di `--top` per le famiglie diverse da `schema` (numero di risultati). Per `schema`
@@ -30,7 +31,7 @@ def search_cmd(
"(default: 12 tabelle per schema, 10 altrimenti).",
),
kind: str = typer.Option(
None, "--kind", help="Filtra per famiglia: evidence | schema | values."
None, "--kind", help="Filtra per famiglia: evidence | schema | values | formula."
),
explain: bool = typer.Option(False, "--explain", help="Mostra anche il testo matchato."),
json_out: bool = typer.Option(
@@ -57,6 +58,30 @@ def search_cmd(
if top is None:
top = cfg.search.top_schema_tables if kind == "schema" else DEFAULT_TOP_FALLBACK
if kind == "formula":
# D14b: recupero formule di concetto dallo store locale (niente LSH/vector).
from tht.cli.evidence_cmd import evidence_root
from tht.evidence.formula_store import search_formulas
formulas = search_formulas(evidence_root(cfg), keyword)[:top]
if json_out:
typer.echo(json.dumps(
[f.model_dump(mode="json") for f in formulas], ensure_ascii=False, indent=2))
return
if not formulas:
typer.secho(f"Nessuna formula per '{keyword}'.", fg=typer.colors.YELLOW)
return
table = Table(title=f"Formule per '{keyword}'")
table.add_column("Concetto")
table.add_column("Status")
table.add_column("Colonne")
table.add_column("SQL")
for f in formulas:
sql_preview = (f.sql[:80] + "…") if len(f.sql) > 80 else f.sql
table.add_row(f.concept, f.status, ", ".join(f.columns), sql_preview)
Console().print(table)
return
lsh_hits = None
try:
lsh, minhashes, meta = load_index(
+16 -8
View File
@@ -163,7 +163,8 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP
from tht.execute.warnings import plan_warnings, runtime_warnings, static_warnings
from tht.report import extract_reviewer_notes, render_validation_report
from tht.session.artifacts import build_evidence_entries
from tht.decisions import list_decisions
from tht.phase import cte_plan as effective_cte_plan
from tht.phase import effective_decisions
from tht.session.models import SchemaLinking
from tht.sqlcheck import validate_sql
@@ -186,8 +187,9 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP
raise typer.Exit(code=5)
# --- gate di ingresso ---
# Vista effective (D15): un sql_approved/cte ritirato o stale post-rollback non conta.
problems = session_problems(cfg, session_id)
decisions = list_decisions(sdir)
decisions = effective_decisions(sdir)
sql_file = sdir / "sql_final.sql"
if not sql_file.exists():
problems.append("sql_final.sql assente")
@@ -196,8 +198,10 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP
"decisione sql_approved assente: la validazione semantica del reviewer "
"e' obbligatoria prima del finalize"
)
cte_dir = sdir / "ctes"
if cte_dir.is_dir():
# I CTE richiesti sono quelli del PIANO effettivo, non i file glob su disco: un
# ctes/*.sql orfano lasciato da un teardown incompleto non deve bloccare il finalize.
plan = effective_cte_plan(sdir)
if plan:
try:
tested = {r.name for r in load_cte_tests(sdir)}
except CteError as e:
@@ -206,9 +210,9 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=3)
for cte_file in sorted(cte_dir.glob("*.sql")):
if cte_file.stem not in tested:
problems.append(f"CTE mai testato: {cte_file.stem}")
for name in plan:
if name not in tested:
problems.append(f"CTE mai testato: {name}")
if problems:
typer.secho(f"Finalize rifiutato per {session_id}:", fg=typer.colors.YELLOW)
for p in problems:
@@ -258,9 +262,13 @@ def finalize_cmd(session_id: str = typer.Argument(...), config: Path = CONFIG_OP
(sdir / "evidence.json").write_text(json.dumps(entries, ensure_ascii=False, indent=2))
# --- manifest + riepilogo ---
from tht.session.store import MANIFEST
from datetime import UTC, datetime
from tht.session.store import MANIFEST, current_author
manifest.status = "finalized"
manifest.updated_at = datetime.now(UTC)
manifest.updated_by = current_author()
manifest.to_yaml(sdir / MANIFEST)
typer.secho(f"OK: sessione {session_id} finalizzata. Artefatti:", fg=typer.colors.GREEN)
for name in ARTIFACT_FILES:
+5
View File
@@ -4,6 +4,11 @@ from tht.config import DatabaseConfig
def make_engine(cfg: DatabaseConfig) -> Engine:
# MVP: il transport `direct` supporta SOLO PostgreSQL (driver psycopg2; introspezione
# su cataloghi pg_*; sqlcheck/EXPLAIN dialetto postgres). Il DWH centrale di
# produzione usa il transport `rest`. Il supporto multi-dialetto (sqlserver/mariadb/
# informix, cfr. Thoth/thoth_sqldb2) e' lavoro futuro separato: aggiungere un campo
# dialect/driver a DatabaseConfig e astrarre introspect/sqlcheck/connection.
url = (
f"postgresql+psycopg2://{cfg.user}:{cfg.password}"
f"@{cfg.host}:{cfg.port}/{cfg.database}"
+12
View File
@@ -56,6 +56,11 @@ class DecisionRecord(BaseModel):
rationale: str = ""
# D15: se type == "decision_retracted", indica quale seq viene ritirata.
retracts: int | None = None
# D15: la fase corrente al momento della scrittura (high-water-mark). Permette a
# effective_decisions() di marcare stale le decisioni il cui subject NON e' "phase:N"
# (es. cte_approved usa il nome del CTE) dopo un rollback. None per record storici
# (pre-fix) o costruiti a mano: in quel caso si ricade sulla logica basata sul subject.
phase: int | None = None
def list_decisions(session_dir: Path) -> list[DecisionRecord]:
@@ -78,6 +83,12 @@ def append_decision(
rationale: str = "",
retracts: int | None = None,
) -> DecisionRecord:
# Fase corrente PRIMA dell'append (high-water-mark D15). Import lazy: phase.py
# importa decisions.py (ciclo). Per i marker di fase (subject "phase:N") il valore
# e' ridondante col subject; per le decisioni sostanziali con subject "a nome"
# (cte_approved, ...) e' l'unico modo per filtrarle dopo un reopen.
from tht.phase import current_phase
record = DecisionRecord(
seq=len(list_decisions(session_dir)) + 1,
ts=datetime.now(UTC),
@@ -86,6 +97,7 @@ def append_decision(
detail=detail,
rationale=rationale,
retracts=retracts,
phase=current_phase(session_dir),
)
path = session_dir / DECISIONS_FILE
path.parent.mkdir(parents=True, exist_ok=True)
+20 -9
View File
@@ -27,7 +27,9 @@ class ConceptFormula(BaseModel):
concept: str
columns: list[str] = []
sql: str
status: Literal["draft", "reviewed"] = "draft"
# auto = sintetizzata dal modello (non ancora rivista); draft = bozza umana;
# reviewed = approvata da un revisore. (spec §4.7.2: status auto/draft/reviewed)
status: Literal["auto", "draft", "reviewed"] = "draft"
sources: list[str] = []
@property
@@ -81,20 +83,29 @@ def save_formula(root: Path | str, formula: ConceptFormula) -> Path:
return path
def retrieve_formula(root: Path | str, concept: str) -> list[ConceptFormula]:
"""All formulas for `concept` under <root>/formulas/. Empty list if none (or if
the dir is absent). Multiple results mean competing drafts/versions for the same
concept -- the caller (gate) lets the reviewer pick."""
root = Path(root)
def _load_all(root: Path) -> list[ConceptFormula]:
formulas_dir = root / FORMULAS_SUBDIR
if not formulas_dir.is_dir():
return []
out: list[ConceptFormula] = []
for f in sorted(formulas_dir.glob("*.sql.md")):
try:
formula = ConceptFormula.parse(f.read_text())
out.append(ConceptFormula.parse(f.read_text()))
except ValueError:
continue # malformed file: skip, don't crash retrieval
if formula.concept == concept:
out.append(formula)
return out
def retrieve_formula(root: Path | str, concept: str) -> list[ConceptFormula]:
"""All formulas matching `concept` exactly under <root>/formulas/. Empty list if
none (or if the dir is absent). Multiple results mean competing drafts/versions for
the same concept -- the caller (gate) lets the reviewer pick."""
return [f for f in _load_all(Path(root)) if f.concept == concept]
def search_formulas(root: Path | str, query: str) -> list[ConceptFormula]:
"""Formulas whose concept contains `query` (case-insensitive). Used by
`tht search find --kind formula` (D14b retrieval, §4.7.2): the reviewer searches a
concept term and gets the candidate formulas to approve before they reach the CTE."""
q = query.strip().lower()
return [f for f in _load_all(Path(root)) if q in f.concept.lower()]
+4
View File
@@ -58,5 +58,9 @@ def load_evidence_dir(root: Path) -> list[EvidenceDoc]:
for f in sorted(root.rglob("*.md")):
if f.name.upper().startswith("README"):
continue
# I file formula (concept->SQL, frontmatter diverso) vivono sotto formulas/ con
# estensione .sql.md: non sono EvidenceDoc, li gestisce formula_store (D14b).
if f.name.endswith(".sql.md"):
continue
docs.append(EvidenceDoc.parse(f.read_text(), path=f))
return docs
+25 -2
View File
@@ -37,6 +37,24 @@ def _inject_limit(sql: str, limit: int) -> tuple[str, bool]:
return ast.limit(limit + 1).sql(dialect="postgres"), True
def assert_read_only(sql: str) -> None:
"""Guard read-only strutturale condiviso dai due codepath di esecuzione (D7).
Difesa in profondita': oltre alla rete server (utente RO / READ ONLY tx / RPC
SELECT-only), ogni esecuzione passa da qui, cosi' anche un riuso diretto delle API
Python (es. dal backend) non puo' bypassare il single-statement + SELECT-only.
Usa sqlcheck.validate_sql senza schema fisico: parse + un solo statement +
read-only per struttura + funzioni vietate.
"""
from tht.sqlcheck import validate_sql
check = validate_sql(sql)
if not check.ok:
raise ExecutionError(
"SQL rifiutato (read-only enforcement): " + "; ".join(check.errors)
)
def _translate_error(e: DBAPIError) -> ExecutionError:
msg = str(e.orig)
if "statement timeout" in msg or "canceling statement" in msg:
@@ -48,7 +66,8 @@ def _translate_error(e: DBAPIError) -> ExecutionError:
def run_controlled(engine: Engine, sql: str, *, limit: int, timeout_ms: int) -> ExecResult:
"""Esecuzione nelle 4 reti: utente RO (a monte), transazione READ ONLY,
statement_timeout, LIMIT iniettato via AST."""
statement_timeout, LIMIT iniettato via AST. Guard read-only client-side a monte."""
assert_read_only(sql)
final_sql, injected = _inject_limit(sql, limit)
start = time.monotonic()
with engine.connect() as conn:
@@ -71,7 +90,11 @@ def run_controlled(engine: Engine, sql: str, *, limit: int, timeout_ms: int) ->
def explain(engine: Engine, sql: str, *, timeout_ms: int) -> PlanSummary:
"""EXPLAIN (FORMAT JSON), mai ANALYZE: il piano si stima, non si esegue."""
"""EXPLAIN (FORMAT JSON), mai ANALYZE: il piano si stima, non si esegue.
assert_read_only a monte: l'EXPLAIN interpola lo SQL, quindi il single-statement +
SELECT-only va garantito anche qui (non solo nel chiamante CLI)."""
assert_read_only(sql)
with engine.connect() as conn:
trans = conn.begin()
try:
+14 -2
View File
@@ -67,9 +67,21 @@ def load_index(directory: Path, name: str) -> tuple[Any, Any, dict]:
def query_index(lsh, minhashes, keyword: str, meta: dict, top_n: int = 10) -> list[LshHit]:
"""Query con score: i parametri MinHash vengono dal meta dell'indice, non dalla config."""
"""Query con score: i parametri MinHash vengono dal meta dell'indice, non dalla config.
Se i due pickle (`*_lsh.pkl` e `*_minhashes.pkl`) sono disallineati (rigenerati
separatamente) `lsh.query` puo' restituire chiavi assenti da `minhashes`: si solleva
LshIndexError azionabile invece di un KeyError opaco."""
qmh = create_minhash(meta["signature_size"], keyword, meta["n_gram"])
scored = [(key, qmh.jaccard(minhashes[key][0])) for key in lsh.query(qmh)]
scored = []
for key in lsh.query(qmh):
entry = minhashes.get(key)
if entry is None:
raise LshIndexError(
"Indice LSH disallineato (lsh.pkl e minhashes.pkl non coerenti): "
"rigenera con `tht lsh build`."
)
scored.append((key, qmh.jaccard(entry[0])))
scored.sort(key=lambda kv: kv[1], reverse=True)
return [
LshHit(table=minhashes[k][1], column=minhashes[k][2], value=minhashes[k][3], score=s)
+16 -4
View File
@@ -6,7 +6,7 @@ from pathlib import Path
from pydantic import BaseModel
from tht.decisions import DecisionRecord, DecisionType, list_decisions
from tht.decisions import DecisionRecord, DecisionType
from tht.session.models import SessionManifest
from tht.vectorstore.records import VectorRecord
@@ -132,7 +132,10 @@ def _compute_promotions(
session_dir: Path, manifest: SessionManifest, *,
seqs: list[int] | None, existing: list[MemoryRecord],
) -> list[MemoryRecord]:
decisions = list_decisions(session_dir)
# Vista effective (D15, §4.8): non promuovere decisioni stale dopo un rollback.
from tht.phase import effective_decisions
decisions = effective_decisions(session_dir)
selected = decisions if seqs is None else [d for d in decisions if d.seq in seqs]
already = {(r.session_id, r.decision_seq) for r in existing}
n = _next_id_num(existing)
@@ -234,7 +237,12 @@ def save_one_memory(
"""Targeted one-row upsert of a promoted decision to pgvector via the writer key
(spec D11). This is NOT a full vectorstore resync: it embeds and pushes a single
record, so a workstation with a writer key can publish one memory without
rebuilding the index. Returns the upsert count (0 if no record matched).
rebuilding the index. Returns the upsert count (0 if no record matched or the
record is already up to date).
Hash dedup client-side (spec §5.4): the SHA-256 of the content is compared with
the writer's existing_vector_hashes; the embedding (Ollama round-trip) and the
upsert are skipped when the content is unchanged. Idempotent by construction.
`writer` is a VectorRestClient (writer key); `embedder` an embeddings client.
The destructive cleanup (sync's delete-stale step) is intentionally absent: it
@@ -246,11 +254,15 @@ def save_one_memory(
record = memory_vector_record_for_decision(records, decision_seq)
if record is None:
return 0
new_hash = content_hash(record.content)
existing = writer.existing_hashes("memory", ["memory"])
if existing.get(record.id) == new_hash:
return 0 # unchanged: skip embedding + upsert
embedding = embedder.embed_documents([record.content])[0]
row = {
"record_key": record.id,
"kind": record.kind,
"content_hash": content_hash(record.content),
"content_hash": new_hash,
"metadata": pack_metadata(record),
"embedding": embedding,
}
+6 -2
View File
@@ -85,8 +85,12 @@ def effective_decisions(session_dir: Path) -> list[DecisionRecord]:
out: list[DecisionRecord] = []
for d in _audit_excluding_retracted(session_dir):
n = _phase_num(d.subject)
# decisioni senza subject di fase (es. concept_clarified puo' avere subject libero):
# le ammettiamo (non sono legate a una fase specifica da filtrare).
# Per le decisioni con subject "a nome" (es. cte_approved -> nome CTE, evidence_*
# -> id evidence) il subject non porta la fase: si usa la fase emittente registrata
# (d.phase, high-water-mark D15). Senza nessuno dei due (record storici) la
# decisione e' ammessa: non c'e' modo di datarla, e il subject non e' di fase.
if n is None:
n = d.phase
if n is None or n <= cur:
out.append(d)
return out
+26 -5
View File
@@ -3,14 +3,31 @@ from tht.sqlcheck import CheckResult
NOTE_HEADING = "## Note del reviewer"
# Bound del report (D16): celle di dati clinici (referti, anamnesi) possono essere
# enormi; senza limiti il report.md esplode (il caso 344KB segnalato) e i '|'/newline
# nei dati corrompono la tabella markdown.
_MAX_CELL_CHARS = 200
_MAX_TABLE_ROWS = 50
def _escape_cell(value: object) -> str:
"""Rende una cella sicura per il markdown: niente '|' o newline grezzi, larghezza
limitata. Evita corruzione della tabella e injection di markdown da dati DB."""
s = str(value).replace("|", "\\|").replace("\r", " ").replace("\n", " ")
if len(s) > _MAX_CELL_CHARS:
s = s[:_MAX_CELL_CHARS] + "…"
return s
def _markdown_table(columns: list[str], rows: list[tuple]) -> str:
lines = [
"| " + " | ".join(columns) + " |",
"| " + " | ".join(_escape_cell(c) for c in columns) + " |",
"|" + "---|" * len(columns),
]
for row in rows:
lines.append("| " + " | ".join(str(v) for v in row) + " |")
for row in rows[:_MAX_TABLE_ROWS]:
lines.append("| " + " | ".join(_escape_cell(v) for v in row) + " |")
if len(rows) > _MAX_TABLE_ROWS:
lines.append(f"| _… {len(rows) - _MAX_TABLE_ROWS} righe in più non mostrate …_ |")
return "\n".join(lines)
@@ -68,7 +85,11 @@ def render_validation_report(
def extract_reviewer_notes(report_md: str) -> str:
"""Estrae il contenuto della sezione note da un report esistente."""
"""Estrae il contenuto della sezione note da un report esistente.
rsplit sull'ULTIMA occorrenza dell'heading: la sezione note e' sempre in coda, e
se per caso la stringa dell'heading comparisse nei dati di preview (ora comunque
con newline escapati, quindi mai a inizio riga) non corromperebbe l'estrazione."""
if NOTE_HEADING not in report_md:
return ""
return report_md.split(NOTE_HEADING, 1)[1].strip()
return report_md.rsplit(NOTE_HEADING, 1)[1].strip()
+11 -1
View File
@@ -9,12 +9,15 @@ il client mantiene solo l'iniezione del LIMIT (per il rilevamento del troncament
import time
from tht.execute import ExecResult, ExecutionError, PlanSummary, _inject_limit
from tht.execute import ExecResult, ExecutionError, PlanSummary, _inject_limit, assert_read_only
from tht.rest.client import RestError
from tht.rest.explain import parse_text_plan
def run_controlled_rest(client, sql: str, *, limit: int) -> ExecResult:
# Guard read-only client-side anche sul path REST (D7): non delegare l'unica verifica
# al server. Stesso check strutturale del path diretto.
assert_read_only(sql)
final_sql, injected = _inject_limit(sql, limit)
start = time.monotonic()
try:
@@ -22,6 +25,12 @@ def run_controlled_rest(client, sql: str, *, limit: int) -> ExecResult:
except RestError as e:
raise ExecutionError(str(e)) from e
elapsed_ms = int((time.monotonic() - start) * 1000)
if not isinstance(rows_dicts, list):
# run_query atteso come lista di righe; una risposta inattesa (dict d'errore,
# scalare) non deve crashare con KeyError/TypeError opaco.
raise ExecutionError(
f"Risposta REST inattesa da run_query (atteso elenco di righe): {type(rows_dicts).__name__}"
)
columns = list(rows_dicts[0].keys()) if rows_dicts else []
rows = [tuple(r.get(c) for c in columns) for r in rows_dicts]
truncated = injected and len(rows) > limit
@@ -31,6 +40,7 @@ def run_controlled_rest(client, sql: str, *, limit: int) -> ExecResult:
def explain_rest(client, sql: str) -> PlanSummary:
assert_read_only(sql)
try:
lines = client.explain_query(sql)
except RestError as e:
+43 -1
View File
@@ -1,3 +1,4 @@
import os
from datetime import UTC, datetime
from pathlib import Path
@@ -7,12 +8,26 @@ from tht.textutil import slugify
MANIFEST = "session_manifest.yaml"
MAX_SLUG_CHARS = 40
MAX_SUMMARY_CHARS = 120
class SessionError(Exception):
pass
def current_author() -> str:
"""Autore della sessione (D6). Nel profilo `none` dell'MVP l'operatore e' sulla
propria macchina: default `dev@local`. Con auth reale il backend passa l'id via
THT_AUTHOR. Una sola fonte per author/updated_by."""
return os.environ.get("THT_AUTHOR", "dev@local").strip() or "dev@local"
def _summarize(question: str) -> str:
"""Domanda sintetica per il manifest/la sidebar (D6 §5.2): prima riga, troncata."""
first = question.strip().splitlines()[0].strip() if question.strip() else ""
return first[:MAX_SUMMARY_CHARS]
def render_question_md(question: str, assumptions: list[str] | None = None) -> str:
"""Rende question.md in modo deterministico: domanda + assunzioni opzionali.
@@ -36,16 +51,28 @@ def _new_id(question: str, sessions_root: Path, stamp: str) -> str:
def create_session(
question: str, db: DatabaseConfig, sessions_root: Path
question: str, db: DatabaseConfig, sessions_root: Path,
*, author: str | None = None, summary: str | None = None,
) -> SessionManifest:
now = datetime.now(UTC)
# stamp con ora/min/sec: identifica univocamente sessioni dello stesso giorno
# sulla stessa domanda. Il contatore -n resta come rete per collisioni nello
# stesso secondo.
session_id = _new_id(question, sessions_root, now.strftime("%Y-%m-%d-%H%M%S"))
who = author or current_author()
# schema_version del workflow usato (D6/§5.3): consente di interpretare il ledger
# secondo la versione anche se il workflow evolve.
try:
from tht.workflow import load_workflow
schema_version = load_workflow().schema_version
except Exception:
schema_version = None
manifest = SessionManifest(
id=session_id, created_at=now, question=question,
database=db.database, schema=db.db_schema,
author=who, summary=summary or _summarize(question),
updated_at=now, updated_by=who, schema_version=schema_version,
)
session_dir = sessions_root / session_id
manifest.to_yaml(session_dir / MANIFEST)
@@ -53,6 +80,18 @@ def create_session(
return manifest
def touch_manifest(
session_id: str, sessions_root: Path, *, updated_by: str | None = None
) -> SessionManifest:
"""Aggiorna updated_at/updated_by del manifest (D6 §5.2). Da chiamare ad ogni
mutazione della sessione (set-question, finalize, close)."""
manifest = load_session(session_id, sessions_root)
manifest.updated_at = datetime.now(UTC)
manifest.updated_by = updated_by or current_author()
manifest.to_yaml(sessions_root / session_id / MANIFEST)
return manifest
def set_question(
session_id: str,
question: str,
@@ -67,6 +106,7 @@ def set_question(
load_session(session_id, sessions_root)
path = sessions_root / session_id / "question.md"
path.write_text(render_question_md(question, assumptions))
touch_manifest(session_id, sessions_root)
return path
@@ -80,5 +120,7 @@ def load_session(session_id: str, sessions_root: Path) -> SessionManifest:
def close_session(session_id: str, sessions_root: Path) -> SessionManifest:
manifest = load_session(session_id, sessions_root)
manifest.status = "closed"
manifest.updated_at = datetime.now(UTC)
manifest.updated_by = current_author()
manifest.to_yaml(sessions_root / session_id / MANIFEST)
return manifest
+33 -2
View File
@@ -10,6 +10,7 @@ post-rollback riflette automaticamente lo stato corretto (le decisioni stale di
"""
from __future__ import annotations
import json
from dataclasses import dataclass
from pathlib import Path
@@ -17,6 +18,7 @@ from tht.phase import effective_decisions
from tht.workflow import load_workflow
MAX_BODY_BYTES = 80_000 # ~20k token (target per task document di una fase)
_TRUNCATION_MARKER = "\n\n[...troncato per il budget di contesto D16...]"
@dataclass
@@ -26,6 +28,27 @@ class TaskDoc:
byte_budget_ok: bool
def _slice_schema_linking(raw: str, promoted_tables: list[str] | None) -> str:
"""Riduce schema_linking.json alle sole tabelle promosse (D16 §4.9: 'solo le
tabelle/colonne promosse, non tutto lo schema'). Senza promoted_tables passa il
contenuto invariato (e' gia' la superficie decisionale, non lo schema fisico).
Su JSON malformato ritorna il raw (il bound piu' sotto lo tronca se enorme)."""
if not promoted_tables:
return raw
try:
data = json.loads(raw)
except json.JSONDecodeError:
return raw
allow = set(promoted_tables)
def table_of(cand: dict) -> str:
return str(cand.get("name", "")).split(".")[0]
if isinstance(data, dict) and isinstance(data.get("candidates"), list):
data["candidates"] = [c for c in data["candidates"] if table_of(c) in allow]
return json.dumps(data, ensure_ascii=False, indent=2)
def generate_task_doc(
session_dir: Path | str,
phase: int,
@@ -48,7 +71,8 @@ def generate_task_doc(
sl = session_dir / "schema_linking.json"
if sl.exists() and phase >= 4:
parts.append("## Schema linking (deciso)\n```json\n" + sl.read_text() + "\n```")
sliced = _slice_schema_linking(sl.read_text(), promoted_tables)
parts.append("## Schema linking (deciso)\n```json\n" + sliced + "\n```")
# Brief decisioni effective (D15-aware)
eff = effective_decisions(session_dir)
@@ -66,4 +90,11 @@ def generate_task_doc(
parts.append(header)
body = "\n\n".join(parts)
return TaskDoc(phase=phase, body=body, byte_budget_ok=len(body.encode()) <= MAX_BODY_BYTES)
# Enforcement del bound (D16): non solo segnalare -- troncare. Un task doc oltre
# budget collasserebbe un 35B; meglio un documento troncato e marcato.
encoded = body.encode()
budget = MAX_BODY_BYTES - len(_TRUNCATION_MARKER.encode())
if len(encoded) > MAX_BODY_BYTES:
body = encoded[:budget].decode("utf-8", "ignore") + _TRUNCATION_MARKER
return TaskDoc(phase=phase, body=body, byte_budget_ok=False)
return TaskDoc(phase=phase, body=body, byte_budget_ok=True)
+7 -1
View File
@@ -46,6 +46,11 @@ class RestSearcher:
for table in tables_for_kinds(kinds):
for row in self.client.search_similar(table, query_vec, top_n):
hits.append(hit_from_metadata(row.get("similarity", 0.0), row.get("metadata")))
# schema_records contiene sia schema_table sia schema_column: la RPC non filtra
# per kind, quindi lo facciamo lato client per parita' col path diretto (#25).
if kinds:
allowed = set(kinds)
hits = [h for h in hits if h.kind in allowed]
return _merge(hits, top_n)
@@ -63,5 +68,6 @@ class DirectSearcher:
hits: list[VectorHit] = []
for table in tables_for_kinds(kinds):
store = VectorStore(self.engine, schema=self.schema, table=table, dim=self.dim)
hits.extend(store.search(query_vec, top_n=top_n))
# passa kinds: dentro schema_records filtra schema_table vs schema_column (#25).
hits.extend(store.search(query_vec, top_n=top_n, kinds=kinds))
return _merge(hits, top_n)
+7 -1
View File
@@ -101,4 +101,10 @@ class VectorRestClient:
return len(rows)
if isinstance(payload, dict):
return int(payload.get("upserted", len(rows)))
return len(payload) if isinstance(payload, list) else len(rows)
# PostgREST puo' incapsulare uno scalar jsonb in una lista [{"upserted": N}]:
# estrai il conteggio dal primo elemento invece di restituire len(lista)=1.
if isinstance(payload, list):
if payload and isinstance(payload[0], dict) and "upserted" in payload[0]:
return int(payload[0]["upserted"])
return len(payload)
return len(rows)
-6
View File
@@ -10,12 +10,6 @@ from tht.vectorstore.rest_client import VectorRestClient
from tht.vectorstore.store import SyncStats, content_hash
KIND_TO_TABLE = {
"schema_table": "schema_records",
"schema_column": "schema_records",
"evidence": "evidence",
"memory": "memory",
}
TABLE_TO_KINDS = {
"schema_records": {"schema_table", "schema_column"},
"evidence": {"evidence"},
+11 -2
View File
@@ -25,6 +25,7 @@ class PhaseSpec:
advance: str
prerequisites: list[Any]
artifacts_out: list[str] = field(default_factory=list)
emits: list[str] = field(default_factory=list)
@dataclass
@@ -46,8 +47,9 @@ class Workflow:
return "?"
def decision_min_phase(self, decision_type: str) -> int:
"""A decision type's min phase = the earliest phase whose prerequisites
reference it (via decision_exists / decision_subject_exists). Defaults to 1."""
"""A decision type's min phase = the earliest phase that EMITS it (workflow.yaml
`emits`), or that references it in a prerequisite. Defaults to 1 (meta types like
phase_approved/decision_retracted are not phase-gated)."""
return self._decision_min_map.get(decision_type, 1)
def schema_linking_phase(self) -> int:
@@ -91,6 +93,12 @@ def _collect_decision_mins(phases: list[PhaseSpec]) -> dict[str, int]:
for p in phases:
scan(p.prerequisites, p.num)
# `emits`: la fase dichiara i decision type che produce -> fonte primaria del
# min phase (copre i tipi non citati nei prerequisites, es. table_promoted,
# value_grounded, concept_formula_*). Vince la fase piu' bassa.
for dtype in p.emits:
if isinstance(dtype, str) and (dtype not in mins or p.num < mins[dtype]):
mins[dtype] = p.num
return mins
@@ -107,6 +115,7 @@ def load_workflow(path: Path | str = _WF_PATH) -> Workflow:
advance=p["advance"],
prerequisites=p.get("prerequisites", []),
artifacts_out=p.get("artifacts_out", []),
emits=p.get("emits", []),
)
)
return Workflow(
+15
View File
@@ -4,34 +4,46 @@
schema_version: 1
# `emits` elenca i decision type sostanziali prodotti da ciascuna fase: e' la fonte
# data-driven di decision_min_phase (la fase minima in cui un tipo e' ammesso),
# usata da require_phase_or_exit per rifiutare decisioni fuori fase. I tipi meta
# cross-fase (phase_approved/auto_approved/reopened/skipped, decision_retracted)
# non vanno elencati: restano ammessi da qualunque fase.
phases:
- id: F1
name: chiarimento
advance: kind:phase
prerequisites: []
artifacts_out: []
emits: [concept_clarified, ambiguity_open]
- id: F2
name: memoria
advance: auto_if_empty
prerequisites: []
artifacts_out: []
emits: [memory_rejected]
- id: F3
name: riscrittura
advance: kind:phase
prerequisites:
- decision_exists: question_rewritten
artifacts_out: [question.md]
emits: [question_rewritten]
- id: F4
name: schema_linking
advance: reviewer_decide
prerequisites: []
artifacts_out: [schema_linking.json]
emits: [table_promoted, table_excluded, column_corrected, join_modified,
evidence_accepted, evidence_rejected, value_grounded,
concept_formula_approved, concept_formula_rejected]
- id: F5
name: sintesi
advance: kind:phase
prerequisites:
- file_validates: [schema_linking.json, SchemaLinking]
artifacts_out: []
emits: []
- id: F6
name: cte
advance: auto_if_empty_or_skipped
@@ -40,12 +52,14 @@ phases:
- decision_subject_exists: [phase_skipped, "phase:6"]
- all_ctes_approved: true
artifacts_out: [cte_plan.json, ctes/, cte_tests.json]
emits: [cte_approved, cte_corrected, cte_rejected]
- id: F7
name: sql_finale
advance: kind:phase
prerequisites:
- decision_exists: sql_approved
artifacts_out: [sql_final.sql]
emits: [sql_revised, sql_approved, sql_rejected]
- id: F8
name: datamart
advance: reviewer_decide
@@ -54,6 +68,7 @@ phases:
- decision_exists: datamart_requested
- decision_exists: datamart_declined
artifacts_out: []
emits: [datamart_requested, datamart_declined]
decision_min_phase: auto
max_phase: auto