# ruff: noqa: BLE001, S110, B008 import logging from pathlib import Path from typing import TypedDict import typer from tht.adapters.factory import build_dwh from tht.cli.config_cmd import CONFIG_OPT from tht.config import ConfigError, load_config from tht.db.sampling import is_text_type from tht.mschema.eligibility import classify_all schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)") logger = logging.getLogger(__name__) def _add_examples(dwh, phys, examples) -> None: for table_name, table in phys.tables.items(): for column_name, column in table.columns.items(): if not is_text_type(column.type): continue try: sampled = dwh.sample_column( table_name, column_name, limit=examples.max_per_column ) except Exception as exc: logger.warning("Campionamento saltato per %s.%s: %s", table_name, column_name, exc) continue column.examples = [str(value) for value in sampled if value not in (None, "")] def _load_config_or_exit(config: Path): try: return load_config(config) except ConfigError as e: typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) def physical_path(cfg) -> Path: from tht.jobs.dwh_pipeline import resolve_dwh_snapshot if not (cfg.paths.artifacts.parent / ".tht-dwh").exists(): return cfg.paths.artifacts / "mschema" / "physical.yaml" return resolve_dwh_snapshot(cfg).physical def annotations_path(cfg) -> Path: return cfg.paths.artifacts / "mschema" / "annotations.yaml" def refresh_catalog(cfg, *, dwh=None, output_path: Path | None = None): """Run the existing catalog algorithm and persist its canonical output.""" target = dwh if dwh is not None else build_dwh(cfg) physical = target.introspect() _add_examples(target, physical, cfg.examples) classify_all(physical, cfg.eligibility) physical.to_yaml(output_path or (cfg.paths.artifacts / "mschema" / "physical.yaml")) return physical @schema_app.command("introspect") def introspect_cmd( config: Path = CONFIG_OPT, refresh: bool = typer.Option( False, "--refresh", help="Forza la re-introspezione del DWH anche se physical.yaml esiste già.", ), ) -> None: """Introspeziona lo schema target e genera artifacts/mschema/physical.yaml. Se physical.yaml esiste già, esce subito (cache); usa --refresh per rigenerarlo. """ cfg = _load_config_or_exit(config) dwh_root = cfg.paths.artifacts.parent / ".tht-dwh" out = cfg.paths.artifacts / "mschema" / "physical.yaml" if dwh_root.exists() or dwh_root.is_symlink(): out = physical_path(cfg) if (dwh_root.exists() or dwh_root.is_symlink()) and out.exists() and not refresh: from datetime import UTC, datetime from tht.mschema.models import PhysicalSchema try: cached = PhysicalSchema.from_yaml(out) except Exception: pass # catalogo illeggibile: procedi con la re-introspezione else: ts = cached.introspected_at if ts.tzinfo is None: ts = ts.replace(tzinfo=UTC) age_days = (datetime.now(UTC) - ts).days typer.secho( f"OK (cache): {out} esistente ({len(cached.tables)} tabelle, " f"età {age_days}g). Re-introspezione solo con --refresh (manutenzione).", fg=typer.colors.GREEN, ) return try: from tht.cli.preprocess_cmd import run_dwh_from_config from tht.mschema.models import PhysicalSchema report = run_dwh_from_config(config, steps=("introspect",)) if report.status != "succeeded": raise RuntimeError("DWH preprocessing failed") out = physical_path(cfg) phys = PhysicalSchema.from_yaml(out) except Exception as e: 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()) n_ignored = sum( 1 for t in phys.tables.values() for c in t.columns.values() if not c.eligible ) typer.secho( f"OK: {len(phys.tables)} tabelle, {n_cols} colonne " f"({n_ignored} ignorate: testo ampio) -> {out}", fg=typer.colors.GREEN, ) class _MachineSchemaError(Exception): """An expected schema CLI failure with a stable public code.""" def __init__(self, code: str, detail: str = ""): self.code = code self.detail = detail super().__init__(code) _MAX_STAGED_SQL_BYTES = 1 << 20 _MAX_STAGED_SQL_TOTAL = 16 << 20 _MAX_STAGED_SQL_FILES = 32 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. ordered = sorted(set(files), key=lambda path: path.resolve().as_posix()) if len(ordered) > _MAX_STAGED_SQL_FILES: raise _MachineSchemaError("staged_sql_too_many") return ordered 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 class CandidateForeignKey(TypedDict): columns: list[str] ref_table: str ref_columns: list[str] class FKCandidate(TypedDict): table: str foreign_keys: list[CandidateForeignKey] class SuggestFksResult(TypedDict): status: str code: str candidates: list[FKCandidate] candidate_count: int candidate_digest: str orphan_count: int staged_sql_count: int mined_join_count: int ambiguous_columns: list[str] class CheckSchemaResult(TypedDict): status: str code: str orphan_count: int orphans: list[str] ignored: list[str] 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) -> SuggestFksResult: """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)) assumed: dict[str, str] = {} for value in assume or []: col, sep, ref = value.partition("=") if not sep or not col or ref not in physical.tables: raise _MachineSchemaError("assumption_invalid", value) assumed[col] = ref 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) -> CheckSchemaResult: """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) ignored = [ 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 ] return { "status": "succeeded" if not orphans else "failed", "code": "ok" if not orphans else "annotation_orphans", "orphan_count": len(orphans), "orphans": orphans, "ignored": ignored, } @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 _render_schema_machine_error(error, config) 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: # ``ignored`` is human diagnostic context, not part of the machine contract. _schema_json({key: payload[key] for key in ("status", "code", "orphan_count", "orphans")}) if payload["status"] != "succeeded": raise typer.Exit(code=3) return if payload["ignored"]: typer.secho( f"Colonne ignorate (testo ampio, {len(payload['ignored'])}):", fg=typer.colors.YELLOW, ) for line in payload["ignored"]: typer.echo(f" - {line}") 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) typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN) def _render_schema_machine_error(error: _MachineSchemaError, config: Path) -> None: """Render expected schema failures for the legacy human command contract.""" cfg = _load_config_or_exit(config) if error.code == "physical_schema_missing": typer.secho( f"ERRORE: {physical_path(cfg)} non trovato. Esegui prima `tht schema introspect`.", fg=typer.colors.RED, err=True, ) elif error.code == "assumption_invalid": detail = getattr(error, "detail", "") typer.secho( f"ERRORE: --assume '{detail}' non valido (atteso col=tabella nel catalogo).", fg=typer.colors.RED, err=True, ) else: typer.secho("ERRORE: impossibile elaborare il catalogo.", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) from None @schema_app.command("suggest-fks") def suggest_fks_cmd( config: Path = CONFIG_OPT, from_sql: list[Path] = typer.Option(None, "--from-sql", help="Directory/file SQL approvati (ripetibile)."), assume: list[str] = typer.Option(None, "--assume", help="Disambigua una PK: col=tabella."), 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."), ) -> None: """Suggerisce FK logiche per la curazione umana in annotations.yaml.""" import yaml as _yaml from tht.mschema.models import Annotations, TableAnnotation if json_output and write: _schema_json({"status": "failed", "code": "write_not_allowed"}) raise typer.Exit(code=2) try: payload = suggest_fks_data(config, from_sql=from_sql, assume=assume) except _MachineSchemaError as error: if json_output: _schema_json({"status": "failed", "code": error.code}) raise typer.Exit(code=1) from None _render_schema_machine_error(error, config) except Exception: if json_output: _schema_json({"status": "failed", "code": "schema_suggestion_failed"}) 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) from None if json_output: _schema_json(payload) return # Human mode remains the original renderer over the pure result. if payload["mined_join_count"]: typer.secho( f"Minati {payload['mined_join_count']} equi-join da {payload['staged_sql_count']} file SQL.", fg=typer.colors.BLUE, err=True, ) if payload["ambiguous_columns"]: typer.secho( "PK ambigue saltate dalla regola same-name (piu' tabelle proprietarie): " + ", ".join(payload["ambiguous_columns"]) + ". Se servono, aggiungile a mano o passa --from-sql.", fg=typer.colors.YELLOW, err=True, ) if not payload["candidates"]: typer.secho("OK: nessuna FK da suggerire.", fg=typer.colors.GREEN) return candidate_tables = { item["table"]: {"foreign_keys": item["foreign_keys"]} for item in payload["candidates"] } if write: cfg = _load_config_or_exit(config) annotations = Annotations.from_yaml(annotations_path(cfg)) for table_name, table_payload in candidate_tables.items(): 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( f"OK: {payload['candidate_count']} FK suggerite aggiunte a {annotations_path(cfg)} " f"({len(candidate_tables)} tabelle). Rivedile a mano prima dell'uso.", fg=typer.colors.GREEN, ) return typer.echo(_yaml.safe_dump({"tables": candidate_tables}, sort_keys=False, allow_unicode=True)) typer.secho( f"{payload['candidate_count']} FK candidate ({len(candidate_tables)} tabelle). " "Usa --write per fonderle in annotations.yaml, poi curale a mano.", 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") def render_cmd( config: Path = CONFIG_OPT, format: str = typer.Option( "markdown", "--format", "-f", help="Formato: markdown | mschema-text | schema-dict" ), tables: list[str] = typer.Option( None, "--table", "-t", help="Limita alle tabelle indicate (ripetibile)." ), output: Path = typer.Option(None, "--output", "-o", help="File di output (default stdout)."), ) -> None: """Serializza mschema (physical + annotations) nel formato richiesto.""" import json from tht.mschema.models import Annotations, PhysicalSchema from tht.mschema.render import to_markdown, to_mschema_text, to_schema_dict cfg = _load_config_or_exit(config) phys_file = physical_path(cfg) if not phys_file.exists(): typer.secho( f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", fg=typer.colors.RED, err=True, ) raise typer.Exit(code=1) physical = PhysicalSchema.from_yaml(phys_file) annotations = Annotations.from_yaml(annotations_path(cfg)) table_filter = list(tables) if tables else None if format == "markdown": out = to_markdown(physical, annotations) elif format == "mschema-text": out = to_mschema_text(physical, annotations, tables=table_filter) elif format == "schema-dict": out = json.dumps(to_schema_dict(physical, annotations), ensure_ascii=False, indent=2) else: typer.secho(f"ERRORE: formato sconosciuto: {format}", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) if output: output.parent.mkdir(parents=True, exist_ok=True) output.write_text(out) typer.secho(f"OK: scritto {output}", fg=typer.colors.GREEN) else: typer.echo(out) @schema_app.command("columns") def columns_cmd( table: str = typer.Argument(..., help="Nome tabella (chiave in physical.yaml)."), json_out: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."), config: Path = CONFIG_OPT, ) -> None: """Elenca nome/descrizione/tipo/pk delle colonne di una tabella dal catalogo.""" import json as _json from tht.mschema.models import PhysicalSchema cfg = _load_config_or_exit(config) phys_file = physical_path(cfg) if not phys_file.exists(): typer.secho( f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", fg=typer.colors.RED, err=True, ) raise typer.Exit(code=1) physical = PhysicalSchema.from_yaml(phys_file) def _payload(name, tbl): return { "table": name, "description": tbl.comment, "columns": [ {"name": n, "description": col.comment, "type": col.type, "pk": col.pk} for n, col in tbl.columns.items() ], } def _emit_human(p): typer.echo(f"{p['table']}: {p['description']}") for c in p["columns"]: typer.echo(f" {'*' if c['pk'] else ' '} {c['name']} ({c['type']}) — {c['description']}") # Glob-friendly: a pattern (containing * ? [) resolves to every matching catalog # table, so a model can ask for a whole family (e.g. fact_sost_impianto_*) in one # call instead of stalling on an unknown wildcard. Exact names keep the original # single-object contract; JSON for a pattern is an array of per-table objects. if any(ch in table for ch in "*?["): from fnmatch import fnmatch matches = sorted(n for n in physical.tables if fnmatch(n, table)) if not matches: typer.secho( f"ERRORE: nessuna tabella corrisponde al pattern: {table}", fg=typer.colors.RED, err=True, ) raise typer.Exit(code=1) payloads = [_payload(n, physical.tables[n]) for n in matches] if json_out: typer.echo(_json.dumps(payloads, ensure_ascii=False)) return for i, p in enumerate(payloads): if i: typer.echo("") _emit_human(p) return tbl = physical.tables.get(table) if tbl is None: # Aid recovery: suggest catalog tables that share the leading segment. prefix = table.rsplit("_", 1)[0] + "_" if "_" in table else table hints = sorted(n for n in physical.tables if n.startswith(prefix))[:12] msg = f"ERRORE: tabella non nel catalogo: {table}" if hints: msg += f" (forse: {', '.join(hints)})" typer.secho(msg, fg=typer.colors.RED, err=True) raise typer.Exit(code=1) payload = _payload(table, tbl) if json_out: typer.echo(_json.dumps(payload, ensure_ascii=False)) return _emit_human(payload)