feat(opt): three efficiency levers for NL→SQL workflow
Lever 1: Join-graph via FK logics in annotations + suggest-fks command
- TableAnnotation.foreign_keys field stores curated logical FKs (DWH has no FK constraints)
- tht schema suggest-fks: mine from approved SQL, heuristics (time_key → dim_time),
same-name discovery + explicit --assume flag for multi-owner PKs
- mschema renders 【Foreign keys】 section populated; validation in merge.py
- SKILL.md F4 now reads FKs from mschema-text, no custom data_time_key logic
Lever 2: Context-pack consolidation at kickoff (tht search pack)
- Single embedding of question, reused for schema + evidence + solved searches
- One command: tht search pack <question> --session <id> → retrieval_pack.md
- Graceful degradation when Ollama/vector store unreachable (exit 0, empty sections)
- SKILL.md F1 prescribes as first call; reduces model thinking turns via pre-retrieval
Lever 3: Phase-summary recap v2 auto-construction from session ledger
- tht session show --json includes full decisions ledger
- tht phase meta --json exports 'emits' (substantive decision types per phase)
- Gate appends deterministic 【Decisioni registrate in questa fase】 section (appendLedgerSection)
- Model authors only summary + checks; recap table comes from persisted state (exact by construction)
- SKILL.md Disciplina 6: brief model output, gate fills the rest
Tests: 358 Python (including 10 FK + 3 pack + 1 session-ledger tests) + 111 JS gate tests, all pass.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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);
|
||||
});
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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 };
|
||||
}
|
||||
|
||||
|
||||
@@ -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 `<workspace>/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 `<workspace>/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 "<term>"` (evidence + schema, LSH
|
||||
over real values) and `tht search find --kind evidence "<term>"`. The LSH exposes
|
||||
EVERY column where a value appears — it does not collapse to a single best match,
|
||||
so a value like "ablazione" may anchor on multiple columns.
|
||||
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 "<original question>" --session <id>` —
|
||||
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 "<term>"` / `tht search find --kind evidence "<term>"` 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 <name>` (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 "<question>" --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
|
||||
|
||||
@@ -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/<id>/ctes/<name>.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/<name>.sql`, appends
|
||||
`SELECT * FROM <last CTE defined in that file>` 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:
|
||||
|
||||
> **<name>** — purpose: <one line>
|
||||
|
||||
@@ -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 = <fact>.data_time_key` and use the dimension's
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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
|
||||
@@ -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"])
|
||||
@@ -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"
|
||||
@@ -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
|
||||
],
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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/<id>/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)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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)})"
|
||||
)
|
||||
|
||||
@@ -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}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user