Files
ThothII/harness/tht/cli/schema_cmd.py
T

670 lines
25 KiB
Python

import logging
import warnings
from itertools import islice
from pathlib import Path
from typing import Literal, TypedDict
import typer
import yaml
from pydantic import ValidationError
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: # noqa: BLE001 - sampling adapter boundary
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 (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError) as exc:
logger.debug("Catalogo cache non leggibile: %s", exc)
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: # noqa: BLE001 - introspection CLI boundary
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 _load_schema_config(config: Path, *, suppress_legacy_warning: bool = False):
"""Load configuration, optionally hiding the legacy-key warning for JSON callers."""
if not suppress_legacy_warning:
return load_config(config)
with warnings.catch_warnings():
warnings.filterwarnings(
"ignore",
message=r"^DEPRECATION: legacy workspace resource keys are deprecated;",
category=FutureWarning,
)
return load_config(config)
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 (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
raise _MachineSchemaError("physical_schema_invalid") from None
def _staged_sql_files(inputs: list[Path] | None) -> list[Path]:
"""Collect staged SQL paths without traversing beyond the file-count bound."""
def candidates():
for item in inputs or []:
if item.is_file():
yield item
elif item.is_dir():
for path in item.rglob("*.sql"):
if path.is_file():
yield path
else:
raise _MachineSchemaError("staged_sql_invalid")
# islice consumes at most the sentinel (33rd) match; unlike a list-producing
# rglob this never enumerates an unbounded directory before rejecting it.
files = list(islice(candidates(), _MAX_STAGED_SQL_FILES + 1))
ordered = sorted(set(files), key=lambda path: path.resolve().as_posix())
if len(ordered) > _MAX_STAGED_SQL_FILES or len(files) > _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:
allowance = min(_MAX_STAGED_SQL_BYTES, _MAX_STAGED_SQL_TOTAL - total)
try:
with path.open("rb") as stream:
raw = stream.read(allowance + 1)
except (OSError, UnicodeError):
raise _MachineSchemaError("staged_sql_invalid") from None
if len(raw) > allowance:
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: Literal["succeeded", "failed"]
code: Literal["ok"]
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: Literal["succeeded", "failed"]
code: Literal["ok", "annotation_orphans"]
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,
suppress_legacy_warning: bool = False,
) -> 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
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineSchemaError("invalid_configuration") from None
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
if column_name in _GENERIC_PK_NAMES:
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:
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, *, suppress_legacy_warning: bool = False
) -> CheckSchemaResult:
"""Validate the physical catalog and imported annotations without writing."""
from tht.mschema.merge import find_orphans
from tht.mschema.models import Annotations
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineSchemaError("invalid_configuration") from None
physical = _physical_or_error(cfg)
try:
annotations = Annotations.from_yaml(annotations_path(cfg))
except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
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, suppress_legacy_warning=json_output)
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: # noqa: BLE001 - JSON CLI boundary
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)."), # noqa: B008
assume: list[str] = typer.Option(None, "--assume", help="Disambigua una PK: col=tabella."), # noqa: B008
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,
suppress_legacy_warning=json_output,
)
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: # noqa: BLE001 - JSON CLI boundary
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["staged_sql_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( # noqa: B008
None, "--table", "-t", help="Limita alle tabelle indicate (ripetibile)."
),
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
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)