From b96d4f13b9b77e18d0016c9a903e08f72c9e6853 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sun, 12 Jul 2026 04:26:24 +0200 Subject: [PATCH] feat(preprocess): publish incremental Evidence corpus --- .superpowers/sdd/evidence-task-5-report.md | 61 +++++++ harness/tests/test_corpus_pipeline.py | 132 ++++++++++++++++ harness/tests/test_corpus_publish.py | 43 +++++ harness/tests/test_preprocess_cli.py | 32 ++++ harness/tht/cli/__init__.py | 2 + harness/tht/cli/preprocess_cmd.py | 65 ++++++++ harness/tht/corpus/pipeline.py | 175 +++++++++++++++++++++ harness/tht/corpus/store.py | 156 ++++++++++++++++++ harness/tht/search/evidence.py | 31 ++++ harness/tht/session/artifacts.py | 11 ++ 10 files changed, 708 insertions(+) create mode 100644 .superpowers/sdd/evidence-task-5-report.md create mode 100644 harness/tests/test_corpus_pipeline.py create mode 100644 harness/tests/test_corpus_publish.py create mode 100644 harness/tests/test_preprocess_cli.py create mode 100644 harness/tht/cli/preprocess_cmd.py create mode 100644 harness/tht/corpus/pipeline.py create mode 100644 harness/tht/corpus/store.py create mode 100644 harness/tht/search/evidence.py diff --git a/.superpowers/sdd/evidence-task-5-report.md b/.superpowers/sdd/evidence-task-5-report.md new file mode 100644 index 00000000..ad45c3de --- /dev/null +++ b/.superpowers/sdd/evidence-task-5-report.md @@ -0,0 +1,61 @@ +# Evidence Task 5 report + +## Outcome + +Implemented an incremental Evidence corpus pipeline with immutable materialized generations, +generation-scoped vector records, and an fsynced atomic `ACTIVE` pointer. Runtime Evidence +artifact lookup reads the active canonical manifest and keeps a legacy source-tree fallback only +when no corpus has been published. + +The CLI is available as `tht preprocess evidence [--dry-run] [--resume RUN_ID] [--json]`. +JSON success and failure output is pristine and failure details are sanitized. + +## Safety and failure model + +- A workspace writer lock serializes preprocess writers; readers never take the lock. +- Generation directories, manifests, materialized files, locks, and `ACTIVE` reject symlink/path + escape cases and use owner-only durable writes. +- Vector records use generation-specific keys and metadata. The active manifest maps each active + document to its valid vector generation, allowing unchanged documents to retain their vectors. +- Runtime retrieval admits only active document IDs and their manifest-selected generations. + Removed documents and partial writes from failed generations are therefore unreachable. +- Embedding count and dimension checks occur before vector upsert; vector write count is checked + before staging/publish. Any failure leaves `ACTIVE` unchanged. +- Dry runs perform discovery/fingerprint planning only and never acquire, embed, write vectors, or + publish. Fully unchanged runs return the active generation without creating a replacement. +- Resume can safely retry idempotent generation-scoped upserts and publish an already staged, + compatibility-checked generation after a crash between staging and pointer replacement. + +## TDD evidence + +Initial focused collection failed because `tht.corpus.pipeline` and `tht.corpus.store` did not +exist. The implemented suite covers incremental skips, removals, model/policy rebuilds, acquire and +partial-vector failures, dry-run isolation, dimension validation, atomic reader snapshots, pointer +validation, symlink defense, and pristine CLI JSON. + +Fresh focused verification: + +```text +18 passed, 3 warnings in 0.39s +``` + +Command: + +```text +.venv/bin/pytest tests/test_corpus_pipeline.py tests/test_corpus_publish.py \ + tests/test_preprocess_cli.py tests/test_search_pack.py tests/test_session_documents.py -q +``` + +Scoped Ruff: `All checks passed!` + +Broader non-Docker/non-packaging run reached `560 passed, 5 deselected`; ten pre-existing HTTP +adapter tests could not bind localhost under the sandbox. The complete suite reached `570 passed, +5 deselected`, with the remaining failures/errors caused by denied Docker socket, localhost bind, +and offline wheel-build access. No task-focused test failed. + +## Remaining operational gate + +Live pgvector integration needs Docker or an authorized local pgvector endpoint. The compensation +strategy is logical isolation rather than destructive cleanup because the shared `VectorStore` +port intentionally exposes no delete/transaction API; unreachable failed generations can be +garbage-collected by a future maintenance job. diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py new file mode 100644 index 00000000..3c41c3ef --- /dev/null +++ b/harness/tests/test_corpus_pipeline.py @@ -0,0 +1,132 @@ +import pytest + +from tht.corpus.chunk import ChunkPolicy +from tht.corpus.pipeline import CorpusPipeline, PipelineError +from tht.corpus.store import CorpusStore +from tht.ports.evidence import AcquiredDocument, SourceObject +from tht.ports.vector import VectorCapabilities + + +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()) + + +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 = [] + + 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 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): + 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", + ) + + +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_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_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 diff --git a/harness/tests/test_corpus_publish.py b/harness/tests/test_corpus_publish.py new file mode 100644 index 00000000..a2fb2268 --- /dev/null +++ b/harness/tests/test_corpus_publish.py @@ -0,0 +1,43 @@ +import pytest + +from tht.corpus.models import CorpusManifest +from tht.corpus.store import CorpusStore, UnsafeCorpusPath + + +def test_publish_switches_active_atomically_and_resolves_materialized_files(tmp_path): + store = CorpusStore(tmp_path / "corpus") + generation = store.stage(CorpusManifest(), {}) + seen = [] + store._replace = lambda source, target: (seen.append(source.read_text()), source.replace(target)) + published = store.publish(generation) + assert published == generation + assert store.active_generation() == generation + assert seen == [generation + "\n"] + + +def test_active_manifest_is_a_consistent_reader_snapshot(tmp_path): + store = CorpusStore(tmp_path / "corpus") + first = store.stage(CorpusManifest(metadata={"name": "first"}), {}) + second = store.stage(CorpusManifest(metadata={"name": "second"}), {}) + store.publish(first) + snapshot = store.active_manifest() + store.publish(second) + assert snapshot.metadata["name"] == "first" + assert store.active_manifest().metadata["name"] == "second" + + +def test_store_rejects_symlinked_generation_root(tmp_path): + outside = tmp_path / "outside" + outside.mkdir() + root = tmp_path / "corpus" + root.symlink_to(outside, target_is_directory=True) + with pytest.raises(UnsafeCorpusPath): + CorpusStore(root) + + +def test_active_pointer_cannot_escape_generation_root(tmp_path): + store = CorpusStore(tmp_path / "corpus") + store.root.mkdir(parents=True, exist_ok=True) + store.active_path.write_text("../outside\n") + with pytest.raises(UnsafeCorpusPath): + store.active_manifest() diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py new file mode 100644 index 00000000..bc48c08e --- /dev/null +++ b/harness/tests/test_preprocess_cli.py @@ -0,0 +1,32 @@ +import json +from types import SimpleNamespace + +from typer.testing import CliRunner + +from tht.cli import app + + +def test_preprocess_evidence_json_is_pristine(monkeypatch, tmp_path): + import tht.cli.preprocess_cmd as command + + result = SimpleNamespace(model_dump=lambda mode=None: { + "status": "succeeded", "generation": "gen:abc", "published": True + }) + monkeypatch.setattr(command, "run_from_config", lambda *args, **kwargs: result) + response = CliRunner().invoke( + app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")] + ) + assert response.exit_code == 0, response.output + assert json.loads(response.output)["generation"] == "gen:abc" + + +def test_preprocess_failure_is_structured_and_nonzero(monkeypatch, tmp_path): + import tht.cli.preprocess_cmd as command + + monkeypatch.setattr(command, "run_from_config", lambda *a, **k: (_ for _ in ()).throw(RuntimeError("secret detail"))) + response = CliRunner().invoke( + app, ["preprocess", "evidence", "--json", "-c", str(tmp_path / "workspace.yaml")] + ) + assert response.exit_code != 0 + assert json.loads(response.output) == {"status": "failed", "error": "preprocessing failed"} + assert "secret detail" not in response.output diff --git a/harness/tht/cli/__init__.py b/harness/tht/cli/__init__.py index afa18f6f..b95cd368 100644 --- a/harness/tht/cli/__init__.py +++ b/harness/tht/cli/__init__.py @@ -48,6 +48,7 @@ from tht.cli.lsh_cmd import lsh_app # noqa: E402 from tht.cli.memory_cmd import memory_app # noqa: E402 from tht.cli.ollama_cmd import ollama_app # noqa: E402 from tht.cli.phase_cmd import phase_app # noqa: E402 +from tht.cli.preprocess_cmd import preprocess_app # noqa: E402 from tht.cli.schema_cmd import schema_app # noqa: E402 from tht.cli.search_cmd import search_app # noqa: E402 from tht.cli.session_cmd import session_app # noqa: E402 @@ -56,6 +57,7 @@ from tht.cli.vector_cmd import vector_app # noqa: E402 import tht.cli.vector_migrate_cmd # noqa: E402, F401 app.add_typer(phase_app, name="phase") +app.add_typer(preprocess_app, name="preprocess") app.add_typer(config_app, name="config") app.add_typer(schema_app, name="schema") app.add_typer(session_app, name="session") diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py new file mode 100644 index 00000000..fbbb330e --- /dev/null +++ b/harness/tht/cli/preprocess_cmd.py @@ -0,0 +1,65 @@ +"""One-shot preprocessing commands.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import typer + +from tht.cli.config_cmd import CONFIG_OPT + + +preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts") + + +def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None): + from tht.adapters.factory import build_evidence_sources, build_vector_store + from tht.cli.schema_cmd import _load_config_or_exit + from tht.cli.vector_cmd import make_embedder + from tht.corpus.chunk import ChunkPolicy + from tht.corpus.pipeline import CorpusPipeline + from tht.corpus.store import CorpusStore + + cfg = _load_config_or_exit(config) + if cfg.embeddings is None: + raise RuntimeError("embeddings are not configured") + corpus_root = cfg.paths.artifacts.parent / "corpus" + generation = None + if resume: + generation = resume if resume.startswith("gen:") else f"gen:{resume}" + pipeline = CorpusPipeline( + store=CorpusStore(corpus_root), sources=build_evidence_sources(cfg), + embedder=make_embedder(cfg.embeddings), + vector_store=build_vector_store(cfg, require_write=True), + embedding_model=cfg.embeddings.model, embedding_dimensions=cfg.embeddings.dim, + chunk_policy=ChunkPolicy(version="chunk-v1", max_chars=cfg.vector.max_chunk_chars), + pipeline_version="evidence-v1", + ) + return pipeline.run(dry_run=dry_run, resume=generation) + + +@preprocess_app.command("evidence") +def evidence_cmd( + config: Path = CONFIG_OPT, + dry_run: bool = typer.Option(False, "--dry-run"), + resume: str | None = typer.Option(None, "--resume"), + json_output: bool = typer.Option(False, "--json"), +) -> None: + try: + result = run_from_config(config, dry_run=dry_run, resume=resume) + except Exception: + payload = {"status": "failed", "error": "preprocessing failed"} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: preprocessing failed", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None + payload = result.model_dump(mode="json") + if json_output: + typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) + else: + typer.echo( + f"OK: generation={payload['generation']} changed={len(payload['changed'])} " + f"unchanged={len(payload['unchanged'])} removed={len(payload['removed'])}" + ) diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py new file mode 100644 index 00000000..5da17151 --- /dev/null +++ b/harness/tht/corpus/pipeline.py @@ -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) diff --git a/harness/tht/corpus/store.py b/harness/tht/corpus/store.py new file mode 100644 index 00000000..5dd60e0d --- /dev/null +++ b/harness/tht/corpus/store.py @@ -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 diff --git a/harness/tht/search/evidence.py b/harness/tht/search/evidence.py new file mode 100644 index 00000000..e58dca74 --- /dev/null +++ b/harness/tht/search/evidence.py @@ -0,0 +1,31 @@ +"""Runtime Evidence lookup bound to the atomically active corpus generation.""" + +from tht.corpus.store import CorpusStore + + +def active_evidence_hits(store: CorpusStore, vector_store, embedding, *, limit: int): + manifest = store.active_manifest() + if manifest is None or manifest.vector_generation is None: + return [] + document_generations = dict(manifest.metadata.get("document_generations", {})) + active_documents = {document.document_id for document in manifest.documents} + hits = vector_store.search(["evidence"], embedding, limit=limit, kinds=["evidence"]) + return [ + hit for hit in hits + if hit.metadata.get("document_id") in active_documents + and hit.metadata.get("vector_generation") + == document_generations.get(hit.metadata.get("document_id"), manifest.vector_generation) + ] + + +def resolve_evidence_file(store: CorpusStore, evidence_id: str) -> str: + manifest = store.active_manifest() + if manifest is None: + return "" + for document in manifest.documents: + frontmatter = document.metadata.get("frontmatter", {}) + identifiers = {document.document_id, document.source_id, str(frontmatter.get("id", ""))} + if evidence_id in identifiers: + path = store.resolve_document(document.document_id) + return str(path) if path else "" + return "" diff --git a/harness/tht/session/artifacts.py b/harness/tht/session/artifacts.py index 2758ba06..6672ce80 100644 --- a/harness/tht/session/artifacts.py +++ b/harness/tht/session/artifacts.py @@ -5,6 +5,17 @@ from tht.session.models import SchemaLinking def _find_evidence_file(evidence_root: Path, evidence_id: str) -> str: + # New deployments resolve only immutable materialized files from ACTIVE. Keep + # the legacy curated-tree fallback for sessions created before a corpus exists. + try: + from tht.corpus.store import CorpusStore + from tht.search.evidence import resolve_evidence_file + + corpus_root = evidence_root.parent.parent / "corpus" + if corpus_root.exists(): + return resolve_evidence_file(CorpusStore(corpus_root), evidence_id) + except (OSError, RuntimeError, ValueError): + pass for match in evidence_root.rglob(f"{evidence_id}.md"): return str(match) return ""