Files
ThothII/harness/tests/test_schema_fk_annotations.py
T

438 lines
17 KiB
Python

# ruff: noqa: DTZ001
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
def test_suggest_fks_json_is_single_deterministic_document_without_writing(tmp_path):
import hashlib
import json
cfg = _write_workspace(tmp_path)
staged = tmp_path / "staged"
staged.mkdir()
(staged / "z.sql").write_text(
"SELECT * FROM fact_ablazione f JOIN dim_patient p ON f.cod_paz = p.cod_paz"
)
(staged / "a.sql").write_text(
"SELECT * FROM fact_ablazione f JOIN dim_time p ON f.data_time_key = p.day_key"
)
response = CliRunner().invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)]
)
assert response.exit_code == 0, response.output
assert response.stderr == ""
assert response.stdout.count("\n") == 1
payload = json.loads(response.stdout)
assert payload["status"] == "succeeded"
assert payload["code"] == "ok"
assert payload["candidate_count"] == 2
assert payload["candidate_digest"].startswith("sha256:")
assert payload["orphan_count"] == 0
# The digest is over the canonical candidate export and --json never writes annotations.
expected = json.dumps(payload["candidates"], ensure_ascii=False, sort_keys=True, separators=(",", ":"))
assert payload["candidate_digest"] == "sha256:" + hashlib.sha256(expected.encode()).hexdigest()
assert not (tmp_path / "artifacts" / "mschema" / "annotations.yaml").exists()
def test_suggest_fks_json_failure_is_safe_and_single_document(tmp_path):
import json
cfg = _write_workspace(tmp_path)
(tmp_path / "artifacts" / "mschema" / "physical.yaml").unlink()
response = CliRunner().invoke(app, ["schema", "suggest-fks", "--json", "-c", str(cfg)])
assert response.exit_code != 0
assert response.stderr == ""
assert response.stdout.count("\n") == 1
payload = json.loads(response.stdout)
assert payload == {"status": "failed", "code": "physical_schema_missing"}
assert str(tmp_path) not in response.stdout
def test_schema_check_json_reports_orphan_count_without_prose(tmp_path):
import json
cfg = _write_workspace(tmp_path)
Annotations(tables={"gone": TableAnnotation(description="x")}).to_yaml(
tmp_path / "artifacts" / "mschema" / "annotations.yaml"
)
response = CliRunner().invoke(app, ["schema", "check", "--json", "-c", str(cfg)])
assert response.exit_code == 3
assert response.stderr == ""
assert response.stdout.count("\n") == 1
assert json.loads(response.stdout) == {
"status": "failed", "code": "annotation_orphans", "orphan_count": 1,
"orphans": ["gone"],
}
def test_suggest_fks_json_rejects_more_than_32_staged_sql_files(tmp_path):
import json
cfg = _write_workspace(tmp_path)
staged = tmp_path / "many"
staged.mkdir()
for i in range(33):
(staged / f"q{i:02d}.sql").write_text("SELECT 1")
response = CliRunner().invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)]
)
assert response.exit_code == 1
assert json.loads(response.stdout) == {"status": "failed", "code": "staged_sql_too_many"}
def test_staged_sql_size_limits_are_inclusive(tmp_path):
import json
cfg = _write_workspace(tmp_path)
staged = tmp_path / "boundary"
staged.mkdir()
# Exactly one MiB per file and exactly 16 MiB in aggregate are accepted.
body = "-- padding\n" + ("x" * (1 << 20))
body = body[: 1 << 20]
for i in range(16):
(staged / f"q{i:02d}.sql").write_bytes(body.encode())
response = CliRunner().invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)]
)
assert response.exit_code == 0, response.output
assert json.loads(response.stdout)["staged_sql_count"] == 16
def test_schema_check_json_success_and_reviewed_annotations_are_preserved(tmp_path):
import json
cfg = _write_workspace(tmp_path)
annotations_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml"
reviewed = _annotations_with_fks()
reviewed.to_yaml(annotations_path)
before = annotations_path.read_text()
check = CliRunner().invoke(app, ["schema", "check", "--json", "-c", str(cfg)])
assert check.exit_code == 0, check.output
assert check.stderr == ""
assert check.stdout.count("\n") == 1
assert json.loads(check.stdout) == {
"status": "succeeded", "code": "ok", "orphan_count": 0, "orphans": []
}
suggest = CliRunner().invoke(app, ["schema", "suggest-fks", "--json", "-c", str(cfg)])
assert suggest.exit_code == 0, suggest.output
assert suggest.stderr == ""
assert json.loads(suggest.stdout)["candidate_count"] == 0
assert annotations_path.read_text() == before
def test_schema_human_check_and_suggest_report_missing_physical_schema(tmp_path):
cfg = _write_workspace(tmp_path)
physical = tmp_path / "artifacts" / "mschema" / "physical.yaml"
physical.unlink()
expected = "physical.yaml non trovato. Esegui prima `tht schema introspect`."
check = CliRunner().invoke(app, ["schema", "check", "-c", str(cfg)])
assert check.exit_code == 1
assert expected in check.output
suggest = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg)])
assert suggest.exit_code == 1
assert expected in suggest.output
def test_schema_human_check_reports_ignored_columns(tmp_path):
cfg = _write_workspace(tmp_path)
physical = _physical()
physical.tables["fact_ablazione"].columns["esito"].eligible = False
physical.tables["fact_ablazione"].columns["esito"].eligibility_reason = "test ignored"
physical.to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml")
response = CliRunner().invoke(app, ["schema", "check", "-c", str(cfg)])
assert response.exit_code == 0, response.output
assert "Colonne ignorate (testo ampio, 1):" in response.output
assert "fact_ablazione.esito (test ignored)" in response.output
def test_staged_sql_file_size_over_one_mib_is_rejected(tmp_path):
import json
cfg = _write_workspace(tmp_path)
staged = tmp_path / "oversize"
staged.mkdir()
(staged / "too-large.sql").write_bytes(b"x" * ((1 << 20) + 1))
response = CliRunner().invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)]
)
assert response.exit_code == 1
assert json.loads(response.stdout) == {"status": "failed", "code": "staged_sql_too_large"}
def test_staged_sql_aggregate_over_16_mib_is_rejected(tmp_path):
import json
cfg = _write_workspace(tmp_path)
staged = tmp_path / "aggregate"
staged.mkdir()
for i in range(16):
(staged / f"q{i:02d}.sql").write_bytes(b"x" * (1 << 20))
(staged / "over.sql").write_bytes(b"x")
response = CliRunner().invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)]
)
assert response.exit_code == 1
assert json.loads(response.stdout) == {"status": "failed", "code": "staged_sql_too_large"}
def test_suggest_fks_candidate_order_is_stable_for_reordered_staged_inputs(tmp_path):
import json
cfg = _write_workspace(tmp_path)
first = tmp_path / "first"
second = tmp_path / "second"
first.mkdir()
second.mkdir()
(first / "join.sql").write_text(
"SELECT * FROM fact_ablazione f JOIN dim_patient p ON f.cod_paz = p.cod_paz"
)
(second / "join.sql").write_text(
"SELECT * FROM fact_ablazione f JOIN dim_time p ON f.data_time_key = p.day_key"
)
runner = CliRunner()
one = runner.invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(first), "--from-sql", str(second)]
)
two = runner.invoke(
app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(second), "--from-sql", str(first)]
)
assert one.exit_code == two.exit_code == 0
assert json.loads(one.stdout) == json.loads(two.stdout)
def test_schema_human_check_missing_config_has_original_error(tmp_path):
response = CliRunner().invoke(app, ["schema", "check", "-c", str(tmp_path / "missing.yaml")])
assert response.exit_code == 1
assert response.output
assert "ERRORE:" in response.output
def test_schema_human_suggest_missing_config_has_original_error(tmp_path):
response = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(tmp_path / "missing.yaml")])
assert response.exit_code == 1
assert response.output
assert "ERRORE:" in response.output