feat(preprocess): publish incremental Evidence corpus

This commit is contained in:
2026-07-12 04:26:24 +02:00
parent 16a8bd9df6
commit b96d4f13b9
10 changed files with 708 additions and 0 deletions
+175
View File
@@ -0,0 +1,175 @@
"""Incremental Evidence preprocessing with generation-isolated vector writes."""
from __future__ import annotations
import hashlib
import json
import uuid
from dataclasses import asdict, dataclass
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
from tht.ports.vector import VectorStore, VectorWriteRecord
from tht.vectorstore.records import VectorRecord
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
def model_dump(self, mode=None):
value = asdict(self)
value["manifest"] = self.manifest.model_dump(mode="json")
return value
def _fingerprint(value) -> str:
payload = json.dumps(value, sort_keys=True, separators=(",", ":"), default=str)
return "sha256:" + hashlib.sha256(payload.encode()).hexdigest()
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,
) -> 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
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():
return self._run(dry_run=dry_run, resume=resume)
def _run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult:
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={
"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) for part, vector in zip(changed_chunks, embeddings, strict=True)]
if records:
written = self.vector_store.upsert("evidence", records)
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)
except PipelineError:
raise
except Exception as error:
raise PipelineError("Evidence preprocessing failed") from error
return PipelineResult("succeeded", generation, True, changed, unchanged, removed, self.store.manifest(generation))
@staticmethod
def _vector_record(chunk: CanonicalChunk, embedding: list[float], generation: str):
record = VectorRecord(
id=f"{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,
"source_uri": chunk.source_uri, "ordinal": chunk.ordinal,
"vector_generation": generation,
},
)
return VectorWriteRecord(record=record, embedding=embedding, content_hash=chunk.content_hash)
+156
View File
@@ -0,0 +1,156 @@
"""Durable immutable corpus generations and an atomic ACTIVE pointer."""
from __future__ import annotations
import json
import fcntl
import os
import re
import stat
import uuid
from pathlib import Path
from contextlib import contextmanager
from tht.corpus.models import CorpusManifest
_GENERATION = re.compile(r"^gen:[0-9a-f]{32}$")
class UnsafeCorpusPath(RuntimeError):
pass
def _atomic_write(path: Path, payload: bytes) -> None:
temporary = path.with_name(f".{path.name}.{uuid.uuid4().hex}.tmp")
fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600)
try:
with os.fdopen(fd, "wb") as stream:
stream.write(payload)
stream.flush()
os.fsync(stream.fileno())
os.replace(temporary, path)
directory = os.open(path.parent, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW)
try:
os.fsync(directory)
finally:
os.close(directory)
except BaseException:
temporary.unlink(missing_ok=True)
raise
class CorpusStore:
def __init__(self, root: Path) -> None:
self.root = Path(root)
self.active_path = self.root / "ACTIVE"
self._replace = os.replace
self._ensure_root()
def _ensure_root(self) -> None:
if self.root.is_symlink():
raise UnsafeCorpusPath("corpus root must not be a symlink")
self.root.mkdir(parents=True, exist_ok=True, mode=0o700)
info = self.root.lstat()
if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid():
raise UnsafeCorpusPath("corpus root is unsafe")
@contextmanager
def writer_lock(self):
lock_path = self.root / ".writer.lock"
fd = os.open(lock_path, os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW | os.O_CLOEXEC, 0o600)
try:
info = os.fstat(fd)
if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() or info.st_nlink != 1:
raise UnsafeCorpusPath("corpus writer lock is unsafe")
fcntl.flock(fd, fcntl.LOCK_EX)
yield
finally:
fcntl.flock(fd, fcntl.LOCK_UN)
os.close(fd)
def generation_path(self, generation: str) -> Path:
if not _GENERATION.fullmatch(generation):
raise UnsafeCorpusPath("invalid corpus generation")
path = self.root / generation.replace(":", "-")
if path.is_symlink():
raise UnsafeCorpusPath("generation must not be a symlink")
return path
def stage(
self,
manifest: CorpusManifest,
materialized: dict[str, str],
*,
generation: str | None = None,
) -> str:
generation = generation or f"gen:{uuid.uuid4().hex}"
path = self.generation_path(generation)
try:
path.mkdir(mode=0o700)
except FileExistsError:
raise UnsafeCorpusPath("generation already exists") from None
documents = path / "documents"
documents.mkdir(mode=0o700)
files: dict[str, str] = {}
for document in manifest.documents:
relative = f"documents/{document.document_id.removeprefix('doc:')}.md"
_atomic_write(path / relative, materialized[document.document_id].encode("utf-8"))
files[document.document_id] = relative
payload = json.loads(manifest.model_dump_json())
metadata = payload["metadata"]
metadata["files"] = files
payload.update({"manifest_id": generation, "metadata": metadata})
staged = CorpusManifest.model_validate(payload)
_atomic_write(path / "manifest.json", (staged.model_dump_json(indent=2) + "\n").encode())
return generation
def publish(self, generation: str) -> str:
manifest = self.manifest(generation)
if manifest.manifest_id != generation:
raise UnsafeCorpusPath("manifest generation mismatch")
temporary = self.active_path.with_name(f".ACTIVE.{uuid.uuid4().hex}.tmp")
_atomic_write(temporary, (generation + "\n").encode())
self._replace(temporary, self.active_path)
directory = os.open(self.root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW)
try:
os.fsync(directory)
finally:
os.close(directory)
return generation
def active_generation(self) -> str | None:
try:
if self.active_path.is_symlink():
raise UnsafeCorpusPath("ACTIVE must not be a symlink")
value = self.active_path.read_text(encoding="ascii").strip()
except FileNotFoundError:
return None
if not _GENERATION.fullmatch(value):
raise UnsafeCorpusPath("ACTIVE contains an invalid generation")
return value
def manifest(self, generation: str) -> CorpusManifest:
path = self.generation_path(generation)
manifest_path = path / "manifest.json"
if manifest_path.is_symlink():
raise UnsafeCorpusPath("manifest must not be a symlink")
return CorpusManifest.model_validate_json(manifest_path.read_text(encoding="utf-8"))
def active_manifest(self) -> CorpusManifest | None:
generation = self.active_generation()
return self.manifest(generation) if generation else None
def resolve_document(self, document_id: str, generation: str | None = None) -> Path | None:
generation = generation or self.active_generation()
if generation is None:
return None
manifest = self.manifest(generation)
relative = manifest.metadata.get("files", {}).get(document_id)
if not isinstance(relative, str):
return None
base = self.generation_path(generation).resolve()
path = (base / relative).resolve()
if not path.is_relative_to(base) or path.is_symlink():
raise UnsafeCorpusPath("materialized document escapes its generation")
return path