test(evidence): harden generation lifecycle boundaries

This commit is contained in:
2026-07-12 05:15:30 +02:00
parent 459ffa0bcd
commit 9e21cce036
7 changed files with 135 additions and 3 deletions
+56
View File
@@ -3,6 +3,7 @@ import pytest
from tht.corpus.chunk import ChunkPolicy
from tht.corpus.pipeline import CorpusPipeline, PipelineError
from tht.corpus.store import CorpusStore
from tht.corpus.models import CorpusManifest
from tht.ports.evidence import AcquiredDocument, SourceObject
from tht.ports.vector import VectorCapabilities
@@ -162,6 +163,61 @@ def test_gc_reconciles_vector_only_generation(tmp_path):
report = candidate.gc(workspace_root=tmp_path)
assert report["evicted"] == [orphan]
assert vectors.list_evidence_generations("evidence") == []
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():
with candidate.store.writer_lock():
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_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path):