From 9e2c2b37b3ebb79159d633401ab19d91694a0449 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 06:09:07 +0200 Subject: [PATCH] fix: harden Task 2 schema quality boundaries --- harness/tests/test_schema_fk_annotations.py | 75 ++++++++++++ harness/tht/cli/schema_cmd.py | 121 ++++++++++++-------- harness/tht/cli/vector_cmd.py | 51 ++++++--- 3 files changed, 183 insertions(+), 64 deletions(-) diff --git a/harness/tests/test_schema_fk_annotations.py b/harness/tests/test_schema_fk_annotations.py index 99eb7926..b93b1c92 100644 --- a/harness/tests/test_schema_fk_annotations.py +++ b/harness/tests/test_schema_fk_annotations.py @@ -1,10 +1,13 @@ # ruff: noqa: DTZ001 +import warnings from datetime import datetime +import pytest import yaml from typer.testing import CliRunner from tht.cli import app +from tht.cli.schema_cmd import check_schema_data from tht.mschema.merge import find_orphans from tht.mschema.models import ( Annotations, @@ -163,6 +166,10 @@ def test_suggest_fks_skips_generic_and_ambiguous_pks(tmp_path): assert res.exit_code == 0, res.output assert "nessuna FK da suggerire" in res.output # id generico, cod_x ambigua assert "cod_x" in res.output # segnalata come ambigua saltata + exact = CliRunner().invoke( + app, ["schema", "suggest-fks", "--json", "-c", str(cfg)] + ) + assert yaml.safe_load(exact.stdout)["ambiguous_columns"] == ["cod_x"] # --assume disambigua la PK multi-proprietario res2 = CliRunner().invoke( @@ -205,6 +212,18 @@ def test_suggest_fks_from_sql_mines_joins(tmp_path): assert "ref_table: dim_patient" in res.output +def test_suggest_fks_reports_staged_file_without_mined_joins(tmp_path): + cfg = _write_workspace(tmp_path) + staged = tmp_path / "approved" + staged.mkdir() + (staged / "no-joins.sql").write_text("SELECT 1") + response = CliRunner().invoke( + app, ["schema", "suggest-fks", "-c", str(cfg), "--from-sql", str(staged)] + ) + assert response.exit_code == 0, response.output + assert "Minati 0 equi-join da 1 file SQL" in response.output + + def test_suggest_fks_write_merges_and_is_idempotent(tmp_path): cfg = _write_workspace(tmp_path) ann_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml" @@ -435,3 +454,59 @@ def test_schema_human_suggest_missing_config_has_original_error(tmp_path): assert response.exit_code == 1 assert response.output assert "ERRORE:" in response.output + + +def test_human_config_keeps_legacy_warning_but_json_suppresses_it(tmp_path): + cfg = _write_workspace(tmp_path) + with pytest.warns(FutureWarning, match="legacy workspace resource keys"): + check_schema_data(cfg) + with warnings.catch_warnings(record=True) as caught: + warnings.simplefilter("always") + check_schema_data(cfg, suppress_legacy_warning=True) + assert not [item for item in caught if issubclass(item.category, FutureWarning)] + + +def test_staged_sql_enumeration_stops_at_count_sentinel(tmp_path, monkeypatch): + from tht.cli.schema_cmd import _MachineSchemaError, _staged_sql_files + + staged = tmp_path / "staged" + staged.mkdir() + files = [staged / f"q{i:03d}.sql" for i in range(100)] + for path in files: + path.write_text("SELECT 1") + yielded = 0 + + def bounded_rglob(_self, _pattern): + nonlocal yielded + for path in files: + yielded += 1 + yield path + + monkeypatch.setattr(type(staged), "rglob", bounded_rglob) + with pytest.raises(_MachineSchemaError, match="staged_sql_too_many"): + _staged_sql_files([staged]) + assert yielded == 33 + + +def test_staged_sql_oversize_read_is_bounded(tmp_path, monkeypatch): + from tht.cli.schema_cmd import _MachineSchemaError, _read_staged_sql + + staged = tmp_path / "oversize.sql" + staged.write_bytes(b"unused") + reads = [] + + class BoundedReader: + def __enter__(self): + return self + + def __exit__(self, *_args): + return False + + def read(self, size): + reads.append(size) + return b"x" * size + + monkeypatch.setattr(type(staged), "open", lambda *_args, **_kwargs: BoundedReader()) + with pytest.raises(_MachineSchemaError, match="staged_sql_too_large"): + _read_staged_sql([staged]) + assert reads == [(1 << 20) + 1] diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 96a86807..45ac53b6 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -1,10 +1,12 @@ -# ruff: noqa: BLE001, S110, B008 import logging import warnings +from itertools import islice from pathlib import Path -from typing import TypedDict +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 @@ -25,7 +27,7 @@ def _add_examples(dwh, phys, examples) -> None: sampled = dwh.sample_column( table_name, column_name, limit=examples.max_per_column ) - except Exception as exc: + except Exception as exc: # noqa: BLE001 - sampling adapter boundary logger.warning("Campionamento saltato per %s.%s: %s", table_name, column_name, exc) continue @@ -87,8 +89,8 @@ def introspect_cmd( try: cached = PhysicalSchema.from_yaml(out) - except Exception: - pass # catalogo illeggibile: procedi con la re-introspezione + 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: @@ -109,7 +111,7 @@ def introspect_cmd( raise RuntimeError("DWH preprocessing failed") out = physical_path(cfg) phys = PhysicalSchema.from_yaml(out) - except Exception as e: + 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()) @@ -143,20 +145,17 @@ def _schema_json(payload: dict) -> None: typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":"))) -def _machine_config(config: Path): - # Legacy resource keys remain accepted, but their deprecation warning is not part of - # the JSON machine contract. Keep this filter scoped to this call so human commands - # (and unrelated warnings/errors) retain their existing behavior. +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, ) - try: - return load_config(config) - except ConfigError: - raise _MachineSchemaError("invalid_configuration") from None + return load_config(config) def _physical_or_error(cfg): @@ -169,22 +168,28 @@ def _physical_or_error(cfg): return PhysicalSchema.from_yaml(path) except _MachineSchemaError: raise - except Exception: + 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]: - files: list[Path] = [] - for item in inputs or []: - if item.is_file(): - files.append(item) - elif item.is_dir(): - files.extend(path for path in item.rglob("*.sql") if path.is_file()) - else: - raise _MachineSchemaError("staged_sql_invalid") - # Resolve only for ordering; content remains read from the caller's staged 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: + if len(ordered) > _MAX_STAGED_SQL_FILES or len(files) > _MAX_STAGED_SQL_FILES: raise _MachineSchemaError("staged_sql_too_many") return ordered @@ -194,11 +199,13 @@ def _read_staged_sql(inputs: list[Path] | None) -> tuple[list[Path], list[str]]: contents: list[str] = [] total = 0 for path in paths: + allowance = min(_MAX_STAGED_SQL_BYTES, _MAX_STAGED_SQL_TOTAL - total) try: - raw = path.read_bytes() + with path.open("rb") as stream: + raw = stream.read(allowance + 1) except (OSError, UnicodeError): raise _MachineSchemaError("staged_sql_invalid") from None - if len(raw) > _MAX_STAGED_SQL_BYTES or total + len(raw) > _MAX_STAGED_SQL_TOTAL: + if len(raw) > allowance: raise _MachineSchemaError("staged_sql_too_large") total += len(raw) try: @@ -220,8 +227,8 @@ class FKCandidate(TypedDict): class SuggestFksResult(TypedDict): - status: str - code: str + status: Literal["succeeded", "failed"] + code: Literal["ok"] candidates: list[FKCandidate] candidate_count: int candidate_digest: str @@ -232,8 +239,8 @@ class SuggestFksResult(TypedDict): class CheckSchemaResult(TypedDict): - status: str - code: str + status: Literal["succeeded", "failed"] + code: Literal["ok", "annotation_orphans"] orphan_count: int orphans: list[str] ignored: list[str] @@ -243,8 +250,13 @@ 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) -> SuggestFksResult: +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 @@ -253,7 +265,10 @@ def suggest_fks_data(config: Path, *, from_sql: list[Path] | None = None, from tht.mschema.merge import find_orphans from tht.mschema.models import Annotations, ForeignKey - cfg = _machine_config(config) + 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)) @@ -316,11 +331,13 @@ def suggest_fks_data(config: Path, *, from_sql: list[Path] | None = None, 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 or column_name in _GENERIC_PK_NAMES: + if not owners: continue add(table_name, column_name, owners[0], column_name) @@ -350,16 +367,21 @@ def suggest_fks_data(config: Path, *, from_sql: list[Path] | None = None, } -def check_schema_data(config: Path) -> CheckSchemaResult: +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 - cfg = _machine_config(config) + 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 Exception: + except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError): raise _MachineSchemaError("annotations_invalid") from None orphans = find_orphans(physical, annotations) ignored = [ @@ -384,13 +406,13 @@ def check_cmd( ) -> None: """Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane.""" try: - payload = check_schema_data(config) + 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: + 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 @@ -441,8 +463,8 @@ def _render_schema_machine_error(error: _MachineSchemaError, config: Path) -> No @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)."), - assume: list[str] = typer.Option(None, "--assume", help="Disambigua una PK: col=tabella."), + 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: @@ -455,13 +477,18 @@ def suggest_fks_cmd( _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) + 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: + 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 @@ -472,7 +499,7 @@ def suggest_fks_cmd( return # Human mode remains the original renderer over the pure result. - if payload["mined_join_count"]: + 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, @@ -524,10 +551,10 @@ def render_cmd( format: str = typer.Option( "markdown", "--format", "-f", help="Formato: markdown | mschema-text | schema-dict" ), - tables: list[str] = typer.Option( + 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)."), + 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 diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index 822d3a08..46d6f974 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -1,6 +1,7 @@ -# ruff: noqa: BLE001 +import logging +from collections.abc import Mapping from pathlib import Path -from typing import TypedDict +from typing import Literal, TypedDict import typer @@ -15,6 +16,7 @@ from tht.ports.vector import VectorWriteRecord from tht.vectorstore.store import SyncStats, content_hash vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)") +logger = logging.getLogger(__name__) def make_embedder(embeddings_cfg): @@ -80,10 +82,19 @@ def sync_canonical_records(collection, records, *, store, embedder): return stats -def _print_stats(stats) -> None: +def _print_stats(stats: SyncStats | Mapping[str, int]) -> None: + if isinstance(stats, SyncStats): + counts = { + "added": stats.added, + "updated": stats.updated, + "deleted": stats.deleted, + "unchanged": stats.unchanged, + } + else: + counts = stats typer.secho( - f"OK: {stats.added} nuovi, {stats.updated} aggiornati, " - f"{stats.deleted} rimossi, {stats.unchanged} invariati", + f"OK: {counts['added']} nuovi, {counts['updated']} aggiornati, " + f"{counts['deleted']} rimossi, {counts['unchanged']} invariati", fg=typer.colors.GREEN, ) @@ -138,8 +149,8 @@ class IndexCounts(TypedDict): class IndexSchemaResult(TypedDict): - status: str - code: str + status: Literal["succeeded", "failed"] + code: Literal["ok"] counts: IndexCounts @@ -158,26 +169,31 @@ def _vector_write_or_error(cfg) -> None: raise _MachineVectorError("vector_write_not_allowed") -def index_schema_data(config: Path) -> IndexSchemaResult: +def index_schema_data( + config: Path, *, suppress_legacy_warning: bool = False +) -> IndexSchemaResult: """Synchronize schema records and return a bounded machine result.""" - from tht.cli.schema_cmd import _machine_config + 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 = _machine_config(config) - except Exception: - # Keep this helper free of Typer rendering while preserving a stable code. + cfg = _load_schema_config(config, suppress_legacy_warning=suppress_legacy_warning) + except ConfigError: 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 _MachineVectorError("physical_schema_missing") + import yaml + from pydantic import ValidationError + try: physical = PhysicalSchema.from_yaml(phys_file) annotations = Annotations.from_yaml(annotations_path(cfg)) - except Exception: + except (OSError, UnicodeError, TypeError, ValueError, yaml.YAMLError, ValidationError): raise _MachineVectorError("schema_artifacts_invalid") from None records = schema_records(physical, annotations) from tht.adapters.factory import build_vector_store @@ -210,23 +226,24 @@ def index_schema_cmd( if json_output: try: - payload = index_schema_data(config) + payload = index_schema_data(config, suppress_legacy_warning=True) except _MachineVectorError as error: typer.echo(json.dumps({"status": "failed", "code": error.code}, sort_keys=True, separators=(",", ":"))) raise typer.Exit(code=1) from None - except Exception: + 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) + payload = index_schema_data(config, suppress_legacy_warning=False) except _MachineVectorError as error: _render_index_schema_error(error, config) except Exception: + logger.exception("Schema indexing failed") 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"])()) + _print_stats(payload["counts"]) def _render_index_schema_error(error: _MachineVectorError, config: Path) -> None: