Files
ThothII/harness/tests/test_corpus_pipeline.py
T

1126 lines
43 KiB
Python

from datetime import UTC, datetime, timedelta
import pytest
from tht.evidence.canonical import CuratedEvidence, dump_curated_markdown
from tht.evidence.contracts import AcquiredDocument, SourceObject
from tht.evidence.corpus.chunk import ChunkPolicy
from tht.evidence.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest
from tht.evidence.corpus.pipeline import CorpusPipeline, PipelineError, PipelineResult
from tht.evidence.corpus.store import CorpusStore
from tht.ports.vector import VectorCapabilities, VectorHealth
class Source:
def __init__(self, documents):
self.documents = documents
self.acquire_calls = []
def discover(self):
return [item[0] for item in self.documents]
def acquire(self, item):
self.acquire_calls.append(item.source_id)
payload = next(payload for source, payload in self.documents if source.source_id == item.source_id)
if isinstance(payload, Exception):
raise payload
return AcquiredDocument(
source=item, content=payload.encode(), media_type=item.metadata.get("media_type")
)
class Embedder:
def __init__(self, dim=3, fail=False):
self.dim = dim
self.fail = fail
self.calls = []
def embed_documents(self, texts):
self.calls.extend(texts)
if self.fail:
raise RuntimeError("embed failed")
return [[float(i) for i in range(self.dim)] for _ in texts]
class Vectors:
capabilities = VectorCapabilities(search=True, existing_hashes=True, upsert=True)
def __init__(self, fail=False):
self.fail = fail
self.records = []
self.dimension = 3
def upsert(self, collection, records):
self.records.extend(records[:1] if self.fail else records)
if self.fail:
raise RuntimeError("partial write")
return len(records)
def existing_hashes(self, collection, kinds):
return {
value.record.id: value.content_hash for value in self.records
}
def health(self):
return VectorHealth(
ok=True, expected_dimension=3, observed_dimensions=(self.dimension,),
dimension_compatible=self.dimension == 3,
)
def delete_generation(self, collection, generation, workspace_id):
self.records = [
value for value in self.records
if not (value.record.metadata["vector_generation"] == generation
and value.record.metadata.get("workspace_id") == workspace_id)
]
return 0
def list_evidence_generations(self, collection, workspace_id):
return sorted({
value.record.metadata["vector_generation"] for value in self.records
if value.record.kind == "evidence"
and value.record.metadata.get("workspace_id") == workspace_id
})
class InterruptingVectors(Vectors):
def __init__(self):
super().__init__()
self.batches = []
self.interrupt = True
def upsert(self, collection, records):
self.batches.append([value.record.id for value in records])
if self.interrupt:
self.interrupt = False
self.records.append(records[0])
raise KeyboardInterrupt("process interruption after partial write")
self.records.extend(records)
return len(records)
def item(name, fingerprint):
return SourceObject(
source_id=f"fs:{name}", uri=f"file:///safe/{name}.md", fingerprint=f"sha256:{fingerprint}"
)
def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", policy=None,
retain=3, candidate_evaluator=None):
return CorpusPipeline(
store=CorpusStore(tmp_path / "corpus"), sources=[source],
embedder=embedder or Embedder(), vector_store=vectors or Vectors(),
embedding_model=model, embedding_dimensions=3,
chunk_policy=policy or ChunkPolicy(version="chunk-v1", max_chars=100),
pipeline_version="evidence-v1",
retain_published_generations=retain,
candidate_evaluator=candidate_evaluator,
)
def test_pipeline_embeds_validated_curated_evidence_as_semantic_fragments(tmp_path):
evidence = CuratedEvidence.model_validate(
{
"schema_version": 3,
"id": "evidence:fascia-pediatrica",
"title": "Fascia pediatrica",
"kind": "formula",
"purposes": ["sql_generation"],
"applies_to": {"columns": ["clinical.patient.birth_date"]},
"language": "it",
"provenance": {
"source_file": "source/paziente.md",
"source_sha256": "sha256:" + "a" * 64,
"supporting_excerpts": ["Pazienti con età inferiore a 18 anni."],
},
"review_items": [],
"payload": {
"concept": "fascia pediatrica",
"columns": ["clinical.patient.birth_date"],
"sql": "CASE WHEN age < 18 THEN 'pediatric' END",
},
}
)
source_item = SourceObject(
source_id="fs:curated-formula",
uri="file:///safe/curated/formula/fascia-pediatrica.md",
fingerprint="sha256:" + "b" * 64,
metadata={"relative_path": "curated/formula/fascia-pediatrica.md"},
)
embedder = Embedder()
vectors = Vectors()
result = pipeline(
tmp_path,
Source([(source_item, dump_curated_markdown(evidence))]),
embedder=embedder,
vectors=vectors,
policy=ChunkPolicy(version="chunk-v1", max_chars=4000),
).run()
assert result.status == "succeeded"
assert len(result.manifest.chunks) == 1
assert "Formula: Fascia pediatrica" in embedder.calls[0]
assert vectors.records[0].record.metadata["evidence_id"] == evidence.id
assert vectors.records[0].record.metadata["provenance"]["source_file"] == "source/paziente.md"
def test_pipeline_exposes_atomic_content_review_code_when_candidate_is_blocked(tmp_path):
evidence = CuratedEvidence.model_validate(
{
"schema_version": 1,
"id": "evidence:formula-lunga",
"title": "Formula lunga",
"kind": "formula",
"purposes": ["sql_generation"],
"language": "it",
"provenance": {
"source_file": "source/paziente.md",
"source_sha256": "sha256:" + "a" * 64,
"supporting_excerpts": ["Una formula molto lunga."],
},
"review_items": [],
"payload": {"concept": "formula lunga", "columns": [], "sql": "x" * 200},
}
)
source_item = SourceObject(
source_id="fs:formula-lunga",
uri="file:///safe/curated/formula/formula-lunga.md",
fingerprint="sha256:" + "b" * 64,
metadata={"relative_path": "curated/formula/formula-lunga.md"},
)
result = pipeline(
tmp_path,
Source([(source_item, dump_curated_markdown(evidence))]),
policy=ChunkPolicy(version="chunk-v1", max_chars=120),
).run()
assert result.status == "blocked"
assert result.published is False
assert [item.code for item in result.review_items] == ["atomic_content_too_large"]
assert result.review_items[0].field == "formula.sql"
def test_job_pipeline_persists_atomic_content_review_item_when_candidate_is_blocked(tmp_path):
evidence = CuratedEvidence.model_validate(
{
"schema_version": 1,
"id": "evidence:formula-lunga-job",
"title": "Formula lunga",
"kind": "formula",
"purposes": ["sql_generation"],
"language": "it",
"provenance": {
"source_file": "source/paziente.md",
"source_sha256": "sha256:" + "a" * 64,
"supporting_excerpts": ["Una formula molto lunga."],
},
"review_items": [],
"payload": {"concept": "formula lunga", "columns": [], "sql": "x" * 200},
}
)
source_item = SourceObject(
source_id="fs:formula-lunga-job",
uri="file:///safe/curated/formula/formula-lunga-job.md",
fingerprint="sha256:" + "b" * 64,
metadata={"relative_path": "curated/formula/formula-lunga-job.md"},
)
result = pipeline(
tmp_path,
Source([(source_item, dump_curated_markdown(evidence))]),
policy=ChunkPolicy(version="chunk-v1", max_chars=120),
).run_as_job(
workspace_id="demo",
workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
assert result.status == "blocked"
assert result.published is False
assert [item.code for item in result.review_items] == ["atomic_content_too_large"]
def test_pipeline_routes_source_io_through_evidence_facade(tmp_path, monkeypatch):
import tht.evidence.acquisition as evidence_acquisition
calls = []
def discover(source):
calls.append(("discover", source))
return source.discover()
def acquire(source, source_item):
calls.append(("acquire", source_item.source_id))
return source.acquire(source_item)
monkeypatch.setattr(evidence_acquisition, "discover", discover)
monkeypatch.setattr(evidence_acquisition, "acquire", acquire)
source = Source([(item("one", "a"), "body")])
result = pipeline(tmp_path, source).run()
assert result.status == "succeeded"
assert calls == [("discover", source), ("acquire", "fs:one")]
def test_retention_bounds_generations_and_purges_vectors_after_publish(tmp_path):
vectors = Vectors()
generations = []
for index in range(4):
result = pipeline(
tmp_path, Source([(item("one", str(index)), f"version {index}")]),
vectors=vectors, retain=2,
).run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + str(index) * 64,
)
generations.append(result.generation)
store = CorpusStore(tmp_path / "corpus")
assert store.list_generations() == generations[-2:]
assert {r.record.metadata["vector_generation"] for r in vectors.records} == set(generations[-2:])
assert store.active_generation() == generations[-1]
def test_retention_keeps_filesystem_when_vector_purge_fails_then_retries(tmp_path):
class FailingDelete(Vectors):
def __init__(self):
super().__init__()
self.fail_delete = True
def delete_generation(self, collection, generation, workspace_id):
if self.fail_delete:
raise RuntimeError("credential secret")
return super().delete_generation(collection, generation, workspace_id)
vectors = FailingDelete()
for index in range(2):
pipeline(tmp_path, Source([(item("one", str(index)), str(index))]), vectors=vectors,
retain=1).run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + str(index) * 64,
)
assert len(CorpusStore(tmp_path / "corpus").list_generations()) == 2
vectors.fail_delete = False
report = pipeline(tmp_path, Source([(item("one", "1"), "1")]), vectors=vectors,
retain=1).gc(workspace_root=tmp_path)
assert report["status"] == "succeeded"
assert len(CorpusStore(tmp_path / "corpus").list_generations()) == 1
def test_gc_reconciles_vector_only_generation(tmp_path):
vectors = Vectors()
orphan = "gen:" + "f" * 32
from tht.ports.vector import VectorWriteRecord
from tht.vectorstore.records import VectorRecord
vectors.records.append(VectorWriteRecord(
record=VectorRecord(id="orphan", kind="evidence", ref="doc:x", title="", content="x",
metadata={"vector_generation": orphan, "workspace_id": "default"}),
embedding=[0.0, 0.0, 0.0], content_hash="sha256:" + "0" * 64,
))
candidate = pipeline(tmp_path, Source([]), vectors=vectors, retain=1)
report = candidate.gc(workspace_root=tmp_path)
assert report["evicted"] == [orphan]
assert vectors.list_evidence_generations("evidence", "default") == []
assert candidate.gc(workspace_root=tmp_path)["evicted"] == []
@pytest.mark.parametrize("status", ["running", "failed"])
def test_gc_protects_generations_referenced_by_resumable_checkpoints(tmp_path, status):
generation = "gen:" + "e" * 32
store = CorpusStore(tmp_path / "corpus")
store.stage(CorpusManifest(), {}, generation=generation)
run = tmp_path / ".tht-jobs" / "evidence" / "runs" / ("a" * 32)
(run / "artifacts").mkdir(parents=True)
(run / "checkpoint.json").write_text(__import__("json").dumps({"status": status}))
(run / "artifacts" / "plan.json").write_text(
__import__("json").dumps({"generation": generation})
)
candidate = pipeline(tmp_path, Source([]), vectors=Vectors(), retain=1)
report = candidate.gc(workspace_root=tmp_path)
assert generation in report["protected"]
assert store.generation_path(generation).exists()
def test_explicit_gc_blocks_while_job_holds_corpus_writer_lock(tmp_path):
import threading
candidate = pipeline(tmp_path, Source([(item("one", "a"), "one")]), vectors=Vectors())
entered = threading.Event()
release = threading.Event()
gc_finished = threading.Event()
def pause(_context, stage):
if stage == "discover":
entered.set()
assert release.wait(5)
job = threading.Thread(target=lambda: candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
after_stage_return=pause,
))
job.start()
assert entered.wait(5)
def collect():
candidate.gc(workspace_root=tmp_path)
gc_finished.set()
gc_thread = threading.Thread(target=collect)
gc_thread.start()
assert not gc_finished.wait(0.1)
release.set()
job.join(5)
gc_thread.join(5)
assert gc_finished.is_set()
assert candidate.store.active_generation() is not None
def test_gc_preserves_vector_dependencies_of_retained_manifests(tmp_path):
vectors = Vectors()
one = item("one", "a")
first = pipeline(tmp_path, Source([(one, "stable")]), vectors=vectors, retain=2).run().generation
second = pipeline(
tmp_path, Source([(one, "stable"), (item("two", "b"), "two")]),
vectors=vectors, retain=2,
).run().generation
third = pipeline(
tmp_path, Source([(one, "stable"), (item("two", "c"), "changed")]),
vectors=vectors, retain=2,
).run().generation
assert CorpusStore(tmp_path / "corpus").list_generations() == [second, third]
assert first in vectors.list_evidence_generations("evidence", "default")
def test_active_searcher_without_active_fails_closed_for_evidence(tmp_path):
from types import SimpleNamespace
from tht.evidence.search import active_searcher
class Delegate:
def search(self, embedding, top_n=10, kinds=None, metadata_filter=None, **kwargs):
return ["legacy"]
cfg = SimpleNamespace(paths=SimpleNamespace(artifacts=tmp_path / "artifacts"))
wrapped = active_searcher(cfg, Delegate())
assert wrapped.search([1.0], kinds=["evidence"]) == []
assert wrapped.search([1.0], kinds=["memory"]) == ["legacy"]
def test_active_searcher_splits_default_and_mixed_kinds_before_global_limit(tmp_path):
from types import SimpleNamespace
from tht.evidence.search import ActiveEvidenceSearcher
store = CorpusStore(tmp_path / "corpus")
generation = store.stage(
CorpusManifest(metadata={"workspace_id": "default"}), {},
generation="gen:" + "a" * 32,
)
store.publish(generation)
calls = []
class Delegate:
def search(self, embedding, top_n=10, kinds=None, metadata_filter=None):
calls.append((kinds, metadata_filter))
if kinds == ["evidence"]:
return [SimpleNamespace(id="active", similarity=0.8)]
return [SimpleNamespace(id="memory", similarity=0.9)]
searcher = ActiveEvidenceSearcher(store, Delegate())
hits = searcher.search([1.0], top_n=1, kinds=["evidence", "memory"])
assert [hit.id for hit in hits] == ["memory"]
assert calls[0] == (["memory"], None)
# Empty manifest means no Evidence query, but the split remains explicit and safe.
assert all(call[0] != ["evidence"] for call in calls)
calls.clear()
searcher.search([1.0], top_n=1)
assert calls[0][0] == ["memory", "schema_column", "schema_table", "solved_question"]
assert all(call[0] is not None for call in calls)
def test_active_evidence_query_holds_lock_against_publish(tmp_path):
import threading
from types import SimpleNamespace
from tht.evidence.search import ActiveEvidenceSearcher
first_pipeline = pipeline(tmp_path, Source([(item("one", "a"), "old")]), vectors=Vectors())
first_pipeline.run()
store = first_pipeline.store
entered = threading.Event()
release = threading.Event()
published = threading.Event()
class Delegate:
def search(self, embedding, top_n=10, kinds=None, metadata_filter=None, **kwargs):
entered.set()
assert release.wait(5)
return [SimpleNamespace(id="active", similarity=1.0)]
search = threading.Thread(
target=lambda: ActiveEvidenceSearcher(store, Delegate()).search(
[1.0], kinds=["evidence"], query_text="old"
)
)
search.start()
assert entered.wait(5)
next_generation = store.stage(CorpusManifest(), {})
def publish():
with store.writer_lock():
store.publish(next_generation)
published.set()
publisher = threading.Thread(target=publish)
publisher.start()
assert not published.wait(0.1)
release.set()
search.join(5)
publisher.join(5)
assert published.is_set()
def test_active_evidence_search_refuses_a_dense_only_fallback(tmp_path):
from tht.evidence.search import ActiveEvidenceSearcher
from tht.ports.vector import VectorStoreError
current = pipeline(tmp_path, Source([(item("one", "a"), "cardiomiopatia")]), vectors=Vectors())
current.run()
with pytest.raises(VectorStoreError, match="hybrid query text"):
ActiveEvidenceSearcher(current.store, object()).search([1.0], kinds=["evidence"])
def test_pipeline_result_dump_does_not_deepcopy_frozen_metadata():
manifest = CorpusManifest(metadata={"nested": {"value": ["safe"]}})
payload = PipelineResult(
"succeeded", None, False, (), (), (), manifest,
).model_dump(mode="json")
assert payload["manifest_id"] is None
assert "manifest" not in payload
def test_pipeline_result_public_dump_is_bounded_and_excludes_evidence_content(tmp_path):
import json
result = pipeline(
tmp_path, Source([(item("one", "a"), "SENSITIVE EVIDENCE CONTENT")])
).run()
payload = result.model_dump(mode="json")
encoded = json.dumps(payload)
assert "SENSITIVE EVIDENCE CONTENT" not in encoded
assert "documents" not in payload and "chunks" not in payload
large = PipelineResult(
"failed", None, False,
tuple(f"fs:item-{index}" for index in range(1000)), (), (), result.manifest,
).model_dump(mode="json")
assert len(large["changed"]) == 100
assert large["counts"]["changed"] == 1000
assert len(json.dumps(large)) < 25_000
def test_pipeline_result_repr_is_bounded_and_excludes_manifest_secrets():
secret = "TOP_SECRET_CONTENT"
manifest = CorpusManifest.model_construct(
manifest_id="gen:" + "a" * 64,
documents=tuple(CanonicalDocument.model_construct(content=secret) for _ in range(1000)),
chunks=tuple(CanonicalChunk.model_construct(content=secret) for _ in range(1000)),
metadata={"password": secret, "credential": "Bearer " + secret},
)
result = PipelineResult(
"succeeded", "gen:" + "a" * 64, True, (), (), (), manifest,
run_id="b" * 32,
)
rendered = repr(result)
assert str(result) == rendered
assert len(rendered) < 1000
assert secret not in rendered
assert "password" not in rendered
assert "credential" not in rendered
assert "manifest" not in rendered.lower()
assert "documents" not in rendered
assert "chunks" not in rendered
def test_reused_corpus_root_rejects_workspace_rename_before_any_mutation(tmp_path):
vectors = Vectors()
first = pipeline(tmp_path, Source([(item("one", "a"), "stable")]), vectors=vectors)
first.run_as_job(
workspace_id="workspace-a", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
active = first.store.active_generation()
records = list(vectors.records)
renamed_source = Source([(item("one", "a"), "stable")])
renamed = pipeline(tmp_path, renamed_source, vectors=vectors)
with pytest.raises(PipelineError, match="different workspace"):
renamed.run_as_job(
workspace_id="workspace-b", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
assert renamed_source.acquire_calls == []
assert renamed.store.active_generation() == active
assert vectors.records == records
def test_gc_rejects_workspace_mismatch_without_deleting(tmp_path):
vectors = Vectors()
owner = pipeline(tmp_path, Source([(item("one", "a"), "stable")]), vectors=vectors)
owner.run_as_job(
workspace_id="workspace-a", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
generations = owner.store.list_generations()
wrong = pipeline(tmp_path, Source([]), vectors=vectors)
wrong.workspace_id = "workspace-b"
with pytest.raises(PipelineError, match="different workspace"):
wrong.gc(workspace_root=tmp_path)
assert wrong.store.list_generations() == generations
@pytest.mark.parametrize("kinds", [None, ["evidence", "memory"], ["memory"]])
def test_active_search_rejects_workspace_mismatch_before_delegate(tmp_path, kinds):
from tht.evidence.search import ActiveEvidenceSearcher, CorpusWorkspaceMismatchError
vectors = Vectors()
owner = pipeline(tmp_path, Source([(item("one", "a"), "stable")]), vectors=vectors)
owner.run_as_job(
workspace_id="workspace-a", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
class Delegate:
def search(self, *args, **kwargs):
raise AssertionError("workspace mismatch reached vector delegate")
with pytest.raises(CorpusWorkspaceMismatchError, match="different workspace"):
ActiveEvidenceSearcher(
owner.store, Delegate(), expected_workspace_id="workspace-b",
).search([1.0], kinds=kinds)
def test_unscoped_active_manifest_is_never_adopted_by_direct_run_or_gc(tmp_path):
store = CorpusStore(tmp_path / "corpus")
generation = store.stage(CorpusManifest(), {})
store.publish(generation)
class ForbiddenSource:
def discover(self):
raise AssertionError("invalid corpus reached source discovery")
class ForbiddenVectors:
def __getattr__(self, name):
raise AssertionError(f"invalid corpus reached vector operation {name}")
candidate = CorpusPipeline(
store=store, sources=[ForbiddenSource()], embedder=Embedder(),
vector_store=ForbiddenVectors(), embedding_model="model", embedding_dimensions=3,
chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=100),
pipeline_version="evidence-v1",
)
with pytest.raises(PipelineError, match="missing or invalid"):
candidate.run()
with pytest.raises(PipelineError, match="missing or invalid"):
candidate.gc(workspace_root=tmp_path)
assert store.active_generation() == generation
def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path):
one = item("one", "a")
first_source = Source([(one, "hello")])
first = pipeline(tmp_path, first_source)
first.run()
second_source = Source([(one, "ignored")])
second_embedder = Embedder()
result = pipeline(tmp_path, second_source, embedder=second_embedder).run()
assert result.unchanged == ("fs:one",)
assert second_source.acquire_calls == []
assert second_embedder.calls == []
def test_unchanged_job_reuses_active_generation_without_new_directory(tmp_path):
source = Source([(
SourceObject(
source_id="fs:one", uri="file:///safe/one.md", fingerprint="sha256:a",
modified_at=datetime(2026, 1, 1, tzinfo=UTC),
metadata={"media_type": "text/markdown", "size": 5, "nested": {"b": 2, "a": 1}},
),
"hello",
)])
candidate = pipeline(tmp_path, source)
args = {"workspace_id": "demo", "workspace_root": tmp_path,
"config_fingerprint": "sha256:" + "1" * 64,
"input_fingerprint": "sha256:" + "2" * 64}
first = candidate.run_as_job(**args)
snapshot = first.manifest.metadata["source_snapshot"]["fs:one"]
assert snapshot == {
"source_id": "fs:one", "uri": "file:///safe/one.md", "fingerprint": "sha256:a",
"modified_at": "2026-01-01T00:00:00Z",
"metadata": {"media_type": "text/markdown", "size": 5,
"nested": {"a": 1, "b": 2}},
"media_type": "text/markdown", "size": 5,
}
count = len(candidate.store.list_generations())
second = candidate.run_as_job(**args)
assert second.generation == first.generation
assert second.published is False
assert len(candidate.store.list_generations()) == count
@pytest.mark.parametrize("field", ["uri", "modified_at", "metadata"])
def test_job_source_snapshot_change_forces_publish_with_same_fingerprint(tmp_path, field):
original = SourceObject(
source_id="fs:one", uri="file:///safe/one.md", fingerprint="sha256:a",
modified_at=datetime(2026, 1, 1, tzinfo=UTC),
metadata={"media_type": "text/markdown", "size": 5, "label": "original"},
)
vectors = Vectors()
args = {"workspace_id": "demo", "workspace_root": tmp_path,
"config_fingerprint": "sha256:" + "1" * 64,
"input_fingerprint": "sha256:" + "2" * 64}
first = pipeline(tmp_path, Source([(original, "hello")]), vectors=vectors).run_as_job(**args)
updates = {
"uri": "file:///safe/renamed.md",
"modified_at": original.modified_at + timedelta(seconds=1),
"metadata": {"media_type": "text/markdown", "size": 5, "label": "changed"},
}
changed = original.model_copy(update={field: updates[field]})
source = Source([(changed, "hello")])
result = pipeline(tmp_path, source, vectors=vectors).run_as_job(**args)
assert result.published is True
assert result.generation != first.generation
assert source.acquire_calls == ["fs:one"]
@pytest.mark.parametrize("fingerprint_name", ["config_fingerprint", "input_fingerprint"])
def test_job_binding_change_forces_publish(tmp_path, fingerprint_name):
vectors = Vectors()
source_object = item("one", "a")
args = {"workspace_id": "demo", "workspace_root": tmp_path,
"config_fingerprint": "sha256:" + "1" * 64,
"input_fingerprint": "sha256:" + "2" * 64}
first = pipeline(tmp_path, Source([(source_object, "hello")]), vectors=vectors).run_as_job(**args)
args[fingerprint_name] = "sha256:" + "3" * 64
source = Source([(source_object, "hello")])
result = pipeline(tmp_path, source, vectors=vectors).run_as_job(**args)
assert result.published is True
assert result.generation != first.generation
assert source.acquire_calls == []
@pytest.mark.parametrize(
"damage", ["legacy_metadata", "corrupt_document_sources", "document", "vector"]
)
def test_job_incomplete_active_contract_never_noops(tmp_path, damage):
import json
vectors = Vectors()
source_object = item("one", "a")
args = {"workspace_id": "demo", "workspace_root": tmp_path,
"config_fingerprint": "sha256:" + "1" * 64,
"input_fingerprint": "sha256:" + "2" * 64}
candidate = pipeline(tmp_path, Source([(source_object, "hello")]), vectors=vectors)
first = candidate.run_as_job(**args)
if damage == "legacy_metadata":
manifest_path = candidate.store.generation_path(first.generation) / "manifest.json"
payload = json.loads(manifest_path.read_text())
payload["metadata"].pop("source_snapshot")
manifest_path.write_text(json.dumps(payload))
elif damage == "corrupt_document_sources":
manifest_path = candidate.store.generation_path(first.generation) / "manifest.json"
payload = json.loads(manifest_path.read_text())
payload["metadata"]["document_sources"] = {}
manifest_path.write_text(json.dumps(payload))
elif damage == "document":
path = candidate.store.resolve_document(first.manifest.documents[0].document_id)
path.unlink()
else:
vectors.records.clear()
source = Source([(source_object, "hello")])
result = pipeline(tmp_path, source, vectors=vectors).run_as_job(**args)
assert result.published is True
assert result.generation != first.generation
assert source.acquire_calls == ["fs:one"]
@pytest.mark.parametrize(
"damage", ["modified_at", "source_metadata", "media_type", "missing_chunk",
"altered_chunk", "extra_chunk", "vector_dimension"]
)
def test_job_corrupt_canonical_document_or_chunk_never_noops(tmp_path, damage):
import hashlib
import json
vectors = Vectors()
source_object = SourceObject(
source_id="fs:one", uri="file:///safe/one.md", fingerprint="sha256:a",
modified_at=datetime(2026, 1, 1, tzinfo=UTC),
metadata={"media_type": "text/markdown", "size": 11, "owner": "docs"},
)
args = {"workspace_id": "demo", "workspace_root": tmp_path,
"config_fingerprint": "sha256:" + "1" * 64,
"input_fingerprint": "sha256:" + "2" * 64}
candidate = pipeline(
tmp_path, Source([(source_object, "hello world")]), vectors=vectors,
policy=ChunkPolicy(version="chunk-v1", max_chars=6),
)
first = candidate.run_as_job(**args)
manifest_path = candidate.store.generation_path(first.generation) / "manifest.json"
payload = json.loads(manifest_path.read_text())
document = payload["documents"][0]
chunks = payload["chunks"]
if damage == "modified_at":
document["modified_at"] = "2026-01-01T00:00:01Z"
elif damage == "source_metadata":
document["metadata"]["source"]["owner"] = "attacker"
elif damage == "media_type":
document["media_type"] = "text/plain"
elif damage == "missing_chunk":
payload["chunks"] = chunks[:-1]
elif damage == "altered_chunk":
chunks[0]["content"] = "HELLO "
chunks[0]["content_hash"] = "sha256:" + hashlib.sha256(b"HELLO ").hexdigest()
chunks[0]["chunk_id"] = "chunk:" + "a" * 64
elif damage == "extra_chunk":
extra = CanonicalChunk(
chunk_id="chunk:" + "b" * 64, document_id=document["document_id"],
ordinal=len(chunks), content="", content_hash="sha256:" + hashlib.sha256(b"").hexdigest(),
source_uri=document["source_uri"], pipeline_version=document["pipeline_version"],
)
chunks.append(extra.model_dump(mode="json"))
else:
vectors.dimension = 4
manifest_path.write_text(json.dumps(payload))
source = Source([(source_object, "hello world")])
result = pipeline(
tmp_path, source, vectors=vectors,
policy=ChunkPolicy(version="chunk-v1", max_chars=6),
).run_as_job(**args)
assert result.published is True
assert result.generation != first.generation
assert source.acquire_calls == ["fs:one"]
def test_removed_documents_are_marked_and_absent_from_new_manifest(tmp_path):
one, two = item("one", "a"), item("two", "b")
pipeline(tmp_path, Source([(one, "one"), (two, "two")])).run()
result = pipeline(tmp_path, Source([(one, "one")])).run()
assert result.removed == ("fs:two",)
assert {doc.source_id for doc in result.manifest.documents} == {"fs:one"}
def test_model_or_chunk_policy_change_forces_full_rebuild(tmp_path):
one = item("one", "a")
pipeline(tmp_path, Source([(one, "hello")])).run()
source = Source([(one, "hello")])
changed = pipeline(tmp_path, source, model="model-b").run()
assert changed.changed == ("fs:one",)
assert source.acquire_calls == ["fs:one"]
def test_partial_vector_failure_never_changes_active_or_exposes_generation(tmp_path):
one = item("one", "a")
good = pipeline(tmp_path, Source([(one, "old")]))
old = good.run().generation
changed = item("one", "b")
vectors = Vectors(fail=True)
broken = pipeline(tmp_path, Source([(changed, "new")]), vectors=vectors)
with pytest.raises(PipelineError):
broken.run()
assert broken.store.active_generation() == old
assert vectors.records[0].record.metadata["vector_generation"] != old
def test_dimension_mismatch_fails_before_vector_write_and_publish(tmp_path):
one = item("one", "a")
vectors = Vectors()
candidate = pipeline(tmp_path, Source([(one, "hello")]), embedder=Embedder(dim=2), vectors=vectors)
with pytest.raises(PipelineError, match="dimension"):
candidate.run()
assert vectors.records == []
assert candidate.store.active_generation() is None
def test_failed_candidate_evaluation_never_switches_the_active_generation(tmp_path):
from types import SimpleNamespace
vectors = Vectors()
active = pipeline(tmp_path, Source([(item("one", "a"), "old")]), vectors=vectors).run().generation
candidate = pipeline(
tmp_path,
Source([(item("one", "b"), "new")]),
vectors=vectors,
candidate_evaluator=lambda manifest: SimpleNamespace(passed=False),
)
with pytest.raises(PipelineError, match="candidate retrieval evaluation failed"):
candidate.run()
assert candidate.store.active_generation() == active
assert {record.record.metadata["vector_generation"] for record in vectors.records} == {active}
def test_job_failed_candidate_evaluation_never_switches_the_active_generation(tmp_path):
from types import SimpleNamespace
vectors = Vectors()
active = pipeline(tmp_path, Source([(item("one", "a"), "old")]), vectors=vectors).run_as_job(
workspace_id="demo",
workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
).generation
candidate = pipeline(
tmp_path,
Source([(item("one", "b"), "new")]),
vectors=vectors,
candidate_evaluator=lambda manifest: SimpleNamespace(passed=False),
)
result = candidate.run_as_job(
workspace_id="demo",
workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "3" * 64,
)
assert result.status == "failed"
assert result.published is False
assert candidate.store.active_generation() == active
assert {record.record.metadata["vector_generation"] for record in vectors.records} == {active}
def test_pipeline_marks_each_evidence_fragment_for_server_side_italian_bm25(tmp_path):
vectors = Vectors()
pipeline(tmp_path, Source([(item("one", "a"), "ricovero cardiologico")]), vectors=vectors).run()
assert [(record.sparse_text, record.sparse_language) for record in vectors.records] == [
("ricovero cardiologico", "italian"),
]
def test_dry_run_and_failed_acquire_never_change_active(tmp_path):
one = item("one", "a")
active = pipeline(tmp_path, Source([(one, "old")])).run().generation
changed = item("one", "b")
dry = pipeline(tmp_path, Source([(changed, "new")])).run(dry_run=True)
assert dry.published is False
assert dry.generation is None
assert dry.manifest.documents[0].content == "old"
with pytest.raises(PipelineError):
pipeline(tmp_path, Source([(changed, RuntimeError("boom"))])).run()
assert CorpusStore(tmp_path / "corpus").active_generation() == active
def test_job_pipeline_uses_ordered_plan_and_returns_run_id(tmp_path):
one = item("one", "a")
candidate = pipeline(tmp_path, Source([(one, "hello")]))
result = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
assert result.status == "succeeded"
assert result.run_id and len(result.run_id) == 32
checkpoint = tmp_path / ".tht-jobs" / "evidence" / "runs" / result.run_id / "checkpoint.json"
payload = __import__("json").loads(checkpoint.read_text())
assert [stage["name"] for stage in payload["stages"]] == [
"discover", "acquire_normalize_chunk", "embed", "vector_upsert",
"stage_validate", "publish", "retention_cleanup",
]
def test_job_pipeline_dry_run_only_discovers_and_reports_changes(tmp_path):
one = item("one", "a")
source = Source([(one, "hello")])
embedder = Embedder()
vectors = Vectors()
result = pipeline(tmp_path, source, embedder=embedder, vectors=vectors).run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
dry_run=True,
)
assert result.changed == ("fs:one",)
assert source.acquire_calls == []
assert embedder.calls == []
assert vectors.records == []
assert result.generation is None and result.published is False
@pytest.mark.parametrize("crash_stage", [
"discover", "acquire_normalize_chunk", "embed", "vector_upsert",
"stage_validate", "publish", "retention_cleanup",
])
def test_job_pipeline_crash_after_each_stage_resumes_without_duplicate_effects(tmp_path, crash_stage):
one = item("one", "a")
source = Source([(one, "hello")])
embedder = Embedder()
vectors = Vectors()
candidate = pipeline(tmp_path, source, embedder=embedder, vectors=vectors)
class Crash(BaseException):
pass
def fault(_context, stage):
if stage == crash_stage:
raise Crash()
with pytest.raises(Crash):
candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
after_stage_return=fault,
)
runs = tmp_path / ".tht-jobs" / "evidence" / "runs"
crashed = next(runs.iterdir()).name
result = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
resume_run_id=crashed,
)
assert result.status == "succeeded"
assert source.acquire_calls == ["fs:one"]
assert len(embedder.calls) == 1
assert len(vectors.records) == 1
def test_job_pipeline_raw_upsert_failure_compensates_and_resumes_with_new_generation(tmp_path):
one = item("one", "a")
vectors = Vectors(fail=True)
candidate = pipeline(tmp_path, Source([(one, "hello")]), vectors=vectors)
first = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
assert first.status == "failed"
assert vectors.records == []
old_generation = first.generation
vectors.fail = False
resumed = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
resume_run_id=first.run_id,
)
assert resumed.status == "succeeded", resumed
assert resumed.generation != old_generation
assert candidate.store.active_generation() == resumed.generation
def test_job_pipeline_raw_stage_failure_compensates_vectors_and_resumes(tmp_path, monkeypatch):
one = item("one", "a")
vectors = Vectors()
candidate = pipeline(tmp_path, Source([(one, "hello")]), vectors=vectors)
real_stage = candidate.store.stage
calls = 0
def fail_once(*args, **kwargs):
nonlocal calls
calls += 1
if calls == 1:
raise OSError("raw stage failure")
return real_stage(*args, **kwargs)
monkeypatch.setattr(candidate.store, "stage", fail_once)
first = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
assert first.status == "failed" and vectors.records == []
resumed = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
resume_run_id=first.run_id,
)
assert resumed.status == "succeeded", resumed
@pytest.mark.parametrize("stage,filename", [
("discover", "plan.json"),
("acquire_normalize_chunk", "manifest.json"),
("embed", "embeddings.json"),
])
@pytest.mark.parametrize("mutation", ["missing", "tampered"])
def test_job_pipeline_rejects_corrupt_required_artifacts_before_resume(
tmp_path, stage, filename, mutation,
):
one = item("one", "a")
candidate = pipeline(tmp_path, Source([(one, "hello")]))
class Crash(BaseException):
pass
with pytest.raises(Crash):
candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
after_stage_return=lambda _context, name: (
(_ for _ in ()).throw(Crash()) if name == stage else None
),
)
runs = tmp_path / ".tht-jobs" / "evidence" / "runs"
crashed = next(runs.iterdir())
target = crashed / "artifacts" / filename
target.unlink() if mutation == "missing" else target.write_text("tampered")
from tht.jobs.runner import CorruptCheckpointError
with pytest.raises(CorruptCheckpointError, match="artifact"):
candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
resume_run_id=crashed.name,
)
def test_vector_intent_is_reconciled_after_process_interruption_without_duplicate_upsert(tmp_path):
one = item("one", "a")
vectors = InterruptingVectors()
candidate = pipeline(
tmp_path, Source([(one, "a" * 250)]), vectors=vectors,
policy=ChunkPolicy(version="chunk-v1", max_chars=100),
)
with pytest.raises(KeyboardInterrupt):
candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
)
runs = tmp_path / ".tht-jobs" / "evidence" / "runs"
interrupted = next(runs.iterdir())
checkpoint = __import__("json").loads((interrupted / "checkpoint.json").read_text())
vector_stage = checkpoint["stages"][3]
assert vector_stage["status"] == "running"
assert vector_stage["effect_state"] == "intent"
first_written = vectors.batches[0][0]
result = candidate.run_as_job(
workspace_id="demo", workspace_root=tmp_path,
config_fingerprint="sha256:" + "1" * 64,
input_fingerprint="sha256:" + "2" * 64,
resume_run_id=interrupted.name,
)
assert result.status == "succeeded" and result.published is True
assert first_written not in vectors.batches[1]
assert len(vectors.records) == 3