fix: preserve Task 2 command loading contracts

This commit is contained in:
2026-08-11 06:31:12 +02:00
parent 9e2c2b37b3
commit f62b4dbc4c
5 changed files with 279 additions and 86 deletions
+3 -5
View File
@@ -1,5 +1,4 @@
"""One-shot preprocessing commands."""
# ruff: noqa: BLE001
from __future__ import annotations
@@ -99,7 +98,6 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None =
def gc_from_config(config: Path, *, dry_run: bool = False):
from tht.adapters.factory import build_evidence_sources, build_vector_store
from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli.vector_cmd import make_embedder
from tht.corpus.chunk import ChunkPolicy
from tht.corpus.pipeline import CorpusPipeline
@@ -134,7 +132,7 @@ def evidence_cmd(
if action == "gc":
try:
payload = gc_from_config(config, dry_run=dry_run)
except Exception:
except (OSError, RuntimeError, ValueError, TypeError, KeyError):
payload = {"status": "failed", "error": "evidence cleanup failed"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
@@ -155,7 +153,7 @@ def evidence_cmd(
raise typer.Exit(code=2)
try:
result = run_from_config(config, dry_run=dry_run, resume=resume)
except Exception:
except (OSError, RuntimeError, ValueError, TypeError, KeyError):
payload = {"status": "failed", "error": "preprocessing failed"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
@@ -206,7 +204,7 @@ def dwh_cmd(
raise typer.Exit(code=2)
try:
result = run_dwh_from_config(config, steps=selected, resume=resume)
except Exception:
except (OSError, RuntimeError, ValueError, TypeError, KeyError):
payload = {"status": "failed", "error": "DWH preprocessing failed"}
if json_output:
typer.echo(json.dumps(payload, sort_keys=True))
+97 -56
View File
@@ -1,6 +1,5 @@
import logging
import warnings
from itertools import islice
from pathlib import Path
from typing import Literal, TypedDict
@@ -10,7 +9,7 @@ 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.config import Config, ConfigError, load_config
from tht.db.sampling import is_text_type
from tht.mschema.eligibility import classify_all
@@ -42,7 +41,7 @@ def _load_config_or_exit(config: Path):
raise typer.Exit(code=1)
def physical_path(cfg) -> Path:
def physical_path(cfg: Config) -> Path:
from tht.jobs.dwh_pipeline import resolve_dwh_snapshot
if not (cfg.paths.artifacts.parent / ".tht-dwh").exists():
@@ -50,7 +49,7 @@ def physical_path(cfg) -> Path:
return resolve_dwh_snapshot(cfg).physical
def annotations_path(cfg) -> Path:
def annotations_path(cfg: Config) -> Path:
return cfg.paths.artifacts / "mschema" / "annotations.yaml"
@@ -134,6 +133,15 @@ class _MachineSchemaError(Exception):
super().__init__(code)
def _annotations_or_error(cfg: Config):
from tht.mschema.models import Annotations
try:
return Annotations.from_yaml(annotations_path(cfg))
except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
raise _MachineSchemaError("annotations_invalid") from None
_MAX_STAGED_SQL_BYTES = 1 << 20
_MAX_STAGED_SQL_TOTAL = 16 << 20
_MAX_STAGED_SQL_FILES = 32
@@ -158,7 +166,7 @@ def _load_schema_config(config: Path, *, suppress_legacy_warning: bool = False):
return load_config(config)
def _physical_or_error(cfg):
def _physical_or_error(cfg: Config):
path = physical_path(cfg)
if not path.exists():
raise _MachineSchemaError("physical_schema_missing")
@@ -173,25 +181,34 @@ def _physical_or_error(cfg):
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")
"""Collect distinct staged SQL paths lazily, bounded by distinct files."""
seen_files: set[Path] = set()
seen_roots: set[Path] = set()
ordered: list[Path] = []
# 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 add(path: Path):
canonical = path.resolve()
if canonical in seen_files:
return
seen_files.add(canonical)
ordered.append(canonical)
if len(ordered) > _MAX_STAGED_SQL_FILES:
raise _MachineSchemaError("staged_sql_too_many")
for item in inputs or []:
canonical_item = item.resolve()
if item.is_file():
add(canonical_item)
elif item.is_dir():
if canonical_item in seen_roots:
continue
seen_roots.add(canonical_item)
for path in item.rglob("*.sql"):
if path.is_file():
add(path)
else:
raise _MachineSchemaError("staged_sql_invalid")
return sorted(ordered, key=Path.as_posix)
def _read_staged_sql(inputs: list[Path] | None) -> tuple[list[Path], list[str]]:
@@ -251,11 +268,13 @@ def _candidate_key(fk) -> tuple:
def suggest_fks_data(
config: Path,
config: Config | Path,
*,
from_sql: list[Path] | None = None,
assume: list[str] | None = None,
suppress_legacy_warning: bool = False,
physical=None,
annotations=None,
) -> SuggestFksResult:
"""Return deterministic FK candidates without reviewing or mutating annotations."""
import hashlib
@@ -263,15 +282,20 @@ def suggest_fks_data(
from tht.mschema.fkmine import mine_join_pairs
from tht.mschema.merge import find_orphans
from tht.mschema.models import Annotations, ForeignKey
from tht.mschema.models import 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)
if isinstance(config, Path):
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineSchemaError("invalid_configuration") from None
else:
cfg = config
if physical is None:
physical = _physical_or_error(cfg)
_, sql_contents = _read_staged_sql(from_sql)
annotations = Annotations.from_yaml(annotations_path(cfg))
if annotations is None:
annotations = _annotations_or_error(cfg)
assumed: dict[str, str] = {}
for value in assume or []:
@@ -368,21 +392,22 @@ def suggest_fks_data(
def check_schema_data(
config: Path, *, suppress_legacy_warning: bool = False
config: Config | Path, *, suppress_legacy_warning: bool = False, physical=None, annotations=None
) -> 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
if isinstance(config, Path):
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineSchemaError("invalid_configuration") from None
else:
cfg = config
if physical is None:
physical = _physical_or_error(cfg)
if annotations is None:
annotations = _annotations_or_error(cfg)
orphans = find_orphans(physical, annotations)
ignored = [
f"{table_name}.{column_name} ({column.eligibility_reason})"
@@ -405,14 +430,25 @@ def check_cmd(
json_output: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."),
) -> None:
"""Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane."""
if json_output:
try:
cfg = _load_schema_config(config, suppress_legacy_warning=True)
except ConfigError:
_schema_json({"status": "failed", "code": "invalid_configuration"})
raise typer.Exit(code=1) from None
else:
cfg = _load_config_or_exit(config)
try:
payload = check_schema_data(config, suppress_legacy_warning=json_output)
physical = _physical_or_error(cfg)
annotations = _annotations_or_error(cfg)
payload = check_schema_data(cfg, physical=physical, annotations=annotations)
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
_render_schema_machine_error(error, cfg)
except Exception:
logger.exception("Schema check failed")
if json_output:
_schema_json({"status": "failed", "code": "schema_check_failed"})
raise typer.Exit(code=1) from None
@@ -439,9 +475,8 @@ def check_cmd(
typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN)
def _render_schema_machine_error(error: _MachineSchemaError, config: Path) -> None:
def _render_schema_machine_error(error: _MachineSchemaError, cfg) -> 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`.",
@@ -471,24 +506,32 @@ def suggest_fks_cmd(
"""Suggerisce FK logiche per la curazione umana in annotations.yaml."""
import yaml as _yaml
from tht.mschema.models import Annotations, TableAnnotation
from tht.mschema.models import TableAnnotation
if json_output and write:
_schema_json({"status": "failed", "code": "write_not_allowed"})
raise typer.Exit(code=2)
if json_output:
try:
cfg = _load_schema_config(config, suppress_legacy_warning=True)
except ConfigError:
_schema_json({"status": "failed", "code": "invalid_configuration"})
raise typer.Exit(code=1) from None
else:
cfg = _load_config_or_exit(config)
try:
physical = _physical_or_error(cfg)
annotations = _annotations_or_error(cfg)
payload = suggest_fks_data(
config,
from_sql=from_sql,
assume=assume,
suppress_legacy_warning=json_output,
cfg, from_sql=from_sql, assume=assume, physical=physical, annotations=annotations
)
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
_render_schema_machine_error(error, cfg)
except Exception:
logger.exception("Schema suggestion failed")
if json_output:
_schema_json({"status": "failed", "code": "schema_suggestion_failed"})
raise typer.Exit(code=1) from None
@@ -519,8 +562,6 @@ def suggest_fks_cmd(
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
+57 -25
View File
@@ -11,7 +11,13 @@ from tht.cli._guards import (
require_vector_write_allowed,
)
from tht.cli.config_cmd import CONFIG_OPT
from tht.cli.schema_cmd import _load_config_or_exit, annotations_path, physical_path
from tht.cli.schema_cmd import (
_load_config_or_exit,
_load_schema_config,
annotations_path,
physical_path,
)
from tht.config import Config, ConfigError
from tht.ports.vector import VectorWriteRecord
from tht.vectorstore.store import SyncStats, content_hash
@@ -154,7 +160,7 @@ class IndexSchemaResult(TypedDict):
counts: IndexCounts
def _vector_cfg_or_error(cfg) -> None:
def _vector_cfg_or_error(cfg: Config) -> None:
missing = []
if cfg.embeddings is None:
missing.append("embeddings")
@@ -164,24 +170,42 @@ def _vector_cfg_or_error(cfg) -> None:
raise _MachineVectorError("vector_configuration_missing")
def _vector_write_or_error(cfg) -> None:
def _vector_write_or_error(cfg: Config) -> None:
if cfg.profile == "workstation" and not has_vector_write_rest(cfg):
raise _MachineVectorError("vector_write_not_allowed")
def _load_schema_artifacts(cfg: Config):
import yaml
from pydantic import ValidationError
from tht.mschema.models import Annotations, PhysicalSchema
phys_file = physical_path(cfg)
if not phys_file.exists():
raise _MachineVectorError("physical_schema_missing")
try:
return (
PhysicalSchema.from_yaml(phys_file),
Annotations.from_yaml(annotations_path(cfg)),
)
except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
raise _MachineVectorError("schema_artifacts_invalid") from None
def index_schema_data(
config: Path, *, suppress_legacy_warning: bool = False
config: Config | Path, *, suppress_legacy_warning: bool = False, physical=None, annotations=None
) -> IndexSchemaResult:
"""Synchronize schema records and return a bounded machine result."""
from tht.cli.schema_cmd import _load_schema_config
from tht.config import ConfigError
from tht.mschema.models import Annotations, PhysicalSchema
from tht.vectorstore.records import schema_records
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineVectorError("invalid_configuration") from None
if isinstance(config, Path):
try:
cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning)
except ConfigError:
raise _MachineVectorError("invalid_configuration") from None
else:
cfg = config
_vector_write_or_error(cfg)
_vector_cfg_or_error(cfg)
phys_file = physical_path(cfg)
@@ -191,8 +215,10 @@ def index_schema_data(
from pydantic import ValidationError
try:
physical = PhysicalSchema.from_yaml(phys_file)
annotations = Annotations.from_yaml(annotations_path(cfg))
if physical is None:
physical = PhysicalSchema.from_yaml(phys_file)
if annotations is None:
annotations = Annotations.from_yaml(annotations_path(cfg))
except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError):
raise _MachineVectorError("schema_artifacts_invalid") from None
records = schema_records(physical, annotations)
@@ -226,29 +252,35 @@ def index_schema_cmd(
if json_output:
try:
payload = index_schema_data(config, suppress_legacy_warning=True)
except _MachineVectorError as error:
cfg = _load_schema_config(config, suppress_legacy_warning=True)
except ConfigError:
typer.echo(json.dumps({"status": "failed", "code": "invalid_configuration"}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
else:
cfg = _load_config_or_exit(config)
try:
physical, annotations = _load_schema_artifacts(cfg)
payload = index_schema_data(cfg, physical=physical, annotations=annotations)
except _MachineVectorError as error:
if json_output:
typer.echo(json.dumps({"status": "failed", "code": error.code}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
except Exception: # noqa: BLE001 - JSON CLI boundary
typer.echo(json.dumps({"status": "failed", "code": "schema_index_failed"}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
typer.echo(json.dumps(payload, sort_keys=True, separators=(",", ":")))
return
try:
payload = index_schema_data(config, suppress_legacy_warning=False)
except _MachineVectorError as error:
_render_index_schema_error(error, config)
_render_index_schema_error(error, cfg)
except Exception:
logger.exception("Schema indexing failed")
if json_output:
typer.echo(json.dumps({"status": "failed", "code": "schema_index_failed"}, sort_keys=True, separators=(",", ":")))
raise typer.Exit(code=1) from None
typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1) from None
if json_output:
typer.echo(json.dumps(payload, sort_keys=True, separators=(",", ":")))
return
_print_stats(payload["counts"])
def _render_index_schema_error(error: _MachineVectorError, config: Path) -> None:
def _render_index_schema_error(error: _MachineVectorError, cfg) -> None:
"""Render expected failures without changing the old human CLI messages."""
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`.",