"""Incremental Evidence preprocessing with generation-isolated vector writes.""" from __future__ import annotations import hashlib import json import re import uuid from collections.abc import Mapping, Sequence from dataclasses import asdict, dataclass, field from datetime import UTC from pathlib import Path from tht.corpus.chunk import ChunkPolicy, chunk from tht.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest from tht.corpus.normalize import normalize from tht.corpus.store import CorpusStore from tht.ports.evidence import EvidenceSource, SourceObject, canonical_provenance_uri from tht.ports.vector import VectorStore, VectorWriteRecord from tht.vectorstore.records import VectorRecord from tht.jobs.models import JobSpec from tht.jobs.runner import JobContext, StageArtifacts, run_job, seal_stage_artifacts EVIDENCE_STAGE_IDS = ( "discover", "acquire_normalize_chunk", "embed", "vector_upsert", "stage_validate", "publish", "retention_cleanup", ) class PipelineError(RuntimeError): """Credential-free failure at the preprocessing boundary.""" @dataclass(frozen=True) class PipelineResult: status: str generation: str | None published: bool changed: tuple[str, ...] unchanged: tuple[str, ...] removed: tuple[str, ...] manifest: CorpusManifest = field(repr=False) run_id: str | None = None resumed_from: str | None = None def __repr__(self) -> str: counts = { "changed": len(self.changed), "unchanged": len(self.unchanged), "removed": len(self.removed), } return ( f"PipelineResult(status={self.status!r}, generation={self.generation!r}, " f"published={self.published!r}, counts={counts!r}, " f"run_id={self.run_id!r}, resumed_from={self.resumed_from!r})" ) def model_dump(self, mode=None): def bounded(values: tuple[str, ...]) -> list[str]: return [value[:200] for value in values[:100]] return { "status": self.status, "generation": self.generation, "published": self.published, "changed": bounded(self.changed), "unchanged": bounded(self.unchanged), "removed": bounded(self.removed), "counts": { "changed": len(self.changed), "unchanged": len(self.unchanged), "removed": len(self.removed), "documents": len(self.manifest.documents), "chunks": len(self.manifest.chunks), }, "manifest_id": self.manifest.manifest_id, "run_id": self.run_id, "resumed_from": self.resumed_from, } def _fingerprint(value) -> str: payload = json.dumps(value, sort_keys=True, separators=(",", ":"), default=str) return "sha256:" + hashlib.sha256(payload.encode()).hexdigest() def _canonical_json(value): if isinstance(value, Mapping): return {str(key): _canonical_json(value[key]) for key in sorted(value)} if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): return [_canonical_json(child) for child in value] return value def _source_snapshot(discovered) -> dict[str, dict]: snapshot = {} for _, item in discovered: modified_at = item.modified_at.astimezone(UTC) if item.modified_at else None metadata = _canonical_json(item.metadata) snapshot[item.source_id] = { "source_id": item.source_id, "uri": item.uri, "fingerprint": item.fingerprint, "modified_at": modified_at.isoformat().replace("+00:00", "Z") if modified_at else None, "metadata": metadata, "media_type": metadata.get("media_type"), "size": metadata.get("size"), } return snapshot 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, retain_published_generations: int = 3, workspace_id: str | None = None, ) -> None: self.store = store self.sources = sources self.embedder = embedder self.vector_store = vector_store self.embedding_model = embedding_model 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 self.workspace_id = workspace_id def _assert_workspace_binding(self) -> None: manifest = self.store.active_manifest() if manifest is None: if self.workspace_id is None: self.workspace_id = "default" return persisted = manifest.metadata.get("workspace_id") if not isinstance(persisted, str) or re.fullmatch( r"[a-z][a-z0-9_-]{0,63}", persisted ) is None: raise PipelineError( "corpus workspace ownership is missing or invalid; use a new corpus root or explicit rebuild" ) if self.workspace_id is None and isinstance(persisted, str): self.workspace_id = persisted return if persisted != self.workspace_id: raise PipelineError( "corpus belongs to a different workspace; use a new corpus root or explicit rebuild" ) 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: with self.store.writer_lock(): return self._gc(workspace_root=workspace_root, dry_run=dry_run) def _gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict: self._assert_workspace_binding() published = self.store.published_generations() list_vectors = getattr(self.vector_store, "list_evidence_generations", None) vector_generations = set(list_vectors("evidence", self.workspace_id)) if list_vectors else set() generations = sorted(set(published) | vector_generations) job_protected = self._protected_generations(workspace_root) active = self.store.active_generation() rollback_count = self.retain_published_generations - 1 rollback = [generation for generation in published if generation != active] keep = ({active} if active else set()) | set(rollback[-rollback_count:] if rollback_count else ()) fs_keep = keep | job_protected vector_protected = set(fs_keep) for generation in fs_keep: try: manifest = self.store.manifest(generation) except (OSError, ValueError): continue vector_protected.update( value for value in manifest.metadata.get("document_generations", {}).values() if isinstance(value, str) ) evicted, failures = [], [] filesystem_generations = set(self.store.list_generations()) for generation in generations: purge_vector = generation not in vector_protected purge_filesystem = generation in filesystem_generations and generation not in fs_keep if not purge_vector and not purge_filesystem: continue if dry_run: evicted.append(generation) continue if purge_vector: try: self.vector_store.delete_generation("evidence", generation, self.workspace_id) except Exception: failures.append({"generation": generation, "error": "vector cleanup failed"}) continue try: if purge_filesystem: 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(vector_protected), "failures": failures} def _discover(self) -> list[tuple[EvidenceSource, SourceObject]]: discovered = [] seen = set() for source in self.sources: for item in source.discover(): if item.source_id in seen: raise PipelineError("duplicate Evidence source identity") seen.add(item.source_id) discovered.append((source, item)) return sorted(discovered, key=lambda pair: pair[1].source_id) def run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult: with self.store.writer_lock(): self._assert_workspace_binding() return self._run(dry_run=dry_run, resume=resume) def run_as_job(self, **kwargs) -> PipelineResult: self.workspace_id = kwargs["workspace_id"] with self.store.writer_lock(): self._assert_workspace_binding() return self._run_as_job(**kwargs) def _run_as_job( self, *, workspace_id: str, workspace_root: Path, config_fingerprint: str, input_fingerprint: str, dry_run: bool = False, resume_run_id: str | None = None, after_stage_return=None, ) -> PipelineResult: """Execute preprocessing through the durable shared job envelope.""" discovered = self._discover() discovered_fingerprint = _fingerprint( {item.source_id: item.fingerprint for _, item in discovered} ) source_snapshot = _source_snapshot(discovered) source_by_id = {item.source_id: (source, item) for source, item in discovered} compatibility = _fingerprint({ "pipeline": self.pipeline_version, "model": self.embedding_model, "dimensions": self.embedding_dimensions, "chunk_policy": asdict(self.chunk_policy), }) job_binding = { "config_fingerprint": config_fingerprint, "input_fingerprint": input_fingerprint, "compatibility_fingerprint": compatibility, "pipeline_version": self.pipeline_version, "chunk_policy_version": self.chunk_policy.version, "embedding_model": self.embedding_model, "embedding_dimensions": self.embedding_dimensions, } previous = self.store.active_manifest() def document_sources(manifest: CorpusManifest) -> dict[str, dict]: return { document.document_id: { "document_id": document.document_id, "source_id": document.source_id, "source_uri": document.source_uri, "source_fingerprint": document.source_fingerprint, "modified_at": ( document.modified_at.isoformat().replace("+00:00", "Z") if document.modified_at else None ), "source_metadata": _canonical_json(document.metadata.get("source")), "media_type": document.media_type, "content_hash": document.content_hash, "pipeline_version": document.pipeline_version, } for document in manifest.documents } def active_assets_are_valid(manifest: CorpusManifest | None) -> bool: if manifest is None or manifest.metadata.get("workspace_id") != workspace_id: return False actual_documents = {document.source_id: document for document in manifest.documents} persisted_snapshot = _canonical_json(manifest.metadata.get("source_snapshot")) if ( not isinstance(persisted_snapshot, dict) or set(persisted_snapshot) != set(actual_documents) or manifest.metadata.get("compatibility_fingerprint") != compatibility or _canonical_json(manifest.metadata.get("document_sources")) != document_sources(manifest) ): return False for source_id, document in actual_documents.items(): source_payload = persisted_snapshot[source_id] source = SourceObject.model_validate({ "source_id": source_payload["source_id"], "uri": source_payload["uri"], "fingerprint": source_payload["fingerprint"], "modified_at": source_payload["modified_at"], "metadata": source_payload["metadata"], }) content = self.store.read_document(document.document_id, manifest.manifest_id) expected_uri = canonical_provenance_uri(source.uri) expected_id = "doc:" + hashlib.sha256( f"{source.source_id}\n{expected_uri}".encode() ).hexdigest() expected_media_type = source_payload.get("media_type") if ( content != document.content or document.document_id != expected_id or document.source_id != source.source_id or document.source_uri != expected_uri or document.source_fingerprint != source.fingerprint or document.modified_at != source.modified_at or _canonical_json(document.metadata.get("source")) != _canonical_json(source.metadata) or ( isinstance(expected_media_type, str) and document.media_type != expected_media_type ) or document.pipeline_version != self.pipeline_version ): return False expected_chunks = tuple( part for document in manifest.documents for part in chunk(document, self.chunk_policy) ) if any(document.content and not chunk(document, self.chunk_policy) for document in manifest.documents): return False if _canonical_json([part.model_dump(mode="json") for part in manifest.chunks]) != ( _canonical_json([part.model_dump(mode="json") for part in expected_chunks]) ): return False generations = manifest.metadata.get("document_generations") if not isinstance(generations, Mapping): return False health = self.vector_store.health() if ( not health.ok or health.dimension_compatible is not True or health.expected_dimension != self.embedding_dimensions or health.observed_dimensions != (self.embedding_dimensions,) ): return False existing = self.vector_store.existing_hashes("evidence", ["evidence"]) for part in expected_chunks: generation = generations.get(part.document_id) if not isinstance(generation, str): return False record_id = f"{workspace_id}:{generation}:{part.chunk_id}" if existing.get(record_id) != part.content_hash: return False return True try: active_assets_valid = active_assets_are_valid(previous) except Exception: active_assets_valid = False reusable = ( active_assets_valid and _canonical_json(previous.metadata.get("source_snapshot")) == source_snapshot and _canonical_json(previous.metadata.get("job_binding")) == job_binding ) if not dry_run and resume_run_id is None and reusable: return PipelineResult( "succeeded", previous.manifest_id, False, (), tuple(sorted(item.source_id for _, item in discovered)), (), previous, ) spec = JobSpec( workspace_id=workspace_id, job_type="evidence", workspace_root=workspace_root, spec_version="jobs-v1", pipeline_version=self.pipeline_version, config_fingerprint=config_fingerprint, input_fingerprint=_fingerprint([input_fingerprint, discovered_fingerprint]), stage_ids=EVIDENCE_STAGE_IDS, dry_run=dry_run, resume_run_id=resume_run_id, ) def artifact(context: JobContext, name: str) -> Path: root = context.run_dir / "artifacts" root.mkdir(exist_ok=True) return root / name def write(context: JobContext, name: str, value) -> None: artifact(context, name).write_text( json.dumps(value, sort_keys=True, separators=(",", ":")), encoding="utf-8" ) def read(context: JobContext, name: str): try: return json.loads(artifact(context, name).read_text(encoding="utf-8")) except (OSError, ValueError) as error: raise PipelineError("preprocessing checkpoint artifact is corrupt") from error def discover_stage(context: JobContext) -> None: previous = self.store.active_manifest() prior = {doc.source_id: doc for doc in previous.documents} if previous else {} fingerprints = {item.source_id: item.fingerprint for _, item in discovered} previous_snapshot = ( _canonical_json(previous.metadata.get("source_snapshot")) if previous else {} ) rebuild = bool(previous and not active_assets_valid) changed = sorted( item.source_id for _, item in discovered if rebuild or item.source_id not in prior or previous_snapshot.get(item.source_id) != source_snapshot[item.source_id] ) unchanged = sorted(set(fingerprints) - set(changed)) removed = sorted(set(prior) - set(fingerprints)) write(context, "plan.json", { "generation": f"gen:{context.run_id}", "compatibility": compatibility, "job_binding": job_binding, "source_snapshot": source_snapshot, "fingerprints": fingerprints, "changed": changed, "unchanged": unchanged, "removed": removed, "previous": previous.model_dump(mode="json") if previous else None, }) return StageArtifacts(("plan.json",)) def acquire_stage(context: JobContext) -> None: if context.dry_run: return StageArtifacts() plan = read(context, "plan.json") previous = CorpusManifest.model_validate(plan["previous"]) if plan["previous"] else None prior = {doc.source_id: doc for doc in previous.documents} if previous else {} documents = [prior[source_id] for source_id in plan["unchanged"]] for source_id in plan["changed"]: source, item = source_by_id[source_id] documents.append(normalize(source.acquire(item), self.pipeline_version)) documents.sort(key=lambda value: value.source_id) chunks = [part for document in documents for part in chunk(document, self.chunk_policy)] previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {} changed = set(plan["changed"]) generations = { document.document_id: ( plan["generation"] if document.source_id in changed else previous_generations.get(document.document_id, previous.vector_generation) ) for document in documents } manifest = CorpusManifest( pipeline_version=self.pipeline_version, embedding_model=self.embedding_model, embedding_dimensions=self.embedding_dimensions, vector_generation=plan["generation"], documents=tuple(documents), chunks=tuple(chunks), metadata={ "workspace_id": self.workspace_id, "compatibility_fingerprint": compatibility, "job_binding": plan["job_binding"], "source_snapshot": plan["source_snapshot"], "document_sources": { document.document_id: { "document_id": document.document_id, "source_id": document.source_id, "source_uri": document.source_uri, "source_fingerprint": document.source_fingerprint, "modified_at": ( document.modified_at.isoformat().replace("+00:00", "Z") if document.modified_at else None ), "source_metadata": _canonical_json( document.metadata.get("source") ), "media_type": document.media_type, "content_hash": document.content_hash, "pipeline_version": document.pipeline_version, } for document in documents }, "fingerprints": plan["fingerprints"], "removed": plan["removed"], "document_generations": generations, }, ) write(context, "manifest.json", manifest.model_dump(mode="json")) return StageArtifacts(("manifest.json",)) def embed_stage(context: JobContext) -> None: if context.dry_run: return StageArtifacts() plan = read(context, "plan.json") manifest = CorpusManifest.model_validate(read(context, "manifest.json")) changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]} parts = [part for part in manifest.chunks if part.document_id in changed_docs] embeddings = self.embedder.embed_documents([part.content for part in parts]) if len(embeddings) != len(parts) or any( len(vector) != self.embedding_dimensions for vector in embeddings ): raise PipelineError("embedding output is incompatible") write(context, "embeddings.json", embeddings) return StageArtifacts(("embeddings.json",)) def records(context: JobContext): plan = read(context, "plan.json") manifest = CorpusManifest.model_validate(read(context, "manifest.json")) changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]} parts = [part for part in manifest.chunks if part.document_id in changed_docs] embeddings = read(context, "embeddings.json") return [self._vector_record(part, vector, plan["generation"], self.workspace_id) for part, vector in zip(parts, embeddings, strict=True)] def compensate(context: JobContext) -> None: generation = read(context, "plan.json")["generation"] if self.store.active_generation() != generation: self.store.discard(generation) try: self.vector_store.delete_generation("evidence", generation, self.workspace_id) except Exception: pass write(context, "compensated.json", {"generation": generation}) def rotate_compensated_generation(context: JobContext) -> None: marker = artifact(context, "compensated.json") if not marker.exists(): return plan = read(context, "plan.json") old = plan["generation"] plan["generation"] = f"gen:{uuid.uuid4().hex}" write(context, "plan.json", plan) manifest = CorpusManifest.model_validate(read(context, "manifest.json")) changed = set(plan["changed"]) generations = dict(manifest.metadata["document_generations"]) for document in manifest.documents: if document.source_id in changed and generations.get(document.document_id) == old: generations[document.document_id] = plan["generation"] manifest_payload = manifest.model_dump(mode="json") manifest_payload["metadata"]["document_generations"] = generations manifest_payload["vector_generation"] = plan["generation"] manifest = CorpusManifest.model_validate(manifest_payload) write(context, "manifest.json", manifest.model_dump(mode="json")) marker.unlink() def vector_stage(context: JobContext) -> None: if context.dry_run: return StageArtifacts() rotate_compensated_generation(context) values = records(context) write(context, "vector-intent.json", { "generation": read(context, "plan.json")["generation"], "records": {value.record.id: value.content_hash for value in values}, }) seal_stage_artifacts( context, "vector_upsert", ("plan.json", "manifest.json", "vector-intent.json"), spec, ) try: existing = self.vector_store.existing_hashes("evidence", ["evidence"]) missing = [ value for value in values if existing.get(value.record.id) != value.content_hash ] if missing and self.vector_store.upsert("evidence", missing) != len(missing): raise PipelineError("vector write count mismatch") except Exception: compensate(context) raise return StageArtifacts(("plan.json", "manifest.json", "vector-intent.json")) def stage_stage(context: JobContext) -> None: if context.dry_run: return StageArtifacts() plan = read(context, "plan.json") manifest = CorpusManifest.model_validate(read(context, "manifest.json")) recovered = False try: if artifact(context, "compensated.json").exists(): recovered = True rotate_compensated_generation(context) values = records(context) existing = self.vector_store.existing_hashes("evidence", ["evidence"]) missing = [value for value in values if existing.get(value.record.id) != value.content_hash] if missing and self.vector_store.upsert("evidence", missing) != len(missing): raise PipelineError("vector write count mismatch") plan = read(context, "plan.json") manifest = CorpusManifest.model_validate(read(context, "manifest.json")) if not self.store.generation_path(plan["generation"]).exists(): self.store.stage( manifest, {doc.document_id: doc.content for doc in manifest.documents}, generation=plan["generation"], ) self.store.manifest(plan["generation"]) except Exception: compensate(context) raise return StageArtifacts( ("plan.json", "manifest.json", "vector-intent.json") if recovered else () ) def publish_stage(context: JobContext) -> None: if context.dry_run: return StageArtifacts() if artifact(context, "compensated.json").exists(): rotate_compensated_generation(context) values = records(context) try: existing = self.vector_store.existing_hashes("evidence", ["evidence"]) missing = [value for value in values if existing.get(value.record.id) != value.content_hash] if missing and self.vector_store.upsert("evidence", missing) != len(missing): raise PipelineError("vector write count mismatch") except Exception: compensate(context) raise manifest = CorpusManifest.model_validate(read(context, "manifest.json")) generation = read(context, "plan.json")["generation"] try: if not self.store.generation_path(generation).exists(): self.store.stage( manifest, {doc.document_id: doc.content for doc in manifest.documents}, generation=generation, ) except Exception: compensate(context) raise generation = read(context, "plan.json")["generation"] try: self.store.publish(generation) except Exception: compensate(context) raise return StageArtifacts(("plan.json", "manifest.json", "vector-intent.json")) def retention_stage(context: JobContext) -> None: if not context.dry_run: self.gc(workspace_root=workspace_root) report = run_job(spec, [ discover_stage, acquire_stage, embed_stage, vector_stage, stage_stage, publish_stage, retention_stage, ], after_stage_return=after_stage_return) run_dir = workspace_root / ".tht-jobs" / "evidence" / "runs" / report.run_id plan = json.loads((run_dir / "artifacts" / "plan.json").read_text()) if dry_run: manifest = self.store.active_manifest() or CorpusManifest(pipeline_version=self.pipeline_version) generation = None published = False elif report.status == "succeeded": generation = plan["generation"] manifest = self.store.manifest(generation) published = True else: generation = plan["generation"] manifest_path = run_dir / "artifacts" / "manifest.json" manifest = (CorpusManifest.model_validate_json(manifest_path.read_text()) if manifest_path.exists() else CorpusManifest(pipeline_version=self.pipeline_version)) published = False return PipelineResult( report.status, generation, published, tuple(plan["changed"]), tuple(plan["unchanged"]), tuple(plan["removed"]), manifest, report.run_id, report.resumed_from, ) def _run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult: generation = None vector_written = False previous = self.store.active_manifest() try: discovered = self._discover() except Exception as error: raise PipelineError("Evidence discovery failed") from error prior_documents = {doc.source_id: doc for doc in previous.documents} if previous else {} fingerprints = {item.source_id: item.fingerprint for _, item in discovered} compatibility = _fingerprint({ "pipeline": self.pipeline_version, "model": self.embedding_model, "dimensions": self.embedding_dimensions, "chunk_policy": asdict(self.chunk_policy), }) previous_compatibility = previous.metadata.get("compatibility_fingerprint") if previous else None rebuild = previous is not None and compatibility != previous_compatibility changed = tuple(item.source_id for _, item in discovered if rebuild or prior_documents.get(item.source_id) is None or prior_documents[item.source_id].source_fingerprint != item.fingerprint) unchanged = tuple(item.source_id for _, item in discovered if item.source_id not in changed) removed = tuple(sorted(set(prior_documents) - set(fingerprints))) if dry_run: manifest = previous or CorpusManifest(pipeline_version=self.pipeline_version) return PipelineResult("succeeded", None, False, changed, unchanged, removed, manifest) if previous is not None and not changed and not removed: return PipelineResult( "succeeded", previous.manifest_id, False, changed, unchanged, removed, previous ) documents: list[CanonicalDocument] = [prior_documents[source_id] for source_id in unchanged] changed_set = set(changed) try: for source, item in discovered: if item.source_id in changed_set: documents.append(normalize(source.acquire(item), self.pipeline_version)) documents.sort(key=lambda document: document.source_id) chunks: list[CanonicalChunk] = [] for document in documents: chunks.extend(chunk(document, self.chunk_policy)) generation = resume or f"gen:{uuid.uuid4().hex}" previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {} document_generations = { document.document_id: ( generation if document.source_id in changed_set else previous_generations.get(document.document_id, previous.vector_generation) ) for document in documents } manifest = CorpusManifest( pipeline_version=self.pipeline_version, embedding_model=self.embedding_model, embedding_dimensions=self.embedding_dimensions, vector_generation=generation, documents=tuple(documents), chunks=tuple(chunks), metadata={ "workspace_id": self.workspace_id, "compatibility_fingerprint": compatibility, "fingerprints": fingerprints, "removed": list(removed), "document_generations": document_generations, }, ) changed_documents = {document.document_id for document in documents if document.source_id in changed_set} changed_chunks = [part for part in chunks if part.document_id in changed_documents] embeddings = self.embedder.embed_documents([part.content for part in changed_chunks]) if len(embeddings) != len(changed_chunks): raise PipelineError("embedding count mismatch") if any(len(vector) != self.embedding_dimensions for vector in embeddings): raise PipelineError("embedding dimension mismatch") records = [self._vector_record(part, vector, generation, self.workspace_id) for part, vector in zip(changed_chunks, embeddings, strict=True)] if records: written = self.vector_store.upsert("evidence", records) vector_written = True if written != len(records): raise PipelineError("vector write count mismatch") generation_path = self.store.generation_path(generation) if resume is not None and generation_path.exists(): staged_manifest = self.store.manifest(generation) expected = manifest.model_dump(mode="json", exclude={"created_at", "manifest_id"}) actual = staged_manifest.model_dump(mode="json", exclude={"created_at", "manifest_id"}) actual["metadata"].pop("files", None) if actual != expected: raise PipelineError("resume generation is incompatible") staged = generation else: staged = self.store.stage( manifest, {document.document_id: document.content for document in documents}, generation=generation, ) self.store.publish(staged) self.gc(workspace_root=self.store.root.parent) except PipelineError: self._compensate(generation, vector_written) raise except Exception as error: self._compensate(generation, vector_written) raise PipelineError("Evidence preprocessing failed") from error return PipelineResult("succeeded", generation, True, changed, unchanged, removed, self.store.manifest(generation)) def _compensate(self, generation: str | None, vector_written: bool) -> None: if generation is None: return try: self.store.discard(generation) except Exception: pass if vector_written: try: self.vector_store.delete_generation("evidence", generation, self.workspace_id) except Exception: pass @staticmethod def _vector_record( chunk: CanonicalChunk, embedding: list[float], generation: str, workspace_id: str, ): record = VectorRecord( id=f"{workspace_id}:{generation}:{chunk.chunk_id}", kind="evidence", ref=chunk.document_id, title=str(chunk.metadata.get("title", "")), content=chunk.content, metadata={ **dict(chunk.metadata), "document_id": chunk.document_id, "workspace_id": workspace_id, "source_uri": chunk.source_uri, "ordinal": chunk.ordinal, "vector_generation": generation, }, ) return VectorWriteRecord(record=record, embedding=embedding, content_hash=chunk.content_hash)