From 554de4ff2f7f353ce130e114b8dc9e0e0f2d11b5 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 05:38:25 +0200 Subject: [PATCH] fix: complete P2 Task 2 schema machine contracts --- harness/tests/test_qdrant_cli_commands.py | 31 ++++ harness/tests/test_schema_fk_annotations.py | 148 ++++++++++++++++++++ harness/tht/cli/schema_cmd.py | 103 +++++++++++--- harness/tht/cli/vector_cmd.py | 104 ++++++++++++-- 4 files changed, 356 insertions(+), 30 deletions(-) diff --git a/harness/tests/test_qdrant_cli_commands.py b/harness/tests/test_qdrant_cli_commands.py index edb17ec8..175b105d 100644 --- a/harness/tests/test_qdrant_cli_commands.py +++ b/harness/tests/test_qdrant_cli_commands.py @@ -221,3 +221,34 @@ def test_vector_index_schema_json_failure_is_safe(monkeypatch, tmp_path): payload = json.loads(response.stdout) assert payload == {"status": "failed", "code": "schema_index_failed"} assert "secret qdrant" not in response.stdout + + + +def test_vector_index_schema_human_missing_physical_has_original_error(tmp_path): + cfg = _qdrant_runtime_config(tmp_path) + _write_schema_artifacts(tmp_path) + physical = tmp_path / "artifacts" / "mschema" / "physical.yaml" + physical.unlink() + response = CliRunner().invoke(app, ["vector", "index-schema", "-c", str(cfg)]) + assert response.exit_code == 1 + assert "physical.yaml non trovato. Esegui prima `tht schema introspect`." in response.output + + +def test_vector_index_schema_json_missing_physical_has_no_stderr_prose(tmp_path): + import json + + cfg = _qdrant_runtime_config(tmp_path) + _write_schema_artifacts(tmp_path) + (tmp_path / "artifacts" / "mschema" / "physical.yaml").unlink() + response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)]) + assert response.exit_code == 1 + assert response.stderr == "" + assert json.loads(response.stdout) == {"status": "failed", "code": "physical_schema_missing"} + + +def test_vector_index_schema_human_missing_config_has_original_error(tmp_path): + cfg = tmp_path / "missing.yaml" + response = CliRunner().invoke(app, ["vector", "index-schema", "-c", str(cfg)]) + assert response.exit_code == 1 + assert response.output + assert "ERRORE:" in response.output diff --git a/harness/tests/test_schema_fk_annotations.py b/harness/tests/test_schema_fk_annotations.py index 185c37e5..7f46c1b6 100644 --- a/harness/tests/test_schema_fk_annotations.py +++ b/harness/tests/test_schema_fk_annotations.py @@ -281,3 +281,151 @@ def test_schema_check_json_reports_orphan_count_without_prose(tmp_path): "status": "failed", "code": "annotation_orphans", "orphan_count": 1, "orphans": ["gone"], } + + +def test_suggest_fks_json_rejects_more_than_32_staged_sql_files(tmp_path): + import json + + cfg = _write_workspace(tmp_path) + staged = tmp_path / "many" + staged.mkdir() + for i in range(33): + (staged / f"q{i:02d}.sql").write_text("SELECT 1") + response = CliRunner().invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)] + ) + assert response.exit_code == 1 + assert json.loads(response.stdout) == {"status": "failed", "code": "staged_sql_too_many"} + + +def test_staged_sql_size_limits_are_inclusive(tmp_path): + import json + + cfg = _write_workspace(tmp_path) + staged = tmp_path / "boundary" + staged.mkdir() + # Exactly one MiB per file and exactly 16 MiB in aggregate are accepted. + body = "-- padding\n" + ("x" * (1 << 20)) + body = body[: 1 << 20] + for i in range(16): + (staged / f"q{i:02d}.sql").write_bytes(body.encode()) + response = CliRunner().invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)] + ) + assert response.exit_code == 0, response.output + assert json.loads(response.stdout)["staged_sql_count"] == 16 + + +def test_schema_check_json_success_and_reviewed_annotations_are_preserved(tmp_path): + import json + + cfg = _write_workspace(tmp_path) + annotations_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml" + reviewed = _annotations_with_fks() + reviewed.to_yaml(annotations_path) + before = annotations_path.read_text() + + check = CliRunner().invoke(app, ["schema", "check", "--json", "-c", str(cfg)]) + assert check.exit_code == 0, check.output + assert json.loads(check.stdout) == { + "status": "succeeded", "code": "ok", "orphan_count": 0, "orphans": [] + } + suggest = CliRunner().invoke(app, ["schema", "suggest-fks", "--json", "-c", str(cfg)]) + assert suggest.exit_code == 0, suggest.output + assert json.loads(suggest.stdout)["candidate_count"] == 0 + assert annotations_path.read_text() == before + + +def test_schema_human_check_and_suggest_report_missing_physical_schema(tmp_path): + cfg = _write_workspace(tmp_path) + physical = tmp_path / "artifacts" / "mschema" / "physical.yaml" + physical.unlink() + expected = "physical.yaml non trovato. Esegui prima `tht schema introspect`." + check = CliRunner().invoke(app, ["schema", "check", "-c", str(cfg)]) + assert check.exit_code == 1 + assert expected in check.output + suggest = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) + assert suggest.exit_code == 1 + assert expected in suggest.output + + +def test_schema_human_check_reports_ignored_columns(tmp_path): + cfg = _write_workspace(tmp_path) + physical = _physical() + physical.tables["fact_ablazione"].columns["esito"].eligible = False + physical.tables["fact_ablazione"].columns["esito"].eligibility_reason = "test ignored" + physical.to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml") + response = CliRunner().invoke(app, ["schema", "check", "-c", str(cfg)]) + assert response.exit_code == 0, response.output + assert "Colonne ignorate (testo ampio, 1):" in response.output + assert "fact_ablazione.esito (test ignored)" in response.output + + + +def test_staged_sql_file_size_over_one_mib_is_rejected(tmp_path): + import json + + cfg = _write_workspace(tmp_path) + staged = tmp_path / "oversize" + staged.mkdir() + (staged / "too-large.sql").write_bytes(b"x" * ((1 << 20) + 1)) + response = CliRunner().invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)] + ) + assert response.exit_code == 1 + assert json.loads(response.stdout) == {"status": "failed", "code": "staged_sql_too_large"} + + +def test_staged_sql_aggregate_over_16_mib_is_rejected(tmp_path): + import json + + cfg = _write_workspace(tmp_path) + staged = tmp_path / "aggregate" + staged.mkdir() + for i in range(16): + (staged / f"q{i:02d}.sql").write_bytes(b"x" * (1 << 20)) + (staged / "over.sql").write_bytes(b"x") + response = CliRunner().invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(staged)] + ) + assert response.exit_code == 1 + assert json.loads(response.stdout) == {"status": "failed", "code": "staged_sql_too_large"} + + +def test_suggest_fks_candidate_order_is_stable_for_reordered_staged_inputs(tmp_path): + import json + + cfg = _write_workspace(tmp_path) + first = tmp_path / "first" + second = tmp_path / "second" + first.mkdir() + second.mkdir() + (first / "join.sql").write_text( + "SELECT * FROM fact_ablazione f JOIN dim_patient p ON f.cod_paz = p.cod_paz" + ) + (second / "join.sql").write_text( + "SELECT * FROM fact_ablazione f JOIN dim_time p ON f.data_time_key = p.day_key" + ) + runner = CliRunner() + one = runner.invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(first), "--from-sql", str(second)] + ) + two = runner.invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg), "--from-sql", str(second), "--from-sql", str(first)] + ) + assert one.exit_code == two.exit_code == 0 + assert json.loads(one.stdout) == json.loads(two.stdout) + + +def test_schema_human_check_missing_config_has_original_error(tmp_path): + response = CliRunner().invoke(app, ["schema", "check", "-c", str(tmp_path / "missing.yaml")]) + assert response.exit_code == 1 + assert response.output + assert "ERRORE:" in response.output + + +def test_schema_human_suggest_missing_config_has_original_error(tmp_path): + response = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(tmp_path / "missing.yaml")]) + assert response.exit_code == 1 + assert response.output + assert "ERRORE:" in response.output diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 9096e895..083ac030 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -1,6 +1,7 @@ # ruff: noqa: BLE001, S110, B008 import logging from pathlib import Path +from typing import TypedDict import typer @@ -124,13 +125,15 @@ def introspect_cmd( class _MachineSchemaError(Exception): """An expected schema CLI failure with a stable public code.""" - def __init__(self, code: str): + 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: @@ -170,7 +173,10 @@ def _staged_sql_files(inputs: list[Path] | None) -> list[Path]: else: raise _MachineSchemaError("staged_sql_invalid") # Resolve only for ordering; content remains read from the caller's staged path. - return sorted(set(files), key=lambda path: path.resolve().as_posix()) + 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]]: @@ -192,12 +198,43 @@ def _read_staged_sql(inputs: list[Path] | None) -> tuple[list[Path], list[str]]: 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) -> dict: + assume: list[str] | None = None) -> SuggestFksResult: """Return deterministic FK candidates without reviewing or mutating annotations.""" import hashlib import json @@ -215,7 +252,7 @@ def suggest_fks_data(config: Path, *, from_sql: list[Path] | None = None, 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") + raise _MachineSchemaError("assumption_invalid", value) assumed[col] = ref def single_pk(table) -> str | None: @@ -303,7 +340,7 @@ def suggest_fks_data(config: Path, *, from_sql: list[Path] | None = None, } -def check_schema_data(config: Path) -> dict: +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 @@ -315,11 +352,18 @@ def check_schema_data(config: Path) -> dict: 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, } @@ -335,8 +379,7 @@ def check_cmd( if json_output: _schema_json({"status": "failed", "code": error.code}) raise typer.Exit(code=1) from None - _load_config_or_exit(config) - 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"}) @@ -344,11 +387,18 @@ def check_cmd( typer.secho("ERRORE: impossibile verificare le annotazioni.", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) from None if json_output: - # Keep the machine response to one object, including expected validation failures. - _schema_json(payload) + # ``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"]: @@ -357,6 +407,27 @@ def check_cmd( 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, @@ -379,12 +450,7 @@ def suggest_fks_cmd( if json_output: _schema_json({"status": "failed", "code": error.code}) raise typer.Exit(code=1) from None - if error.code == "assumption_invalid": - typer.secho("ERRORE: --assume non valido (atteso col=tabella nel catalogo).", fg=typer.colors.RED, err=True) - else: - _load_config_or_exit(config) - typer.secho("ERRORE: impossibile elaborare il catalogo.", fg=typer.colors.RED, err=True) - 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"}) @@ -403,7 +469,9 @@ def suggest_fks_cmd( ) if payload["ambiguous_columns"]: typer.secho( - "PK ambigue saltate dalla regola same-name: " + ", ".join(payload["ambiguous_columns"]), + "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"]: @@ -422,7 +490,8 @@ def suggest_fks_cmd( 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"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 diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index dc019a0a..822d3a08 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -1,9 +1,14 @@ # ruff: noqa: BLE001 from pathlib import Path +from typing import TypedDict import typer -from tht.cli._guards import require_server_profile, require_vector_write_allowed +from tht.cli._guards import ( + has_vector_write_rest, + require_server_profile, + 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.ports.vector import VectorWriteRecord @@ -117,20 +122,63 @@ def init_cmd( ) -def index_schema_data(config: Path) -> dict: +class _MachineVectorError(Exception): + """Expected vector CLI failure with a stable machine-readable code.""" + + def __init__(self, code: str): + self.code = code + super().__init__(code) + + +class IndexCounts(TypedDict): + added: int + updated: int + deleted: int + unchanged: int + + +class IndexSchemaResult(TypedDict): + status: str + code: str + counts: IndexCounts + + +def _vector_cfg_or_error(cfg) -> None: + missing = [] + if cfg.embeddings is None: + missing.append("embeddings") + if cfg.vectors is None: + missing.append("vectors") + if missing: + raise _MachineVectorError("vector_configuration_missing") + + +def _vector_write_or_error(cfg) -> None: + if cfg.profile == "workstation" and not has_vector_write_rest(cfg): + raise _MachineVectorError("vector_write_not_allowed") + + +def index_schema_data(config: Path) -> IndexSchemaResult: """Synchronize schema records and return a bounded machine result.""" from tht.cli.schema_cmd import _machine_config from tht.mschema.models import Annotations, PhysicalSchema from tht.vectorstore.records import schema_records - cfg = _machine_config(config) - require_vector_write_allowed(cfg, "vector index-schema") - require_vector_cfg(cfg) + try: + cfg = _machine_config(config) + except Exception: + # Keep this helper free of Typer rendering while preserving a stable code. + raise _MachineVectorError("invalid_configuration") from None + _vector_write_or_error(cfg) + _vector_cfg_or_error(cfg) phys_file = physical_path(cfg) if not phys_file.exists(): - raise ValueError("physical schema missing") - physical = PhysicalSchema.from_yaml(phys_file) - annotations = Annotations.from_yaml(annotations_path(cfg)) + raise _MachineVectorError("physical_schema_missing") + try: + physical = PhysicalSchema.from_yaml(phys_file) + annotations = Annotations.from_yaml(annotations_path(cfg)) + except Exception: + raise _MachineVectorError("schema_artifacts_invalid") from None records = schema_records(physical, annotations) from tht.adapters.factory import build_vector_store @@ -158,14 +206,44 @@ def index_schema_cmd( json_output: bool = typer.Option(False, "--json", help="Emetti JSON puro su stdout."), ) -> None: """Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store.""" + import json + if json_output: try: payload = index_schema_data(config) - except Exception: - payload = {"status": "failed", "code": "schema_index_failed"} - typer.echo(__import__("json").dumps(payload, sort_keys=True, separators=(",", ":"))) + except _MachineVectorError as error: + typer.echo(json.dumps({"status": "failed", "code": error.code}, sort_keys=True, separators=(",", ":"))) raise typer.Exit(code=1) from None - typer.echo(__import__("json").dumps(payload, sort_keys=True, separators=(",", ":"))) + except Exception: + 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 - payload = index_schema_data(config) + try: + payload = index_schema_data(config) + except _MachineVectorError as error: + _render_index_schema_error(error, config) + except Exception: + typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None _print_stats(type("Stats", (), payload["counts"])()) + + +def _render_index_schema_error(error: _MachineVectorError, config: Path) -> 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`.", + fg=typer.colors.RED, + err=True, + ) + elif error.code == "vector_configuration_missing": + require_vector_cfg(cfg) + elif error.code == "vector_write_not_allowed": + require_vector_write_allowed(cfg, "vector index-schema") + elif error.code == "schema_artifacts_invalid": + typer.secho("ERRORE: schema artifacts non validi.", fg=typer.colors.RED, err=True) + else: + typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None