feat: add P2 harness machine contracts
This commit is contained in:
@@ -1,4 +1,5 @@
|
|||||||
import json
|
import json
|
||||||
|
from pathlib import Path
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
|
|
||||||
from typer.testing import CliRunner
|
from typer.testing import CliRunner
|
||||||
@@ -51,9 +52,10 @@ def test_preprocess_failed_job_report_is_sanitized_json_and_nonzero(monkeypatch,
|
|||||||
|
|
||||||
|
|
||||||
def test_preprocess_real_failed_stage_result_exits_nonzero(monkeypatch, tmp_path):
|
def test_preprocess_real_failed_stage_result_exits_nonzero(monkeypatch, tmp_path):
|
||||||
import tht.cli.preprocess_cmd as command
|
|
||||||
from test_corpus_pipeline import Source, item, pipeline
|
from test_corpus_pipeline import Source, item, pipeline
|
||||||
|
|
||||||
|
import tht.cli.preprocess_cmd as command
|
||||||
|
|
||||||
result = pipeline(
|
result = pipeline(
|
||||||
tmp_path, Source([(item("one", "a"), RuntimeError("SENSITIVE EVIDENCE secret"))])
|
tmp_path, Source([(item("one", "a"), RuntimeError("SENSITIVE EVIDENCE secret"))])
|
||||||
).run_as_job(
|
).run_as_job(
|
||||||
@@ -129,3 +131,32 @@ def test_preprocess_evidence_gc_json_is_pristine(monkeypatch, tmp_path):
|
|||||||
)
|
)
|
||||||
assert response.exit_code == 0, response.output
|
assert response.exit_code == 0, response.output
|
||||||
assert json.loads(response.output)["dry_run"] is True
|
assert json.loads(response.output)["dry_run"] is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_preprocess_evidence_uses_runtime_identity_for_dev_fd_config(monkeypatch, tmp_path):
|
||||||
|
import tht.cli.preprocess_cmd as command
|
||||||
|
|
||||||
|
cfg = SimpleNamespace(
|
||||||
|
runtime_identity=SimpleNamespace(workspace_id="runtime-workspace"),
|
||||||
|
embeddings=SimpleNamespace(model="m", dim=4),
|
||||||
|
vector=SimpleNamespace(max_chunk_chars=10, retain_published_generations=1),
|
||||||
|
paths=SimpleNamespace(artifacts=tmp_path / "artifacts"),
|
||||||
|
model_dump_json=lambda: "{}",
|
||||||
|
)
|
||||||
|
captured = {}
|
||||||
|
|
||||||
|
class FakePipeline:
|
||||||
|
def __init__(self, **kwargs):
|
||||||
|
captured.update(kwargs)
|
||||||
|
|
||||||
|
def run_as_job(self, **kwargs):
|
||||||
|
captured.update(kwargs)
|
||||||
|
return SimpleNamespace(model_dump=lambda mode=None: {"status": "succeeded"})
|
||||||
|
|
||||||
|
monkeypatch.setattr(command, "_load_config_or_exit", lambda _: cfg)
|
||||||
|
monkeypatch.setattr("tht.adapters.factory.build_evidence_sources", lambda _: [])
|
||||||
|
monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda *_args, **_kwargs: object())
|
||||||
|
monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: object())
|
||||||
|
monkeypatch.setattr("tht.corpus.pipeline.CorpusPipeline", FakePipeline)
|
||||||
|
command.run_from_config(Path("/dev/fd/3"))
|
||||||
|
assert captured["workspace_id"] == "runtime-workspace"
|
||||||
|
|||||||
@@ -187,3 +187,37 @@ def test_memory_solved_index_help_uses_semantic_store_wording():
|
|||||||
assert res.exit_code == 0, res.output
|
assert res.exit_code == 0, res.output
|
||||||
assert "semantic" in res.output.lower() or "qdrant" in res.output.lower()
|
assert "semantic" in res.output.lower() or "qdrant" in res.output.lower()
|
||||||
assert "vectordb" not in res.output.lower()
|
assert "vectordb" not in res.output.lower()
|
||||||
|
|
||||||
|
|
||||||
|
def test_vector_index_schema_json_is_single_document(monkeypatch, tmp_path):
|
||||||
|
import json
|
||||||
|
|
||||||
|
cfg = _qdrant_runtime_config(tmp_path)
|
||||||
|
_write_schema_artifacts(tmp_path)
|
||||||
|
store = _FakeVectorStore()
|
||||||
|
monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store)
|
||||||
|
monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: _FakeEmbedder())
|
||||||
|
response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)])
|
||||||
|
assert response.exit_code == 0, response.output
|
||||||
|
assert response.stdout.count("\n") == 1
|
||||||
|
payload = json.loads(response.stdout)
|
||||||
|
assert payload["status"] == "succeeded"
|
||||||
|
assert payload["code"] == "ok"
|
||||||
|
assert payload["counts"]["added"] == 2
|
||||||
|
|
||||||
|
|
||||||
|
def test_vector_index_schema_json_failure_is_safe(monkeypatch, tmp_path):
|
||||||
|
import json
|
||||||
|
|
||||||
|
cfg = _qdrant_runtime_config(tmp_path)
|
||||||
|
_write_schema_artifacts(tmp_path)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
"tht.adapters.factory.build_vector_store",
|
||||||
|
lambda cfg, require_write: (_ for _ in ()).throw(RuntimeError("secret qdrant endpoint")),
|
||||||
|
)
|
||||||
|
response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)])
|
||||||
|
assert response.exit_code != 0
|
||||||
|
assert response.stdout.count("\n") == 1
|
||||||
|
payload = json.loads(response.stdout)
|
||||||
|
assert payload == {"status": "failed", "code": "schema_index_failed"}
|
||||||
|
assert "secret qdrant" not in response.stdout
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
|
# ruff: noqa: DTZ001
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
|
|
||||||
import yaml
|
import yaml
|
||||||
@@ -221,3 +222,62 @@ def test_suggest_fks_write_merges_and_is_idempotent(tmp_path):
|
|||||||
assert "nessuna FK da suggerire" in res2.output
|
assert "nessuna FK da suggerire" in res2.output
|
||||||
ann2 = Annotations.from_yaml(ann_path)
|
ann2 = Annotations.from_yaml(ann_path)
|
||||||
assert len(ann2.tables["fact_ablazione"].foreign_keys) == 2
|
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.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.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.stdout.count("\n") == 1
|
||||||
|
assert json.loads(response.stdout) == {
|
||||||
|
"status": "failed", "code": "annotation_orphans", "orphan_count": 1,
|
||||||
|
"orphans": ["gone"],
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,17 +1,18 @@
|
|||||||
"""One-shot preprocessing commands."""
|
"""One-shot preprocessing commands."""
|
||||||
|
# ruff: noqa: BLE001
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import hashlib
|
||||||
import json
|
import json
|
||||||
import re
|
import re
|
||||||
import hashlib
|
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import typer
|
import typer
|
||||||
|
|
||||||
from tht.cli.config_cmd import CONFIG_OPT
|
from tht.cli.config_cmd import CONFIG_OPT
|
||||||
from tht.config import workspace_id_from_path
|
from tht.cli.schema_cmd import _load_config_or_exit
|
||||||
|
from tht.config import workspace_id_for_config
|
||||||
|
|
||||||
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
|
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
|
||||||
|
|
||||||
@@ -22,7 +23,8 @@ def run_dwh_from_config(
|
|||||||
from tht.cli.lsh_cmd import build_lsh_artifacts
|
from tht.cli.lsh_cmd import build_lsh_artifacts
|
||||||
from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog
|
from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog
|
||||||
from tht.jobs.dwh_pipeline import (
|
from tht.jobs.dwh_pipeline import (
|
||||||
DwhPreprocessPipeline, config_dwh_binding,
|
DwhPreprocessPipeline,
|
||||||
|
config_dwh_binding,
|
||||||
)
|
)
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
cfg = _load_config_or_exit(config)
|
||||||
@@ -64,7 +66,6 @@ def _parse_dwh_steps(value: str) -> tuple[str, ...]:
|
|||||||
|
|
||||||
def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None):
|
def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None):
|
||||||
from tht.adapters.factory import build_evidence_sources, build_vector_store
|
from tht.adapters.factory import build_evidence_sources, build_vector_store
|
||||||
from tht.cli.schema_cmd import _load_config_or_exit
|
|
||||||
from tht.cli.vector_cmd import make_embedder
|
from tht.cli.vector_cmd import make_embedder
|
||||||
from tht.corpus.chunk import ChunkPolicy
|
from tht.corpus.chunk import ChunkPolicy
|
||||||
from tht.corpus.pipeline import CorpusPipeline
|
from tht.corpus.pipeline import CorpusPipeline
|
||||||
@@ -87,7 +88,7 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None =
|
|||||||
return "sha256:" + hashlib.sha256(value.encode()).hexdigest()
|
return "sha256:" + hashlib.sha256(value.encode()).hexdigest()
|
||||||
|
|
||||||
return pipeline.run_as_job(
|
return pipeline.run_as_job(
|
||||||
workspace_id=workspace_id_from_path(config),
|
workspace_id=workspace_id_for_config(cfg, config),
|
||||||
workspace_root=corpus_root.parent,
|
workspace_root=corpus_root.parent,
|
||||||
config_fingerprint=fingerprint(cfg.model_dump_json()),
|
config_fingerprint=fingerprint(cfg.model_dump_json()),
|
||||||
input_fingerprint=fingerprint(config.resolve().as_posix()),
|
input_fingerprint=fingerprint(config.resolve().as_posix()),
|
||||||
@@ -116,7 +117,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False):
|
|||||||
pipeline_version="evidence-v1",
|
pipeline_version="evidence-v1",
|
||||||
retain_published_generations=cfg.vector.retain_published_generations,
|
retain_published_generations=cfg.vector.retain_published_generations,
|
||||||
)
|
)
|
||||||
pipeline.workspace_id = workspace_id_from_path(config)
|
pipeline.workspace_id = workspace_id_for_config(cfg, config)
|
||||||
return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run)
|
return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+285
-172
@@ -1,7 +1,9 @@
|
|||||||
from pathlib import Path
|
# ruff: noqa: BLE001, S110, B008
|
||||||
import logging
|
import logging
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
import typer
|
import typer
|
||||||
|
|
||||||
from tht.adapters.factory import build_dwh
|
from tht.adapters.factory import build_dwh
|
||||||
from tht.cli.config_cmd import CONFIG_OPT
|
from tht.cli.config_cmd import CONFIG_OPT
|
||||||
from tht.config import ConfigError, load_config
|
from tht.config import ConfigError, load_config
|
||||||
@@ -119,213 +121,324 @@ def introspect_cmd(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@schema_app.command("check")
|
class _MachineSchemaError(Exception):
|
||||||
def check_cmd(config: Path = CONFIG_OPT) -> None:
|
"""An expected schema CLI failure with a stable public code."""
|
||||||
"""Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane."""
|
|
||||||
from tht.mschema.merge import find_orphans
|
|
||||||
from tht.mschema.models import Annotations, PhysicalSchema
|
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
def __init__(self, code: str):
|
||||||
phys_file = physical_path(cfg)
|
self.code = code
|
||||||
if not phys_file.exists():
|
super().__init__(code)
|
||||||
typer.secho(
|
|
||||||
f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.",
|
|
||||||
fg=typer.colors.RED, err=True,
|
_MAX_STAGED_SQL_BYTES = 1 << 20
|
||||||
)
|
_MAX_STAGED_SQL_TOTAL = 16 << 20
|
||||||
raise typer.Exit(code=1)
|
|
||||||
physical = PhysicalSchema.from_yaml(phys_file)
|
|
||||||
|
def _schema_json(payload: dict) -> None:
|
||||||
|
import json
|
||||||
|
|
||||||
|
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":")))
|
||||||
|
|
||||||
|
|
||||||
|
def _machine_config(config: Path):
|
||||||
|
try:
|
||||||
|
return load_config(config)
|
||||||
|
except ConfigError:
|
||||||
|
raise _MachineSchemaError("invalid_configuration") from None
|
||||||
|
|
||||||
|
|
||||||
|
def _physical_or_error(cfg):
|
||||||
|
path = physical_path(cfg)
|
||||||
|
if not path.exists():
|
||||||
|
raise _MachineSchemaError("physical_schema_missing")
|
||||||
|
try:
|
||||||
|
from tht.mschema.models import PhysicalSchema
|
||||||
|
|
||||||
|
return PhysicalSchema.from_yaml(path)
|
||||||
|
except _MachineSchemaError:
|
||||||
|
raise
|
||||||
|
except Exception:
|
||||||
|
raise _MachineSchemaError("physical_schema_invalid") from None
|
||||||
|
|
||||||
|
|
||||||
|
def _staged_sql_files(inputs: list[Path] | None) -> list[Path]:
|
||||||
|
files: list[Path] = []
|
||||||
|
for item in inputs or []:
|
||||||
|
if item.is_file():
|
||||||
|
files.append(item)
|
||||||
|
elif item.is_dir():
|
||||||
|
files.extend(path for path in item.rglob("*.sql") if path.is_file())
|
||||||
|
else:
|
||||||
|
raise _MachineSchemaError("staged_sql_invalid")
|
||||||
|
# Resolve only for ordering; content remains read from the caller's staged path.
|
||||||
|
return sorted(set(files), key=lambda path: path.resolve().as_posix())
|
||||||
|
|
||||||
|
|
||||||
|
def _read_staged_sql(inputs: list[Path] | None) -> tuple[list[Path], list[str]]:
|
||||||
|
paths = _staged_sql_files(inputs)
|
||||||
|
contents: list[str] = []
|
||||||
|
total = 0
|
||||||
|
for path in paths:
|
||||||
|
try:
|
||||||
|
raw = path.read_bytes()
|
||||||
|
except (OSError, UnicodeError):
|
||||||
|
raise _MachineSchemaError("staged_sql_invalid") from None
|
||||||
|
if len(raw) > _MAX_STAGED_SQL_BYTES or total + len(raw) > _MAX_STAGED_SQL_TOTAL:
|
||||||
|
raise _MachineSchemaError("staged_sql_too_large")
|
||||||
|
total += len(raw)
|
||||||
|
try:
|
||||||
|
contents.append(raw.decode("utf-8"))
|
||||||
|
except UnicodeDecodeError:
|
||||||
|
raise _MachineSchemaError("staged_sql_invalid") from None
|
||||||
|
return paths, contents
|
||||||
|
|
||||||
|
|
||||||
|
def _candidate_key(fk) -> tuple:
|
||||||
|
return (tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))
|
||||||
|
|
||||||
|
|
||||||
|
def suggest_fks_data(config: Path, *, from_sql: list[Path] | None = None,
|
||||||
|
assume: list[str] | None = None) -> dict:
|
||||||
|
"""Return deterministic FK candidates without reviewing or mutating annotations."""
|
||||||
|
import hashlib
|
||||||
|
import json
|
||||||
|
|
||||||
|
from tht.mschema.fkmine import mine_join_pairs
|
||||||
|
from tht.mschema.merge import find_orphans
|
||||||
|
from tht.mschema.models import Annotations, ForeignKey
|
||||||
|
|
||||||
|
cfg = _machine_config(config)
|
||||||
|
physical = _physical_or_error(cfg)
|
||||||
|
_, sql_contents = _read_staged_sql(from_sql)
|
||||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||||
|
|
||||||
ignored = [
|
assumed: dict[str, str] = {}
|
||||||
f"{t}.{c} ({col.eligibility_reason})"
|
for value in assume or []:
|
||||||
for t, table in physical.tables.items()
|
col, sep, ref = value.partition("=")
|
||||||
for c, col in table.columns.items()
|
if not sep or not col or ref not in physical.tables:
|
||||||
if not col.eligible
|
raise _MachineSchemaError("assumption_invalid")
|
||||||
]
|
assumed[col] = ref
|
||||||
if ignored:
|
|
||||||
typer.secho(
|
|
||||||
f"Colonne ignorate (testo ampio, {len(ignored)}):", fg=typer.colors.YELLOW
|
|
||||||
)
|
|
||||||
for line in ignored:
|
|
||||||
typer.echo(f" - {line}")
|
|
||||||
|
|
||||||
|
def single_pk(table) -> str | None:
|
||||||
|
pks = [name for name, column in table.columns.items() if column.pk]
|
||||||
|
return pks[0] if len(pks) == 1 else None
|
||||||
|
|
||||||
|
pk_owners: dict[str, list[str]] = {}
|
||||||
|
for table_name in sorted(physical.tables):
|
||||||
|
pk = single_pk(physical.tables[table_name])
|
||||||
|
if pk:
|
||||||
|
pk_owners.setdefault(pk, []).append(table_name)
|
||||||
|
dim_time_pk = single_pk(physical.tables["dim_time"]) if "dim_time" in physical.tables else None
|
||||||
|
|
||||||
|
known: dict[str, set] = {}
|
||||||
|
for table_name in physical.tables:
|
||||||
|
keys = {
|
||||||
|
_candidate_key(fk)
|
||||||
|
for fk in physical.tables[table_name].foreign_keys
|
||||||
|
}
|
||||||
|
annotation = annotations.tables.get(table_name)
|
||||||
|
if annotation:
|
||||||
|
keys.update(_candidate_key(fk) for fk in annotation.foreign_keys)
|
||||||
|
known[table_name] = keys
|
||||||
|
|
||||||
|
suggested: dict[str, list[ForeignKey]] = {}
|
||||||
|
|
||||||
|
def add(table_name: str, column: str, ref_table: str, ref_column: str) -> None:
|
||||||
|
key = ((column,), ref_table, (ref_column,))
|
||||||
|
if key in known[table_name]:
|
||||||
|
return
|
||||||
|
known[table_name].add(key)
|
||||||
|
suggested.setdefault(table_name, []).append(
|
||||||
|
ForeignKey(columns=[column], ref_table=ref_table, ref_columns=[ref_column])
|
||||||
|
)
|
||||||
|
|
||||||
|
mined_total = 0
|
||||||
|
for sql_text in sql_contents:
|
||||||
|
pairs = mine_join_pairs(sql_text, physical)
|
||||||
|
mined_total += sum(pairs.values())
|
||||||
|
for src_table, src_column, ref_table, ref_column in sorted(pairs):
|
||||||
|
add(src_table, src_column, ref_table, ref_column)
|
||||||
|
|
||||||
|
ambiguous_columns: set[str] = set()
|
||||||
|
for table_name in sorted(physical.tables):
|
||||||
|
table = physical.tables[table_name]
|
||||||
|
for column_name in sorted(table.columns):
|
||||||
|
if dim_time_pk and column_name.endswith("time_key") and table_name != "dim_time":
|
||||||
|
add(table_name, column_name, "dim_time", dim_time_pk)
|
||||||
|
continue
|
||||||
|
if column_name in assumed:
|
||||||
|
if assumed[column_name] != table_name:
|
||||||
|
add(table_name, column_name, assumed[column_name], column_name)
|
||||||
|
continue
|
||||||
|
owners = [owner for owner in pk_owners.get(column_name, []) if owner != table_name]
|
||||||
|
if len(pk_owners.get(column_name, [])) > 1:
|
||||||
|
ambiguous_columns.add(column_name)
|
||||||
|
continue
|
||||||
|
if not owners or column_name in _GENERIC_PK_NAMES:
|
||||||
|
continue
|
||||||
|
add(table_name, column_name, owners[0], column_name)
|
||||||
|
|
||||||
|
candidates = [
|
||||||
|
{
|
||||||
|
"table": table_name,
|
||||||
|
"foreign_keys": [
|
||||||
|
fk.model_dump(exclude_defaults=True)
|
||||||
|
for fk in sorted(fks, key=_candidate_key)
|
||||||
|
],
|
||||||
|
}
|
||||||
|
for table_name, fks in sorted(suggested.items())
|
||||||
|
]
|
||||||
|
canonical = json.dumps(candidates, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
|
||||||
|
digest = "sha256:" + hashlib.sha256(canonical.encode("utf-8")).hexdigest()
|
||||||
|
orphan_count = len(find_orphans(physical, annotations))
|
||||||
|
return {
|
||||||
|
"status": "succeeded",
|
||||||
|
"code": "ok",
|
||||||
|
"candidates": candidates,
|
||||||
|
"candidate_count": sum(len(item["foreign_keys"]) for item in candidates),
|
||||||
|
"candidate_digest": digest,
|
||||||
|
"orphan_count": orphan_count,
|
||||||
|
"staged_sql_count": len(sql_contents),
|
||||||
|
"mined_join_count": mined_total,
|
||||||
|
"ambiguous_columns": sorted(ambiguous_columns),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def check_schema_data(config: Path) -> dict:
|
||||||
|
"""Validate the physical catalog and imported annotations without writing."""
|
||||||
|
from tht.mschema.merge import find_orphans
|
||||||
|
from tht.mschema.models import Annotations
|
||||||
|
|
||||||
|
cfg = _machine_config(config)
|
||||||
|
physical = _physical_or_error(cfg)
|
||||||
|
try:
|
||||||
|
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||||
|
except Exception:
|
||||||
|
raise _MachineSchemaError("annotations_invalid") from None
|
||||||
orphans = find_orphans(physical, annotations)
|
orphans = find_orphans(physical, annotations)
|
||||||
if orphans:
|
return {
|
||||||
typer.secho(f"ATTENZIONE: {len(orphans)} annotazioni orfane:", fg=typer.colors.YELLOW)
|
"status": "succeeded" if not orphans else "failed",
|
||||||
for o in orphans:
|
"code": "ok" if not orphans else "annotation_orphans",
|
||||||
typer.echo(f" - {o}")
|
"orphan_count": len(orphans),
|
||||||
|
"orphans": orphans,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@schema_app.command("check")
|
||||||
|
def check_cmd(
|
||||||
|
config: Path = CONFIG_OPT,
|
||||||
|
json_output: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."),
|
||||||
|
) -> None:
|
||||||
|
"""Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane."""
|
||||||
|
try:
|
||||||
|
payload = check_schema_data(config)
|
||||||
|
except _MachineSchemaError as error:
|
||||||
|
if json_output:
|
||||||
|
_schema_json({"status": "failed", "code": error.code})
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
_load_config_or_exit(config)
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
except Exception:
|
||||||
|
if json_output:
|
||||||
|
_schema_json({"status": "failed", "code": "schema_check_failed"})
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
typer.secho("ERRORE: impossibile verificare le annotazioni.", fg=typer.colors.RED, err=True)
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
if json_output:
|
||||||
|
# Keep the machine response to one object, including expected validation failures.
|
||||||
|
_schema_json(payload)
|
||||||
|
if payload["status"] != "succeeded":
|
||||||
|
raise typer.Exit(code=3)
|
||||||
|
return
|
||||||
|
if payload["orphan_count"]:
|
||||||
|
typer.secho(f"ATTENZIONE: {payload['orphan_count']} annotazioni orfane:", fg=typer.colors.YELLOW)
|
||||||
|
for orphan in payload["orphans"]:
|
||||||
|
typer.echo(f" - {orphan}")
|
||||||
raise typer.Exit(code=3)
|
raise typer.Exit(code=3)
|
||||||
typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN)
|
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")
|
@schema_app.command("suggest-fks")
|
||||||
def suggest_fks_cmd(
|
def suggest_fks_cmd(
|
||||||
config: Path = CONFIG_OPT,
|
config: Path = CONFIG_OPT,
|
||||||
from_sql: list[Path] = typer.Option(
|
from_sql: list[Path] = typer.Option(None, "--from-sql", help="Directory/file SQL approvati (ripetibile)."),
|
||||||
None, "--from-sql",
|
assume: list[str] = typer.Option(None, "--assume", help="Disambigua una PK: col=tabella."),
|
||||||
help="Directory di .sql approvati da cui minare i join reali (ripetibile).",
|
write: bool = typer.Option(False, "--write", help="Fonde i suggerimenti in annotations.yaml."),
|
||||||
),
|
json_output: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."),
|
||||||
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:
|
) -> None:
|
||||||
"""Suggerisce FK logiche per la curazione umana in annotations.yaml.
|
"""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
|
import yaml as _yaml
|
||||||
|
|
||||||
from tht.mschema.fkmine import mine_join_pairs
|
from tht.mschema.models import Annotations, TableAnnotation
|
||||||
from tht.mschema.models import Annotations, ForeignKey, PhysicalSchema, TableAnnotation
|
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
if json_output and write:
|
||||||
phys_file = physical_path(cfg)
|
_schema_json({"status": "failed", "code": "write_not_allowed"})
|
||||||
if not phys_file.exists():
|
raise typer.Exit(code=2)
|
||||||
typer.secho(
|
try:
|
||||||
f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.",
|
payload = suggest_fks_data(config, from_sql=from_sql, assume=assume)
|
||||||
fg=typer.colors.RED, err=True,
|
except _MachineSchemaError as error:
|
||||||
)
|
if json_output:
|
||||||
raise typer.Exit(code=1)
|
_schema_json({"status": "failed", "code": error.code})
|
||||||
physical = PhysicalSchema.from_yaml(phys_file)
|
raise typer.Exit(code=1) from None
|
||||||
ann_path = annotations_path(cfg)
|
if error.code == "assumption_invalid":
|
||||||
annotations = Annotations.from_yaml(ann_path)
|
typer.secho("ERRORE: --assume non valido (atteso col=tabella nel catalogo).", fg=typer.colors.RED, err=True)
|
||||||
|
else:
|
||||||
assumed: dict[str, str] = {}
|
_load_config_or_exit(config)
|
||||||
for a in assume or []:
|
typer.secho("ERRORE: impossibile elaborare il catalogo.", fg=typer.colors.RED, err=True)
|
||||||
col, _, ref = a.partition("=")
|
raise typer.Exit(code=1) from None
|
||||||
if not ref or ref not in physical.tables:
|
except Exception:
|
||||||
typer.secho(
|
if json_output:
|
||||||
f"ERRORE: --assume '{a}' non valido (atteso col=tabella nel catalogo).",
|
_schema_json({"status": "failed", "code": "schema_suggestion_failed"})
|
||||||
fg=typer.colors.RED, err=True,
|
raise typer.Exit(code=1) from None
|
||||||
)
|
typer.secho("ERRORE: impossibile elaborare il catalogo.", fg=typer.colors.RED, err=True)
|
||||||
raise typer.Exit(code=1)
|
raise typer.Exit(code=1) from None
|
||||||
assumed[col] = ref
|
if json_output:
|
||||||
|
_schema_json(payload)
|
||||||
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
|
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.
|
# Human mode remains the original renderer over the pure result.
|
||||||
n_sql_files = 0
|
if payload["mined_join_count"]:
|
||||||
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(
|
typer.secho(
|
||||||
f"Minati {mined_total} equi-join da {n_sql_files} file SQL.",
|
f"Minati {payload['mined_join_count']} equi-join da {payload['staged_sql_count']} file SQL.",
|
||||||
fg=typer.colors.BLUE, err=True,
|
fg=typer.colors.BLUE, err=True,
|
||||||
)
|
)
|
||||||
if ambiguous_skipped:
|
if payload["ambiguous_columns"]:
|
||||||
typer.secho(
|
typer.secho(
|
||||||
"PK ambigue saltate dalla regola same-name (piu' tabelle proprietarie): "
|
"PK ambigue saltate dalla regola same-name: " + ", ".join(payload["ambiguous_columns"]),
|
||||||
+ ", ".join(sorted(ambiguous_skipped))
|
|
||||||
+ ". Se servono, aggiungile a mano o passa --from-sql.",
|
|
||||||
fg=typer.colors.YELLOW, err=True,
|
fg=typer.colors.YELLOW, err=True,
|
||||||
)
|
)
|
||||||
|
if not payload["candidates"]:
|
||||||
n_fks = sum(len(v) for v in suggested.values())
|
|
||||||
if not suggested:
|
|
||||||
typer.secho("OK: nessuna FK da suggerire.", fg=typer.colors.GREEN)
|
typer.secho("OK: nessuna FK da suggerire.", fg=typer.colors.GREEN)
|
||||||
return
|
return
|
||||||
|
candidate_tables = {
|
||||||
|
item["table"]: {"foreign_keys": item["foreign_keys"]}
|
||||||
|
for item in payload["candidates"]
|
||||||
|
}
|
||||||
if write:
|
if write:
|
||||||
for tname, fks in suggested.items():
|
cfg = _load_config_or_exit(config)
|
||||||
ann = annotations.tables.setdefault(tname, TableAnnotation())
|
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||||
ann.foreign_keys.extend(fks)
|
for table_name, table_payload in candidate_tables.items():
|
||||||
annotations.to_yaml(ann_path)
|
ann = annotations.tables.setdefault(table_name, TableAnnotation())
|
||||||
|
from tht.mschema.models import ForeignKey
|
||||||
|
ann.foreign_keys.extend(ForeignKey(**fk) for fk in table_payload["foreign_keys"])
|
||||||
|
annotations.to_yaml(annotations_path(cfg))
|
||||||
typer.secho(
|
typer.secho(
|
||||||
f"OK: {n_fks} FK suggerite aggiunte a {ann_path} "
|
f"OK: {payload['candidate_count']} FK suggerite aggiunte a {annotations_path(cfg)}.",
|
||||||
f"({len(suggested)} tabelle). Rivedile a mano prima dell'uso.",
|
|
||||||
fg=typer.colors.GREEN,
|
fg=typer.colors.GREEN,
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
typer.echo(_yaml.safe_dump({"tables": candidate_tables}, sort_keys=False, allow_unicode=True))
|
||||||
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(
|
typer.secho(
|
||||||
f"{n_fks} FK candidate ({len(suggested)} tabelle). "
|
f"{payload['candidate_count']} FK candidate ({len(candidate_tables)} tabelle). "
|
||||||
f"Usa --write per fonderle in annotations.yaml, poi curale a mano.",
|
"Usa --write per fonderle in annotations.yaml, poi curale a mano.",
|
||||||
fg=typer.colors.YELLOW,
|
fg=typer.colors.YELLOW,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# 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("render")
|
@schema_app.command("render")
|
||||||
def render_cmd(
|
def render_cmd(
|
||||||
config: Path = CONFIG_OPT,
|
config: Path = CONFIG_OPT,
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
|
# ruff: noqa: BLE001
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import typer
|
import typer
|
||||||
@@ -116,22 +117,18 @@ def init_cmd(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@vector_app.command("index-schema")
|
def index_schema_data(config: Path) -> dict:
|
||||||
def index_schema_cmd(config: Path = CONFIG_OPT) -> None:
|
"""Synchronize schema records and return a bounded machine result."""
|
||||||
"""Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store."""
|
from tht.cli.schema_cmd import _machine_config
|
||||||
from tht.mschema.models import Annotations, PhysicalSchema
|
from tht.mschema.models import Annotations, PhysicalSchema
|
||||||
from tht.vectorstore.records import schema_records
|
from tht.vectorstore.records import schema_records
|
||||||
|
|
||||||
cfg = _load_config_or_exit(config)
|
cfg = _machine_config(config)
|
||||||
require_vector_write_allowed(cfg, "vector index-schema")
|
require_vector_write_allowed(cfg, "vector index-schema")
|
||||||
require_vector_cfg(cfg)
|
require_vector_cfg(cfg)
|
||||||
phys_file = physical_path(cfg)
|
phys_file = physical_path(cfg)
|
||||||
if not phys_file.exists():
|
if not phys_file.exists():
|
||||||
typer.secho(
|
raise ValueError("physical schema missing")
|
||||||
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)
|
physical = PhysicalSchema.from_yaml(phys_file)
|
||||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||||
records = schema_records(physical, annotations)
|
records = schema_records(physical, annotations)
|
||||||
@@ -143,4 +140,32 @@ def index_schema_cmd(config: Path = CONFIG_OPT) -> None:
|
|||||||
store=build_vector_store(cfg, require_write=True),
|
store=build_vector_store(cfg, require_write=True),
|
||||||
embedder=make_embedder(cfg.embeddings),
|
embedder=make_embedder(cfg.embeddings),
|
||||||
)
|
)
|
||||||
_print_stats(stats)
|
return {
|
||||||
|
"status": "succeeded",
|
||||||
|
"code": "ok",
|
||||||
|
"counts": {
|
||||||
|
"added": stats.added,
|
||||||
|
"updated": stats.updated,
|
||||||
|
"deleted": stats.deleted,
|
||||||
|
"unchanged": stats.unchanged,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@vector_app.command("index-schema")
|
||||||
|
def index_schema_cmd(
|
||||||
|
config: Path = CONFIG_OPT,
|
||||||
|
json_output: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."),
|
||||||
|
) -> None:
|
||||||
|
"""Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store."""
|
||||||
|
if json_output:
|
||||||
|
try:
|
||||||
|
payload = index_schema_data(config)
|
||||||
|
except Exception:
|
||||||
|
payload = {"status": "failed", "code": "schema_index_failed"}
|
||||||
|
typer.echo(__import__("json").dumps(payload, sort_keys=True, separators=(",", ":")))
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
typer.echo(__import__("json").dumps(payload, sort_keys=True, separators=(",", ":")))
|
||||||
|
return
|
||||||
|
payload = index_schema_data(config)
|
||||||
|
_print_stats(type("Stats", (), payload["counts"])())
|
||||||
|
|||||||
Reference in New Issue
Block a user