diff --git a/.superpowers/sdd/evidence-task-5c-report.md b/.superpowers/sdd/evidence-task-5c-report.md new file mode 100644 index 00000000..62a6a7b8 --- /dev/null +++ b/.superpowers/sdd/evidence-task-5c-report.md @@ -0,0 +1,31 @@ +# Evidence Task 5C report + +## Delivered + +- Added `vector.retain_published_generations` (default `3`, validation minimum `1`). +- Retention runs only after publication. It keeps ACTIVE, the newest configured generations, + and generations referenced by running or resumable failed job checkpoints. +- Cleanup deletes the exact Evidence generation from the vector store before removing its + immutable filesystem directory. Vector failures retain filesystem metadata for retry and + produce credential-free partial reports. +- Added idempotent `tht preprocess evidence gc [--dry-run] --json` reconciliation with pristine + JSON output. +- Materialized document reads now open generation/documents components with directory file + descriptors and `O_NOFOLLOW`, require a regular file owned by the process with one link, and + hash the bytes read from the same descriptor against the canonical manifest. +- HTTP generation deletion is pinned to `delete_vector_generation` with exact + table/kind/generation arguments. Legacy 404 responses fail closed with an actionable, + sanitized migration message. + +## Evidence + +- Focused retention, safe-read, CLI, and HTTP contract tests: `51 passed` (Docker-backed direct + parametrizations excluded from that focused invocation). +- Real Docker pgvector adapter suites: `33 passed`. +- Full harness suite, including Docker-backed tests: `668 passed, 5 deselected`. +- Changed-file Ruff: clean. +- `git diff --check`: clean. + +The five deselected tests are the repository's opt-in `l2` tests requiring external services; +they are not local pgvector tests. Test output retains pre-existing Pydantic serialization and +legacy-config deprecation warnings. diff --git a/harness/tests/l0/test_vector_adapter_parity.py b/harness/tests/l0/test_vector_adapter_parity.py index deab75a7..b1655d0f 100644 --- a/harness/tests/l0/test_vector_adapter_parity.py +++ b/harness/tests/l0/test_vector_adapter_parity.py @@ -8,7 +8,7 @@ from tht.adapters.vector.pgvector import PgVectorStore from tht.adapters.vector.thoth_http import ThothHttpVectorStore from tht.config import DatabaseConfig, RestConfig from tht.ports.vector import VectorRecord, VectorStoreError, VectorWriteRecord -from tht.vectorstore.rest_client import VectorRestClient +from tht.vectorstore.rest_client import VectorRestClient, VectorRestError def _write(record_id, kind, embedding, content_hash): @@ -226,3 +226,27 @@ def test_http_adapter_legacy_fallback_preserves_kind_semantics(monkeypatch): ) assert [hit.id for hit in hits] == ["right"] assert "kinds" in calls[0] and "kinds" not in calls[1] + + +def test_http_delete_generation_uses_exact_allowlisted_rpc_payload(monkeypatch): + calls = [] + monkeypatch.setattr( + "tht.vectorstore.rest_client.requests.post", + lambda url, json, **kwargs: calls.append((url, json)) or Response({"deleted": 2}), + ) + client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) + assert client.delete_generation("evidence", "gen:" + "a" * 32) == 2 + assert calls == [("https://vectors.test/rpc/delete_vector_generation", { + "table_name": "evidence", "kind": "evidence", "generation": "gen:" + "a" * 32, + })] + + +def test_http_delete_generation_legacy_404_fails_closed_without_body_leak(monkeypatch): + monkeypatch.setattr( + "tht.vectorstore.rest_client.requests.post", + lambda *args, **kwargs: Response({"message": "secret legacy endpoint detail"}, status=404), + ) + client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) + with pytest.raises(VectorRestError, match="delete_vector_generation RPC is unavailable") as error: + client.delete_generation("evidence", "gen:" + "a" * 32) + assert "secret" not in str(error.value) diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 58cd25f4..373f1738 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -84,16 +84,64 @@ def item(name, fingerprint): ) -def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", policy=None): +def pipeline(tmp_path, source, *, embedder=None, vectors=None, model="model-a", policy=None, + retain=3): 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", + retain_published_generations=retain, ) +def test_retention_bounds_generations_and_purges_vectors_after_publish(tmp_path): + vectors = Vectors() + generations = [] + for index in range(4): + result = pipeline( + tmp_path, Source([(item("one", str(index)), f"version {index}")]), + vectors=vectors, retain=2, + ).run_as_job( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + str(index) * 64, + ) + generations.append(result.generation) + store = CorpusStore(tmp_path / "corpus") + assert store.list_generations() == generations[-2:] + assert {r.record.metadata["vector_generation"] for r in vectors.records} == set(generations[-2:]) + assert store.active_generation() == generations[-1] + + +def test_retention_keeps_filesystem_when_vector_purge_fails_then_retries(tmp_path): + class FailingDelete(Vectors): + def __init__(self): + super().__init__() + self.fail_delete = True + + def delete_generation(self, collection, generation): + if self.fail_delete: + raise RuntimeError("credential secret") + return super().delete_generation(collection, generation) + + vectors = FailingDelete() + for index in range(2): + pipeline(tmp_path, Source([(item("one", str(index)), str(index))]), vectors=vectors, + retain=1).run_as_job( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + str(index) * 64, + ) + assert len(CorpusStore(tmp_path / "corpus").list_generations()) == 2 + vectors.fail_delete = False + report = pipeline(tmp_path, Source([(item("one", "1"), "1")]), vectors=vectors, + retain=1).gc(workspace_root=tmp_path) + assert report["status"] == "succeeded" + assert len(CorpusStore(tmp_path / "corpus").list_generations()) == 1 + + def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path): one = item("one", "a") first_source = Source([(one, "hello")]) diff --git a/harness/tests/test_corpus_publish.py b/harness/tests/test_corpus_publish.py index 44ab343d..12c4ad61 100644 --- a/harness/tests/test_corpus_publish.py +++ b/harness/tests/test_corpus_publish.py @@ -1,4 +1,5 @@ import pytest +import os from tht.corpus.models import CorpusManifest from tht.corpus.store import CorpusStore, UnsafeCorpusPath @@ -56,3 +57,43 @@ def test_publish_restores_previous_active_when_directory_fsync_fails_after_repla with pytest.raises(OSError, match="post replace"): store.publish(second) assert store.active_generation() == first + + +def test_read_document_rejects_symlink_hardlink_and_hash_mismatch(tmp_path): + from tht.corpus.models import CanonicalDocument + + content = "trusted" + digest = "sha256:" + __import__("hashlib").sha256(content.encode()).hexdigest() + document = CanonicalDocument( + document_id="doc:" + "a" * 64, source_id="fs:one", source_uri="file:///one", + source_fingerprint="sha256:" + "b" * 64, content_hash=digest, content=content, + pipeline_version="evidence-v1", + ) + store = CorpusStore(tmp_path / "corpus") + generation = store.stage(CorpusManifest(documents=(document,)), {document.document_id: content}) + path = store.resolve_document(document.document_id, generation) + assert store.read_document(document.document_id, generation) == content + + path.unlink() + path.symlink_to(tmp_path / "outside") + (tmp_path / "outside").write_text(content) + with pytest.raises(UnsafeCorpusPath): + store.read_document(document.document_id, generation) + + path.unlink() + os.link(tmp_path / "outside", path) + with pytest.raises(UnsafeCorpusPath): + store.read_document(document.document_id, generation) + + path.unlink() + path.write_text("tampered") + with pytest.raises(UnsafeCorpusPath): + store.read_document(document.document_id, generation) + + +def test_generation_inventory_is_validated_and_sorted(tmp_path): + store = CorpusStore(tmp_path / "corpus") + first = store.stage(CorpusManifest(), {}, generation="gen:" + "1" * 32) + second = store.stage(CorpusManifest(), {}, generation="gen:" + "2" * 32) + (store.root / "unrelated").mkdir() + assert store.list_generations() == [first, second] diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index 59929cef..e0a757e5 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -54,3 +54,17 @@ def test_preprocess_resume_rejects_generation_id_before_configuration(monkeypatc "status": "failed", "error": "resume requires a preprocessing run id" } assert called is False + + +def test_preprocess_evidence_gc_json_is_pristine(monkeypatch, tmp_path): + import tht.cli.preprocess_cmd as command + + monkeypatch.setattr(command, "gc_from_config", lambda *a, **k: { + "status": "succeeded", "dry_run": True, "evicted": [], "failures": [], + }) + response = CliRunner().invoke( + app, ["preprocess", "evidence", "gc", "--dry-run", "--json", "-c", + str(tmp_path / "workspace.yaml")] + ) + assert response.exit_code == 0, response.output + assert json.loads(response.output)["dry_run"] is True diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 5465e76c..be070eb7 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -34,6 +34,7 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = 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", + retain_published_generations=cfg.vector.retain_published_generations, ) def fingerprint(value: str) -> str: return "sha256:" + hashlib.sha256(value.encode()).hexdigest() @@ -48,13 +49,55 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = ) +def gc_from_config(config: Path, *, dry_run: bool = False): + 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" + 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", + retain_published_generations=cfg.vector.retain_published_generations, + ) + with pipeline.store.writer_lock(): + return pipeline.gc(workspace_root=corpus_root.parent, dry_run=dry_run) + + @preprocess_app.command("evidence") def evidence_cmd( + action: str | None = typer.Argument(None), 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: + if action is not None and action != "gc": + raise typer.BadParameter("only the optional 'gc' action is supported") + if action == "gc": + try: + payload = gc_from_config(config, dry_run=dry_run) + except Exception: + payload = {"status": "failed", "error": "evidence cleanup failed"} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: evidence cleanup failed", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) from None + if json_output: + typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) + else: + typer.echo(f"OK: evicted={len(payload['evicted'])} failures={len(payload['failures'])}") + return if resume is not None and re.fullmatch(r"[0-9a-f]{32}", resume) is None: payload = {"status": "failed", "error": "resume requires a preprocessing run id"} if json_output: diff --git a/harness/tht/config.py b/harness/tht/config.py index cc2caab9..e98cbf55 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -228,6 +228,8 @@ class EmbeddingsConfig(BaseModel): class VectorConfig(BaseModel): max_chunk_chars: int = 4000 + # ACTIVE plus the two most recent rollback generations by default. + retain_published_generations: int = Field(default=3, ge=1) class SearchConfig(BaseModel): diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py index 3ac39e0a..5afcee21 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -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 diff --git a/harness/tht/corpus/store.py b/harness/tht/corpus/store.py index 44d4d534..4b2fe831 100644 --- a/harness/tht/corpus/store.py +++ b/harness/tht/corpus/store.py @@ -9,6 +9,7 @@ import re import stat import shutil import uuid +import hashlib from pathlib import Path from contextlib import contextmanager @@ -170,6 +171,14 @@ class CorpusStore: generation = self.active_generation() return self.manifest(generation) if generation else None + def list_generations(self) -> list[str]: + values = [] + for entry in self.root.iterdir(): + match = re.fullmatch(r"gen-([0-9a-f]{32})", entry.name) + if match and not entry.is_symlink() and stat.S_ISDIR(entry.lstat().st_mode): + values.append(f"gen:{match.group(1)}") + return sorted(values, key=lambda value: self.generation_path(value).stat().st_mtime_ns) + def resolve_document(self, document_id: str, generation: str | None = None) -> Path | None: generation = generation or self.active_generation() if generation is None: @@ -178,8 +187,37 @@ class CorpusStore: 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 + parts = Path(relative).parts + if Path(relative).is_absolute() or parts[:1] != ("documents",) or len(parts) != 2: + raise UnsafeCorpusPath("materialized document path is unsafe") + return self.generation_path(generation) / relative + + def read_document(self, document_id: str, generation: str | None = None) -> str | None: + generation = generation or self.active_generation() + if generation is None: + return None + manifest = self.manifest(generation) + path = self.resolve_document(document_id, generation) + document = next((item for item in manifest.documents if item.document_id == document_id), None) + if path is None or document is None: + return None + generation_fd = os.open(self.generation_path(generation), os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + documents_fd = fd = None + try: + documents_fd = os.open("documents", os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=generation_fd) + fd = os.open(path.name, os.O_RDONLY | os.O_NOFOLLOW | os.O_CLOEXEC, dir_fd=documents_fd) + 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("materialized document is unsafe") + payload = os.read(fd, info.st_size + 1) + if len(payload) != info.st_size or "sha256:" + hashlib.sha256(payload).hexdigest() != document.content_hash: + raise UnsafeCorpusPath("materialized document hash mismatch") + return payload.decode("utf-8") + except (OSError, UnicodeError) as error: + raise UnsafeCorpusPath("materialized document read failed") from error + finally: + if fd is not None: + os.close(fd) + if documents_fd is not None: + os.close(documents_fd) + os.close(generation_fd) diff --git a/harness/tht/search/evidence.py b/harness/tht/search/evidence.py index a42cb72f..d2529966 100644 --- a/harness/tht/search/evidence.py +++ b/harness/tht/search/evidence.py @@ -64,6 +64,9 @@ def resolve_evidence_file(store: CorpusStore, evidence_id: str) -> str: frontmatter = document.metadata.get("frontmatter", {}) identifiers = {document.document_id, document.source_id, str(frontmatter.get("id", ""))} if evidence_id in identifiers: + # Validate the immutable bytes through the dirfd/O_NOFOLLOW reader before + # handing the path to legacy session-artifact consumers. + store.read_document(document.document_id) path = store.resolve_document(document.document_id) return str(path) if path else "" return "" diff --git a/harness/tht/vectorstore/rest_client.py b/harness/tht/vectorstore/rest_client.py index f78db770..e9c44e77 100644 --- a/harness/tht/vectorstore/rest_client.py +++ b/harness/tht/vectorstore/rest_client.py @@ -128,10 +128,17 @@ class VectorRestClient: return len(rows) def delete_generation(self, table_name: str, generation: str) -> int: - payload = self._call( - "delete_vector_generation", - {"table_name": table_name, "kind": "evidence", "generation": generation}, - ) + try: + payload = self._call( + "delete_vector_generation", + {"table_name": table_name, "kind": "evidence", "generation": generation}, + ) + except VectorRestError as error: + if "HTTP 404" in str(error): + raise VectorRestError( + "delete_vector_generation RPC is unavailable; deploy the cleanup migration" + ) from None + raise if isinstance(payload, dict): return int(payload.get("deleted", 0)) return 0