diff --git a/harness/.pi/extensions/gate/__tests__/enrich.test.js b/harness/.pi/extensions/gate/__tests__/enrich.test.js index b5394f61..5002b76a 100644 --- a/harness/.pi/extensions/gate/__tests__/enrich.test.js +++ b/harness/.pi/extensions/gate/__tests__/enrich.test.js @@ -199,3 +199,32 @@ test("enrichPhaseSummaryV2 preserves open_questions/checks untouched", () => { assert.deepEqual(out.checks, [{ label: "x", status: "ok" }]); assert.deepEqual(out.open_questions, ["domanda aperta"]); }); + +// --- appendLedgerSection ----------------------------------------------------------- + +const { appendLedgerSection } = require("../enrich.js"); + +test("appendLedgerSection appends only the phase's substantive decisions", () => { + const data = { schema_version: 2, summary: "s", sections: [{ title: "Criteri", items: [] }] }; + const decisions = [ + { seq: 1, type: "concept_clarified", subject: "ablazione", detail: "solo transcatetere", rationale: "r1" }, + { seq: 2, type: "phase_approved", subject: "phase:1", detail: "", rationale: "" }, // meta: non in emits + { seq: 3, type: "table_promoted", subject: "dim_patient", detail: "", rationale: "" }, // altra fase + ]; + const out = appendLedgerSection(data, decisions, ["concept_clarified", "ambiguity_open"]); + assert.equal(out.sections.length, 2); + const ledger = out.sections[1]; + assert.equal(ledger.title, "Decisioni registrate in questa fase (dal ledger)"); + assert.equal(ledger.items.length, 1); + assert.equal(ledger.items[0].label, "concept_clarified: ablazione"); + assert.equal(ledger.items[0].value, "solo transcatetere"); + // input non mutato + assert.equal(data.sections.length, 1); +}); + +test("appendLedgerSection is a no-op without matching decisions or with empty ledger", () => { + const data = { schema_version: 2, summary: "s" }; + assert.equal(appendLedgerSection(data, [], ["concept_clarified"]), data); + assert.equal(appendLedgerSection(data, [{ type: "sql_approved", subject: "x" }], ["concept_clarified"]), data); + assert.equal(appendLedgerSection(data, undefined, undefined), data); +}); diff --git a/harness/.pi/extensions/gate/enrich.js b/harness/.pi/extensions/gate/enrich.js index f90e47bf..ca73035b 100644 --- a/harness/.pi/extensions/gate/enrich.js +++ b/harness/.pi/extensions/gate/enrich.js @@ -151,10 +151,32 @@ function enrichCteResultColumns(payload, planTables, getColumns) { return { ...payload, columns }; } +// Contract C bis. Appends a deterministic "decisions of this phase" section built +// from the session ledger, so the model's recap can stay thin (Discipline 6 exactness +// comes from persisted state, not model prose). `decisions` is session-show's ledger +// dump; `emits` is the phase's substantive decision-type list from workflow.yaml. +// No matching decisions -> data returned unchanged. Pure: returns a NEW object. +function appendLedgerSection(data, decisions, emits) { + const emitSet = new Set(emits || []); + const rows = (decisions || []).filter((d) => d && emitSet.has(d.type)); + if (rows.length === 0) return data; + const section = { + title: "Decisioni registrate in questa fase (dal ledger)", + items: rows.map((d) => ({ + label: `${d.type}: ${d.subject}`, + value: d.detail || "", + kind: "decision", + rationale: d.rationale || "", + })), + }; + return { ...data, sections: [...(data.sections || []), section] }; +} + module.exports = { memoizeGetColumns, enrichCtePlanV2, buildCteResultV2, enrichCteResultColumns, enrichPhaseSummaryV2, + appendLedgerSection, }; diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index f3726377..9d6798c0 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -42,6 +42,7 @@ import { buildCteResultV2, enrichCteResultColumns, enrichPhaseSummaryV2, + appendLedgerSection, } from "./gate/enrich.js"; import { isReserved } from "./reserved-labels.mjs"; @@ -898,11 +899,26 @@ export default function (pi) { "Riepilogo di fase v2 non valido:\n- " + v.errors.join("\n- "), ); } - const enriched = enrichPhaseSummaryV2( + let enriched = enrichPhaseSummaryV2( data, phaseMetaForNum(ctx, curNum), makeGetColumns(ctx), ); + // Recap deterministico: appende la sezione "decisioni di questa fase" + // dal ledger (best-effort: un errore di lettura non blocca il gate). + try { + const show = JSON.parse( + tht(ctx, ["session", "show", session, "--json"]), + ); + const meta = phaseMeta(ctx).phases.find((p) => p.num === curNum); + enriched = appendLedgerSection( + enriched, + show.decisions, + meta ? meta.emits : [], + ); + } catch { + // ledger non leggibile: il riepilogo resta quello del modello + } artifact = { ...artifact, data: enriched }; } diff --git a/harness/.pi/skills/tht-sessione/SKILL.md b/harness/.pi/skills/tht-sessione/SKILL.md index 57755389..8d2690d2 100644 --- a/harness/.pi/skills/tht-sessione/SKILL.md +++ b/harness/.pi/skills/tht-sessione/SKILL.md @@ -81,7 +81,10 @@ substantive decisions. For a phase-closing gate (`reviewer_confirm kind:"phase"`) prefer the **structured v2 recap** `artifact:{kind:"phase", data:{schema_version:2, …}}`: you author `summary` (1-3 sentence markdown), `checks[]`, `sections[]` and `tables[]`; the gate - fills `phase` (from workflow meta) and every `description` from the catalog. Every + fills `phase` (from workflow meta) and every `description` from the catalog, and + APPENDS a deterministic section "Decisioni registrate in questa fase (dal ledger)" + — do NOT re-enumerate the phase's recorded decisions yourself: author only the + summary, the checks and the context the ledger cannot express. Every `sections[].items[]` MUST cite the concrete **table**, **column** and the **value** that motivates the choice (booleans, time windows, thresholds) — not just prose. Compact example: @@ -124,6 +127,13 @@ substantive decisions. 11. **Rollback (D15).** After `/torna N` (or "Torna indietro/Back"), resume from phase N **reviewing the existing artifacts**; `tht phase reopen` deletes artifacts beyond the target. Do NOT re-run `tht` commands for artifacts that are still valid. +12. **This skill is the complete contract.** Every command, flag and behavior you + need is named in this skill and its reference docs (`rewriting.md`, `cte.md`, + `sql-generation.md`). Do NOT run `--help`, do NOT read the harness source + (`tht/`, `.pi/extensions/`, tests) to figure out how a command works, and do NOT + explore the filesystem with `find`/`grep`/`cat` for that purpose. If something + genuinely seems missing or a command behaves unexpectedly, say so to the reviewer + instead of reverse-engineering the tooling. ## Phase 0 — Resume (cold start) @@ -148,21 +158,25 @@ it is complete. (The backend already refuses resume for finalized/archived sessi Prerequisite: you must already be in Phase 1. -**F1 toolbox.** The only commands you need here are `tht search find` and `tht schema -render` — both fast, read-only lookups over workspace artifacts already on disk. -Evidence lives in `/evidence/**` and is what `tht search find --kind evidence` -returns — do not browse it with `find`/`cat`. Do NOT run `tht schema introspect`: it is -a maintenance command that re-reads the remote DWH (~3 minutes); the catalog -`artifacts/mschema/physical.yaml` is already in the workspace. Do NOT explore with -`--help` or ad-hoc shell commands — every command you need is named in this skill. +**F1 toolbox.** The only commands you need here are `tht search pack`, `tht search +find` and `tht schema render` — all fast, read-only lookups over workspace artifacts +already on disk. Evidence lives in `/evidence/**` and is what `tht search +find --kind evidence` returns — do not browse it with `find`/`cat`. Do NOT run `tht +schema introspect`: it is a maintenance command that re-reads the remote DWH (~3 +minutes); the catalog `artifacts/mschema/physical.yaml` is already in the workspace. +Do NOT explore with `--help` or ad-hoc shell commands — every command you need is +named in this skill. -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. - First list ALL the ambiguous terms in the question, then run the `tht search find` - calls for every term in ONE batch (a single message with multiple shell invocations) - — not one lookup per turn. +1. **First call, one shot:** `tht search pack "" --session ` — + it bundles candidate tables, relevant evidence and similar solved questions for the + WHOLE question in a single command (one embedding, three searches) and persists + `retrieval_pack.md` in the session. Read it before anything else; it usually + answers "which tables/evidence matter here" without further exploration. + Then ground the individual ambiguous terms: list ALL of them and run the + `tht search find ""` / `tht search find --kind evidence ""` calls in + ONE batch (a single message with multiple shell invocations) — not one lookup per + turn. 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 the candidate interpretations (`recommended:true` on the best) + "Altro". Pick the widget by the question's shape: @@ -244,6 +258,10 @@ Prerequisite: Phase 3 closed. 1. `tht schema render --format mschema-text` for the schema context (the catalog `artifacts/mschema/physical.yaml` is already in the workspace; only if render fails with `physical.yaml non trovato`, run `tht schema introspect` once, then render). + To inspect specific tables use `--table ` (repeatable: `-t t1 -t t2`) — + do NOT dump the full catalog or slice it with `awk`/`grep`. The session's + `retrieval_pack.md` (built in F1) already lists the candidate tables for the + question — start from those. Copy table/column names EXACTLY from it — never invent objects. Also run `tht memory solved-search "" --json`: similar already-solved questions show which tables comparable questions used. Cite relevant precedents @@ -262,7 +280,11 @@ Prerequisite: Phase 3 closed. to reference other columns as join keys or filter predicates when the query requires them. Propose joins separately in `reviewer_decide(advance:false)`, registering - `join_modified`. + `join_modified`. Ground them in the `【Foreign keys】` section of the mschema-text + render: it lists the curated logical FKs of the workspace (e.g. + `fact_x.cod_paz=dim_patient.cod_paz`, `*_time_key=dim_time.day_key`) — prefer + those to joins you derive yourself, and flag to the reviewer any join you need + that is NOT in the list. 3. **Value grounding (D14a).** If a cited value (e.g. "ablazione") matches MULTIPLE columns (a boolean flag + a free-text patologia field), present a `reviewer_decide` with a `value_grounded` option for each candidate column (the LSH exposes all of diff --git a/harness/.pi/skills/tht-sessione/cte.md b/harness/.pi/skills/tht-sessione/cte.md index 2fb54314..1ce695b2 100644 --- a/harness/.pi/skills/tht-sessione/cte.md +++ b/harness/.pi/skills/tht-sessione/cte.md @@ -13,8 +13,8 @@ Rules (from the AV-SQL discipline, hold verbatim): 3. Better one column too many than one too few: if unsure, include it. 4. One CTE = one informative subset with a clear purpose (e.g. "ricoveri with ablazione in 2025"), named in a speaking snake_case. -5. CTEs can chain-reference each other; the last one in the file is the one that - `tht cte test` will query. +5. CTEs can chain-reference each other **within the same file**; the last one in the + file is the one that `tht cte test` will query (see the execution contract below). 6. Each file in `sessions//ctes/.sql` contains ONLY the `WITH ... AS (...)` block (multi-CTE allowed), WITHOUT a trailing SELECT. A `SELECT ...` line after the WITH block causes an error in `tht cte test`: never add it. CTEs are tested @@ -24,6 +24,24 @@ Rules (from the AV-SQL discipline, hold verbatim): 7. Filters: use field values verified with `tht search` (LSH match on real values), not imagined values. +## How `tht cte test` executes (complete contract — do not read the harness source) + +- Each CTE file is **standalone**: the test reads ONLY `ctes/.sql`, appends + `SELECT * FROM ` and runs it read-only against the + DWH with an injected LIMIT and statement timeout. Cross-file references are NOT + resolved: to build on an earlier CTE, repeat its definition in the same `WITH` + chain (that is why multi-CTE files are allowed). A test normally completes in + well under a second. +- Before execution the SQL is validated statically: parsable, a single statement, + read-only by structure, no blacklisted functions, and every referenced table must + exist in the catalog (names defined in the `WITH` chain are exempt). Tables outside + the promoted perimeter produce warnings. Unqualified columns are not statically + checked — the DWH will catch them at run time. +- Test order is enforced by the CLI from the approved plan (exit 5 names the CTE + whose turn it is). Every outcome (ok or error) is appended to `cte_tests.json`; + the gate reads the persisted file + last test via `tht cte info`, so never paste + SQL, columns or preview rows into the gate text. + Presentation to the reviewer, for each CTE in the plan: > **** — purpose: diff --git a/harness/.pi/skills/tht-sessione/sql-generation.md b/harness/.pi/skills/tht-sessione/sql-generation.md index ed8a33a0..763025fc 100644 --- a/harness/.pi/skills/tht-sessione/sql-generation.md +++ b/harness/.pi/skills/tht-sessione/sql-generation.md @@ -21,9 +21,9 @@ copy-pasteable: no rationale comments (that lives in the audit artifacts). ## Time dimension (analysis by year/month/quarter) Fact tables have `data_time_key` (`integer`, format `YYYYMMDD`): it is the FK to -`dim_time.day_key`. **This FK is NOT declared** in the DWH (facts have -`foreign_keys: []`), so it will NOT appear in `schema_linking.json`: you must add it -by hand to the join. +`dim_time.day_key`. The DWH does not declare it, but the workspace annotations do: +it appears in the `【Foreign keys】` section of the mschema-text render (every +`*_time_key` column maps to `dim_time.day_key`) — take it from there for the join. - To extract year, month, quarter, semester etc. do `JOIN dim_time dt ON dt.day_key = .data_time_key` and use the dimension's diff --git a/harness/tests/test_cli_phase_meta.py b/harness/tests/test_cli_phase_meta.py index ebcfc5a3..42ce122b 100644 --- a/harness/tests/test_cli_phase_meta.py +++ b/harness/tests/test_cli_phase_meta.py @@ -30,3 +30,11 @@ def test_phase_meta_each_phase_carries_id_and_name(): assert p["num"] == i assert p["id"], f"phase {i} missing id" assert p["name"], f"phase {i} missing name" + + +def test_phase_meta_exposes_emits(): + # Il gate filtra il ledger per fase con `emits` (recap deterministico v2). + result = runner.invoke(app, ["phase", "meta", "--json"]) + data = json.loads(result.stdout) + f1 = data["phases"][0] + assert f1["emits"] == ["concept_clarified", "ambiguity_open"] diff --git a/harness/tests/test_schema_fk_annotations.py b/harness/tests/test_schema_fk_annotations.py new file mode 100644 index 00000000..88c037c8 --- /dev/null +++ b/harness/tests/test_schema_fk_annotations.py @@ -0,0 +1,223 @@ +from datetime import datetime + +import yaml +from typer.testing import CliRunner + +from tht.cli import app +from tht.mschema.merge import find_orphans +from tht.mschema.models import ( + Annotations, + ColumnPhysical, + ForeignKey, + PhysicalSchema, + TableAnnotation, + TablePhysical, +) +from tht.mschema.render import to_mschema_text, to_schema_dict + + +def _physical(): + return PhysicalSchema( + database="d", schema="s", introspected_at=datetime(2026, 1, 1), + tables={ + "dim_patient": TablePhysical( + columns={"cod_paz": ColumnPhysical(type="bigint", pk=True)}, + ), + "dim_time": TablePhysical( + columns={"day_key": ColumnPhysical(type="integer", pk=True)}, + ), + "fact_ablazione": TablePhysical( + columns={ + "cod_paz": ColumnPhysical(type="bigint"), + "data_time_key": ColumnPhysical(type="integer"), + "esito": ColumnPhysical(type="text"), + }, + ), + }, + ) + + +def _annotations_with_fks(): + return Annotations( + tables={ + "fact_ablazione": TableAnnotation( + foreign_keys=[ + ForeignKey(columns=["cod_paz"], ref_table="dim_patient", + ref_columns=["cod_paz"]), + ForeignKey(columns=["data_time_key"], ref_table="dim_time", + ref_columns=["day_key"]), + ], + ) + } + ) + + +def test_mschema_text_renders_annotation_fks(): + text = to_mschema_text(_physical(), _annotations_with_fks()) + assert "fact_ablazione.cod_paz=dim_patient.cod_paz" in text + assert "fact_ablazione.data_time_key=dim_time.day_key" in text + + +def test_schema_dict_merges_annotation_fks(): + d = to_schema_dict(_physical(), _annotations_with_fks()) + fks = d["fact_ablazione"]["foreign_keys"] + assert {"columns": ["cod_paz"], "ref_table": "dim_patient", + "ref_columns": ["cod_paz"]} in fks + + +def test_find_orphans_flags_broken_annotation_fk(): + ann = Annotations( + tables={ + "fact_ablazione": TableAnnotation( + foreign_keys=[ + ForeignKey(columns=["cod_paz"], ref_table="dim_sparita", + ref_columns=["x"]), + ForeignKey(columns=["colonna_sparita"], ref_table="dim_time", + ref_columns=["day_key"]), + ], + ) + } + ) + orphans = find_orphans(_physical(), ann) + assert "fact_ablazione.fk(cod_paz)->dim_sparita" in orphans + assert "fact_ablazione.fk(colonna_sparita)->dim_time" in orphans + + +def test_find_orphans_ok_with_valid_fk(): + assert find_orphans(_physical(), _annotations_with_fks()) == [] + + +def _write_workspace(tmp_path): + _physical().to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml") + cfg = tmp_path / "workspace.yaml" + cfg.write_text( + "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" + f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, sessions: {tmp_path/'s'}}}\n" + ) + return cfg + + +def test_suggest_fks_prints_candidates(tmp_path): + cfg = _write_workspace(tmp_path) + res = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) + assert res.exit_code == 0, res.output + data = yaml.safe_load(res.output.rsplit("\n", 2)[0].split("FK candidate")[0]) + fks = data["tables"]["fact_ablazione"]["foreign_keys"] + assert {"columns": ["cod_paz"], "ref_table": "dim_patient", + "ref_columns": ["cod_paz"]} in fks + assert {"columns": ["data_time_key"], "ref_table": "dim_time", + "ref_columns": ["day_key"]} in fks + + +def test_mine_join_pairs_from_approved_sql(): + from tht.mschema.fkmine import mine_join_pairs + + sql = """ + WITH abl AS ( + SELECT sea.cod_paz, dt.year + FROM datawarehouse.fact_ablazione AS sea + JOIN datawarehouse.dim_time AS dt ON sea.data_time_key = dt.day_key + ) + SELECT * FROM abl JOIN abl b ON abl.year = b.year; + """ + pairs = mine_join_pairs(sql, _physical()) + assert pairs[("fact_ablazione", "data_time_key", "dim_time", "day_key")] == 1 + # il join CTE-CTE (abl.year=b.year) non produce coppie + assert len(pairs) == 1 + + +def test_mine_join_pairs_ignores_non_pk_pairs_and_bad_sql(): + from tht.mschema.fkmine import mine_join_pairs + + # esito=esito: nessun lato e' PK -> scartato + sql = ("SELECT * FROM fact_ablazione a JOIN fact_ablazione b " + "ON a.esito = b.esito") + assert len(mine_join_pairs(sql, _physical())) == 0 + assert len(mine_join_pairs("WITH broken (", _physical())) == 0 + + +def test_suggest_fks_skips_generic_and_ambiguous_pks(tmp_path): + phys = PhysicalSchema( + database="d", schema="s", introspected_at=datetime(2026, 1, 1), + tables={ + "dim_a": TablePhysical(columns={"id": ColumnPhysical(type="int", pk=True)}), + "dim_b": TablePhysical(columns={"id": ColumnPhysical(type="int", pk=True)}), + "dim_c1": TablePhysical(columns={"cod_x": ColumnPhysical(type="int", pk=True)}), + "dim_c2": TablePhysical(columns={"cod_x": ColumnPhysical(type="int", pk=True)}), + "fact_f": TablePhysical( + columns={ + "id": ColumnPhysical(type="int"), + "cod_x": ColumnPhysical(type="int"), + }, + ), + }, + ) + phys.to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml") + cfg = tmp_path / "workspace.yaml" + cfg.write_text( + "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" + f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, sessions: {tmp_path/'s'}}}\n" + ) + res = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) + assert res.exit_code == 0, res.output + assert "nessuna FK da suggerire" in res.output # id generico, cod_x ambigua + assert "cod_x" in res.output # segnalata come ambigua saltata + + # --assume disambigua la PK multi-proprietario + res2 = CliRunner().invoke( + app, ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=dim_c1"] + ) + assert res2.exit_code == 0, res2.output + yaml_text = "\n".join( + line for line in res2.output.splitlines() if "FK candidate" not in line + ) + data = yaml.safe_load(yaml_text) + fact_fks = data["tables"]["fact_f"]["foreign_keys"] + assert {"columns": ["cod_x"], "ref_table": "dim_c1", + "ref_columns": ["cod_x"]} in fact_fks + # dim_c2.cod_x -> dim_c1 (estensione 1:1), ma NON dim_c1 -> se stessa + assert "dim_c1" not in data["tables"] or all( + fk["ref_table"] != "dim_c1" for fk in data["tables"].get("dim_c1", {}).get("foreign_keys", []) + ) + + # --assume con tabella inesistente -> errore chiaro + res3 = CliRunner().invoke( + app, ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=nope"] + ) + assert res3.exit_code == 1 + assert "non valido" in res3.output + + +def test_suggest_fks_from_sql_mines_joins(tmp_path): + cfg = _write_workspace(tmp_path) + sqldir = tmp_path / "approved" + sqldir.mkdir() + (sqldir / "q1.sql").write_text( + "SELECT f.esito FROM datawarehouse.fact_ablazione f " + "JOIN datawarehouse.dim_patient p ON f.cod_paz = p.cod_paz" + ) + res = CliRunner().invoke( + app, ["schema", "suggest-fks", "-c", str(cfg), "--from-sql", str(sqldir)] + ) + assert res.exit_code == 0, res.output + assert "Minati 1 equi-join da 1 file SQL" in res.output + assert "ref_table: dim_patient" in res.output + + +def test_suggest_fks_write_merges_and_is_idempotent(tmp_path): + cfg = _write_workspace(tmp_path) + ann_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml" + Annotations( + tables={"fact_ablazione": TableAnnotation(description="Ablazioni")} + ).to_yaml(ann_path) + + res = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"]) + assert res.exit_code == 0, res.output + ann = Annotations.from_yaml(ann_path) + assert ann.tables["fact_ablazione"].description == "Ablazioni" # non distrutta + assert len(ann.tables["fact_ablazione"].foreign_keys) == 2 + + res2 = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"]) + assert "nessuna FK da suggerire" in res2.output + ann2 = Annotations.from_yaml(ann_path) + assert len(ann2.tables["fact_ablazione"].foreign_keys) == 2 diff --git a/harness/tests/test_search_pack.py b/harness/tests/test_search_pack.py new file mode 100644 index 00000000..5f246219 --- /dev/null +++ b/harness/tests/test_search_pack.py @@ -0,0 +1,115 @@ +import json +from datetime import datetime +from types import SimpleNamespace + +from typer.testing import CliRunner + +from tht.cli import app +from tht.mschema.models import ColumnPhysical, PhysicalSchema, TablePhysical +from tht.vectorstore.embeddings import EmbeddingsError + + +class _FakeEmbedder: + def __init__(self): + self.calls = 0 + + def embed_query(self, text): + self.calls += 1 + return [0.1, 0.2, 0.3] + + +class _FakeSearcher: + def search(self, vec, top_n, kinds=None): + if kinds == ["solved_question"]: + return [SimpleNamespace( + kind="memory", ref="s-1", id="m1", title="q solved", + similarity=0.91, content="quanti pazienti nel 2024?", + metadata={"session_id": "2026-01-01-000000-x", "sql": "SELECT 1", + "tables": ["fact_ablazione"], "question": "quanti pazienti nel 2024?"}, + )] + if kinds == ["schema_table", "schema_column"]: + return [SimpleNamespace( + kind="schema_table", ref="fact_ablazione", id="t1", + title="Tabella fact_ablazione", similarity=0.88, + content="Tabella fact_ablazione", metadata={}, + )] + if kinds == ["evidence"]: + return [SimpleNamespace( + kind="evidence", ref="ev1", id="ev1", title="Dominio ablazione", + similarity=0.8, content="L'ablazione e' una procedura...", + metadata={"status": "approved"}, + )] + return [] + + +def _workspace(tmp_path, with_session=None): + PhysicalSchema( + database="d", schema="s", introspected_at=datetime(2026, 1, 1), + tables={"fact_ablazione": TablePhysical( + comment="Ablazioni", columns={"cod_paz": ColumnPhysical(type="bigint")})}, + ).to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml") + cfg = tmp_path / "workspace.yaml" + cfg.write_text( + "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" + "vector_db: {database: v, schema: public, user: u, password: p}\n" + "embeddings: {base_url: 'http://localhost:11434', model: nomic-embed-text, dim: 8}\n" + f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, " + f"sessions: {tmp_path/'sessions'}}}\n" + ) + if with_session: + sdir = tmp_path / "sessions" / with_session + sdir.mkdir(parents=True) + (sdir / "session_manifest.yaml").write_text( + f"id: {with_session}\nquestion: q\ndatabase: d\nschema: s\n" + "created_at: 2026-01-01T00:00:00+00:00\nstatus: open\n" + ) + return cfg + + +def _patch(monkeypatch, embedder, searcher): + import tht.cli.vector_cmd as vc + + monkeypatch.setattr(vc, "make_embedder", lambda _cfg: embedder) + monkeypatch.setattr(vc, "open_searcher", lambda _cfg: searcher) + + +def test_pack_single_embed_and_sections(tmp_path, monkeypatch): + cfg = _workspace(tmp_path) + emb = _FakeEmbedder() + _patch(monkeypatch, emb, _FakeSearcher()) + res = CliRunner().invoke(app, ["search", "pack", "quanti pazienti", "-c", str(cfg)]) + assert res.exit_code == 0, res.output + assert emb.calls == 1 # UN solo embedding per le tre ricerche + assert "fact_ablazione" in res.output and "Ablazioni" in res.output + assert "Dominio ablazione" in res.output + assert "SELECT 1" in res.output + + +def test_pack_json_and_session_file(tmp_path, monkeypatch): + sid = "2026-01-01-000000-test" + cfg = _workspace(tmp_path, with_session=sid) + _patch(monkeypatch, _FakeEmbedder(), _FakeSearcher()) + res = CliRunner().invoke( + app, ["search", "pack", "q", "-c", str(cfg), "--session", sid, "--json"] + ) + assert res.exit_code == 0, res.output + data = json.loads(res.output) + assert data["tables"][0]["name"] == "fact_ablazione" + pack = tmp_path / "sessions" / sid / "retrieval_pack.md" + assert pack.exists() + assert "Retrieval pack" in pack.read_text() + + +def test_pack_degrades_gracefully(tmp_path, monkeypatch): + cfg = _workspace(tmp_path) + + class _Broken: + def embed_query(self, text): + raise EmbeddingsError("ollama down") + + _patch(monkeypatch, _Broken(), _FakeSearcher()) + res = CliRunner().invoke(app, ["search", "pack", "q", "-c", str(cfg), "--json"]) + assert res.exit_code == 0, res.output + data = json.loads(res.output[res.output.index("{"):]) + assert data["tables"] == [] and data["evidence"] == [] and data["solved"] == [] + assert any("retrieval non disponibile" in w for w in data["warnings"]) diff --git a/harness/tests/test_session_show_decisions.py b/harness/tests/test_session_show_decisions.py new file mode 100644 index 00000000..1daf4b72 --- /dev/null +++ b/harness/tests/test_session_show_decisions.py @@ -0,0 +1,34 @@ +import json + +from typer.testing import CliRunner + +from tht.cli import app + + +def test_session_show_json_includes_ledger(tmp_path): + sid = "2026-01-01-000000-test" + sdir = tmp_path / "sessions" / sid + sdir.mkdir(parents=True) + (sdir / "session_manifest.yaml").write_text( + f"id: {sid}\nquestion: q\ndatabase: d\nschema: s\n" + "created_at: 2026-01-01T00:00:00+00:00\nstatus: open\n" + ) + (sdir / "review_decisions.jsonl").write_text( + '{"seq": 1, "ts": "2026-01-01T00:01:00+00:00", "type": "concept_clarified", ' + '"subject": "ablazione", "detail": "solo transcatetere", ' + '"rationale": "scelta reviewer", "phase": 1}\n' + ) + cfg = tmp_path / "workspace.yaml" + cfg.write_text( + "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" + f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, " + f"sessions: {tmp_path/'sessions'}}}\n" + ) + res = CliRunner().invoke(app, ["session", "show", sid, "--json", "-c", str(cfg)]) + assert res.exit_code == 0, res.output + data = json.loads(res.output) + assert len(data["decisions"]) == 1 + d = data["decisions"][0] + assert d["type"] == "concept_clarified" + assert d["subject"] == "ablazione" + assert d["detail"] == "solo transcatetere" diff --git a/harness/tht/cli/phase_cmd.py b/harness/tht/cli/phase_cmd.py index 47aef7ba..85b5051a 100644 --- a/harness/tht/cli/phase_cmd.py +++ b/harness/tht/cli/phase_cmd.py @@ -71,6 +71,7 @@ def meta_cmd( "name": p.name, "advance": p.advance, "artifacts_out": p.artifacts_out, + "emits": p.emits, } for p in wf.phases ], diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index e860abf6..e2e671a1 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -141,6 +141,174 @@ def check_cmd(config: Path = CONFIG_OPT) -> None: typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN) +# PK con questi nomi sono identificatori generici: la regola same-name non si applica +# (nel DWH reale `id` e' la PK di ~50 tabelle e produrrebbe migliaia di falsi positivi). +_GENERIC_PK_NAMES = {"id", "key", "code"} + + +@schema_app.command("suggest-fks") +def suggest_fks_cmd( + config: Path = CONFIG_OPT, + from_sql: list[Path] = typer.Option( + None, "--from-sql", + help="Directory di .sql approvati da cui minare i join reali (ripetibile).", + ), + assume: list[str] = typer.Option( + None, "--assume", + help="Disambigua una PK con piu' proprietari: col=tabella_ref " + "(es. cod_paz=dim_patient). Ripetibile.", + ), + write: bool = typer.Option( + False, "--write", + help="Fonde i suggerimenti in annotations.yaml (aggiunge solo FK mancanti).", + ), +) -> None: + """Suggerisce FK logiche per la curazione umana in annotations.yaml. + + Tre regole, in ordine di confidenza: (1) equi-join minati dall'SQL gia' + approvato (--from-sql); (2) colonna `*time_key` verso la PK di dim_time; + (3) colonna con lo stesso nome della PK di UN'ALTRA tabella, solo se quel + nome ha un unico proprietario e non e' generico (id/key/code) — salvo + disambiguazione esplicita con --assume. + """ + import yaml as _yaml + + from tht.mschema.fkmine import mine_join_pairs + from tht.mschema.models import Annotations, ForeignKey, PhysicalSchema, TableAnnotation + + cfg = _load_config_or_exit(config) + phys_file = physical_path(cfg) + if not phys_file.exists(): + typer.secho( + f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + physical = PhysicalSchema.from_yaml(phys_file) + ann_path = annotations_path(cfg) + annotations = Annotations.from_yaml(ann_path) + + assumed: dict[str, str] = {} + for a in assume or []: + col, _, ref = a.partition("=") + if not ref or ref not in physical.tables: + typer.secho( + f"ERRORE: --assume '{a}' non valido (atteso col=tabella nel catalogo).", + fg=typer.colors.RED, err=True, + ) + raise typer.Exit(code=1) + assumed[col] = ref + + def _single_pk(table) -> str | None: + pks = [c for c, col in table.columns.items() if col.pk] + return pks[0] if len(pks) == 1 else None + + pk_owners: dict[str, list[str]] = {} + for tname, table in physical.tables.items(): + pk = _single_pk(table) + if pk: + pk_owners.setdefault(pk, []).append(tname) + + dim_time_pk = None + if "dim_time" in physical.tables: + dim_time_pk = _single_pk(physical.tables["dim_time"]) + + def _known(tname: str) -> set: + keys = set() + for fk in physical.tables[tname].foreign_keys: + keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))) + ann = annotations.tables.get(tname) + if ann: + for fk in ann.foreign_keys: + keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))) + return keys + + known_by_table: dict[str, set] = {t: _known(t) for t in physical.tables} + suggested: dict[str, list[ForeignKey]] = {} + + def _add(tname: str, col: str, ref_table: str, ref_col: str) -> None: + key = ((col,), ref_table, (ref_col,)) + if key in known_by_table[tname]: + return + known_by_table[tname].add(key) + suggested.setdefault(tname, []).append( + ForeignKey(columns=[col], ref_table=ref_table, ref_columns=[ref_col]) + ) + + # Regola 1: join minati dall'SQL approvato. + n_sql_files = 0 + mined_total = 0 + for d in from_sql or []: + for sql_file in sorted(d.rglob("*.sql")): + n_sql_files += 1 + pairs = mine_join_pairs(sql_file.read_text(), physical) + mined_total += sum(pairs.values()) + for (src_t, src_c, ref_t, ref_c) in pairs: + _add(src_t, src_c, ref_t, ref_c) + + # Regole 2 e 3: convenzioni di naming. + ambiguous_skipped: set[str] = set() + for tname, table in physical.tables.items(): + for cname in table.columns: + if dim_time_pk and cname.endswith("time_key") and tname != "dim_time": + _add(tname, cname, "dim_time", dim_time_pk) + continue + if cname in assumed: + if assumed[cname] != tname: + _add(tname, cname, assumed[cname], cname) + continue + owners = [o for o in pk_owners.get(cname, []) if o != tname] + if not owners or cname in _GENERIC_PK_NAMES: + continue + if len(pk_owners[cname]) > 1: + ambiguous_skipped.add(cname) + continue + _add(tname, cname, owners[0], cname) + + if n_sql_files: + typer.secho( + f"Minati {mined_total} equi-join da {n_sql_files} file SQL.", + fg=typer.colors.BLUE, err=True, + ) + if ambiguous_skipped: + typer.secho( + "PK ambigue saltate dalla regola same-name (piu' tabelle proprietarie): " + + ", ".join(sorted(ambiguous_skipped)) + + ". Se servono, aggiungile a mano o passa --from-sql.", + fg=typer.colors.YELLOW, err=True, + ) + + n_fks = sum(len(v) for v in suggested.values()) + if not suggested: + typer.secho("OK: nessuna FK da suggerire.", fg=typer.colors.GREEN) + return + + if write: + for tname, fks in suggested.items(): + ann = annotations.tables.setdefault(tname, TableAnnotation()) + ann.foreign_keys.extend(fks) + annotations.to_yaml(ann_path) + typer.secho( + f"OK: {n_fks} FK suggerite aggiunte a {ann_path} " + f"({len(suggested)} tabelle). Rivedile a mano prima dell'uso.", + fg=typer.colors.GREEN, + ) + return + + payload = { + "tables": { + tname: {"foreign_keys": [fk.model_dump(exclude_defaults=True) for fk in fks]} + for tname, fks in suggested.items() + } + } + typer.echo(_yaml.safe_dump(payload, sort_keys=False, allow_unicode=True)) + typer.secho( + f"{n_fks} FK candidate ({len(suggested)} tabelle). " + f"Usa --write per fonderle in annotations.yaml, poi curale a mano.", + fg=typer.colors.YELLOW, + ) + + @schema_app.command("render") def render_cmd( config: Path = CONFIG_OPT, diff --git a/harness/tht/cli/search_cmd.py b/harness/tht/cli/search_cmd.py index e0ef8755..082bc6f6 100644 --- a/harness/tht/cli/search_cmd.py +++ b/harness/tht/cli/search_cmd.py @@ -205,3 +205,156 @@ def search_cmd( row.append((r.content[:120] + "…") if len(r.content) > 120 else r.content) table.add_row(*row) Console().print(table) + + +# Dimensioni fisse del pack (niente config: il pack deve restare piccolo perche' +# entra nel contesto del modello in un turno solo). +PACK_EVIDENCE_TOP = 5 +PACK_SOLVED_TOP = 3 +PACK_EXCERPT_CHARS = 400 + + +@search_app.command("pack") +def pack_cmd( + question: str = typer.Argument(..., help="La domanda in linguaggio naturale."), + config: Path = CONFIG_OPT, + session: str = typer.Option( + None, "--session", help="Scrive il pack in sessions//retrieval_pack.md." + ), + json_out: bool = typer.Option(False, "--json", help="Output JSON (per Pi)."), +) -> None: + """Context-pack F1: tabelle candidate + evidence + domande risolte in UNA chiamata. + + Un solo embedding della domanda, riusato per le tre ricerche vettoriali. + Degrado gentile: se Ollama/vectordb non rispondono, le sezioni restano vuote + con un'avvertenza (exit 0) — la sessione prosegue con le ricerche live. + """ + from sqlalchemy.exc import OperationalError + + from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg + from tht.search import combined_search, schema_tables + from tht.solved import SOLVED_KIND + from tht.vectorstore.embeddings import EmbeddingsError + from tht.vectorstore.rest_client import VectorRestError + + cfg = _load_config_or_exit(config) + require_vector_cfg(cfg) + + tables: list[dict] = [] + evidence: list[dict] = [] + solved: list[dict] = [] + warnings: list[str] = [] + degrade = (VectorRestError, EmbeddingsError, OperationalError) + + vec = None + searcher = embedder = None + try: + searcher = open_searcher(cfg) + embedder = make_embedder(cfg.embeddings) + vec = embedder.embed_query(question) + except degrade as e: + warnings.append(f"retrieval non disponibile ({e}): prosegui con le ricerche live") + + if vec is not None: + from tht.cli.schema_cmd import physical_path + + descriptions: dict[str, str] = {} + phys_file = physical_path(cfg) + if phys_file.exists(): + from tht.mschema.models import PhysicalSchema + + phys = PhysicalSchema.from_yaml(phys_file) + descriptions = {t: tab.comment for t, tab in phys.tables.items()} + try: + cand = combined_search( + keyword=question, lsh_hits=None, store=searcher, embedder=embedder, + top=cfg.search.schema_chunk_pool, rrf_k=cfg.search.rrf_k, + kinds=KIND_MAP["schema"], query_vec=vec, + ) + tables = [ + {"name": n, "rrf": round(s, 6), "description": descriptions.get(n, "")} + for n, s in schema_tables(cand, top_tables=cfg.search.top_schema_tables) + ] + except degrade as e: + warnings.append(f"ricerca schema fallita ({e})") + try: + ev = combined_search( + keyword=question, lsh_hits=None, store=searcher, embedder=embedder, + top=PACK_EVIDENCE_TOP, rrf_k=cfg.search.rrf_k, + kinds=KIND_MAP["evidence"], query_vec=vec, + ) + evidence = [ + {"title": r.label, "status": r.status, + "excerpt": r.content[:PACK_EXCERPT_CHARS]} + for r in ev + ] + except degrade as e: + warnings.append(f"ricerca evidence fallita ({e})") + try: + hits = searcher.search(vec, top_n=PACK_SOLVED_TOP, kinds=[SOLVED_KIND]) + solved = [ + { + "session_id": h.metadata.get("session_id", h.ref), + "question": h.metadata.get("question", h.content), + "sql": h.metadata.get("sql", ""), + "tables": h.metadata.get("tables", []), + "score": round(h.similarity, 4), + } + for h in hits + ] + except degrade as e: + warnings.append(f"solved-search fallita ({e})") + + for w in warnings: + typer.secho(f"ATTENZIONE: {w}", fg=typer.colors.YELLOW, err=True) + + md_lines = ["# Retrieval pack", "", f"Domanda: {question}", ""] + md_lines += [f"## Tabelle candidate (top {len(tables)}, vettoriale sull'intera domanda)", ""] + if tables: + for i, t in enumerate(tables, 1): + desc = f" — {t['description']}" if t["description"] else "" + md_lines.append(f"{i}. **{t['name']}**{desc} (rrf {t['rrf']})") + else: + md_lines.append("_nessuna (retrieval non disponibile o nessun match)_") + md_lines += ["", "## Evidence rilevanti", ""] + if evidence: + for e in evidence: + status = f" [{e['status']}]" if e["status"] else "" + md_lines.append(f"- **{e['title']}**{status}: {e['excerpt']}") + else: + md_lines.append("_nessuna_") + md_lines += ["", "## Domande risolte simili (exemplar di riferimento, NON decisioni)", ""] + if solved: + for s in solved: + md_lines.append( + f"### {s['question']} \n(sessione `{s['session_id']}`; " + f"tabelle: {', '.join(s['tables']) or '-'})" + ) + if s["sql"]: + md_lines += ["", "```sql", s["sql"], "```", ""] + else: + md_lines.append("_nessuna_") + if warnings: + md_lines += ["", "## Avvertenze", ""] + [f"- {w}" for w in warnings] + md = "\n".join(md_lines) + "\n" + + if session: + from tht.cli.session_cmd import load_session_or_exit, session_dir + + load_session_or_exit(cfg, session) + out = session_dir(cfg, session) / "retrieval_pack.md" + out.write_text(md) + if not json_out: + typer.secho( + f"OK: retrieval pack scritto in {out} " + f"({len(tables)} tabelle, {len(evidence)} evidence, {len(solved)} solved).", + fg=typer.colors.GREEN, + ) + if json_out: + typer.echo(json.dumps( + {"question": question, "tables": tables, "evidence": evidence, + "solved": solved, "warnings": warnings}, + ensure_ascii=False, indent=2, + )) + elif not session: + typer.echo(md) diff --git a/harness/tht/cli/session_cmd.py b/harness/tht/cli/session_cmd.py index a0994768..01c73ad9 100644 --- a/harness/tht/cli/session_cmd.py +++ b/harness/tht/cli/session_cmd.py @@ -181,6 +181,11 @@ def show_cmd( data = manifest.model_dump(mode="json", by_alias=True) data["phase"] = phase data["has_schema_linking"] = has_schema_linking + # Ledger integrale: il gate lo usa per costruire deterministicamente il + # recap delle decisioni nei riepiloghi di fase (v2). + data["decisions"] = [ + d.model_dump(mode="json") for d in list_decisions(sdir) + ] typer.echo(json.dumps(data, ensure_ascii=False, indent=2)) return diff --git a/harness/tht/mschema/fkmine.py b/harness/tht/mschema/fkmine.py new file mode 100644 index 00000000..2bdc958e --- /dev/null +++ b/harness/tht/mschema/fkmine.py @@ -0,0 +1,58 @@ +"""Mining dei join reali dall'SQL approvato: coppie equi-join -> FK logiche candidate. + +La fonte di verita' sono le query gia' validate da un umano (sql_final.sql, ctes/*.sql +delle sessioni approvate): un equi-join ricorrente tra due tabelle del catalogo, con +una delle due colonne PK della propria tabella, e' una FK logica ad alta confidenza. +""" + +from collections import Counter + +import sqlglot +from sqlglot import exp + +from tht.mschema.models import PhysicalSchema + +JoinPair = tuple[str, str, str, str] # (src_table, src_col, ref_table, ref_col) + + +def mine_join_pairs(sql_text: str, physical: PhysicalSchema) -> Counter: + """Estrae le coppie equi-join tra tabelle del catalogo da un testo SQL. + + Ritorna un Counter {(src_table, src_col, ref_table, ref_col): occorrenze}. + Il lato ref e' quello la cui colonna e' PK della propria tabella; coppie in cui + nessuno o entrambi i lati sono PK vengono scartate (non FK-like). Alias e CTE + vengono risolti; i riferimenti a CTE (non nel catalogo) sono ignorati. + """ + pairs: Counter = Counter() + try: + statements = sqlglot.parse(sql_text, read="postgres") + except sqlglot.errors.ParseError: + return pairs + for stmt in statements: + if stmt is None: + continue + alias_map: dict[str, str] = {} + for t in stmt.find_all(exp.Table): + alias_map[t.alias_or_name] = t.name + for eq in stmt.find_all(exp.EQ): + left, right = eq.left, eq.right + if not (isinstance(left, exp.Column) and isinstance(right, exp.Column)): + continue + if not (left.table and right.table): + continue + lt = alias_map.get(left.table, left.table) + rt = alias_map.get(right.table, right.table) + if lt == rt or lt not in physical.tables or rt not in physical.tables: + continue + lc, rc = left.name, right.name + if lc not in physical.tables[lt].columns or rc not in physical.tables[rt].columns: + continue + l_pk = physical.tables[lt].columns[lc].pk + r_pk = physical.tables[rt].columns[rc].pk + if l_pk == r_pk: # nessuna o entrambe PK: non FK-like + continue + if r_pk: + pairs[(lt, lc, rt, rc)] += 1 + else: + pairs[(rt, rc, lt, lc)] += 1 + return pairs diff --git a/harness/tht/mschema/merge.py b/harness/tht/mschema/merge.py index bc9a6b7d..f9c72949 100644 --- a/harness/tht/mschema/merge.py +++ b/harness/tht/mschema/merge.py @@ -12,4 +12,14 @@ def find_orphans(physical: PhysicalSchema, annotations: Annotations) -> list[str for column_name in table_ann.columns: if column_name not in table.columns: orphans.append(f"{table_name}.{column_name}") + for fk in table_ann.foreign_keys: + label = f"{table_name}.fk({','.join(fk.columns)})->{fk.ref_table}" + ref = physical.tables.get(fk.ref_table) + if ref is None: + orphans.append(label) + continue + missing = [c for c in fk.columns if c not in table.columns] + missing += [c for c in fk.ref_columns if c not in ref.columns] + if missing: + orphans.append(label) return orphans diff --git a/harness/tht/mschema/models.py b/harness/tht/mschema/models.py index 256ffb0b..2de69e7d 100644 --- a/harness/tht/mschema/models.py +++ b/harness/tht/mschema/models.py @@ -78,6 +78,10 @@ class TableAnnotation(BaseModel): concepts: list[str] = [] notes: str = "" columns: dict[str, ColumnAnnotation] = {} + # FK "logiche" curate a mano: il DWH non dichiara vincoli, quindi i join + # noti (es. data_time_key -> dim_time.day_key) vivono qui e vengono fusi + # con le FK fisiche in tutte le viste renderizzate. + foreign_keys: list[ForeignKey] = [] class Annotations(_YamlModel): diff --git a/harness/tht/mschema/render.py b/harness/tht/mschema/render.py index 81e606d9..eb7b41a0 100644 --- a/harness/tht/mschema/render.py +++ b/harness/tht/mschema/render.py @@ -1,11 +1,25 @@ from typing import Any from tht.mschema.eligibility import effective_eligibility -from tht.mschema.models import Annotations, ColumnAnnotation, PhysicalSchema +from tht.mschema.models import Annotations, ColumnAnnotation, ForeignKey, PhysicalSchema MAX_EXAMPLES_IN_PROMPT = 5 +def table_foreign_keys( + physical: PhysicalSchema, annotations: Annotations, table: str +) -> list[ForeignKey]: + """FK fisiche + FK logiche dalle annotations (dedup su columns/ref).""" + fks = list(physical.tables[table].foreign_keys) + ann = annotations.tables.get(table) + if ann: + seen = {(tuple(f.columns), f.ref_table, tuple(f.ref_columns)) for f in fks} + for fk in ann.foreign_keys: + if (tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns)) not in seen: + fks.append(fk) + return fks + + def _ann_col(annotations: Annotations, table: str, column: str) -> ColumnAnnotation | None: ann = annotations.tables.get(table) if ann is None: @@ -59,7 +73,7 @@ def to_mschema_text( shown = ", ".join(column.examples[:MAX_EXAMPLES_IN_PROMPT]) lines.append(f" -- Examples: {shown}") lines.append(");") - for fk in table.foreign_keys: + for fk in table_foreign_keys(physical, annotations, table_name): for src, dst in zip(fk.columns, fk.ref_columns): fk_lines.append(f"{table_name}.{src}={fk.ref_table}.{dst}") lines.extend(["", "【Foreign keys】", *fk_lines]) @@ -89,7 +103,7 @@ def to_schema_dict( "primary_keys": [c for c in cols if table.columns[c].pk], "foreign_keys": [ {"columns": fk.columns, "ref_table": fk.ref_table, "ref_columns": fk.ref_columns} - for fk in table.foreign_keys + for fk in table_foreign_keys(physical, annotations, table_name) ], } return out @@ -132,9 +146,10 @@ def to_markdown(physical: PhysicalSchema, annotations: Annotations | None = None f"| {column_name} | {column.type} | {'sì' if column.nullable else 'no'} " f"| {'sì' if column.pk else ''} | {cdesc} | {examples} |" ) - if table.foreign_keys: + fks = table_foreign_keys(physical, annotations, table_name) + if fks: lines += ["", "Foreign keys:"] - for fk in table.foreign_keys: + for fk in fks: lines.append( f"- ({', '.join(fk.columns)}) → {fk.ref_table} ({', '.join(fk.ref_columns)})" ) diff --git a/harness/tht/search/__init__.py b/harness/tht/search/__init__.py index 469e0ecc..fb1bf713 100644 --- a/harness/tht/search/__init__.py +++ b/harness/tht/search/__init__.py @@ -97,8 +97,12 @@ def combined_search( top: int, rrf_k: int, kinds: list[str] | None, + query_vec: list[float] | None = None, ) -> list[SearchResult]: - """Fonde LSH (valori di campo) e pgvector con Reciprocal Rank Fusion.""" + """Fonde LSH (valori di campo) e pgvector con Reciprocal Rank Fusion. + + `query_vec` permette di riusare un embedding gia' calcolato della stessa + keyword (es. `tht search pack`, che fa piu' ricerche sulla stessa domanda).""" rankings: dict[str, list[tuple[str, float]]] = {} lsh_values: dict[str, str] = {} if lsh_hits: @@ -106,7 +110,9 @@ def combined_search( rankings["lsh"] = [(key, score) for key, score, _ in aggregated] lsh_values = {key: value for key, _, value in aggregated} - vector_hits = store.search(embedder.embed_query(keyword), top_n=top * 2, kinds=kinds) + if query_vec is None: + query_vec = embedder.embed_query(keyword) + vector_hits = store.search(query_vec, top_n=top * 2, kinds=kinds) rankings["vector"] = [(_vector_key(h), h.similarity) for h in vector_hits] by_key = {_vector_key(h): h for h in vector_hits}