From f62b4dbc4c6f8b745a0f69f6350dc4364fb20051 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 06:29:09 +0200 Subject: [PATCH] fix: preserve Task 2 command loading contracts --- harness/tests/test_preprocess_cli.py | 31 ++++ harness/tests/test_schema_fk_annotations.py | 91 ++++++++++++ harness/tht/cli/preprocess_cmd.py | 8 +- harness/tht/cli/schema_cmd.py | 153 +++++++++++++------- harness/tht/cli/vector_cmd.py | 82 +++++++---- 5 files changed, 279 insertions(+), 86 deletions(-) diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index 5ad69cec..38ad21c3 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -160,3 +160,34 @@ def test_preprocess_evidence_uses_runtime_identity_for_dev_fd_config(monkeypatch monkeypatch.setattr("tht.corpus.pipeline.CorpusPipeline", FakePipeline) command.run_from_config(Path("/dev/fd/3")) assert captured["workspace_id"] == "runtime-workspace" + + +def test_preprocess_gc_uses_runtime_identity_for_dev_fd_config(monkeypatch, tmp_path): + import tht.cli.preprocess_cmd as command + + cfg = SimpleNamespace( + runtime_identity=SimpleNamespace(workspace_id="runtime-workspace"), + embeddings=SimpleNamespace(model="m", dim=4), + vector=SimpleNamespace(max_chunk_chars=10, retain_published_generations=1), + paths=SimpleNamespace(artifacts=tmp_path / "artifacts"), + ) + captured = {} + + class FakePipeline: + def __init__(self, **kwargs): + captured.update(kwargs) + + def gc(self, **kwargs): + captured.update(kwargs) + captured["workspace_id"] = self.workspace_id + return {"status": "succeeded", "dry_run": True, "evicted": [], "failures": []} + + monkeypatch.setattr(command, "_load_config_or_exit", lambda _: cfg) + monkeypatch.setattr("tht.adapters.factory.build_evidence_sources", lambda _: []) + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda *_args, **_kwargs: object()) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: object()) + monkeypatch.setattr("tht.corpus.pipeline.CorpusPipeline", FakePipeline) + + command.gc_from_config(Path("/dev/fd/3"), dry_run=True) + assert captured["dry_run"] is True + assert captured["workspace_id"] == "runtime-workspace" diff --git a/harness/tests/test_schema_fk_annotations.py b/harness/tests/test_schema_fk_annotations.py index b93b1c92..35747845 100644 --- a/harness/tests/test_schema_fk_annotations.py +++ b/harness/tests/test_schema_fk_annotations.py @@ -1,4 +1,8 @@ # ruff: noqa: DTZ001 +import json +import subprocess +import sys +import textwrap import warnings from datetime import datetime @@ -510,3 +514,90 @@ def test_staged_sql_oversize_read_is_bounded(tmp_path, monkeypatch): with pytest.raises(_MachineSchemaError, match="staged_sql_too_large"): _read_staged_sql([staged]) assert reads == [(1 << 20) + 1] + + +def test_staged_sql_deduplicates_overlapping_roots_and_repeated_files(tmp_path): + from tht.cli.schema_cmd import _staged_sql_files + + root = tmp_path / "approved" + nested = root / "nested" + nested.mkdir(parents=True) + first = root / "first.sql" + second = nested / "second.sql" + first.write_text("SELECT 1") + second.write_text("SELECT 1") + + files = _staged_sql_files([root, nested, first, root]) + assert files == sorted({first.resolve(), second.resolve()}) + + +def test_staged_sql_file_limit_counts_distinct_paths_only(tmp_path): + from tht.cli.schema_cmd import _MachineSchemaError, _staged_sql_files + + path = tmp_path / "same.sql" + path.write_text("SELECT 1") + assert _staged_sql_files([path] * 100) == [path.resolve()] + + root = tmp_path / "many" + root.mkdir() + for index in range(33): + (root / f"q{index:02d}.sql").write_text("SELECT 1") + with pytest.raises(_MachineSchemaError, match="staged_sql_too_many"): + _staged_sql_files([root]) + + +def test_fresh_process_human_warning_cardinality_is_one_across_failure_and_write_paths(tmp_path): + cfg = _write_workspace(tmp_path) + physical = tmp_path / "artifacts" / "mschema" / "physical.yaml" + annotations = tmp_path / "artifacts" / "mschema" / "annotations.yaml" + annotations.parent.mkdir(parents=True, exist_ok=True) + Annotations().to_yaml(annotations) + probe = textwrap.dedent( + """ + import json, sys, warnings + from typer.testing import CliRunner + from tht.cli import app + + with warnings.catch_warnings(record=True) as caught: + warnings.simplefilter("always") + result = CliRunner().invoke(app, sys.argv[1:]) + print(json.dumps({ + "warnings": sum(issubclass(w.category, FutureWarning) for w in caught), + "exit": result.exit_code, + })) + """ + ) + + + physical.unlink() + branches = [ + ["schema", "check"], + ["schema", "suggest-fks"], + ["vector", "index-schema"], + ] + for branch in branches: + response = subprocess.run( + [sys.executable, "-c", probe, *branch, "-c", str(cfg)], + check=True, capture_output=True, text=True, + ) + assert json.loads(response.stdout) == {"warnings": 1, "exit": 1} + + _physical().to_yaml(physical) + response = subprocess.run( + [sys.executable, "-c", probe, "schema", "suggest-fks", "--write", "-c", str(cfg)], + check=True, capture_output=True, text=True, + ) + assert json.loads(response.stdout) == {"warnings": 1, "exit": 0} + + physical.unlink() + for branch in branches: + response = subprocess.run( + [sys.executable, "-c", probe, *branch, "--json", "-c", str(cfg)], + check=True, capture_output=True, text=True, + ) + assert json.loads(response.stdout)["warnings"] == 0 + response = subprocess.run( + [sys.executable, "-c", probe, "schema", "suggest-fks", "--write", "--json", "-c", str(cfg)], + check=True, capture_output=True, text=True, + ) + assert json.loads(response.stdout)["warnings"] == 0 diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 532d1508..dd15c40a 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -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)) diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 45ac53b6..879f1110 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -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 diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index 46d6f974..b0acb55f 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -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`.",