fix: harden task2 preprocessing and vector guards
This commit is contained in:
@@ -2,9 +2,13 @@ import json
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
|
|
||||||
|
import pytest
|
||||||
from typer.testing import CliRunner
|
from typer.testing import CliRunner
|
||||||
|
|
||||||
from tht.cli import app
|
from tht.cli import app
|
||||||
|
from tht.ports.evidence import EvidenceSourceError, EvidenceSourceErrorCategory
|
||||||
|
from tht.ports.vector import VectorStoreError
|
||||||
|
from tht.vectorstore.embeddings import EmbeddingsError
|
||||||
|
|
||||||
|
|
||||||
def test_preprocess_evidence_json_is_pristine(monkeypatch, tmp_path):
|
def test_preprocess_evidence_json_is_pristine(monkeypatch, tmp_path):
|
||||||
@@ -191,3 +195,64 @@ def test_preprocess_gc_uses_runtime_identity_for_dev_fd_config(monkeypatch, tmp_
|
|||||||
command.gc_from_config(Path("/dev/fd/3"), dry_run=True)
|
command.gc_from_config(Path("/dev/fd/3"), dry_run=True)
|
||||||
assert captured["dry_run"] is True
|
assert captured["dry_run"] is True
|
||||||
assert captured["workspace_id"] == "runtime-workspace"
|
assert captured["workspace_id"] == "runtime-workspace"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"error",
|
||||||
|
[
|
||||||
|
pytest.param(
|
||||||
|
EvidenceSourceError("secret source", category=EvidenceSourceErrorCategory.PERMANENT),
|
||||||
|
id="evidence-source",
|
||||||
|
),
|
||||||
|
pytest.param(VectorStoreError("secret vector"), id="vector-store"),
|
||||||
|
pytest.param(EmbeddingsError("secret embeddings"), id="embeddings"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_preprocess_evidence_json_catches_domain_failures_without_stderr(monkeypatch, tmp_path, error):
|
||||||
|
import tht.cli.preprocess_cmd as command
|
||||||
|
|
||||||
|
monkeypatch.setattr(command, "run_from_config", lambda *a, **k: (_ for _ in ()).throw(error))
|
||||||
|
response = CliRunner().invoke(
|
||||||
|
app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")]
|
||||||
|
)
|
||||||
|
|
||||||
|
assert response.exit_code == 1
|
||||||
|
assert json.loads(response.stdout) == {"status": "failed", "error": "preprocessing failed"}
|
||||||
|
assert response.stderr == ""
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"error",
|
||||||
|
[
|
||||||
|
pytest.param(
|
||||||
|
EvidenceSourceError("secret source", category=EvidenceSourceErrorCategory.PERMANENT),
|
||||||
|
id="evidence-source",
|
||||||
|
),
|
||||||
|
pytest.param(VectorStoreError("secret vector"), id="vector-store"),
|
||||||
|
pytest.param(EmbeddingsError("secret embeddings"), id="embeddings"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_preprocess_evidence_gc_json_catches_domain_failures_without_stderr(monkeypatch, tmp_path, error):
|
||||||
|
import tht.cli.preprocess_cmd as command
|
||||||
|
|
||||||
|
monkeypatch.setattr(command, "gc_from_config", lambda *a, **k: (_ for _ in ()).throw(error))
|
||||||
|
response = CliRunner().invoke(
|
||||||
|
app, ["preprocess", "evidence", "gc", "--json", "-c", str(tmp_path / "workspace.yaml")]
|
||||||
|
)
|
||||||
|
|
||||||
|
assert response.exit_code == 1
|
||||||
|
assert json.loads(response.stdout) == {"status": "failed", "error": "evidence cleanup failed"}
|
||||||
|
assert response.stderr == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_preprocess_evidence_json_unexpected_failure_has_safe_boundary(monkeypatch, tmp_path):
|
||||||
|
import tht.cli.preprocess_cmd as command
|
||||||
|
|
||||||
|
monkeypatch.setattr(command, "run_from_config", lambda *a, **k: (_ for _ in ()).throw(Exception("secret unexpected")))
|
||||||
|
response = CliRunner().invoke(
|
||||||
|
app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")]
|
||||||
|
)
|
||||||
|
|
||||||
|
assert response.exit_code == 1
|
||||||
|
assert json.loads(response.stdout) == {"status": "failed", "error": "preprocessing failed"}
|
||||||
|
assert response.stderr == ""
|
||||||
|
|||||||
@@ -225,7 +225,7 @@ def test_vector_index_schema_json_failure_is_safe(monkeypatch, tmp_path):
|
|||||||
_write_schema_artifacts(tmp_path)
|
_write_schema_artifacts(tmp_path)
|
||||||
monkeypatch.setattr(
|
monkeypatch.setattr(
|
||||||
"tht.adapters.factory.build_vector_store",
|
"tht.adapters.factory.build_vector_store",
|
||||||
lambda cfg, require_write: (_ for _ in ()).throw(RuntimeError("secret qdrant endpoint")),
|
lambda cfg, require_write: (_ for _ in ()).throw(Exception("secret qdrant endpoint")),
|
||||||
)
|
)
|
||||||
response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)])
|
response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)])
|
||||||
assert response.exit_code != 0
|
assert response.exit_code != 0
|
||||||
@@ -233,6 +233,7 @@ def test_vector_index_schema_json_failure_is_safe(monkeypatch, tmp_path):
|
|||||||
payload = json.loads(response.stdout)
|
payload = json.loads(response.stdout)
|
||||||
assert payload == {"status": "failed", "code": "schema_index_failed"}
|
assert payload == {"status": "failed", "code": "schema_index_failed"}
|
||||||
assert "secret qdrant" not in response.stdout
|
assert "secret qdrant" not in response.stdout
|
||||||
|
assert response.stderr == ""
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
@@ -294,3 +295,39 @@ def test_vector_index_schema_json_legacy_config_has_no_stderr_on_failure(tmp_pat
|
|||||||
assert json.loads(response.stdout) == {
|
assert json.loads(response.stdout) == {
|
||||||
"status": "failed", "code": "physical_schema_missing"
|
"status": "failed", "code": "physical_schema_missing"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_vector_index_schema_guards_before_artifact_access(tmp_path, monkeypatch):
|
||||||
|
import tht.cli.vector_cmd as command
|
||||||
|
|
||||||
|
cfg = _legacy_qdrant_runtime_config(tmp_path)
|
||||||
|
text = cfg.read_text()
|
||||||
|
text = text.replace("vectors:\n type: qdrant\n base_url: http://qdrant:6333\n collection: psd-clinical\n", "")
|
||||||
|
text = text.replace("embeddings:\n provider: ollama_internal\n base_url: http://embedding:11434\n model: qwen3-embedding:0.6b\n dim: 1024\n", "")
|
||||||
|
cfg.write_text(text)
|
||||||
|
monkeypatch.setattr(command, "_load_schema_artifacts", lambda cfg: (_ for _ in ()).throw(AssertionError("artifact access")))
|
||||||
|
response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)])
|
||||||
|
|
||||||
|
assert response.exit_code == 1
|
||||||
|
assert json.loads(response.stdout) == {"status": "failed", "code": "vector_configuration_missing"}
|
||||||
|
assert response.stderr == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_vector_index_schema_core_reuses_injected_artifacts_without_path_resolution(tmp_path, monkeypatch):
|
||||||
|
import tht.cli.vector_cmd as command
|
||||||
|
from tht.config import load_config
|
||||||
|
from tht.mschema.models import Annotations, PhysicalSchema
|
||||||
|
|
||||||
|
cfg_path = _qdrant_runtime_config(tmp_path)
|
||||||
|
_write_schema_artifacts(tmp_path)
|
||||||
|
physical = PhysicalSchema.from_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml")
|
||||||
|
annotations = Annotations.from_yaml(tmp_path / "artifacts" / "mschema" / "annotations.yaml")
|
||||||
|
store = _FakeVectorStore()
|
||||||
|
monkeypatch.setattr(command, "physical_path", lambda cfg: (_ for _ in ()).throw(AssertionError("physical path")))
|
||||||
|
monkeypatch.setattr(command, "annotations_path", lambda cfg: (_ for _ in ()).throw(AssertionError("annotations path")))
|
||||||
|
monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store)
|
||||||
|
monkeypatch.setattr(command, "make_embedder", lambda _: _FakeEmbedder())
|
||||||
|
|
||||||
|
payload = command.index_schema_data(load_config(cfg_path), physical=physical, annotations=annotations)
|
||||||
|
|
||||||
|
assert payload["status"] == "succeeded"
|
||||||
|
|||||||
@@ -228,6 +228,24 @@ def test_suggest_fks_reports_staged_file_without_mined_joins(tmp_path):
|
|||||||
assert "Minati 0 equi-join da 1 file SQL" in response.output
|
assert "Minati 0 equi-join da 1 file SQL" in response.output
|
||||||
|
|
||||||
|
|
||||||
|
def test_suggest_fks_human_write_reads_one_annotation_snapshot(tmp_path, monkeypatch):
|
||||||
|
cfg = _write_workspace(tmp_path)
|
||||||
|
ann_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml"
|
||||||
|
Annotations().to_yaml(ann_path)
|
||||||
|
reads = []
|
||||||
|
original = Annotations.from_yaml
|
||||||
|
|
||||||
|
def tracked(path):
|
||||||
|
reads.append(path)
|
||||||
|
return original(path)
|
||||||
|
|
||||||
|
monkeypatch.setattr(Annotations, "from_yaml", tracked)
|
||||||
|
response = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"])
|
||||||
|
|
||||||
|
assert response.exit_code == 0, response.output
|
||||||
|
assert reads == [ann_path]
|
||||||
|
|
||||||
|
|
||||||
def test_suggest_fks_write_merges_and_is_idempotent(tmp_path):
|
def test_suggest_fks_write_merges_and_is_idempotent(tmp_path):
|
||||||
cfg = _write_workspace(tmp_path)
|
cfg = _write_workspace(tmp_path)
|
||||||
ann_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml"
|
ann_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml"
|
||||||
@@ -601,3 +619,27 @@ def test_fresh_process_human_warning_cardinality_is_one_across_failure_and_write
|
|||||||
check=True, capture_output=True, text=True,
|
check=True, capture_output=True, text=True,
|
||||||
)
|
)
|
||||||
assert json.loads(response.stdout)["warnings"] == 0
|
assert json.loads(response.stdout)["warnings"] == 0
|
||||||
|
|
||||||
|
|
||||||
|
def test_schema_check_json_unexpected_failure_has_no_stderr_secret(monkeypatch, tmp_path):
|
||||||
|
import tht.cli.schema_cmd as command
|
||||||
|
|
||||||
|
cfg = _write_workspace(tmp_path)
|
||||||
|
monkeypatch.setattr(command, "_physical_or_error", lambda cfg: (_ for _ in ()).throw(Exception("secret schema adapter")))
|
||||||
|
response = CliRunner().invoke(app, ["schema", "check", "--json", "-c", str(cfg)])
|
||||||
|
|
||||||
|
assert response.exit_code == 1
|
||||||
|
assert json.loads(response.stdout) == {"status": "failed", "code": "schema_check_failed"}
|
||||||
|
assert response.stderr == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_schema_suggest_json_unexpected_failure_has_no_stderr_secret(monkeypatch, tmp_path):
|
||||||
|
import tht.cli.schema_cmd as command
|
||||||
|
|
||||||
|
cfg = _write_workspace(tmp_path)
|
||||||
|
monkeypatch.setattr(command, "_physical_or_error", lambda cfg: (_ for _ in ()).throw(Exception("secret schema adapter")))
|
||||||
|
response = CliRunner().invoke(app, ["schema", "suggest-fks", "--json", "-c", str(cfg)])
|
||||||
|
|
||||||
|
assert response.exit_code == 1
|
||||||
|
assert json.loads(response.stdout) == {"status": "failed", "code": "schema_suggestion_failed"}
|
||||||
|
assert response.stderr == ""
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import hashlib
|
import hashlib
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
import re
|
import re
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -12,8 +13,17 @@ import typer
|
|||||||
from tht.cli.config_cmd import CONFIG_OPT
|
from tht.cli.config_cmd import CONFIG_OPT
|
||||||
from tht.cli.schema_cmd import _load_config_or_exit
|
from tht.cli.schema_cmd import _load_config_or_exit
|
||||||
from tht.config import workspace_id_for_config
|
from tht.config import workspace_id_for_config
|
||||||
|
from tht.ports.evidence import EvidenceSourceError
|
||||||
|
from tht.ports.vector import VectorStoreError
|
||||||
|
from tht.vectorstore.embeddings import EmbeddingsError
|
||||||
|
|
||||||
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
|
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
_PREPROCESS_EXPECTED_ERRORS = (
|
||||||
|
OSError, RuntimeError, ValueError, TypeError, KeyError,
|
||||||
|
EvidenceSourceError, VectorStoreError, EmbeddingsError,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def run_dwh_from_config(
|
def run_dwh_from_config(
|
||||||
@@ -132,7 +142,16 @@ def evidence_cmd(
|
|||||||
if action == "gc":
|
if action == "gc":
|
||||||
try:
|
try:
|
||||||
payload = gc_from_config(config, dry_run=dry_run)
|
payload = gc_from_config(config, dry_run=dry_run)
|
||||||
except (OSError, RuntimeError, ValueError, TypeError, KeyError):
|
except _PREPROCESS_EXPECTED_ERRORS:
|
||||||
|
payload = {"status": "failed", "error": "evidence cleanup failed"}
|
||||||
|
if json_output:
|
||||||
|
typer.echo(json.dumps(payload, sort_keys=True))
|
||||||
|
else:
|
||||||
|
typer.secho("ERRORE: evidence cleanup failed", fg=typer.colors.RED, err=True)
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
except Exception:
|
||||||
|
if not json_output:
|
||||||
|
logger.exception("Evidence cleanup failed")
|
||||||
payload = {"status": "failed", "error": "evidence cleanup failed"}
|
payload = {"status": "failed", "error": "evidence cleanup failed"}
|
||||||
if json_output:
|
if json_output:
|
||||||
typer.echo(json.dumps(payload, sort_keys=True))
|
typer.echo(json.dumps(payload, sort_keys=True))
|
||||||
@@ -153,7 +172,16 @@ def evidence_cmd(
|
|||||||
raise typer.Exit(code=2)
|
raise typer.Exit(code=2)
|
||||||
try:
|
try:
|
||||||
result = run_from_config(config, dry_run=dry_run, resume=resume)
|
result = run_from_config(config, dry_run=dry_run, resume=resume)
|
||||||
except (OSError, RuntimeError, ValueError, TypeError, KeyError):
|
except _PREPROCESS_EXPECTED_ERRORS:
|
||||||
|
payload = {"status": "failed", "error": "preprocessing failed"}
|
||||||
|
if json_output:
|
||||||
|
typer.echo(json.dumps(payload, sort_keys=True))
|
||||||
|
else:
|
||||||
|
typer.secho("ERRORE: preprocessing failed", fg=typer.colors.RED, err=True)
|
||||||
|
raise typer.Exit(code=1) from None
|
||||||
|
except Exception:
|
||||||
|
if not json_output:
|
||||||
|
logger.exception("Evidence preprocessing failed")
|
||||||
payload = {"status": "failed", "error": "preprocessing failed"}
|
payload = {"status": "failed", "error": "preprocessing failed"}
|
||||||
if json_output:
|
if json_output:
|
||||||
typer.echo(json.dumps(payload, sort_keys=True))
|
typer.echo(json.dumps(payload, sort_keys=True))
|
||||||
|
|||||||
@@ -448,10 +448,10 @@ def check_cmd(
|
|||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
_render_schema_machine_error(error, cfg)
|
_render_schema_machine_error(error, cfg)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Schema check failed")
|
|
||||||
if json_output:
|
if json_output:
|
||||||
_schema_json({"status": "failed", "code": "schema_check_failed"})
|
_schema_json({"status": "failed", "code": "schema_check_failed"})
|
||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
|
logger.exception("Schema check failed")
|
||||||
typer.secho("ERRORE: impossibile verificare le annotazioni.", fg=typer.colors.RED, err=True)
|
typer.secho("ERRORE: impossibile verificare le annotazioni.", fg=typer.colors.RED, err=True)
|
||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
if json_output:
|
if json_output:
|
||||||
@@ -531,10 +531,10 @@ def suggest_fks_cmd(
|
|||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
_render_schema_machine_error(error, cfg)
|
_render_schema_machine_error(error, cfg)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Schema suggestion failed")
|
|
||||||
if json_output:
|
if json_output:
|
||||||
_schema_json({"status": "failed", "code": "schema_suggestion_failed"})
|
_schema_json({"status": "failed", "code": "schema_suggestion_failed"})
|
||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
|
logger.exception("Schema suggestion failed")
|
||||||
typer.secho("ERRORE: impossibile elaborare il catalogo.", fg=typer.colors.RED, err=True)
|
typer.secho("ERRORE: impossibile elaborare il catalogo.", fg=typer.colors.RED, err=True)
|
||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
if json_output:
|
if json_output:
|
||||||
|
|||||||
@@ -208,9 +208,10 @@ def index_schema_data(
|
|||||||
cfg = config
|
cfg = config
|
||||||
_vector_write_or_error(cfg)
|
_vector_write_or_error(cfg)
|
||||||
_vector_cfg_or_error(cfg)
|
_vector_cfg_or_error(cfg)
|
||||||
phys_file = physical_path(cfg)
|
if physical is None:
|
||||||
if not phys_file.exists():
|
phys_file = physical_path(cfg)
|
||||||
raise _MachineVectorError("physical_schema_missing")
|
if not phys_file.exists():
|
||||||
|
raise _MachineVectorError("physical_schema_missing")
|
||||||
import yaml
|
import yaml
|
||||||
from pydantic import ValidationError
|
from pydantic import ValidationError
|
||||||
|
|
||||||
@@ -259,6 +260,9 @@ def index_schema_cmd(
|
|||||||
else:
|
else:
|
||||||
cfg = _load_config_or_exit(config)
|
cfg = _load_config_or_exit(config)
|
||||||
try:
|
try:
|
||||||
|
# Authorization and configuration are checked before touching any artifacts.
|
||||||
|
_vector_write_or_error(cfg)
|
||||||
|
_vector_cfg_or_error(cfg)
|
||||||
physical, annotations = _load_schema_artifacts(cfg)
|
physical, annotations = _load_schema_artifacts(cfg)
|
||||||
payload = index_schema_data(cfg, physical=physical, annotations=annotations)
|
payload = index_schema_data(cfg, physical=physical, annotations=annotations)
|
||||||
except _MachineVectorError as error:
|
except _MachineVectorError as error:
|
||||||
@@ -267,10 +271,10 @@ def index_schema_cmd(
|
|||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
_render_index_schema_error(error, cfg)
|
_render_index_schema_error(error, cfg)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Schema indexing failed")
|
|
||||||
if json_output:
|
if json_output:
|
||||||
typer.echo(json.dumps({"status": "failed", "code": "schema_index_failed"}, sort_keys=True, separators=(",", ":")))
|
typer.echo(json.dumps({"status": "failed", "code": "schema_index_failed"}, sort_keys=True, separators=(",", ":")))
|
||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
|
logger.exception("Schema indexing failed")
|
||||||
typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True)
|
typer.secho("ERRORE: impossibile indicizzare lo schema.", fg=typer.colors.RED, err=True)
|
||||||
raise typer.Exit(code=1) from None
|
raise typer.Exit(code=1) from None
|
||||||
if json_output:
|
if json_output:
|
||||||
|
|||||||
Reference in New Issue
Block a user