feat: pristine harness JSON interfaces and require-existing semantic mode (P2)
This commit is contained in:
@@ -2,27 +2,57 @@
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import re
|
||||
import hashlib
|
||||
from pathlib import Path
|
||||
|
||||
import typer
|
||||
|
||||
from tht.cli.config_cmd import CONFIG_OPT
|
||||
from tht.config import workspace_id_from_path
|
||||
|
||||
|
||||
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
|
||||
|
||||
|
||||
def _evidence_json_context(config: Path):
|
||||
from tht.cli.schema_cmd import _load_config_or_exit
|
||||
|
||||
return _load_config_or_exit(config)
|
||||
|
||||
|
||||
def _evidence_json_payload(cfg, payload: dict, *, code: str, error: str | None = None) -> dict:
|
||||
value = {
|
||||
**payload,
|
||||
"schemaVersion": 1,
|
||||
"status": payload.get("status", "failed"),
|
||||
"code": code,
|
||||
"operation": "preprocess_evidence",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
}
|
||||
if error is not None:
|
||||
value["error"] = error
|
||||
return value
|
||||
|
||||
|
||||
def _simple_json_payload(*, code: str, error: str) -> dict:
|
||||
return {
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"code": code,
|
||||
"operation": "preprocess_evidence",
|
||||
"error": error,
|
||||
}
|
||||
|
||||
|
||||
def run_dwh_from_config(
|
||||
config: Path, *, steps: tuple[str, ...], resume: str | None = None,
|
||||
):
|
||||
from tht.cli.lsh_cmd import build_lsh_artifacts
|
||||
from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog
|
||||
from tht.jobs.dwh_pipeline import (
|
||||
DwhPreprocessPipeline, config_dwh_binding,
|
||||
DwhPreprocessPipeline,
|
||||
config_dwh_binding,
|
||||
)
|
||||
|
||||
cfg = _load_config_or_exit(config)
|
||||
@@ -87,7 +117,7 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None =
|
||||
return "sha256:" + hashlib.sha256(value.encode()).hexdigest()
|
||||
|
||||
return pipeline.run_as_job(
|
||||
workspace_id=workspace_id_from_path(config),
|
||||
workspace_id=cfg._workspace_id,
|
||||
workspace_root=corpus_root.parent,
|
||||
config_fingerprint=fingerprint(cfg.model_dump_json()),
|
||||
input_fingerprint=fingerprint(config.resolve().as_posix()),
|
||||
@@ -116,7 +146,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False):
|
||||
pipeline_version="evidence-v1",
|
||||
retain_published_generations=cfg.vector.retain_published_generations,
|
||||
)
|
||||
pipeline.workspace_id = workspace_id_from_path(config)
|
||||
pipeline.workspace_id = cfg._workspace_id
|
||||
return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run)
|
||||
|
||||
|
||||
@@ -133,7 +163,7 @@ def evidence_cmd(
|
||||
if action == "gc":
|
||||
try:
|
||||
payload = gc_from_config(config, dry_run=dry_run)
|
||||
except Exception:
|
||||
except Exception: # noqa: BLE001
|
||||
payload = {"status": "failed", "error": "evidence cleanup failed"}
|
||||
if json_output:
|
||||
typer.echo(json.dumps(payload, sort_keys=True))
|
||||
@@ -145,8 +175,12 @@ def evidence_cmd(
|
||||
else:
|
||||
typer.echo(f"OK: evicted={len(payload['evicted'])} failures={len(payload['failures'])}")
|
||||
return
|
||||
cfg = _evidence_json_context(config) if json_output else None
|
||||
if resume is not None and re.fullmatch(r"[0-9a-f]{32}", resume) is None:
|
||||
payload = {"status": "failed", "error": "resume requires a preprocessing run id"}
|
||||
payload = _simple_json_payload(
|
||||
code="invalid_resume",
|
||||
error="resume requires a preprocessing run id",
|
||||
)
|
||||
if json_output:
|
||||
typer.echo(json.dumps(payload, sort_keys=True))
|
||||
else:
|
||||
@@ -154,23 +188,43 @@ def evidence_cmd(
|
||||
raise typer.Exit(code=2)
|
||||
try:
|
||||
result = run_from_config(config, dry_run=dry_run, resume=resume)
|
||||
except Exception:
|
||||
payload = {"status": "failed", "error": "preprocessing failed"}
|
||||
except Exception: # noqa: BLE001
|
||||
payload = {"status": "failed"}
|
||||
if json_output:
|
||||
typer.echo(json.dumps(payload, sort_keys=True))
|
||||
typer.echo(json.dumps(
|
||||
_evidence_json_payload(
|
||||
cfg,
|
||||
payload,
|
||||
code="preprocessing_failed",
|
||||
error="preprocessing failed",
|
||||
),
|
||||
sort_keys=True,
|
||||
))
|
||||
else:
|
||||
typer.secho("ERRORE: preprocessing failed", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1) from None
|
||||
payload = result.model_dump(mode="json")
|
||||
if payload.get("status") != "succeeded":
|
||||
payload["error"] = "preprocessing job failed"
|
||||
if json_output:
|
||||
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
||||
typer.echo(json.dumps(
|
||||
_evidence_json_payload(
|
||||
cfg,
|
||||
payload,
|
||||
code="preprocessing_failed",
|
||||
error="preprocessing job failed",
|
||||
),
|
||||
ensure_ascii=False,
|
||||
sort_keys=True,
|
||||
))
|
||||
else:
|
||||
typer.secho("ERRORE: preprocessing job failed", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1)
|
||||
if json_output:
|
||||
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
||||
typer.echo(json.dumps(
|
||||
_evidence_json_payload(cfg, payload, code="ok"),
|
||||
ensure_ascii=False,
|
||||
sort_keys=True,
|
||||
))
|
||||
else:
|
||||
counts = payload["counts"]
|
||||
typer.echo(
|
||||
@@ -205,7 +259,7 @@ def dwh_cmd(
|
||||
raise typer.Exit(code=2)
|
||||
try:
|
||||
result = run_dwh_from_config(config, steps=selected, resume=resume)
|
||||
except Exception:
|
||||
except Exception: # noqa: BLE001
|
||||
payload = {"status": "failed", "error": "DWH preprocessing failed"}
|
||||
if json_output:
|
||||
typer.echo(json.dumps(payload, sort_keys=True))
|
||||
|
||||
+344
-130
@@ -1,7 +1,11 @@
|
||||
from pathlib import Path
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
from pathlib import Path
|
||||
|
||||
import typer
|
||||
import yaml
|
||||
|
||||
from tht.adapters.factory import build_dwh
|
||||
from tht.cli.config_cmd import CONFIG_OPT
|
||||
from tht.config import ConfigError, load_config
|
||||
@@ -21,7 +25,7 @@ def _add_examples(dwh, phys, examples) -> None:
|
||||
sampled = dwh.sample_column(
|
||||
table_name, column_name, limit=examples.max_per_column
|
||||
)
|
||||
except Exception as exc:
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.warning("Campionamento saltato per %s.%s: %s",
|
||||
table_name, column_name, exc)
|
||||
continue
|
||||
@@ -83,7 +87,7 @@ def introspect_cmd(
|
||||
|
||||
try:
|
||||
cached = PhysicalSchema.from_yaml(out)
|
||||
except Exception:
|
||||
except Exception: # noqa: BLE001,S110
|
||||
pass # catalogo illeggibile: procedi con la re-introspezione
|
||||
else:
|
||||
ts = cached.introspected_at
|
||||
@@ -105,7 +109,7 @@ def introspect_cmd(
|
||||
raise RuntimeError("DWH preprocessing failed")
|
||||
out = physical_path(cfg)
|
||||
phys = PhysicalSchema.from_yaml(out)
|
||||
except Exception as e:
|
||||
except Exception as e: # noqa: BLE001
|
||||
typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1)
|
||||
n_cols = sum(len(t.columns) for t in phys.tables.values())
|
||||
@@ -119,8 +123,188 @@ def introspect_cmd(
|
||||
)
|
||||
|
||||
|
||||
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 _emit_json(payload: dict) -> None:
|
||||
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
||||
|
||||
|
||||
def _sorted_fk_payloads(foreign_keys) -> list[dict]:
|
||||
payloads = [fk.model_dump(mode="json", exclude_defaults=True) for fk in foreign_keys]
|
||||
return sorted(
|
||||
payloads,
|
||||
key=lambda payload: (
|
||||
tuple(payload.get("columns", [])),
|
||||
payload.get("ref_table", ""),
|
||||
tuple(payload.get("ref_columns", [])),
|
||||
payload.get("name", ""),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _sorted_annotations_payload(annotations) -> dict:
|
||||
tables = {}
|
||||
for table_name in sorted(annotations.tables):
|
||||
table = annotations.tables[table_name]
|
||||
payload = {}
|
||||
if table.description:
|
||||
payload["description"] = table.description
|
||||
if table.concepts:
|
||||
payload["concepts"] = table.concepts
|
||||
if table.notes:
|
||||
payload["notes"] = table.notes
|
||||
if table.columns:
|
||||
payload["columns"] = {
|
||||
name: value.model_dump(mode="json", exclude_defaults=True)
|
||||
for name, value in sorted(table.columns.items())
|
||||
}
|
||||
if table.foreign_keys:
|
||||
payload["foreign_keys"] = _sorted_fk_payloads(table.foreign_keys)
|
||||
tables[table_name] = payload
|
||||
return {"tables": tables}
|
||||
|
||||
|
||||
def _suggested_fk_payload(annotations_by_table: dict) -> dict:
|
||||
return {
|
||||
"tables": {
|
||||
table_name: {"foreign_keys": _sorted_fk_payloads(foreign_keys)}
|
||||
for table_name, foreign_keys in sorted(annotations_by_table.items())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
def _load_sql_inputs(entries: list[Path] | None) -> list[tuple[str, str]]:
|
||||
max_file_bytes = 1024 * 1024
|
||||
max_total_bytes = 16 * 1024 * 1024
|
||||
total_bytes = 0
|
||||
sql_files: list[Path] = []
|
||||
for entry in entries or []:
|
||||
if entry.is_dir():
|
||||
sql_files.extend(sorted(path for path in entry.rglob("*.sql") if path.is_file()))
|
||||
continue
|
||||
sql_files.append(entry)
|
||||
loaded = []
|
||||
for sql_file in sorted(sql_files, key=lambda candidate: candidate.as_posix()):
|
||||
if not sql_file.exists() or not sql_file.is_file() or sql_file.is_symlink():
|
||||
raise ValueError("invalid SQL input")
|
||||
size = sql_file.stat().st_size
|
||||
total_bytes += size
|
||||
if size > max_file_bytes or total_bytes > max_total_bytes:
|
||||
raise ValueError("invalid SQL input")
|
||||
loaded.append((sql_file.as_posix(), sql_file.read_text(encoding="utf-8")))
|
||||
return loaded
|
||||
|
||||
|
||||
def _suggest_fk_result(physical, annotations, *, sql_inputs: list[tuple[str, str]], assume: list[str] | None):
|
||||
from tht.mschema.fkmine import mine_join_pairs
|
||||
from tht.mschema.models import ForeignKey
|
||||
|
||||
assumed: dict[str, str] = {}
|
||||
for value in assume or []:
|
||||
col, _, ref = value.partition("=")
|
||||
if not ref or ref not in physical.tables:
|
||||
raise ValueError("invalid assume mapping")
|
||||
assumed[col] = ref
|
||||
|
||||
def _single_pk(table) -> str | None:
|
||||
pks = [column_name for column_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, table in physical.tables.items():
|
||||
pk = _single_pk(table)
|
||||
if pk:
|
||||
pk_owners.setdefault(pk, []).append(table_name)
|
||||
|
||||
dim_time_pk = None
|
||||
if "dim_time" in physical.tables:
|
||||
dim_time_pk = _single_pk(physical.tables["dim_time"])
|
||||
|
||||
def _known(table_name: str) -> set:
|
||||
keys = set()
|
||||
for fk in physical.tables[table_name].foreign_keys:
|
||||
keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns)))
|
||||
annotation = annotations.tables.get(table_name)
|
||||
if annotation:
|
||||
for fk in annotation.foreign_keys:
|
||||
keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns)))
|
||||
return keys
|
||||
|
||||
known_by_table: dict[str, set] = {table_name: _known(table_name) for table_name in physical.tables}
|
||||
suggested: dict[str, list[ForeignKey]] = {}
|
||||
|
||||
def _add(table_name: str, column_name: str, ref_table: str, ref_column: str) -> None:
|
||||
key = ((column_name,), ref_table, (ref_column,))
|
||||
if key in known_by_table[table_name]:
|
||||
return
|
||||
known_by_table[table_name].add(key)
|
||||
suggested.setdefault(table_name, []).append(
|
||||
ForeignKey(columns=[column_name], ref_table=ref_table, ref_columns=[ref_column])
|
||||
)
|
||||
|
||||
mined_total = 0
|
||||
for _name, sql_text in sql_inputs:
|
||||
pairs = mine_join_pairs(sql_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)
|
||||
|
||||
ambiguous_skipped: set[str] = set()
|
||||
for table_name, table in physical.tables.items():
|
||||
for column_name in 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 not owners or column_name in _GENERIC_PK_NAMES:
|
||||
continue
|
||||
if len(pk_owners[column_name]) > 1:
|
||||
ambiguous_skipped.add(column_name)
|
||||
continue
|
||||
_add(table_name, column_name, owners[0], column_name)
|
||||
|
||||
candidate_annotations = _suggested_fk_payload(suggested)
|
||||
candidate_yaml = yaml.safe_dump(candidate_annotations, sort_keys=False, allow_unicode=True)
|
||||
counts = {
|
||||
"ambiguousColumns": len(ambiguous_skipped),
|
||||
"candidateTables": len(candidate_annotations["tables"]),
|
||||
"candidates": sum(len(value["foreign_keys"]) for value in candidate_annotations["tables"].values()),
|
||||
"minedJoins": mined_total,
|
||||
"sqlFiles": len(sql_inputs),
|
||||
}
|
||||
candidate_document = {
|
||||
"annotations": candidate_annotations,
|
||||
"counts": {
|
||||
"candidateTables": counts["candidateTables"],
|
||||
"candidates": counts["candidates"],
|
||||
},
|
||||
"schemaVersion": 1,
|
||||
}
|
||||
return {
|
||||
"ambiguous": sorted(ambiguous_skipped),
|
||||
"candidate_count": counts["candidates"],
|
||||
"candidateDigest": "sha256:" + hashlib.sha256(candidate_yaml.encode("utf-8")).hexdigest(),
|
||||
"candidateDocument": candidate_document,
|
||||
"candidate_yaml": candidate_yaml,
|
||||
"counts": counts,
|
||||
"suggested": suggested,
|
||||
}
|
||||
|
||||
|
||||
@schema_app.command("check")
|
||||
def check_cmd(config: Path = CONFIG_OPT) -> None:
|
||||
def check_cmd(
|
||||
config: Path = CONFIG_OPT,
|
||||
annotations: Path | None = typer.Option(None, "--annotations"), # noqa: B008
|
||||
reviewed_candidates: str | None = typer.Option(None, "--reviewed-candidates"),
|
||||
json_output: bool = typer.Option(False, "--json"),
|
||||
) -> None:
|
||||
"""Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane."""
|
||||
from tht.mschema.merge import find_orphans
|
||||
from tht.mschema.models import Annotations, PhysicalSchema
|
||||
@@ -128,20 +312,80 @@ def check_cmd(config: Path = CONFIG_OPT) -> None:
|
||||
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,
|
||||
)
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": "schema_missing",
|
||||
"error": "physical schema is missing",
|
||||
"operation": "schema_check",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
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)
|
||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||
annotations_file = annotations or annotations_path(cfg)
|
||||
try:
|
||||
loaded_annotations = Annotations.from_yaml(annotations_file)
|
||||
except Exception: # noqa: BLE001
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": "annotation_invalid",
|
||||
"error": "annotations are invalid",
|
||||
"operation": "schema_check",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
typer.secho("ERRORE: annotations non valide.", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1) from None
|
||||
|
||||
ignored = [
|
||||
f"{t}.{c} ({col.eligibility_reason})"
|
||||
for t, table in physical.tables.items()
|
||||
for c, col in table.columns.items()
|
||||
if not col.eligible
|
||||
f"{table_name}.{column_name} ({column.eligibility_reason})"
|
||||
for table_name, table in physical.tables.items()
|
||||
for column_name, column in table.columns.items()
|
||||
if not column.eligible
|
||||
]
|
||||
orphans = sorted(find_orphans(physical, loaded_annotations))
|
||||
if json_output:
|
||||
annotations_payload = _sorted_annotations_payload(loaded_annotations)
|
||||
payload = {
|
||||
"annotationsDigest": _json_sha({"annotations": annotations_payload, "schemaVersion": 1}),
|
||||
"code": "ok" if not orphans else "annotation_invalid",
|
||||
"counts": {
|
||||
"annotationTables": len(annotations_payload["tables"]),
|
||||
"foreignKeys": sum(
|
||||
len(table_payload.get("foreign_keys", []))
|
||||
for table_payload in annotations_payload["tables"].values()
|
||||
),
|
||||
"orphans": len(orphans),
|
||||
},
|
||||
"operation": "schema_check",
|
||||
"orphan_count": len(orphans),
|
||||
"orphans": orphans,
|
||||
"schemaVersion": 1,
|
||||
"status": "succeeded" if not orphans else "blocked",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
"zeroOrphans": not orphans,
|
||||
}
|
||||
payload["annotations_digest"] = "sha256:" + hashlib.sha256(Path(annotations_file).read_bytes()).hexdigest()
|
||||
if reviewed_candidates is not None:
|
||||
payload["reviewedCandidates"] = reviewed_candidates
|
||||
payload["reviewed_candidates_digest"] = reviewed_candidates
|
||||
_emit_json(payload)
|
||||
if orphans:
|
||||
raise typer.Exit(code=3)
|
||||
return
|
||||
|
||||
if ignored:
|
||||
typer.secho(
|
||||
f"Colonne ignorate (testo ampio, {len(ignored)}):", fg=typer.colors.YELLOW
|
||||
@@ -149,11 +393,10 @@ def check_cmd(config: Path = CONFIG_OPT) -> None:
|
||||
for line in ignored:
|
||||
typer.echo(f" - {line}")
|
||||
|
||||
orphans = find_orphans(physical, annotations)
|
||||
if orphans:
|
||||
typer.secho(f"ATTENZIONE: {len(orphans)} annotazioni orfane:", fg=typer.colors.YELLOW)
|
||||
for o in orphans:
|
||||
typer.echo(f" - {o}")
|
||||
for orphan in orphans:
|
||||
typer.echo(f" - {orphan}")
|
||||
raise typer.Exit(code=3)
|
||||
typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN)
|
||||
|
||||
@@ -166,11 +409,11 @@ _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(
|
||||
from_sql: list[Path] = typer.Option( # noqa: B008
|
||||
None, "--from-sql",
|
||||
help="Directory di .sql approvati da cui minare i join reali (ripetibile).",
|
||||
help="Directory o file .sql approvati da cui minare i join reali (ripetibile).",
|
||||
),
|
||||
assume: list[str] = typer.Option(
|
||||
assume: list[str] = typer.Option( # noqa: B008
|
||||
None, "--assume",
|
||||
help="Disambigua una PK con piu' proprietari: col=tabella_ref "
|
||||
"(es. cod_paz=dim_patient). Ripetibile.",
|
||||
@@ -179,132 +422,103 @@ def suggest_fks_cmd(
|
||||
False, "--write",
|
||||
help="Fonde i suggerimenti in annotations.yaml (aggiunge solo FK mancanti).",
|
||||
),
|
||||
json_output: bool = typer.Option(False, "--json"),
|
||||
) -> 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.
|
||||
"""
|
||||
"""Suggerisce FK logiche per la curazione umana in annotations.yaml."""
|
||||
import yaml as _yaml
|
||||
|
||||
from tht.mschema.fkmine import mine_join_pairs
|
||||
from tht.mschema.models import Annotations, ForeignKey, PhysicalSchema, TableAnnotation
|
||||
from tht.mschema.models import Annotations, 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,
|
||||
)
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": "schema_missing",
|
||||
"error": "physical schema is missing",
|
||||
"operation": "schema_suggest_fks",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
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,
|
||||
loaded_annotations = Annotations.from_yaml(ann_path)
|
||||
try:
|
||||
sql_inputs = _load_sql_inputs(from_sql)
|
||||
result = _suggest_fk_result(physical, loaded_annotations, sql_inputs=sql_inputs, assume=assume)
|
||||
except ValueError as exc:
|
||||
code = "invalid_argument"
|
||||
error = str(exc)
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": code,
|
||||
"error": error,
|
||||
"operation": "schema_suggest_fks",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
human_error = (
|
||||
"--assume non valido (atteso col=tabella nel catalogo)."
|
||||
if error == "invalid assume mapping"
|
||||
else error
|
||||
)
|
||||
raise typer.Exit(code=1)
|
||||
assumed[col] = ref
|
||||
typer.secho(f"ERRORE: {human_error}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1) from None
|
||||
|
||||
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
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"candidate_count": result["candidate_count"],
|
||||
"candidateDigest": result["candidateDigest"],
|
||||
"candidateDocument": result["candidateDocument"],
|
||||
"candidate_yaml": result["candidate_yaml"],
|
||||
"code": "ok",
|
||||
"counts": result["counts"],
|
||||
"operation": "schema_suggest_fks",
|
||||
"schemaVersion": 1,
|
||||
"status": "succeeded",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
return
|
||||
|
||||
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:
|
||||
if result["counts"]["sqlFiles"]:
|
||||
typer.secho(
|
||||
f"Minati {mined_total} equi-join da {n_sql_files} file SQL.",
|
||||
fg=typer.colors.BLUE, err=True,
|
||||
f"Minati {result['counts']['minedJoins']} equi-join da {result['counts']['sqlFiles']} file SQL.",
|
||||
fg=typer.colors.BLUE,
|
||||
err=True,
|
||||
)
|
||||
if ambiguous_skipped:
|
||||
if result["ambiguous"]:
|
||||
typer.secho(
|
||||
"PK ambigue saltate dalla regola same-name (piu' tabelle proprietarie): "
|
||||
+ ", ".join(sorted(ambiguous_skipped))
|
||||
+ ", ".join(result["ambiguous"])
|
||||
+ ". Se servono, aggiungile a mano o passa --from-sql.",
|
||||
fg=typer.colors.YELLOW, err=True,
|
||||
fg=typer.colors.YELLOW,
|
||||
err=True,
|
||||
)
|
||||
|
||||
n_fks = sum(len(v) for v in suggested.values())
|
||||
suggested = result["suggested"]
|
||||
n_fks = result["counts"]["candidates"]
|
||||
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)
|
||||
for table_name, foreign_keys in suggested.items():
|
||||
annotation = loaded_annotations.tables.setdefault(table_name, TableAnnotation())
|
||||
annotation.foreign_keys.extend(foreign_keys)
|
||||
loaded_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.",
|
||||
@@ -312,13 +526,13 @@ def suggest_fks_cmd(
|
||||
)
|
||||
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.echo(
|
||||
_yaml.safe_dump(
|
||||
result["candidateDocument"]["annotations"],
|
||||
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.",
|
||||
@@ -332,10 +546,10 @@ def render_cmd(
|
||||
format: str = typer.Option(
|
||||
"markdown", "--format", "-f", help="Formato: markdown | mschema-text | schema-dict"
|
||||
),
|
||||
tables: list[str] = typer.Option(
|
||||
tables: list[str] = typer.Option( # noqa: B008
|
||||
None, "--table", "-t", help="Limita alle tabelle indicate (ripetibile)."
|
||||
),
|
||||
output: Path = typer.Option(None, "--output", "-o", help="File di output (default stdout)."),
|
||||
output: Path = typer.Option(None, "--output", "-o", help="File di output (default stdout)."), # noqa: B008
|
||||
) -> None:
|
||||
"""Serializza mschema (physical + annotations) nel formato richiesto."""
|
||||
import json
|
||||
|
||||
@@ -1,3 +1,5 @@
|
||||
import hashlib
|
||||
import json
|
||||
from pathlib import Path
|
||||
|
||||
import typer
|
||||
@@ -11,6 +13,14 @@ from tht.vectorstore.store import SyncStats, content_hash
|
||||
vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)")
|
||||
|
||||
|
||||
def _artifact_digest(path: Path) -> str:
|
||||
return "sha256:" + hashlib.sha256(path.read_bytes()).hexdigest()
|
||||
|
||||
|
||||
def _emit_json(payload: dict) -> None:
|
||||
typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True))
|
||||
|
||||
|
||||
def make_embedder(embeddings_cfg):
|
||||
"""Factory del client embeddings (monkeypatchabile nei test)."""
|
||||
from tht.vectorstore.embeddings import OllamaEmbeddings
|
||||
@@ -117,9 +127,14 @@ def init_cmd(
|
||||
|
||||
|
||||
@vector_app.command("index-schema")
|
||||
def index_schema_cmd(config: Path = CONFIG_OPT) -> None:
|
||||
def index_schema_cmd(
|
||||
config: Path = CONFIG_OPT,
|
||||
json_output: bool = typer.Option(False, "--json"),
|
||||
) -> None:
|
||||
"""Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store."""
|
||||
from tht.adapters.factory import build_vector_store
|
||||
from tht.mschema.models import Annotations, PhysicalSchema
|
||||
from tht.ports.vector import VectorStoreError
|
||||
from tht.vectorstore.records import schema_records
|
||||
|
||||
cfg = _load_config_or_exit(config)
|
||||
@@ -127,20 +142,69 @@ def index_schema_cmd(config: Path = CONFIG_OPT) -> None:
|
||||
require_vector_cfg(cfg)
|
||||
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,
|
||||
)
|
||||
message = f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`."
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": "schema_missing",
|
||||
"error": "physical schema is missing",
|
||||
"operation": "index_schema",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
typer.secho(message, fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1)
|
||||
physical = PhysicalSchema.from_yaml(phys_file)
|
||||
annotations = Annotations.from_yaml(annotations_path(cfg))
|
||||
annotations_file = annotations_path(cfg)
|
||||
annotations = Annotations.from_yaml(annotations_file)
|
||||
records = schema_records(physical, annotations)
|
||||
from tht.adapters.factory import build_vector_store
|
||||
|
||||
stats = sync_canonical_records(
|
||||
"schema_records",
|
||||
records,
|
||||
store=build_vector_store(cfg, require_write=True),
|
||||
embedder=make_embedder(cfg.embeddings),
|
||||
)
|
||||
try:
|
||||
stats = sync_canonical_records(
|
||||
"schema_records",
|
||||
records,
|
||||
store=build_vector_store(cfg, require_write=True),
|
||||
embedder=make_embedder(cfg.embeddings),
|
||||
)
|
||||
except VectorStoreError as exc:
|
||||
code = str(exc)
|
||||
error = "semantic index incompatible" if code == "semantic_index_incompatible" else "schema indexing failed"
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"code": code,
|
||||
"error": error,
|
||||
"operation": "index_schema",
|
||||
"schemaVersion": 1,
|
||||
"status": "failed",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
else:
|
||||
typer.secho(f"ERRORE: {error}", fg=typer.colors.RED, err=True)
|
||||
raise typer.Exit(code=1) from None
|
||||
if json_output:
|
||||
_emit_json({
|
||||
"artifactIdentities": [
|
||||
{"digest": _artifact_digest(annotations_file), "kind": "schema_annotations"},
|
||||
{"digest": _artifact_digest(phys_file), "kind": "physical_schema"},
|
||||
],
|
||||
"code": "ok",
|
||||
"collection": cfg.vectors.collection,
|
||||
"counts": {
|
||||
"added": stats.added,
|
||||
"columns": sum(len(table.columns) for table in physical.tables.values()),
|
||||
"deleted": stats.deleted,
|
||||
"records": len(records),
|
||||
"tables": len(physical.tables),
|
||||
"unchanged": stats.unchanged,
|
||||
"updated": stats.updated,
|
||||
},
|
||||
"operation": "index_schema",
|
||||
"schemaVersion": 1,
|
||||
"status": "succeeded",
|
||||
"workspaceId": cfg._workspace_id,
|
||||
"workspaceRevision": cfg._workspace_revision,
|
||||
})
|
||||
return
|
||||
_print_stats(stats)
|
||||
|
||||
Reference in New Issue
Block a user