471 lines
16 KiB
Python
471 lines
16 KiB
Python
import hashlib
|
|
import json
|
|
from datetime import UTC, 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
|
|
|
|
RUNNER = CliRunner()
|
|
|
|
|
|
def _json_sha(value) -> str:
|
|
payload = json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
|
|
return "sha256:" + hashlib.sha256(payload.encode("utf-8")).hexdigest()
|
|
|
|
|
|
def _physical():
|
|
return PhysicalSchema(
|
|
database="d",
|
|
schema="s",
|
|
introspected_at=datetime(2026, 1, 1, tzinfo=UTC),
|
|
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(
|
|
f"""
|
|
runtime_identity:
|
|
workspace_id: demo
|
|
workspace_revision: {'a' * 40}
|
|
database: {{database: d, schema: s, user: u, password: p, transport: direct}}
|
|
paths: {{artifacts: {tmp_path / 'artifacts'}, indexes: {tmp_path / 'i'}, sessions: {tmp_path / 's'}}}
|
|
"""
|
|
)
|
|
return cfg
|
|
|
|
|
|
def test_suggest_fks_prints_candidates(tmp_path):
|
|
cfg = _write_workspace(tmp_path)
|
|
res = RUNNER.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_suggest_fks_json_is_pristine_and_stable(tmp_path):
|
|
cfg = _write_workspace(tmp_path)
|
|
first = tmp_path / "second.sql"
|
|
first.write_text(
|
|
"SELECT f.esito FROM datawarehouse.fact_ablazione f "
|
|
"JOIN datawarehouse.dim_patient p ON f.cod_paz = p.cod_paz"
|
|
)
|
|
second = tmp_path / "first.sql"
|
|
second.write_text(
|
|
"SELECT dt.year FROM datawarehouse.fact_ablazione f "
|
|
"JOIN datawarehouse.dim_time dt ON f.data_time_key = dt.day_key"
|
|
)
|
|
|
|
response = RUNNER.invoke(
|
|
app,
|
|
[
|
|
"schema",
|
|
"suggest-fks",
|
|
"-c",
|
|
str(cfg),
|
|
"--from-sql",
|
|
str(first),
|
|
"--from-sql",
|
|
str(second),
|
|
"--json",
|
|
],
|
|
)
|
|
|
|
assert response.exit_code == 0, response.output
|
|
assert response.stderr == ""
|
|
payload = json.loads(response.stdout)
|
|
assert payload["schemaVersion"] == 1
|
|
assert payload["status"] == "succeeded"
|
|
assert payload["code"] == "ok"
|
|
assert payload["operation"] == "schema_suggest_fks"
|
|
assert payload["workspaceId"] == "demo"
|
|
assert payload["workspaceRevision"] == "a" * 40
|
|
assert payload["counts"] == {
|
|
"ambiguousColumns": 0,
|
|
"candidateTables": 1,
|
|
"candidates": 2,
|
|
"minedJoins": 2,
|
|
"sqlFiles": 2,
|
|
}
|
|
assert payload["candidateDocument"] == {
|
|
"annotations": {
|
|
"tables": {
|
|
"fact_ablazione": {
|
|
"foreign_keys": [
|
|
{
|
|
"columns": ["cod_paz"],
|
|
"ref_columns": ["cod_paz"],
|
|
"ref_table": "dim_patient",
|
|
},
|
|
{
|
|
"columns": ["data_time_key"],
|
|
"ref_columns": ["day_key"],
|
|
"ref_table": "dim_time",
|
|
},
|
|
]
|
|
}
|
|
}
|
|
},
|
|
"counts": {"candidateTables": 1, "candidates": 2},
|
|
"schemaVersion": 1,
|
|
}
|
|
assert payload["candidate_count"] == payload["counts"]["candidates"]
|
|
candidate_yaml = payload["candidate_yaml"]
|
|
assert payload["candidateDigest"] == "sha256:" + hashlib.sha256(candidate_yaml.encode("utf-8")).hexdigest()
|
|
assert yaml.safe_load(candidate_yaml) == payload["candidateDocument"]["annotations"]
|
|
|
|
rerun = RUNNER.invoke(
|
|
app,
|
|
[
|
|
"schema",
|
|
"suggest-fks",
|
|
"-c",
|
|
str(cfg),
|
|
"--from-sql",
|
|
str(second),
|
|
"--from-sql",
|
|
str(first),
|
|
"--json",
|
|
],
|
|
)
|
|
assert rerun.exit_code == 0, rerun.output
|
|
assert json.loads(rerun.stdout) == payload
|
|
|
|
|
|
def test_suggest_fks_json_rejects_invalid_assume_without_prose(tmp_path):
|
|
cfg = _write_workspace(tmp_path)
|
|
|
|
response = RUNNER.invoke(
|
|
app,
|
|
["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=nope", "--json"],
|
|
)
|
|
|
|
assert response.exit_code == 1
|
|
assert response.stderr == ""
|
|
payload = json.loads(response.stdout)
|
|
assert payload["schemaVersion"] == 1
|
|
assert payload["status"] == "failed"
|
|
assert payload["code"] == "invalid_argument"
|
|
assert payload["error"] == "invalid assume mapping"
|
|
|
|
|
|
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
|
|
assert len(pairs) == 1
|
|
|
|
|
|
def test_mine_join_pairs_ignores_non_pk_pairs_and_bad_sql():
|
|
from tht.mschema.fkmine import mine_join_pairs
|
|
|
|
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, tzinfo=UTC),
|
|
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 = RUNNER.invoke(app, ["schema", "suggest-fks", "-c", str(cfg)])
|
|
assert res.exit_code == 0, res.output
|
|
assert "nessuna FK da suggerire" in res.output
|
|
assert "cod_x" in res.output
|
|
|
|
res2 = RUNNER.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
|
|
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", [])
|
|
)
|
|
|
|
res3 = RUNNER.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)
|
|
sql_file = tmp_path / "approved.sql"
|
|
sql_file.write_text(
|
|
"SELECT f.esito FROM datawarehouse.fact_ablazione f "
|
|
"JOIN datawarehouse.dim_patient p ON f.cod_paz = p.cod_paz"
|
|
)
|
|
res = RUNNER.invoke(
|
|
app,
|
|
["schema", "suggest-fks", "-c", str(cfg), "--from-sql", str(sql_file)],
|
|
)
|
|
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_schema_check_json_validates_staged_annotations_without_mutating_runtime(tmp_path):
|
|
cfg = _write_workspace(tmp_path)
|
|
runtime_annotations = tmp_path / "artifacts" / "mschema" / "annotations.yaml"
|
|
runtime_annotations.write_text("tables: {}\n")
|
|
reviewed = tmp_path / "reviewed.yaml"
|
|
_annotations_with_fks().to_yaml(reviewed)
|
|
|
|
response = RUNNER.invoke(
|
|
app,
|
|
[
|
|
"schema",
|
|
"check",
|
|
"-c",
|
|
str(cfg),
|
|
"--annotations",
|
|
str(reviewed),
|
|
"--reviewed-candidates",
|
|
"sha256:" + "b" * 64,
|
|
"--json",
|
|
],
|
|
)
|
|
|
|
assert response.exit_code == 0, response.output
|
|
assert response.stderr == ""
|
|
payload = json.loads(response.stdout)
|
|
assert payload["orphan_count"] == 0
|
|
assert payload["reviewed_candidates_digest"] == "sha256:" + "b" * 64
|
|
assert payload["annotations_digest"] == "sha256:" + hashlib.sha256(reviewed.read_bytes()).hexdigest()
|
|
assert payload == {
|
|
"annotationsDigest": _json_sha(
|
|
{
|
|
"annotations": {
|
|
"tables": {
|
|
"fact_ablazione": {
|
|
"foreign_keys": [
|
|
{
|
|
"columns": ["cod_paz"],
|
|
"ref_columns": ["cod_paz"],
|
|
"ref_table": "dim_patient",
|
|
},
|
|
{
|
|
"columns": ["data_time_key"],
|
|
"ref_columns": ["day_key"],
|
|
"ref_table": "dim_time",
|
|
},
|
|
]
|
|
}
|
|
}
|
|
},
|
|
"schemaVersion": 1,
|
|
}
|
|
),
|
|
"annotations_digest": "sha256:" + hashlib.sha256(reviewed.read_bytes()).hexdigest(),
|
|
"code": "ok",
|
|
"orphan_count": 0,
|
|
"reviewed_candidates_digest": "sha256:" + "b" * 64,
|
|
"counts": {"annotationTables": 1, "foreignKeys": 2, "orphans": 0},
|
|
"operation": "schema_check",
|
|
"orphans": [],
|
|
"reviewedCandidates": "sha256:" + "b" * 64,
|
|
"schemaVersion": 1,
|
|
"status": "succeeded",
|
|
"workspaceId": "demo",
|
|
"workspaceRevision": "a" * 40,
|
|
"zeroOrphans": True,
|
|
}
|
|
assert runtime_annotations.read_text() == "tables: {}\n"
|
|
|
|
|
|
def test_schema_check_json_reports_orphans_without_prose(tmp_path):
|
|
cfg = _write_workspace(tmp_path)
|
|
reviewed = tmp_path / "reviewed.yaml"
|
|
Annotations(
|
|
tables={
|
|
"fact_ablazione": TableAnnotation(
|
|
foreign_keys=[
|
|
ForeignKey(
|
|
columns=["cod_paz"],
|
|
ref_table="dim_missing",
|
|
ref_columns=["cod_paz"],
|
|
)
|
|
]
|
|
)
|
|
}
|
|
).to_yaml(reviewed)
|
|
|
|
response = RUNNER.invoke(
|
|
app,
|
|
["schema", "check", "-c", str(cfg), "--annotations", str(reviewed), "--json"],
|
|
)
|
|
|
|
assert response.exit_code == 3
|
|
assert response.stderr == ""
|
|
payload = json.loads(response.stdout)
|
|
assert payload["schemaVersion"] == 1
|
|
assert payload["status"] == "blocked"
|
|
assert payload["code"] == "annotation_invalid"
|
|
assert payload["zeroOrphans"] is False
|
|
assert payload["counts"]["orphans"] == 1
|
|
assert payload["orphans"] == ["fact_ablazione.fk(cod_paz)->dim_missing"]
|
|
|
|
|
|
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 = RUNNER.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"
|
|
assert len(ann.tables["fact_ablazione"].foreign_keys) == 2
|
|
|
|
res2 = RUNNER.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
|