Files
ThothII/harness/tests/test_p3_dwh_binding.py
T
Codex cffa60772e
Publish documentation / publish (push) Successful in 2m12s
feat: complete catalog-driven preprocessing
2026-09-06 17:49:35 +02:00

192 lines
7.0 KiB
Python

import pytest
def _cfg(tmp_path, **overrides):
base = {
"dwh": {
"type": "thoth_rest",
"database": {"database": "warehouse", "schema": "dw"},
"endpoint": {"base_url": "http://dwh.example.invalid", "api_key": "secret"},
},
"vectors": {"type": "qdrant", "base_url": "http://qdrant:6333", "collection": "psd", "collection_lifecycle": "self_heal"},
"embeddings": {"provider": "ollama_internal", "base_url": "http://embedding:11434", "model": "qwen3-embedding:0.6b", "dimensions": 1024},
"roots": {"artifacts": str(tmp_path / "artifacts"), "indexes": str(tmp_path / "indexes")},
"paths": {"artifacts": str(tmp_path / "artifacts"), "indexes": str(tmp_path / "indexes"), "sessions": str(tmp_path / "sessions")},
}
base.update(overrides)
return base
def _write_cfg(tmp_path, raw):
import yaml as _yaml
from tht.config import load_config
path = tmp_path / "config.yaml"
path.write_text(_yaml.safe_dump(raw), encoding="utf-8")
cfg = load_config(path)
cfg._workspace_id = "psd"
cfg._config_source = "workspace://psd"
return cfg
@pytest.fixture()
def binding(tmp_path):
from tht.jobs.dwh_pipeline import config_dwh_binding
return config_dwh_binding(_write_cfg(tmp_path, _cfg(tmp_path)))
def test_binding_has_versioned_fingerprints(binding):
assert binding["workspace_id"] == "psd"
assert binding["config_fingerprint"].startswith("sha256:")
assert binding["input_fingerprint"].startswith("sha256:")
assert len(binding["config_fingerprint"]) == 71
def test_content_only_change_keeps_binding(tmp_path):
from tht.jobs.dwh_pipeline import config_dwh_binding
cfg1 = _write_cfg(tmp_path, _cfg(tmp_path))
import copy
raw = copy.deepcopy(_cfg(tmp_path))
raw["vectors"]["collection_lifecycle"] = "require_existing"
cfg3 = _write_cfg(tmp_path, raw)
assert config_dwh_binding(cfg1)["config_fingerprint"] == config_dwh_binding(cfg3)["config_fingerprint"]
def test_endpoint_change_changes_binding(tmp_path):
from tht.jobs.dwh_pipeline import config_dwh_binding
cfg1 = _write_cfg(tmp_path, _cfg(tmp_path))
raw = _cfg(tmp_path)
raw["dwh"] = {"type": "thoth_rest", "database": {"database": "warehouse", "schema": "dw"}, "endpoint": {"base_url": "http://other.example.invalid", "api_key": "secret"}}
cfg2 = _write_cfg(tmp_path, raw)
assert config_dwh_binding(cfg1)["config_fingerprint"] != config_dwh_binding(cfg2)["config_fingerprint"]
def test_catalog_database_and_metadata_revision_are_explicit_lsh_binding_inputs(tmp_path):
import json
from tht.jobs.dwh_pipeline import config_dwh_binding
snapshot = tmp_path / "catalog-metadata.json"
snapshot.write_text(json.dumps({
"schemaVersion": 1,
"workspaceId": "psd",
"databaseId": "database-1",
"databaseName": "warehouse",
"schemaName": "dw",
"metadataContentRevision": 7,
"tables": [],
"relationships": [],
}), encoding="utf-8")
raw = _cfg(tmp_path)
raw["paths"]["catalog_metadata_snapshot"] = str(snapshot)
binding = config_dwh_binding(_write_cfg(tmp_path, raw))
assert binding["workspace_id"] == "psd"
assert binding["catalog_database_id"] == "database-1"
assert binding["metadata_content_revision"] == 7
def test_catalog_pipeline_materializes_schema_and_bound_lsh_generation(monkeypatch, tmp_path):
import json
from tht.cli.preprocess_cmd import run_catalog_dwh_from_config
snapshot = tmp_path / "catalog-metadata.json"
snapshot.write_text(json.dumps({
"schemaVersion": 1,
"workspaceId": "psd",
"databaseId": "database-1",
"databaseName": "warehouse",
"schemaName": "dw",
"metadataContentRevision": 7,
"tables": [{
"id": "table-orders",
"name": "orders",
"description": "Orders from the Catalog",
"descriptionSource": "curated",
"columns": [{
"id": "column-id",
"name": "id",
"ordinalPosition": 1,
"dataType": "bigint",
"isNullable": False,
"defaultExpression": None,
"primaryKeyPosition": 1,
"sensitive": False,
"description": "Order identifier",
"descriptionSource": "generated",
}],
}],
"relationships": [],
}), encoding="utf-8")
raw = _cfg(tmp_path)
raw["paths"]["catalog_metadata_snapshot"] = str(snapshot)
cfg = _write_cfg(tmp_path, raw)
def build_lsh(_cfg, *, physical_file, output_dir):
assert "Orders from the Catalog" in physical_file.read_text(encoding="utf-8")
for name, value in (
("dw_lsh.pkl", b"lsh"),
("dw_minhashes.pkl", b"minhashes"),
("dw_meta.json", b"{}"),
):
(output_dir / name).write_bytes(value)
return ({"orders.id": "hash"}, [], [], {})
monkeypatch.setattr("tht.cli.lsh_cmd.build_lsh_artifacts", build_lsh)
report, result = run_catalog_dwh_from_config(cfg)
assert report.status == "succeeded"
assert result[0] == {"orders.id": "hash"}
owner = json.loads((tmp_path / ".tht-dwh" / "OWNER.json").read_text(encoding="utf-8"))
assert owner["schema_version"] == 2
assert owner["binding"]["workspace_id"] == "psd"
assert owner["binding"]["catalog_database_id"] == "database-1"
assert owner["binding"]["metadata_content_revision"] == 7
def test_catalog_pipeline_replaces_prior_derived_binding(monkeypatch, tmp_path):
import json
from tht.cli.preprocess_cmd import run_catalog_dwh_from_config
snapshot = tmp_path / "catalog-metadata.json"
def publish_snapshot(revision: int) -> None:
snapshot.write_text(json.dumps({
"schemaVersion": 1,
"workspaceId": "psd",
"databaseId": "database-1",
"databaseName": "warehouse",
"schemaName": "dw",
"metadataContentRevision": revision,
"tables": [],
"relationships": [],
}), encoding="utf-8")
publish_snapshot(7)
raw = _cfg(tmp_path)
raw["paths"]["catalog_metadata_snapshot"] = str(snapshot)
cfg = _write_cfg(tmp_path, raw)
def build_lsh(_cfg, *, physical_file, output_dir):
assert physical_file.is_file()
for name, value in (
("dw_lsh.pkl", b"lsh"),
("dw_minhashes.pkl", b"minhashes"),
("dw_meta.json", b"{}"),
):
(output_dir / name).write_bytes(value)
return ({}, [], [], {})
monkeypatch.setattr("tht.cli.lsh_cmd.build_lsh_artifacts", build_lsh)
first, _ = run_catalog_dwh_from_config(cfg)
assert first.status == "succeeded"
publish_snapshot(8)
second, _ = run_catalog_dwh_from_config(cfg)
assert second.status == "succeeded"
owner = json.loads((tmp_path / ".tht-dwh" / "OWNER.json").read_text(encoding="utf-8"))
assert owner["binding"]["metadata_content_revision"] == 8