diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index 8fb87e6b..f6d54838 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -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)); } diff --git a/harness/.pi/skills/tht-sessione/SKILL.md b/harness/.pi/skills/tht-sessione/SKILL.md index 4f1502cd..97abc8cb 100644 --- a/harness/.pi/skills/tht-sessione/SKILL.md +++ b/harness/.pi/skills/tht-sessione/SKILL.md @@ -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 ""` (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 ""` (evidence + schema, LSH - over real values) and `tht search --kind evidence ""`. The LSH exposes +1. Explore the DWH and knowledge base: `tht search find ""` (evidence + schema, LSH + over real values) and `tht search find --kind evidence ""`. 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 + ""` (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). diff --git a/harness/.pi/skills/tht-sessione/sql-generation.md b/harness/.pi/skills/tht-sessione/sql-generation.md index 76d1156e..ed8a33a0 100644 --- a/harness/.pi/skills/tht-sessione/sql-generation.md +++ b/harness/.pi/skills/tht-sessione/sql-generation.md @@ -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 ""` (match on real values). + Verify the filter values with `tht search find ""` (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? diff --git a/harness/README.md b/harness/README.md index e2887a92..299b8e60 100644 --- a/harness/README.md +++ b/harness/README.md @@ -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 diff --git a/harness/tests/integration/__init__.py b/harness/tests/integration/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/harness/tests/integration/test_gate_cli_signatures.py b/harness/tests/integration/test_gate_cli_signatures.py new file mode 100644 index 00000000..beca9d63 --- /dev/null +++ b/harness/tests/integration/test_gate_cli_signatures.py @@ -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 --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}" + ) diff --git a/harness/tests/test_blocco6_robustness.py b/harness/tests/test_blocco6_robustness.py new file mode 100644 index 00000000..d6db9f19 --- /dev/null +++ b/harness/tests/test_blocco6_robustness.py @@ -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 diff --git a/harness/tests/test_decision_min_phase.py b/harness/tests/test_decision_min_phase.py new file mode 100644 index 00000000..5b3fe112 --- /dev/null +++ b/harness/tests/test_decision_min_phase.py @@ -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 diff --git a/harness/tests/test_decision_retract_cli.py b/harness/tests/test_decision_retract_cli.py new file mode 100644 index 00000000..4991d2aa --- /dev/null +++ b/harness/tests/test_decision_retract_cli.py @@ -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 diff --git a/harness/tests/test_effective_name_subject_stale.py b/harness/tests/test_effective_name_subject_stale.py new file mode 100644 index 00000000..a0982257 --- /dev/null +++ b/harness/tests/test_effective_name_subject_stale.py @@ -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 diff --git a/harness/tests/test_formula_wiring.py b/harness/tests/test_formula_wiring.py new file mode 100644 index 00000000..d670315a --- /dev/null +++ b/harness/tests/test_formula_wiring.py @@ -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) diff --git a/harness/tests/test_manifest_fields.py b/harness/tests/test_manifest_fields.py new file mode 100644 index 00000000..5493bb95 --- /dev/null +++ b/harness/tests/test_manifest_fields.py @@ -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" diff --git a/harness/tests/test_readonly_guard.py b/harness/tests/test_readonly_guard.py new file mode 100644 index 00000000..01d4d138 --- /dev/null +++ b/harness/tests/test_readonly_guard.py @@ -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 diff --git a/harness/tht/cli/__init__.py b/harness/tht/cli/__init__.py index edea3b7c..397de75f 100644 --- a/harness/tht/cli/__init__.py +++ b/harness/tht/cli/__init__.py @@ -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") diff --git a/harness/tht/cli/decision_cmd.py b/harness/tht/cli/decision_cmd.py index fa9cd33b..de7ea37d 100644 --- a/harness/tht/cli/decision_cmd.py +++ b/harness/tht/cli/decision_cmd.py @@ -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 (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"), diff --git a/harness/tht/cli/formula_cmd.py b/harness/tht/cli/formula_cmd.py new file mode 100644 index 00000000..b7253537 --- /dev/null +++ b/harness/tht/cli/formula_cmd.py @@ -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 ""`. 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 oppure --sql \"\".", + 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) diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py index 17de2345..1faf0ea1 100644 --- a/harness/tht/cli/memory_cmd.py +++ b/harness/tht/cli/memory_cmd.py @@ -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."), diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py index 12e2fc3e..e0ef8755 100644 --- a/harness/tht/cli/search_cmd.py +++ b/harness/tht/cli/search_cmd.py @@ -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( diff --git a/harness/tht/cli/session_cmd.py b/harness/tht/cli/session_cmd.py index a3e2d7ab..1ed0ff47 100644 --- a/harness/tht/cli/session_cmd.py +++ b/harness/tht/cli/session_cmd.py @@ -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: diff --git a/harness/tht/db/connection.py b/harness/tht/db/connection.py index 04abedcd..2f7a9b71 100644 --- a/harness/tht/db/connection.py +++ b/harness/tht/db/connection.py @@ -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}" diff --git a/harness/tht/decisions.py b/harness/tht/decisions.py index 3be675ae..0b2cb090 100644 --- a/harness/tht/decisions.py +++ b/harness/tht/decisions.py @@ -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) diff --git a/harness/tht/evidence/formula_store.py b/harness/tht/evidence/formula_store.py index 47c8117e..06a93d58 100644 --- a/harness/tht/evidence/formula_store.py +++ b/harness/tht/evidence/formula_store.py @@ -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 /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 /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()] diff --git a/harness/tht/evidence/model.py b/harness/tht/evidence/model.py index c21180eb..6d0b4300 100644 --- a/harness/tht/evidence/model.py +++ b/harness/tht/evidence/model.py @@ -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 diff --git a/harness/tht/execute/__init__.py b/harness/tht/execute/__init__.py index f3240e58..6d3dd14c 100644 --- a/harness/tht/execute/__init__.py +++ b/harness/tht/execute/__init__.py @@ -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: diff --git a/harness/tht/lshindex/__init__.py b/harness/tht/lshindex/__init__.py index 6f3bbc03..be12b8bd 100644 --- a/harness/tht/lshindex/__init__.py +++ b/harness/tht/lshindex/__init__.py @@ -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) diff --git a/harness/tht/memory.py b/harness/tht/memory.py index cc213636..23b24f4d 100644 --- a/harness/tht/memory.py +++ b/harness/tht/memory.py @@ -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, } diff --git a/harness/tht/phase.py b/harness/tht/phase.py index c95c84bd..dbcc8d7a 100644 --- a/harness/tht/phase.py +++ b/harness/tht/phase.py @@ -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 diff --git a/harness/tht/report.py b/harness/tht/report.py index 09c4408f..7a23ce4f 100644 --- a/harness/tht/report.py +++ b/harness/tht/report.py @@ -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() diff --git a/harness/tht/rest/execute.py b/harness/tht/rest/execute.py index 07626ec5..223f460f 100644 --- a/harness/tht/rest/execute.py +++ b/harness/tht/rest/execute.py @@ -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: diff --git a/harness/tht/session/store.py b/harness/tht/session/store.py index 3019abf2..ad8b00bd 100644 --- a/harness/tht/session/store.py +++ b/harness/tht/session/store.py @@ -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 diff --git a/harness/tht/taskdoc.py b/harness/tht/taskdoc.py index ad615c20..d5a6f32b 100644 --- a/harness/tht/taskdoc.py +++ b/harness/tht/taskdoc.py @@ -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) diff --git a/harness/tht/vectorstore/reader.py b/harness/tht/vectorstore/reader.py index 083521d1..972ec85e 100644 --- a/harness/tht/vectorstore/reader.py +++ b/harness/tht/vectorstore/reader.py @@ -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) diff --git a/harness/tht/vectorstore/rest_client.py b/harness/tht/vectorstore/rest_client.py index 4e041910..d5ff5576 100644 --- a/harness/tht/vectorstore/rest_client.py +++ b/harness/tht/vectorstore/rest_client.py @@ -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) diff --git a/harness/tht/vectorstore/rest_writer.py b/harness/tht/vectorstore/rest_writer.py index d83d84bf..3d2ac028 100644 --- a/harness/tht/vectorstore/rest_writer.py +++ b/harness/tht/vectorstore/rest_writer.py @@ -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"}, diff --git a/harness/tht/workflow.py b/harness/tht/workflow.py index 975c2724..882f8024 100644 --- a/harness/tht/workflow.py +++ b/harness/tht/workflow.py @@ -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( diff --git a/harness/workflow.yaml b/harness/workflow.yaml index 0efeed04..5fb9cd12 100644 --- a/harness/workflow.yaml +++ b/harness/workflow.yaml @@ -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