192 lines
7.0 KiB
Python
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
|