diff --git a/harness/tests/test_local_compose_contract.py b/harness/tests/test_local_compose_contract.py index 525b4ea2..8066766c 100644 --- a/harness/tests/test_local_compose_contract.py +++ b/harness/tests/test_local_compose_contract.py @@ -8,7 +8,11 @@ def test_local_compose_uses_the_generic_external_endpoint_contract(): compose = yaml.safe_load((root / "compose.yaml").read_text()) local = yaml.safe_load((root / "deploy/compose.local.yaml").read_text()) - assert set(compose["services"]) == {"core", "frontend", "qdrant", "embedding", "embedding-model-init"} + assert set(compose["services"]) == { + "core", "frontend", "qdrant", "embedding", "embedding-model-init", "workspace-maintenance", + } + # workspace-maintenance is profile-gated: it must not be part of the default local startup. + assert compose["services"]["workspace-maintenance"].get("profiles") == ["workspace-maintenance"] assert local["services"]["core"]["environment"]["AUTH_MODE"] == "none" assert local["services"]["core"]["ports"] == ["127.0.0.1:${THOTH_CORE_HTTP_PORT:-8787}:8787"] assert local["services"]["frontend"]["ports"] == ["127.0.0.1:${THOTH_HTTP_PORT:-8080}:8080"] diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index 8040865e..792111eb 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -1,4 +1,5 @@ import json +from pathlib import Path from types import SimpleNamespace from typer.testing import CliRunner @@ -6,68 +7,152 @@ from typer.testing import CliRunner from tht.cli import app +def _runtime_config(tmp_path: Path, name: str = "workspace.yaml") -> Path: + path = tmp_path / name + (tmp_path / "evidence").mkdir(exist_ok=True) + path.write_text( + f""" +runtime_identity: + workspace_id: psd-clinical + workspace_revision: {'a' * 40} +dwh: + type: postgres_direct + connection: {{database: analytics, schema: mart, user: reader, password: secret}} +vectors: + type: qdrant + base_url: http://qdrant:6333 + collection: psd-clinical +embeddings: + provider: ollama_internal + base_url: http://embedding:11434 + model: qwen3-embedding:0.6b + dim: 1024 +evidence: + sources: + - type: filesystem + root: {tmp_path / 'evidence'} +roots: + sessions: {tmp_path / 'sessions'} + artifacts: {tmp_path / 'artifacts'} + indexes: {tmp_path / 'indexes'} +""" + ) + return path + + def test_preprocess_evidence_json_is_pristine(monkeypatch, tmp_path): import tht.cli.preprocess_cmd as command - result = SimpleNamespace(model_dump=lambda mode=None: { - "status": "succeeded", "generation": "gen:abc", "published": True - }) - monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) - response = CliRunner().invoke( - app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")] + config = _runtime_config(tmp_path) + result = SimpleNamespace( + model_dump=lambda mode=None: { + "status": "succeeded", + "generation": "gen:abc", + "published": True, + "counts": {"changed": 0, "unchanged": 0, "removed": 0, "documents": 0, "chunks": 0}, + "changed": [], + "unchanged": [], + "removed": [], + "manifest_id": "manifest-1", + "run_id": "a" * 32, + "resumed_from": None, + } ) + monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) + response = CliRunner().invoke(app, ["preprocess", "evidence", "--json", "-c", str(config)]) assert response.exit_code == 0, response.output - assert json.loads(response.output)["generation"] == "gen:abc" + assert response.stderr == "" + assert json.loads(response.stdout) == { + "changed": [], + "code": "ok", + "counts": {"changed": 0, "chunks": 0, "documents": 0, "removed": 0, "unchanged": 0}, + "generation": "gen:abc", + "manifest_id": "manifest-1", + "operation": "preprocess_evidence", + "published": True, + "removed": [], + "resumed_from": None, + "run_id": "a" * 32, + "schemaVersion": 1, + "status": "succeeded", + "unchanged": [], + "workspaceId": "psd-clinical", + "workspaceRevision": "a" * 40, + } def test_preprocess_failure_is_structured_and_nonzero(monkeypatch, tmp_path): import tht.cli.preprocess_cmd as command - monkeypatch.setattr(command, "run_from_config", lambda *a, **k: (_ for _ in ()).throw(RuntimeError("secret detail"))) - response = CliRunner().invoke( - app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")] + config = _runtime_config(tmp_path) + monkeypatch.setattr( + command, + "run_from_config", + lambda *a, **k: (_ for _ in ()).throw(RuntimeError("secret detail")), ) + response = CliRunner().invoke(app, ["preprocess", "evidence", "--json", "-c", str(config)]) assert response.exit_code != 0 - assert json.loads(response.output) == {"status": "failed", "error": "preprocessing failed"} + payload = json.loads(response.stdout) + assert payload == { + "code": "preprocessing_failed", + "error": "preprocessing failed", + "operation": "preprocess_evidence", + "schemaVersion": 1, + "status": "failed", + "workspaceId": "psd-clinical", + "workspaceRevision": "a" * 40, + } assert "secret detail" not in response.output def test_preprocess_failed_job_report_is_sanitized_json_and_nonzero(monkeypatch, tmp_path): import tht.cli.preprocess_cmd as command - result = SimpleNamespace(model_dump=lambda mode=None: { - "status": "failed", "run_id": "a" * 32, "published": False, - "generation": "gen:" + "b" * 32, "changed": ["fs:one"], - }) - monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) - response = CliRunner().invoke( - app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")] + config = _runtime_config(tmp_path) + result = SimpleNamespace( + model_dump=lambda mode=None: { + "status": "failed", + "run_id": "a" * 32, + "published": False, + "generation": "gen:" + "b" * 32, + "changed": ["fs:one"], + "unchanged": [], + "removed": [], + "counts": {"changed": 1, "unchanged": 0, "removed": 0, "documents": 1, "chunks": 1}, + "manifest_id": "manifest-1", + "resumed_from": None, + } ) + monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) + response = CliRunner().invoke(app, ["preprocess", "evidence", "--json", "-c", str(config)]) assert response.exit_code == 1 - payload = json.loads(response.output) + payload = json.loads(response.stdout) assert payload["status"] == "failed" assert payload["error"] == "preprocessing job failed" + assert payload["workspaceId"] == "psd-clinical" assert "traceback" not in response.output.lower() def test_preprocess_real_failed_stage_result_exits_nonzero(monkeypatch, tmp_path): - import tht.cli.preprocess_cmd as command from test_corpus_pipeline import Source, item, pipeline + import tht.cli.preprocess_cmd as command + + config = _runtime_config(tmp_path) result = pipeline( - tmp_path, Source([(item("one", "a"), RuntimeError("SENSITIVE EVIDENCE secret"))]) + tmp_path, + Source([(item("one", "a"), RuntimeError("SENSITIVE EVIDENCE secret"))]), ).run_as_job( - workspace_id="demo", workspace_root=tmp_path, + workspace_id="demo", + workspace_root=tmp_path, config_fingerprint="sha256:" + "1" * 64, input_fingerprint="sha256:" + "2" * 64, ) assert result.status == "failed" monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) - response = CliRunner().invoke( - app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")] - ) + response = CliRunner().invoke(app, ["preprocess", "evidence", "--json", "-c", str(config)]) assert response.exit_code == 1 - assert json.loads(response.output)["status"] == "failed" + assert json.loads(response.stdout)["status"] == "failed" assert "SENSITIVE EVIDENCE" not in response.output assert "secret" not in response.output @@ -75,19 +160,24 @@ def test_preprocess_real_failed_stage_result_exits_nonzero(monkeypatch, tmp_path def test_preprocess_evidence_text_uses_uncapped_result_counts(monkeypatch, tmp_path): import tht.cli.preprocess_cmd as command - result = SimpleNamespace(model_dump=lambda mode=None: { - "status": "succeeded", "run_id": "a" * 32, - "generation": "gen:" + "b" * 64, "published": True, - "changed": ["fs:item"] * 100, - "unchanged": ["fs:item"] * 100, - "removed": ["fs:item"] * 100, - "counts": {"changed": 1001, "unchanged": 902, "removed": 803}, - }) + config = _runtime_config(tmp_path) + result = SimpleNamespace( + model_dump=lambda mode=None: { + "status": "succeeded", + "run_id": "a" * 32, + "generation": "gen:" + "b" * 64, + "published": True, + "changed": ["fs:item"] * 100, + "unchanged": ["fs:item"] * 100, + "removed": ["fs:item"] * 100, + "counts": {"changed": 1001, "unchanged": 902, "removed": 803}, + "manifest_id": "manifest-1", + "resumed_from": None, + } + ) monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) - response = CliRunner().invoke( - app, ["preprocess", "evidence", "-c", str(tmp_path / "workspace.yaml")] - ) + response = CliRunner().invoke(app, ["preprocess", "evidence", "-c", str(config)]) assert response.exit_code == 0, response.output assert "changed=1001 unchanged=902 removed=803" in response.output @@ -96,6 +186,7 @@ def test_preprocess_evidence_text_uses_uncapped_result_counts(monkeypatch, tmp_p def test_preprocess_resume_rejects_generation_id_before_configuration(monkeypatch, tmp_path): import tht.cli.preprocess_cmd as command + config = _runtime_config(tmp_path) called = False def forbidden(*args, **kwargs): @@ -106,26 +197,80 @@ def test_preprocess_resume_rejects_generation_id_before_configuration(monkeypatc response = CliRunner().invoke( app, [ - "preprocess", "evidence", "--resume", "gen:" + "a" * 32, - "--json", "-c", str(tmp_path / "workspace.yaml"), + "preprocess", + "evidence", + "--resume", + "gen:" + "a" * 32, + "--json", + "-c", + str(config), ], ) assert response.exit_code != 0 assert json.loads(response.output) == { - "status": "failed", "error": "resume requires a preprocessing run id" + "code": "invalid_resume", + "error": "resume requires a preprocessing run id", + "operation": "preprocess_evidence", + "schemaVersion": 1, + "status": "failed", } assert called is False +def test_run_from_config_uses_runtime_identity_workspace_id(monkeypatch, tmp_path): + import tht.cli.preprocess_cmd as command + + config = _runtime_config(tmp_path, name="3") + calls = {} + + class FakePipeline: + def __init__( + self, + *, + store, + sources, + embedder, + vector_store, + embedding_model, + embedding_dimensions, + chunk_policy, + pipeline_version, + retain_published_generations, + ): + calls["init"] = { + "embedding_model": embedding_model, + "embedding_dimensions": embedding_dimensions, + "pipeline_version": pipeline_version, + } + + def run_as_job(self, **kwargs): + calls["run_as_job"] = kwargs + return SimpleNamespace(model_dump=lambda mode=None: {"status": "succeeded"}) + + monkeypatch.setattr("tht.adapters.factory.build_evidence_sources", lambda cfg: []) + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: object()) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda cfg: object()) + monkeypatch.setattr("tht.corpus.pipeline.CorpusPipeline", FakePipeline) + + command.run_from_config(config) + + assert calls["run_as_job"]["workspace_id"] == "psd-clinical" + assert calls["run_as_job"]["input_fingerprint"] != calls["run_as_job"]["config_fingerprint"] + + def test_preprocess_evidence_gc_json_is_pristine(monkeypatch, tmp_path): import tht.cli.preprocess_cmd as command + config = _runtime_config(tmp_path) monkeypatch.setattr(command, "gc_from_config", lambda *a, **k: { - "status": "succeeded", "dry_run": True, "evicted": [], "failures": [], + "status": "succeeded", + "dry_run": True, + "evicted": [], + "failures": [], }) response = CliRunner().invoke( - app, ["preprocess", "evidence", "gc", "--dry-run", "--json", "-c", - str(tmp_path / "workspace.yaml")] + app, + ["preprocess", "evidence", "gc", "--dry-run", "--json", "-c", str(config)], ) assert response.exit_code == 0, response.output assert json.loads(response.output)["dry_run"] is True diff --git a/harness/tests/test_qdrant_cli_commands.py b/harness/tests/test_qdrant_cli_commands.py index 4ad73271..40a740c4 100644 --- a/harness/tests/test_qdrant_cli_commands.py +++ b/harness/tests/test_qdrant_cli_commands.py @@ -1,5 +1,6 @@ from __future__ import annotations +import hashlib import json from datetime import UTC, datetime from pathlib import Path @@ -9,6 +10,7 @@ from typer.testing import CliRunner from tht.cli import app from tht.memory import MemoryRecord, save_registry +from tht.ports.vector import VectorStoreError class _FakeEmbedder: @@ -33,6 +35,10 @@ class _FakeVectorStore: return 3 +def _sha_file(path: Path) -> str: + return "sha256:" + hashlib.sha256(path.read_bytes()).hexdigest() + + def _qdrant_runtime_config(tmp_path: Path) -> Path: cfg = tmp_path / "workspace.yaml" cfg.write_text( @@ -76,9 +82,7 @@ tables: type: bigint """ ) - (tmp_path / "artifacts" / "mschema" / "annotations.yaml").write_text( - "tables: {}\n" - ) + (tmp_path / "artifacts" / "mschema" / "annotations.yaml").write_text("tables: {}\n") def _memory_record() -> MemoryRecord: @@ -111,6 +115,73 @@ def test_vector_index_schema_accepts_qdrant_only_runtime_config(tmp_path, monkey assert store.upserts +def test_vector_index_schema_json_is_pristine_and_reports_artifacts(tmp_path, monkeypatch): + cfg = _qdrant_runtime_config(tmp_path) + _write_schema_artifacts(tmp_path) + store = _FakeVectorStore() + + monkeypatch.setattr("tht.adapters.factory.build_vector_store", lambda cfg, require_write: store) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: _FakeEmbedder()) + + response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)]) + + assert response.exit_code == 0, response.output + assert response.stderr == "" + assert json.loads(response.stdout) == { + "artifactIdentities": [ + { + "digest": _sha_file(tmp_path / "artifacts" / "mschema" / "annotations.yaml"), + "kind": "schema_annotations", + }, + { + "digest": _sha_file(tmp_path / "artifacts" / "mschema" / "physical.yaml"), + "kind": "physical_schema", + }, + ], + "code": "ok", + "collection": "psd-clinical", + "counts": { + "added": 2, + "columns": 1, + "deleted": 0, + "records": 2, + "tables": 1, + "unchanged": 0, + "updated": 0, + }, + "operation": "index_schema", + "schemaVersion": 1, + "status": "succeeded", + "workspaceId": "psd-clinical", + "workspaceRevision": "a" * 40, + } + + +def test_vector_index_schema_json_failure_is_pristine(tmp_path, monkeypatch): + cfg = _qdrant_runtime_config(tmp_path) + _write_schema_artifacts(tmp_path) + + def boom(cfg, require_write): + raise VectorStoreError("semantic_index_incompatible") + + monkeypatch.setattr("tht.adapters.factory.build_vector_store", boom) + monkeypatch.setattr("tht.cli.vector_cmd.make_embedder", lambda _: _FakeEmbedder()) + + response = CliRunner().invoke(app, ["vector", "index-schema", "--json", "-c", str(cfg)]) + + assert response.exit_code == 1 + assert response.stderr == "" + assert json.loads(response.stdout) == { + "code": "semantic_index_incompatible", + "error": "semantic index incompatible", + "operation": "index_schema", + "schemaVersion": 1, + "status": "failed", + "workspaceId": "psd-clinical", + "workspaceRevision": "a" * 40, + } + + def test_memory_promote_accepts_qdrant_only_runtime_config(tmp_path, monkeypatch): cfg = _qdrant_runtime_config(tmp_path) store = _FakeVectorStore() diff --git a/harness/tests/test_qdrant_vector_store.py b/harness/tests/test_qdrant_vector_store.py index fe36104d..45b6e334 100644 --- a/harness/tests/test_qdrant_vector_store.py +++ b/harness/tests/test_qdrant_vector_store.py @@ -39,6 +39,7 @@ class FakeQdrantHttp: self.malformed_query = False self.malformed_scroll = False self.scroll_pages: list[dict] | None = None + self.drop_collection_on_points = False def request(self, method, url, *, json=None, timeout=None): self.calls.append((method, url, json)) @@ -75,6 +76,9 @@ class FakeQdrantHttp: return FakeResponse(200, {"status": "ok"}) if method == "PUT" and path == "/collections/workspace-semantic/points": + if self.drop_collection_on_points: + self.collection = None + return FakeResponse(404, {"status": "error"}) for point in json["points"]: self.points[point["id"]] = point return FakeResponse(200, {"result": {"status": "acknowledged"}}) @@ -161,13 +165,16 @@ def _write_record(record_id: str, kind: str, *, metadata=None): ) -def _store(fake: FakeQdrantHttp) -> QdrantVectorStore: +def _store( + fake: FakeQdrantHttp, *, collection_lifecycle: str = "self_heal" +) -> QdrantVectorStore: return QdrantVectorStore( base_url="http://qdrant:6333", collection="workspace-semantic", workspace_id="demo", workspace_revision="a" * 40, expected_dimension=1024, + collection_lifecycle=collection_lifecycle, request=fake.request, ) @@ -212,6 +219,87 @@ def test_upsert_refuses_collection_dimension_or_distance_mismatch_without_recrea assert creates == [] +def test_upsert_require_existing_refuses_missing_collection_without_creating(): + fake = FakeQdrantHttp() + store = _store(fake, collection_lifecycle="require_existing") + + with pytest.raises(VectorStoreError, match="semantic_index_incompatible"): + store.upsert("memory", [_write_record("memory:1", "memory")]) + + assert fake.collection is None + creates = [ + call + for call in fake.calls + if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic") + ] + assert creates == [] + + +def test_upsert_require_existing_refuses_incompatible_collection_without_mutating(): + fake = FakeQdrantHttp(dimension=384, distance="Dot") + fake.collection = {"vectors": {"size": 384, "distance": "Dot"}} + store = _store(fake, collection_lifecycle="require_existing") + + with pytest.raises(VectorStoreError, match="semantic_index_incompatible"): + store.upsert("memory", [_write_record("memory:1", "memory")]) + + assert fake.payload_indexes == set() + mutating = [call for call in fake.calls if call[0] == "PUT"] + assert mutating == [] + + +def test_upsert_require_existing_writes_to_existing_compatible_collection(): + fake = FakeQdrantHttp() + fake.collection = {"vectors": {"size": 1024, "distance": "Cosine"}} + fake.payload_indexes = { + "content_hash", + "document_id", + "kind", + "record_key", + "record_kind", + "vector_generation", + "workspace_id", + "workspace_revision", + } + store = _store(fake, collection_lifecycle="require_existing") + + assert store.upsert("memory", [_write_record("memory:1", "memory")]) == 1 + + create_or_index = [ + call + for call in fake.calls + if call[0] == "PUT" and not call[1].endswith("/points?wait=true") + ] + assert create_or_index == [] + + +def test_upsert_require_existing_fails_if_collection_disappears_after_preflight(): + fake = FakeQdrantHttp() + fake.collection = {"vectors": {"size": 1024, "distance": "Cosine"}} + fake.payload_indexes = { + "content_hash", + "document_id", + "kind", + "record_key", + "record_kind", + "vector_generation", + "workspace_id", + "workspace_revision", + } + fake.drop_collection_on_points = True + store = _store(fake, collection_lifecycle="require_existing") + + with pytest.raises(VectorStoreError, match="HTTP 404"): + store.upsert("memory", [_write_record("memory:1", "memory")]) + + creates = [ + call + for call in fake.calls + if call[0] == "PUT" and call[1].endswith("/collections/workspace-semantic") + ] + assert creates == [] + + def test_health_fails_when_the_bound_collection_is_missing(): fake = FakeQdrantHttp() diff --git a/harness/tests/test_registry_evidence_config.py b/harness/tests/test_registry_evidence_config.py index f972974b..358061fd 100644 --- a/harness/tests/test_registry_evidence_config.py +++ b/harness/tests/test_registry_evidence_config.py @@ -113,6 +113,61 @@ def test_signed_http_file_resolves_in_memory_and_preserves_provenance_order(tmp_ assert_no_canaries(repr(adapter)) +def test_qdrant_runtime_config_keeps_require_existing_and_http_policy_fields(tmp_path): + path = tmp_path / "runtime.yaml" + path.write_text( + yaml.safe_dump( + { + "runtime_identity": { + "workspace_id": "psd-clinical", + "workspace_revision": "a" * 40, + }, + "dwh": { + "type": "postgres_direct", + "connection": { + "database": "analytics", + "schema": "public", + "user": "reader", + "password": "secret", + }, + }, + "vectors": { + "type": "qdrant", + "base_url": "http://qdrant:6333", + "collection": "psd-clinical", + "collection_lifecycle": "require_existing", + }, + "embeddings": { + "provider": "ollama_internal", + "base_url": "http://embedding:11434", + "model": "qwen3-embedding:0.6b", + "dim": 1024, + }, + "evidence": { + "sources": [ + { + "type": "http", + "urls": ["https://evidence.example.test/guide.md"], + "allow_private_hosts": True, + "max_redirects": 2, + "max_cache_bytes": 1234, + } + ] + }, + } + ) + ) + + cfg = load_config(path) + rendered = cfg.model_dump(mode="json") + + assert cfg.vectors.collection_lifecycle == "require_existing" + assert rendered["vectors"]["collection_lifecycle"] == "require_existing" + assert rendered["evidence"]["sources"][0]["allow_private_hosts"] is True + assert rendered["evidence"]["sources"][0]["max_redirects"] == 2 + assert rendered["evidence"]["sources"][0]["max_cache_bytes"] == 1234 + + def test_signed_http_file_requires_explicit_provenance_urls(tmp_path): secret_file = tmp_path / "signed-urls.json" secret_file.write_text(json.dumps([ diff --git a/harness/tests/test_schema_fk_annotations.py b/harness/tests/test_schema_fk_annotations.py index 88c037c8..7a2c5153 100644 --- a/harness/tests/test_schema_fk_annotations.py +++ b/harness/tests/test_schema_fk_annotations.py @@ -1,4 +1,6 @@ -from datetime import datetime +import hashlib +import json +from datetime import UTC, datetime import yaml from typer.testing import CliRunner @@ -15,10 +17,19 @@ from tht.mschema.models import ( ) from tht.mschema.render import to_mschema_text, to_schema_dict +RUNNER = CliRunner() + + +def _json_sha(value) -> str: + payload = json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + return "sha256:" + hashlib.sha256(payload.encode("utf-8")).hexdigest() + def _physical(): return PhysicalSchema( - database="d", schema="s", introspected_at=datetime(2026, 1, 1), + database="d", + schema="s", + introspected_at=datetime(2026, 1, 1, tzinfo=UTC), tables={ "dim_patient": TablePhysical( columns={"cod_paz": ColumnPhysical(type="bigint", pk=True)}, @@ -42,10 +53,16 @@ def _annotations_with_fks(): tables={ "fact_ablazione": TableAnnotation( foreign_keys=[ - ForeignKey(columns=["cod_paz"], ref_table="dim_patient", - ref_columns=["cod_paz"]), - ForeignKey(columns=["data_time_key"], ref_table="dim_time", - ref_columns=["day_key"]), + ForeignKey( + columns=["cod_paz"], + ref_table="dim_patient", + ref_columns=["cod_paz"], + ), + ForeignKey( + columns=["data_time_key"], + ref_table="dim_time", + ref_columns=["day_key"], + ), ], ) } @@ -61,8 +78,11 @@ def test_mschema_text_renders_annotation_fks(): def test_schema_dict_merges_annotation_fks(): d = to_schema_dict(_physical(), _annotations_with_fks()) fks = d["fact_ablazione"]["foreign_keys"] - assert {"columns": ["cod_paz"], "ref_table": "dim_patient", - "ref_columns": ["cod_paz"]} in fks + assert { + "columns": ["cod_paz"], + "ref_table": "dim_patient", + "ref_columns": ["cod_paz"], + } in fks def test_find_orphans_flags_broken_annotation_fk(): @@ -70,10 +90,16 @@ def test_find_orphans_flags_broken_annotation_fk(): tables={ "fact_ablazione": TableAnnotation( foreign_keys=[ - ForeignKey(columns=["cod_paz"], ref_table="dim_sparita", - ref_columns=["x"]), - ForeignKey(columns=["colonna_sparita"], ref_table="dim_time", - ref_columns=["day_key"]), + ForeignKey( + columns=["cod_paz"], + ref_table="dim_sparita", + ref_columns=["x"], + ), + ForeignKey( + columns=["colonna_sparita"], + ref_table="dim_time", + ref_columns=["day_key"], + ), ], ) } @@ -91,22 +117,139 @@ def _write_workspace(tmp_path): _physical().to_yaml(tmp_path / "artifacts" / "mschema" / "physical.yaml") cfg = tmp_path / "workspace.yaml" cfg.write_text( - "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" - f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, sessions: {tmp_path/'s'}}}\n" + f""" +runtime_identity: + workspace_id: demo + workspace_revision: {'a' * 40} +database: {{database: d, schema: s, user: u, password: p, transport: direct}} +paths: {{artifacts: {tmp_path / 'artifacts'}, indexes: {tmp_path / 'i'}, sessions: {tmp_path / 's'}}} +""" ) return cfg def test_suggest_fks_prints_candidates(tmp_path): cfg = _write_workspace(tmp_path) - res = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) + res = RUNNER.invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) assert res.exit_code == 0, res.output data = yaml.safe_load(res.output.rsplit("\n", 2)[0].split("FK candidate")[0]) fks = data["tables"]["fact_ablazione"]["foreign_keys"] - assert {"columns": ["cod_paz"], "ref_table": "dim_patient", - "ref_columns": ["cod_paz"]} in fks - assert {"columns": ["data_time_key"], "ref_table": "dim_time", - "ref_columns": ["day_key"]} in fks + assert { + "columns": ["cod_paz"], + "ref_table": "dim_patient", + "ref_columns": ["cod_paz"], + } in fks + assert { + "columns": ["data_time_key"], + "ref_table": "dim_time", + "ref_columns": ["day_key"], + } in fks + + +def test_suggest_fks_json_is_pristine_and_stable(tmp_path): + cfg = _write_workspace(tmp_path) + first = tmp_path / "second.sql" + first.write_text( + "SELECT f.esito FROM datawarehouse.fact_ablazione f " + "JOIN datawarehouse.dim_patient p ON f.cod_paz = p.cod_paz" + ) + second = tmp_path / "first.sql" + second.write_text( + "SELECT dt.year FROM datawarehouse.fact_ablazione f " + "JOIN datawarehouse.dim_time dt ON f.data_time_key = dt.day_key" + ) + + response = RUNNER.invoke( + app, + [ + "schema", + "suggest-fks", + "-c", + str(cfg), + "--from-sql", + str(first), + "--from-sql", + str(second), + "--json", + ], + ) + + assert response.exit_code == 0, response.output + assert response.stderr == "" + payload = json.loads(response.stdout) + assert payload["schemaVersion"] == 1 + assert payload["status"] == "succeeded" + assert payload["code"] == "ok" + assert payload["operation"] == "schema_suggest_fks" + assert payload["workspaceId"] == "demo" + assert payload["workspaceRevision"] == "a" * 40 + assert payload["counts"] == { + "ambiguousColumns": 0, + "candidateTables": 1, + "candidates": 2, + "minedJoins": 2, + "sqlFiles": 2, + } + assert payload["candidateDocument"] == { + "annotations": { + "tables": { + "fact_ablazione": { + "foreign_keys": [ + { + "columns": ["cod_paz"], + "ref_columns": ["cod_paz"], + "ref_table": "dim_patient", + }, + { + "columns": ["data_time_key"], + "ref_columns": ["day_key"], + "ref_table": "dim_time", + }, + ] + } + } + }, + "counts": {"candidateTables": 1, "candidates": 2}, + "schemaVersion": 1, + } + assert payload["candidate_count"] == payload["counts"]["candidates"] + candidate_yaml = payload["candidate_yaml"] + assert payload["candidateDigest"] == "sha256:" + hashlib.sha256(candidate_yaml.encode("utf-8")).hexdigest() + assert yaml.safe_load(candidate_yaml) == payload["candidateDocument"]["annotations"] + + rerun = RUNNER.invoke( + app, + [ + "schema", + "suggest-fks", + "-c", + str(cfg), + "--from-sql", + str(second), + "--from-sql", + str(first), + "--json", + ], + ) + assert rerun.exit_code == 0, rerun.output + assert json.loads(rerun.stdout) == payload + + +def test_suggest_fks_json_rejects_invalid_assume_without_prose(tmp_path): + cfg = _write_workspace(tmp_path) + + response = RUNNER.invoke( + app, + ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=nope", "--json"], + ) + + assert response.exit_code == 1 + assert response.stderr == "" + payload = json.loads(response.stdout) + assert payload["schemaVersion"] == 1 + assert payload["status"] == "failed" + assert payload["code"] == "invalid_argument" + assert payload["error"] == "invalid assume mapping" def test_mine_join_pairs_from_approved_sql(): @@ -122,23 +265,22 @@ def test_mine_join_pairs_from_approved_sql(): """ pairs = mine_join_pairs(sql, _physical()) assert pairs[("fact_ablazione", "data_time_key", "dim_time", "day_key")] == 1 - # il join CTE-CTE (abl.year=b.year) non produce coppie assert len(pairs) == 1 def test_mine_join_pairs_ignores_non_pk_pairs_and_bad_sql(): from tht.mschema.fkmine import mine_join_pairs - # esito=esito: nessun lato e' PK -> scartato - sql = ("SELECT * FROM fact_ablazione a JOIN fact_ablazione b " - "ON a.esito = b.esito") + sql = "SELECT * FROM fact_ablazione a JOIN fact_ablazione b ON a.esito = b.esito" assert len(mine_join_pairs(sql, _physical())) == 0 assert len(mine_join_pairs("WITH broken (", _physical())) == 0 def test_suggest_fks_skips_generic_and_ambiguous_pks(tmp_path): phys = PhysicalSchema( - database="d", schema="s", introspected_at=datetime(2026, 1, 1), + database="d", + schema="s", + introspected_at=datetime(2026, 1, 1, tzinfo=UTC), tables={ "dim_a": TablePhysical(columns={"id": ColumnPhysical(type="int", pk=True)}), "dim_b": TablePhysical(columns={"id": ColumnPhysical(type="int", pk=True)}), @@ -158,14 +300,14 @@ def test_suggest_fks_skips_generic_and_ambiguous_pks(tmp_path): "database: {database: d, schema: s, user: u, password: p, transport: direct}\n" f"paths: {{artifacts: {tmp_path/'artifacts'}, indexes: {tmp_path/'i'}, sessions: {tmp_path/'s'}}}\n" ) - res = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) + res = RUNNER.invoke(app, ["schema", "suggest-fks", "-c", str(cfg)]) 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 + assert "nessuna FK da suggerire" in res.output + assert "cod_x" in res.output - # --assume disambigua la PK multi-proprietario - res2 = CliRunner().invoke( - app, ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=dim_c1"] + res2 = RUNNER.invoke( + app, + ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=dim_c1"], ) assert res2.exit_code == 0, res2.output yaml_text = "\n".join( @@ -173,16 +315,19 @@ def test_suggest_fks_skips_generic_and_ambiguous_pks(tmp_path): ) data = yaml.safe_load(yaml_text) fact_fks = data["tables"]["fact_f"]["foreign_keys"] - assert {"columns": ["cod_x"], "ref_table": "dim_c1", - "ref_columns": ["cod_x"]} in fact_fks - # dim_c2.cod_x -> dim_c1 (estensione 1:1), ma NON dim_c1 -> se stessa + assert { + "columns": ["cod_x"], + "ref_table": "dim_c1", + "ref_columns": ["cod_x"], + } in fact_fks assert "dim_c1" not in data["tables"] or all( - fk["ref_table"] != "dim_c1" for fk in data["tables"].get("dim_c1", {}).get("foreign_keys", []) + fk["ref_table"] != "dim_c1" + for fk in data["tables"].get("dim_c1", {}).get("foreign_keys", []) ) - # --assume con tabella inesistente -> errore chiaro - res3 = CliRunner().invoke( - app, ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=nope"] + res3 = RUNNER.invoke( + app, + ["schema", "suggest-fks", "-c", str(cfg), "--assume", "cod_x=nope"], ) assert res3.exit_code == 1 assert "non valido" in res3.output @@ -190,20 +335,122 @@ def test_suggest_fks_skips_generic_and_ambiguous_pks(tmp_path): def test_suggest_fks_from_sql_mines_joins(tmp_path): cfg = _write_workspace(tmp_path) - sqldir = tmp_path / "approved" - sqldir.mkdir() - (sqldir / "q1.sql").write_text( + sql_file = tmp_path / "approved.sql" + sql_file.write_text( "SELECT f.esito FROM datawarehouse.fact_ablazione f " "JOIN datawarehouse.dim_patient p ON f.cod_paz = p.cod_paz" ) - res = CliRunner().invoke( - app, ["schema", "suggest-fks", "-c", str(cfg), "--from-sql", str(sqldir)] + res = RUNNER.invoke( + app, + ["schema", "suggest-fks", "-c", str(cfg), "--from-sql", str(sql_file)], ) assert res.exit_code == 0, res.output assert "Minati 1 equi-join da 1 file SQL" in res.output assert "ref_table: dim_patient" in res.output +def test_schema_check_json_validates_staged_annotations_without_mutating_runtime(tmp_path): + cfg = _write_workspace(tmp_path) + runtime_annotations = tmp_path / "artifacts" / "mschema" / "annotations.yaml" + runtime_annotations.write_text("tables: {}\n") + reviewed = tmp_path / "reviewed.yaml" + _annotations_with_fks().to_yaml(reviewed) + + response = RUNNER.invoke( + app, + [ + "schema", + "check", + "-c", + str(cfg), + "--annotations", + str(reviewed), + "--reviewed-candidates", + "sha256:" + "b" * 64, + "--json", + ], + ) + + assert response.exit_code == 0, response.output + assert response.stderr == "" + payload = json.loads(response.stdout) + assert payload["orphan_count"] == 0 + assert payload["reviewed_candidates_digest"] == "sha256:" + "b" * 64 + assert payload["annotations_digest"] == "sha256:" + hashlib.sha256(reviewed.read_bytes()).hexdigest() + assert payload == { + "annotationsDigest": _json_sha( + { + "annotations": { + "tables": { + "fact_ablazione": { + "foreign_keys": [ + { + "columns": ["cod_paz"], + "ref_columns": ["cod_paz"], + "ref_table": "dim_patient", + }, + { + "columns": ["data_time_key"], + "ref_columns": ["day_key"], + "ref_table": "dim_time", + }, + ] + } + } + }, + "schemaVersion": 1, + } + ), + "annotations_digest": "sha256:" + hashlib.sha256(reviewed.read_bytes()).hexdigest(), + "code": "ok", + "orphan_count": 0, + "reviewed_candidates_digest": "sha256:" + "b" * 64, + "counts": {"annotationTables": 1, "foreignKeys": 2, "orphans": 0}, + "operation": "schema_check", + "orphans": [], + "reviewedCandidates": "sha256:" + "b" * 64, + "schemaVersion": 1, + "status": "succeeded", + "workspaceId": "demo", + "workspaceRevision": "a" * 40, + "zeroOrphans": True, + } + assert runtime_annotations.read_text() == "tables: {}\n" + + +def test_schema_check_json_reports_orphans_without_prose(tmp_path): + cfg = _write_workspace(tmp_path) + reviewed = tmp_path / "reviewed.yaml" + Annotations( + tables={ + "fact_ablazione": TableAnnotation( + foreign_keys=[ + ForeignKey( + columns=["cod_paz"], + ref_table="dim_missing", + ref_columns=["cod_paz"], + ) + ] + ) + } + ).to_yaml(reviewed) + + response = RUNNER.invoke( + app, + ["schema", "check", "-c", str(cfg), "--annotations", str(reviewed), "--json"], + ) + + assert response.exit_code == 3 + assert response.stderr == "" + payload = json.loads(response.stdout) + assert payload["schemaVersion"] == 1 + assert payload["status"] == "blocked" + assert payload["code"] == "annotation_invalid" + assert payload["zeroOrphans"] is False + assert payload["counts"]["orphans"] == 1 + assert payload["orphans"] == ["fact_ablazione.fk(cod_paz)->dim_missing"] + + def test_suggest_fks_write_merges_and_is_idempotent(tmp_path): cfg = _write_workspace(tmp_path) ann_path = tmp_path / "artifacts" / "mschema" / "annotations.yaml" @@ -211,13 +458,13 @@ def test_suggest_fks_write_merges_and_is_idempotent(tmp_path): tables={"fact_ablazione": TableAnnotation(description="Ablazioni")} ).to_yaml(ann_path) - res = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"]) + res = RUNNER.invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"]) assert res.exit_code == 0, res.output ann = Annotations.from_yaml(ann_path) - assert ann.tables["fact_ablazione"].description == "Ablazioni" # non distrutta + assert ann.tables["fact_ablazione"].description == "Ablazioni" assert len(ann.tables["fact_ablazione"].foreign_keys) == 2 - res2 = CliRunner().invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"]) + res2 = RUNNER.invoke(app, ["schema", "suggest-fks", "-c", str(cfg), "--write"]) assert "nessuna FK da suggerire" in res2.output ann2 = Annotations.from_yaml(ann_path) assert len(ann2.tables["fact_ablazione"].foreign_keys) == 2 diff --git a/harness/tht/adapters/factory.py b/harness/tht/adapters/factory.py index 3d594ed7..a9f52936 100644 --- a/harness/tht/adapters/factory.py +++ b/harness/tht/adapters/factory.py @@ -38,6 +38,7 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto workspace_id=cfg._workspace_id, workspace_revision=cfg._workspace_revision, expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None, + collection_lifecycle=resource.collection_lifecycle, ) case other: # pragma: no cover - Pydantic's discriminator rejects this first. raise ConfigError(f"Adapter vector non supportato: {other}") diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index f95611f6..12513a86 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -55,6 +55,7 @@ class QdrantVectorStore: workspace_id: str, workspace_revision: str | None = None, expected_dimension: int | None = None, + collection_lifecycle: str = "self_heal", request: Callable[..., object] | None = None, connect_timeout: float = 2.0, read_timeout: float = 10.0, @@ -64,6 +65,7 @@ class QdrantVectorStore: self._workspace_id = workspace_id self._workspace_revision = workspace_revision self._expected_dimension = expected_dimension + self._collection_lifecycle = collection_lifecycle self._request = request or requests.request self._timeout = (connect_timeout, read_timeout) @@ -306,6 +308,8 @@ class QdrantVectorStore: if response is None: if not strict: raise VectorStoreError("Qdrant collection is missing") + if self._collection_lifecycle == "require_existing": + raise VectorStoreError("semantic_index_incompatible") self._call( "PUT", f"/collections/{self._collection}", @@ -328,9 +332,13 @@ class QdrantVectorStore: self._expected_dimension is not None and (size != self._expected_dimension or distance != "Cosine") ): + if strict and self._collection_lifecycle == "require_existing": + raise VectorStoreError("semantic_index_incompatible") raise VectorStoreError("Qdrant collection configuration mismatch") for field_name in _KEYWORD_INDEXES: if field_name not in result.get("payload_schema", {}): + if strict and self._collection_lifecycle == "require_existing": + raise VectorStoreError("semantic_index_incompatible") if not strict: raise VectorStoreError("Qdrant collection payload indexes mismatch") self._call( diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index dc937737..3bd008c6 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -2,27 +2,57 @@ from __future__ import annotations +import hashlib import json import re -import hashlib from pathlib import Path import typer from tht.cli.config_cmd import CONFIG_OPT -from tht.config import workspace_id_from_path - preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts") +def _evidence_json_context(config: Path): + from tht.cli.schema_cmd import _load_config_or_exit + + return _load_config_or_exit(config) + + +def _evidence_json_payload(cfg, payload: dict, *, code: str, error: str | None = None) -> dict: + value = { + **payload, + "schemaVersion": 1, + "status": payload.get("status", "failed"), + "code": code, + "operation": "preprocess_evidence", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + } + if error is not None: + value["error"] = error + return value + + +def _simple_json_payload(*, code: str, error: str) -> dict: + return { + "schemaVersion": 1, + "status": "failed", + "code": code, + "operation": "preprocess_evidence", + "error": error, + } + + def run_dwh_from_config( config: Path, *, steps: tuple[str, ...], resume: str | None = None, ): from tht.cli.lsh_cmd import build_lsh_artifacts from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog from tht.jobs.dwh_pipeline import ( - DwhPreprocessPipeline, config_dwh_binding, + DwhPreprocessPipeline, + config_dwh_binding, ) cfg = _load_config_or_exit(config) @@ -87,7 +117,7 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = return "sha256:" + hashlib.sha256(value.encode()).hexdigest() return pipeline.run_as_job( - workspace_id=workspace_id_from_path(config), + workspace_id=cfg._workspace_id, workspace_root=corpus_root.parent, config_fingerprint=fingerprint(cfg.model_dump_json()), input_fingerprint=fingerprint(config.resolve().as_posix()), @@ -116,7 +146,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False): pipeline_version="evidence-v1", retain_published_generations=cfg.vector.retain_published_generations, ) - pipeline.workspace_id = workspace_id_from_path(config) + pipeline.workspace_id = cfg._workspace_id return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run) @@ -133,7 +163,7 @@ def evidence_cmd( if action == "gc": try: payload = gc_from_config(config, dry_run=dry_run) - except Exception: + except Exception: # noqa: BLE001 payload = {"status": "failed", "error": "evidence cleanup failed"} if json_output: typer.echo(json.dumps(payload, sort_keys=True)) @@ -145,8 +175,12 @@ def evidence_cmd( else: typer.echo(f"OK: evicted={len(payload['evicted'])} failures={len(payload['failures'])}") return + cfg = _evidence_json_context(config) if json_output else None if resume is not None and re.fullmatch(r"[0-9a-f]{32}", resume) is None: - payload = {"status": "failed", "error": "resume requires a preprocessing run id"} + payload = _simple_json_payload( + code="invalid_resume", + error="resume requires a preprocessing run id", + ) if json_output: typer.echo(json.dumps(payload, sort_keys=True)) else: @@ -154,23 +188,43 @@ def evidence_cmd( raise typer.Exit(code=2) try: result = run_from_config(config, dry_run=dry_run, resume=resume) - except Exception: - payload = {"status": "failed", "error": "preprocessing failed"} + except Exception: # noqa: BLE001 + payload = {"status": "failed"} if json_output: - typer.echo(json.dumps(payload, sort_keys=True)) + typer.echo(json.dumps( + _evidence_json_payload( + cfg, + payload, + code="preprocessing_failed", + error="preprocessing failed", + ), + sort_keys=True, + )) else: typer.secho("ERRORE: preprocessing failed", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) from None payload = result.model_dump(mode="json") if payload.get("status") != "succeeded": - payload["error"] = "preprocessing job failed" if json_output: - typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) + typer.echo(json.dumps( + _evidence_json_payload( + cfg, + payload, + code="preprocessing_failed", + error="preprocessing job failed", + ), + ensure_ascii=False, + sort_keys=True, + )) else: typer.secho("ERRORE: preprocessing job failed", fg=typer.colors.RED, err=True) raise typer.Exit(code=1) if json_output: - typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) + typer.echo(json.dumps( + _evidence_json_payload(cfg, payload, code="ok"), + ensure_ascii=False, + sort_keys=True, + )) else: counts = payload["counts"] typer.echo( @@ -205,7 +259,7 @@ def dwh_cmd( raise typer.Exit(code=2) try: result = run_dwh_from_config(config, steps=selected, resume=resume) - except Exception: + except Exception: # noqa: BLE001 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 9110b40c..742ae146 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -1,7 +1,11 @@ -from pathlib import Path +import hashlib +import json import logging +from pathlib import Path import typer +import yaml + from tht.adapters.factory import build_dwh from tht.cli.config_cmd import CONFIG_OPT from tht.config import ConfigError, load_config @@ -21,7 +25,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 logger.warning("Campionamento saltato per %s.%s: %s", table_name, column_name, exc) continue @@ -83,7 +87,7 @@ def introspect_cmd( try: cached = PhysicalSchema.from_yaml(out) - except Exception: + except Exception: # noqa: BLE001,S110 pass # catalogo illeggibile: procedi con la re-introspezione else: ts = cached.introspected_at @@ -105,7 +109,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 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()) @@ -119,8 +123,188 @@ def introspect_cmd( ) +def _json_sha(value) -> str: + payload = json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + return "sha256:" + hashlib.sha256(payload.encode("utf-8")).hexdigest() + + +def _emit_json(payload: dict) -> None: + typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) + + +def _sorted_fk_payloads(foreign_keys) -> list[dict]: + payloads = [fk.model_dump(mode="json", exclude_defaults=True) for fk in foreign_keys] + return sorted( + payloads, + key=lambda payload: ( + tuple(payload.get("columns", [])), + payload.get("ref_table", ""), + tuple(payload.get("ref_columns", [])), + payload.get("name", ""), + ), + ) + + +def _sorted_annotations_payload(annotations) -> dict: + tables = {} + for table_name in sorted(annotations.tables): + table = annotations.tables[table_name] + payload = {} + if table.description: + payload["description"] = table.description + if table.concepts: + payload["concepts"] = table.concepts + if table.notes: + payload["notes"] = table.notes + if table.columns: + payload["columns"] = { + name: value.model_dump(mode="json", exclude_defaults=True) + for name, value in sorted(table.columns.items()) + } + if table.foreign_keys: + payload["foreign_keys"] = _sorted_fk_payloads(table.foreign_keys) + tables[table_name] = payload + return {"tables": tables} + + +def _suggested_fk_payload(annotations_by_table: dict) -> dict: + return { + "tables": { + table_name: {"foreign_keys": _sorted_fk_payloads(foreign_keys)} + for table_name, foreign_keys in sorted(annotations_by_table.items()) + } + } + + +def _load_sql_inputs(entries: list[Path] | None) -> list[tuple[str, str]]: + max_file_bytes = 1024 * 1024 + max_total_bytes = 16 * 1024 * 1024 + total_bytes = 0 + sql_files: list[Path] = [] + for entry in entries or []: + if entry.is_dir(): + sql_files.extend(sorted(path for path in entry.rglob("*.sql") if path.is_file())) + continue + sql_files.append(entry) + loaded = [] + for sql_file in sorted(sql_files, key=lambda candidate: candidate.as_posix()): + if not sql_file.exists() or not sql_file.is_file() or sql_file.is_symlink(): + raise ValueError("invalid SQL input") + size = sql_file.stat().st_size + total_bytes += size + if size > max_file_bytes or total_bytes > max_total_bytes: + raise ValueError("invalid SQL input") + loaded.append((sql_file.as_posix(), sql_file.read_text(encoding="utf-8"))) + return loaded + + +def _suggest_fk_result(physical, annotations, *, sql_inputs: list[tuple[str, str]], assume: list[str] | None): + from tht.mschema.fkmine import mine_join_pairs + from tht.mschema.models import ForeignKey + + assumed: dict[str, str] = {} + for value in assume or []: + col, _, ref = value.partition("=") + if not ref or ref not in physical.tables: + raise ValueError("invalid assume mapping") + assumed[col] = ref + + def _single_pk(table) -> str | None: + pks = [column_name for column_name, column in table.columns.items() if column.pk] + return pks[0] if len(pks) == 1 else None + + pk_owners: dict[str, list[str]] = {} + for table_name, table in physical.tables.items(): + pk = _single_pk(table) + if pk: + pk_owners.setdefault(pk, []).append(table_name) + + dim_time_pk = None + if "dim_time" in physical.tables: + dim_time_pk = _single_pk(physical.tables["dim_time"]) + + def _known(table_name: str) -> set: + keys = set() + for fk in physical.tables[table_name].foreign_keys: + keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))) + annotation = annotations.tables.get(table_name) + if annotation: + for fk in annotation.foreign_keys: + keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))) + return keys + + known_by_table: dict[str, set] = {table_name: _known(table_name) for table_name in physical.tables} + suggested: dict[str, list[ForeignKey]] = {} + + def _add(table_name: str, column_name: str, ref_table: str, ref_column: str) -> None: + key = ((column_name,), ref_table, (ref_column,)) + if key in known_by_table[table_name]: + return + known_by_table[table_name].add(key) + suggested.setdefault(table_name, []).append( + ForeignKey(columns=[column_name], ref_table=ref_table, ref_columns=[ref_column]) + ) + + mined_total = 0 + for _name, sql_text in sql_inputs: + pairs = mine_join_pairs(sql_text, physical) + mined_total += sum(pairs.values()) + for src_t, src_c, ref_t, ref_c in pairs: + _add(src_t, src_c, ref_t, ref_c) + + ambiguous_skipped: set[str] = set() + for table_name, table in physical.tables.items(): + for column_name in table.columns: + if dim_time_pk and column_name.endswith("time_key") and table_name != "dim_time": + _add(table_name, column_name, "dim_time", dim_time_pk) + continue + if column_name in assumed: + if assumed[column_name] != table_name: + _add(table_name, column_name, assumed[column_name], column_name) + continue + owners = [owner for owner in pk_owners.get(column_name, []) if owner != table_name] + if not owners or column_name in _GENERIC_PK_NAMES: + continue + if len(pk_owners[column_name]) > 1: + ambiguous_skipped.add(column_name) + continue + _add(table_name, column_name, owners[0], column_name) + + candidate_annotations = _suggested_fk_payload(suggested) + candidate_yaml = yaml.safe_dump(candidate_annotations, sort_keys=False, allow_unicode=True) + counts = { + "ambiguousColumns": len(ambiguous_skipped), + "candidateTables": len(candidate_annotations["tables"]), + "candidates": sum(len(value["foreign_keys"]) for value in candidate_annotations["tables"].values()), + "minedJoins": mined_total, + "sqlFiles": len(sql_inputs), + } + candidate_document = { + "annotations": candidate_annotations, + "counts": { + "candidateTables": counts["candidateTables"], + "candidates": counts["candidates"], + }, + "schemaVersion": 1, + } + return { + "ambiguous": sorted(ambiguous_skipped), + "candidate_count": counts["candidates"], + "candidateDigest": "sha256:" + hashlib.sha256(candidate_yaml.encode("utf-8")).hexdigest(), + "candidateDocument": candidate_document, + "candidate_yaml": candidate_yaml, + "counts": counts, + "suggested": suggested, + } + + @schema_app.command("check") -def check_cmd(config: Path = CONFIG_OPT) -> None: +def check_cmd( + config: Path = CONFIG_OPT, + annotations: Path | None = typer.Option(None, "--annotations"), # noqa: B008 + reviewed_candidates: str | None = typer.Option(None, "--reviewed-candidates"), + json_output: bool = typer.Option(False, "--json"), +) -> None: """Confronta physical.yaml e annotations.yaml; segnala annotazioni orfane.""" from tht.mschema.merge import find_orphans from tht.mschema.models import Annotations, PhysicalSchema @@ -128,20 +312,80 @@ def check_cmd(config: Path = CONFIG_OPT) -> None: cfg = _load_config_or_exit(config) phys_file = physical_path(cfg) if not phys_file.exists(): - typer.secho( - f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", - fg=typer.colors.RED, err=True, - ) + if json_output: + _emit_json({ + "code": "schema_missing", + "error": "physical schema is missing", + "operation": "schema_check", + "schemaVersion": 1, + "status": "failed", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + else: + typer.secho( + f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", + fg=typer.colors.RED, + err=True, + ) raise typer.Exit(code=1) physical = PhysicalSchema.from_yaml(phys_file) - annotations = Annotations.from_yaml(annotations_path(cfg)) + annotations_file = annotations or annotations_path(cfg) + try: + loaded_annotations = Annotations.from_yaml(annotations_file) + except Exception: # noqa: BLE001 + if json_output: + _emit_json({ + "code": "annotation_invalid", + "error": "annotations are invalid", + "operation": "schema_check", + "schemaVersion": 1, + "status": "failed", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + else: + typer.secho("ERRORE: annotations non valide.", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None ignored = [ - f"{t}.{c} ({col.eligibility_reason})" - for t, table in physical.tables.items() - for c, col in table.columns.items() - if not col.eligible + 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 ] + orphans = sorted(find_orphans(physical, loaded_annotations)) + if json_output: + annotations_payload = _sorted_annotations_payload(loaded_annotations) + payload = { + "annotationsDigest": _json_sha({"annotations": annotations_payload, "schemaVersion": 1}), + "code": "ok" if not orphans else "annotation_invalid", + "counts": { + "annotationTables": len(annotations_payload["tables"]), + "foreignKeys": sum( + len(table_payload.get("foreign_keys", [])) + for table_payload in annotations_payload["tables"].values() + ), + "orphans": len(orphans), + }, + "operation": "schema_check", + "orphan_count": len(orphans), + "orphans": orphans, + "schemaVersion": 1, + "status": "succeeded" if not orphans else "blocked", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + "zeroOrphans": not orphans, + } + payload["annotations_digest"] = "sha256:" + hashlib.sha256(Path(annotations_file).read_bytes()).hexdigest() + if reviewed_candidates is not None: + payload["reviewedCandidates"] = reviewed_candidates + payload["reviewed_candidates_digest"] = reviewed_candidates + _emit_json(payload) + if orphans: + raise typer.Exit(code=3) + return + if ignored: typer.secho( f"Colonne ignorate (testo ampio, {len(ignored)}):", fg=typer.colors.YELLOW @@ -149,11 +393,10 @@ def check_cmd(config: Path = CONFIG_OPT) -> None: for line in ignored: typer.echo(f" - {line}") - orphans = find_orphans(physical, annotations) if orphans: typer.secho(f"ATTENZIONE: {len(orphans)} annotazioni orfane:", fg=typer.colors.YELLOW) - for o in orphans: - typer.echo(f" - {o}") + for orphan in orphans: + typer.echo(f" - {orphan}") raise typer.Exit(code=3) typer.secho("OK: nessuna annotazione orfana.", fg=typer.colors.GREEN) @@ -166,11 +409,11 @@ _GENERIC_PK_NAMES = {"id", "key", "code"} @schema_app.command("suggest-fks") def suggest_fks_cmd( config: Path = CONFIG_OPT, - from_sql: list[Path] = typer.Option( + from_sql: list[Path] = typer.Option( # noqa: B008 None, "--from-sql", - help="Directory di .sql approvati da cui minare i join reali (ripetibile).", + help="Directory o file .sql approvati da cui minare i join reali (ripetibile).", ), - assume: list[str] = typer.Option( + assume: list[str] = typer.Option( # noqa: B008 None, "--assume", help="Disambigua una PK con piu' proprietari: col=tabella_ref " "(es. cod_paz=dim_patient). Ripetibile.", @@ -179,132 +422,103 @@ def suggest_fks_cmd( False, "--write", help="Fonde i suggerimenti in annotations.yaml (aggiunge solo FK mancanti).", ), + json_output: bool = typer.Option(False, "--json"), ) -> None: - """Suggerisce FK logiche per la curazione umana in annotations.yaml. - - Tre regole, in ordine di confidenza: (1) equi-join minati dall'SQL gia' - approvato (--from-sql); (2) colonna `*time_key` verso la PK di dim_time; - (3) colonna con lo stesso nome della PK di UN'ALTRA tabella, solo se quel - nome ha un unico proprietario e non e' generico (id/key/code) — salvo - disambiguazione esplicita con --assume. - """ + """Suggerisce FK logiche per la curazione umana in annotations.yaml.""" import yaml as _yaml - from tht.mschema.fkmine import mine_join_pairs - from tht.mschema.models import Annotations, ForeignKey, PhysicalSchema, TableAnnotation + from tht.mschema.models import Annotations, PhysicalSchema, TableAnnotation cfg = _load_config_or_exit(config) phys_file = physical_path(cfg) if not phys_file.exists(): - typer.secho( - f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", - fg=typer.colors.RED, err=True, - ) + if json_output: + _emit_json({ + "code": "schema_missing", + "error": "physical schema is missing", + "operation": "schema_suggest_fks", + "schemaVersion": 1, + "status": "failed", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + else: + typer.secho( + f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", + fg=typer.colors.RED, + err=True, + ) raise typer.Exit(code=1) physical = PhysicalSchema.from_yaml(phys_file) ann_path = annotations_path(cfg) - annotations = Annotations.from_yaml(ann_path) - - assumed: dict[str, str] = {} - for a in assume or []: - col, _, ref = a.partition("=") - if not ref or ref not in physical.tables: - typer.secho( - f"ERRORE: --assume '{a}' non valido (atteso col=tabella nel catalogo).", - fg=typer.colors.RED, err=True, + loaded_annotations = Annotations.from_yaml(ann_path) + try: + sql_inputs = _load_sql_inputs(from_sql) + result = _suggest_fk_result(physical, loaded_annotations, sql_inputs=sql_inputs, assume=assume) + except ValueError as exc: + code = "invalid_argument" + error = str(exc) + if json_output: + _emit_json({ + "code": code, + "error": error, + "operation": "schema_suggest_fks", + "schemaVersion": 1, + "status": "failed", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + else: + human_error = ( + "--assume non valido (atteso col=tabella nel catalogo)." + if error == "invalid assume mapping" + else error ) - raise typer.Exit(code=1) - assumed[col] = ref + typer.secho(f"ERRORE: {human_error}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None - def _single_pk(table) -> str | None: - pks = [c for c, col in table.columns.items() if col.pk] - return pks[0] if len(pks) == 1 else None + if json_output: + _emit_json({ + "candidate_count": result["candidate_count"], + "candidateDigest": result["candidateDigest"], + "candidateDocument": result["candidateDocument"], + "candidate_yaml": result["candidate_yaml"], + "code": "ok", + "counts": result["counts"], + "operation": "schema_suggest_fks", + "schemaVersion": 1, + "status": "succeeded", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + return - pk_owners: dict[str, list[str]] = {} - for tname, table in physical.tables.items(): - pk = _single_pk(table) - if pk: - pk_owners.setdefault(pk, []).append(tname) - - dim_time_pk = None - if "dim_time" in physical.tables: - dim_time_pk = _single_pk(physical.tables["dim_time"]) - - def _known(tname: str) -> set: - keys = set() - for fk in physical.tables[tname].foreign_keys: - keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))) - ann = annotations.tables.get(tname) - if ann: - for fk in ann.foreign_keys: - keys.add((tuple(fk.columns), fk.ref_table, tuple(fk.ref_columns))) - return keys - - known_by_table: dict[str, set] = {t: _known(t) for t in physical.tables} - suggested: dict[str, list[ForeignKey]] = {} - - def _add(tname: str, col: str, ref_table: str, ref_col: str) -> None: - key = ((col,), ref_table, (ref_col,)) - if key in known_by_table[tname]: - return - known_by_table[tname].add(key) - suggested.setdefault(tname, []).append( - ForeignKey(columns=[col], ref_table=ref_table, ref_columns=[ref_col]) - ) - - # Regola 1: join minati dall'SQL approvato. - n_sql_files = 0 - mined_total = 0 - for d in from_sql or []: - for sql_file in sorted(d.rglob("*.sql")): - n_sql_files += 1 - pairs = mine_join_pairs(sql_file.read_text(), physical) - mined_total += sum(pairs.values()) - for (src_t, src_c, ref_t, ref_c) in pairs: - _add(src_t, src_c, ref_t, ref_c) - - # Regole 2 e 3: convenzioni di naming. - ambiguous_skipped: set[str] = set() - for tname, table in physical.tables.items(): - for cname in table.columns: - if dim_time_pk and cname.endswith("time_key") and tname != "dim_time": - _add(tname, cname, "dim_time", dim_time_pk) - continue - if cname in assumed: - if assumed[cname] != tname: - _add(tname, cname, assumed[cname], cname) - continue - owners = [o for o in pk_owners.get(cname, []) if o != tname] - if not owners or cname in _GENERIC_PK_NAMES: - continue - if len(pk_owners[cname]) > 1: - ambiguous_skipped.add(cname) - continue - _add(tname, cname, owners[0], cname) - - if n_sql_files: + if result["counts"]["sqlFiles"]: typer.secho( - f"Minati {mined_total} equi-join da {n_sql_files} file SQL.", - fg=typer.colors.BLUE, err=True, + f"Minati {result['counts']['minedJoins']} equi-join da {result['counts']['sqlFiles']} file SQL.", + fg=typer.colors.BLUE, + err=True, ) - if ambiguous_skipped: + if result["ambiguous"]: typer.secho( "PK ambigue saltate dalla regola same-name (piu' tabelle proprietarie): " - + ", ".join(sorted(ambiguous_skipped)) + + ", ".join(result["ambiguous"]) + ". Se servono, aggiungile a mano o passa --from-sql.", - fg=typer.colors.YELLOW, err=True, + fg=typer.colors.YELLOW, + err=True, ) - n_fks = sum(len(v) for v in suggested.values()) + suggested = result["suggested"] + n_fks = result["counts"]["candidates"] if not suggested: typer.secho("OK: nessuna FK da suggerire.", fg=typer.colors.GREEN) return if write: - for tname, fks in suggested.items(): - ann = annotations.tables.setdefault(tname, TableAnnotation()) - ann.foreign_keys.extend(fks) - annotations.to_yaml(ann_path) + for table_name, foreign_keys in suggested.items(): + annotation = loaded_annotations.tables.setdefault(table_name, TableAnnotation()) + annotation.foreign_keys.extend(foreign_keys) + loaded_annotations.to_yaml(ann_path) typer.secho( f"OK: {n_fks} FK suggerite aggiunte a {ann_path} " f"({len(suggested)} tabelle). Rivedile a mano prima dell'uso.", @@ -312,13 +526,13 @@ def suggest_fks_cmd( ) return - payload = { - "tables": { - tname: {"foreign_keys": [fk.model_dump(exclude_defaults=True) for fk in fks]} - for tname, fks in suggested.items() - } - } - typer.echo(_yaml.safe_dump(payload, sort_keys=False, allow_unicode=True)) + typer.echo( + _yaml.safe_dump( + result["candidateDocument"]["annotations"], + sort_keys=False, + allow_unicode=True, + ) + ) typer.secho( f"{n_fks} FK candidate ({len(suggested)} tabelle). " f"Usa --write per fonderle in annotations.yaml, poi curale a mano.", @@ -332,10 +546,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 2bbe7578..3b3754e4 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -1,3 +1,5 @@ +import hashlib +import json from pathlib import Path import typer @@ -11,6 +13,14 @@ from tht.vectorstore.store import SyncStats, content_hash vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)") +def _artifact_digest(path: Path) -> str: + return "sha256:" + hashlib.sha256(path.read_bytes()).hexdigest() + + +def _emit_json(payload: dict) -> None: + typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) + + def make_embedder(embeddings_cfg): """Factory del client embeddings (monkeypatchabile nei test).""" from tht.vectorstore.embeddings import OllamaEmbeddings @@ -117,9 +127,14 @@ def init_cmd( @vector_app.command("index-schema") -def index_schema_cmd(config: Path = CONFIG_OPT) -> None: +def index_schema_cmd( + config: Path = CONFIG_OPT, + json_output: bool = typer.Option(False, "--json"), +) -> None: """Embedda e sincronizza i record schema (tabelle e colonne) nel semantic store.""" + from tht.adapters.factory import build_vector_store from tht.mschema.models import Annotations, PhysicalSchema + from tht.ports.vector import VectorStoreError from tht.vectorstore.records import schema_records cfg = _load_config_or_exit(config) @@ -127,20 +142,69 @@ def index_schema_cmd(config: Path = CONFIG_OPT) -> None: require_vector_cfg(cfg) phys_file = physical_path(cfg) if not phys_file.exists(): - typer.secho( - f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`.", - fg=typer.colors.RED, err=True, - ) + message = f"ERRORE: {phys_file} non trovato. Esegui prima `tht schema introspect`." + if json_output: + _emit_json({ + "code": "schema_missing", + "error": "physical schema is missing", + "operation": "index_schema", + "schemaVersion": 1, + "status": "failed", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + else: + typer.secho(message, fg=typer.colors.RED, err=True) raise typer.Exit(code=1) physical = PhysicalSchema.from_yaml(phys_file) - annotations = Annotations.from_yaml(annotations_path(cfg)) + annotations_file = annotations_path(cfg) + annotations = Annotations.from_yaml(annotations_file) records = schema_records(physical, annotations) - from tht.adapters.factory import build_vector_store - - stats = sync_canonical_records( - "schema_records", - records, - store=build_vector_store(cfg, require_write=True), - embedder=make_embedder(cfg.embeddings), - ) + try: + stats = sync_canonical_records( + "schema_records", + records, + store=build_vector_store(cfg, require_write=True), + embedder=make_embedder(cfg.embeddings), + ) + except VectorStoreError as exc: + code = str(exc) + error = "semantic index incompatible" if code == "semantic_index_incompatible" else "schema indexing failed" + if json_output: + _emit_json({ + "code": code, + "error": error, + "operation": "index_schema", + "schemaVersion": 1, + "status": "failed", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + else: + typer.secho(f"ERRORE: {error}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None + if json_output: + _emit_json({ + "artifactIdentities": [ + {"digest": _artifact_digest(annotations_file), "kind": "schema_annotations"}, + {"digest": _artifact_digest(phys_file), "kind": "physical_schema"}, + ], + "code": "ok", + "collection": cfg.vectors.collection, + "counts": { + "added": stats.added, + "columns": sum(len(table.columns) for table in physical.tables.values()), + "deleted": stats.deleted, + "records": len(records), + "tables": len(physical.tables), + "unchanged": stats.unchanged, + "updated": stats.updated, + }, + "operation": "index_schema", + "schemaVersion": 1, + "status": "succeeded", + "workspaceId": cfg._workspace_id, + "workspaceRevision": cfg._workspace_revision, + }) + return _print_stats(stats) diff --git a/harness/tht/config.py b/harness/tht/config.py index 6916d13d..877fe7c6 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -222,6 +222,7 @@ class QdrantConfig(BaseModel): type: Literal["qdrant"] base_url: str collection: str = Field(min_length=1) + collection_lifecycle: Literal["self_heal", "require_existing"] = "self_heal" VectorResourceConfig = Annotated[