feat(evidence): add safe retention and materialized reads
This commit is contained in:
@@ -61,7 +61,7 @@ class CorpusPipeline:
|
||||
def __init__(
|
||||
self, *, store: CorpusStore, sources: list[EvidenceSource], embedder,
|
||||
vector_store: VectorStore, embedding_model: str, embedding_dimensions: int,
|
||||
chunk_policy: ChunkPolicy, pipeline_version: str,
|
||||
chunk_policy: ChunkPolicy, pipeline_version: str, retain_published_generations: int = 3,
|
||||
) -> None:
|
||||
self.store = store
|
||||
self.sources = sources
|
||||
@@ -71,6 +71,50 @@ class CorpusPipeline:
|
||||
self.embedding_dimensions = embedding_dimensions
|
||||
self.chunk_policy = chunk_policy
|
||||
self.pipeline_version = pipeline_version
|
||||
if isinstance(retain_published_generations, bool) or retain_published_generations < 1:
|
||||
raise ValueError("retain_published_generations must be at least 1")
|
||||
self.retain_published_generations = retain_published_generations
|
||||
|
||||
def _protected_generations(self, workspace_root: Path) -> set[str]:
|
||||
protected = {value for value in (self.store.active_generation(),) if value}
|
||||
runs = workspace_root / ".tht-jobs" / "evidence" / "runs"
|
||||
for checkpoint in runs.glob("*/checkpoint.json") if runs.exists() else ():
|
||||
try:
|
||||
state = json.loads(checkpoint.read_text(encoding="utf-8"))
|
||||
if state.get("status") not in {"running", "failed"}:
|
||||
continue
|
||||
plan = checkpoint.parent / "artifacts" / "plan.json"
|
||||
generation = json.loads(plan.read_text(encoding="utf-8")).get("generation")
|
||||
if isinstance(generation, str):
|
||||
protected.add(generation)
|
||||
except (OSError, ValueError):
|
||||
continue
|
||||
return protected
|
||||
|
||||
def gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict:
|
||||
generations = self.store.list_generations()
|
||||
protected = self._protected_generations(workspace_root)
|
||||
keep = set(generations[-self.retain_published_generations:]) | protected
|
||||
evicted, failures = [], []
|
||||
for generation in generations:
|
||||
if generation in keep:
|
||||
continue
|
||||
if dry_run:
|
||||
evicted.append(generation)
|
||||
continue
|
||||
try:
|
||||
self.vector_store.delete_generation("evidence", generation)
|
||||
except Exception:
|
||||
failures.append({"generation": generation, "error": "vector cleanup failed"})
|
||||
continue
|
||||
try:
|
||||
self.store.discard(generation)
|
||||
evicted.append(generation)
|
||||
except Exception:
|
||||
failures.append({"generation": generation, "error": "filesystem cleanup failed"})
|
||||
return {"status": "partial" if failures else "succeeded", "dry_run": dry_run,
|
||||
"active_generation": self.store.active_generation(), "evicted": evicted,
|
||||
"protected": sorted(protected), "failures": failures}
|
||||
|
||||
def _discover(self) -> list[tuple[EvidenceSource, SourceObject]]:
|
||||
discovered = []
|
||||
@@ -343,8 +387,8 @@ class CorpusPipeline:
|
||||
return StageArtifacts(("plan.json", "manifest.json", "vector-intent.json"))
|
||||
|
||||
def retention_stage(context: JobContext) -> None:
|
||||
# Retention policy is intentionally a stable no-op until configured.
|
||||
return
|
||||
if not context.dry_run:
|
||||
self.gc(workspace_root=workspace_root)
|
||||
|
||||
report = run_job(spec, [
|
||||
discover_stage, acquire_stage, embed_stage, vector_stage,
|
||||
@@ -459,6 +503,7 @@ class CorpusPipeline:
|
||||
generation=generation,
|
||||
)
|
||||
self.store.publish(staged)
|
||||
self.gc(workspace_root=self.store.root.parent)
|
||||
except PipelineError:
|
||||
self._compensate(generation, vector_written)
|
||||
raise
|
||||
|
||||
Reference in New Issue
Block a user