681 lines
25 KiB
Python
681 lines
25 KiB
Python
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
|
|
from tht.db.sampling import is_text_type
|
|
from tht.mschema.context import (
|
|
SchemaContextError,
|
|
annotations_path,
|
|
load_schema_context,
|
|
physical_path,
|
|
)
|
|
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
|
|
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 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: # noqa: BLE001,S110
|
|
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: # 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())
|
|
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,
|
|
)
|
|
|
|
|
|
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,
|
|
annotations: Path | None = typer.Option(None, "--annotations"),
|
|
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
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
phys_file = physical_path(cfg)
|
|
if not phys_file.exists():
|
|
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_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"{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
|
|
)
|
|
for line in ignored:
|
|
typer.echo(f" - {line}")
|
|
|
|
if orphans:
|
|
typer.secho(f"ATTENZIONE: {len(orphans)} annotazioni orfane:", fg=typer.colors.YELLOW)
|
|
for orphan in orphans:
|
|
typer.echo(f" - {orphan}")
|
|
raise typer.Exit(code=3)
|
|
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")
|
|
def suggest_fks_cmd(
|
|
config: Path = CONFIG_OPT,
|
|
from_sql: list[Path] = typer.Option(
|
|
None, "--from-sql",
|
|
help="Directory o file .sql approvati da cui minare i join reali (ripetibile).",
|
|
),
|
|
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).",
|
|
),
|
|
json_output: bool = typer.Option(False, "--json"),
|
|
) -> None:
|
|
"""Suggerisce FK logiche per la curazione umana in annotations.yaml."""
|
|
import yaml as _yaml
|
|
|
|
from tht.mschema.models import Annotations, PhysicalSchema, TableAnnotation
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
if write and cfg.paths.effective_relationships is not None:
|
|
error = "effective relationships are managed by the catalog; --write is disabled"
|
|
if json_output:
|
|
_emit_json({
|
|
"code": "catalog_relationships_managed",
|
|
"error": error,
|
|
"operation": "schema_suggest_fks",
|
|
"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)
|
|
phys_file = physical_path(cfg)
|
|
if not phys_file.exists():
|
|
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)
|
|
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
|
|
)
|
|
typer.secho(f"ERRORE: {human_error}", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from 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
|
|
|
|
if result["counts"]["sqlFiles"]:
|
|
typer.secho(
|
|
f"Minati {result['counts']['minedJoins']} equi-join da {result['counts']['sqlFiles']} file SQL.",
|
|
fg=typer.colors.BLUE,
|
|
err=True,
|
|
)
|
|
if result["ambiguous"]:
|
|
typer.secho(
|
|
"PK ambigue saltate dalla regola same-name (piu' tabelle proprietarie): "
|
|
+ ", ".join(result["ambiguous"])
|
|
+ ". Se servono, aggiungile a mano o passa --from-sql.",
|
|
fg=typer.colors.YELLOW,
|
|
err=True,
|
|
)
|
|
|
|
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 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.",
|
|
fg=typer.colors.GREEN,
|
|
)
|
|
return
|
|
|
|
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.",
|
|
fg=typer.colors.YELLOW,
|
|
)
|
|
|
|
|
|
@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:
|
|
"""Serialize the current PostgreSQL Catalog projection."""
|
|
import json
|
|
|
|
from tht.mschema.render import to_markdown, to_mschema_text, to_schema_dict
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
try:
|
|
context = load_schema_context(cfg)
|
|
except SchemaContextError as exc:
|
|
typer.secho(f"ERRORE: {exc}", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1) from None
|
|
table_filter = list(tables) if tables else None
|
|
|
|
if format == "markdown":
|
|
out = to_markdown(
|
|
context.physical,
|
|
context.annotations,
|
|
effective_relationships=context.effective_relationships,
|
|
)
|
|
elif format == "mschema-text":
|
|
out = to_mschema_text(
|
|
context.physical,
|
|
context.annotations,
|
|
tables=table_filter,
|
|
effective_relationships=context.effective_relationships,
|
|
)
|
|
elif format == "schema-dict":
|
|
out = json.dumps(
|
|
to_schema_dict(
|
|
context.physical,
|
|
context.annotations,
|
|
effective_relationships=context.effective_relationships,
|
|
),
|
|
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
|
|
|
|
cfg = _load_config_or_exit(config)
|
|
try:
|
|
physical = load_schema_context(cfg).physical
|
|
except SchemaContextError as exc:
|
|
typer.secho(f"ERRORE: {exc}", fg=typer.colors.RED, err=True)
|
|
raise typer.Exit(code=1)
|
|
|
|
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)
|