From 459ffa0bcda3b7245a815dc3e1e2a6711974d381 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sun, 12 Jul 2026 05:10:19 +0200 Subject: [PATCH] fix(evidence): reconcile generations safely --- .superpowers/sdd/evidence-task-5c-report.md | 24 ++++++++++++ harness/scripts/create_vector_writer_rpc.sql | 23 ++++++++++++ harness/tests/l0/test_pgvector_store.py | 22 +++++++++++ .../tests/l0/test_vector_adapter_parity.py | 14 +++++++ harness/tests/test_corpus_pipeline.py | 22 +++++++++++ harness/tests/test_corpus_publish.py | 11 ++++++ harness/tht/adapters/vector/pgvector.py | 22 +++++++++++ harness/tht/adapters/vector/thoth_http.py | 9 +++++ harness/tht/corpus/pipeline.py | 20 ++++++++-- harness/tht/corpus/store.py | 37 +++++++++++++++++++ harness/tht/ports/vector.py | 3 ++ harness/tht/search/evidence.py | 11 +++--- harness/tht/session/artifacts.py | 5 ++- harness/tht/vectorstore/rest_client.py | 14 +++++++ 14 files changed, 227 insertions(+), 10 deletions(-) diff --git a/.superpowers/sdd/evidence-task-5c-report.md b/.superpowers/sdd/evidence-task-5c-report.md index 62a6a7b8..e77c1cc8 100644 --- a/.superpowers/sdd/evidence-task-5c-report.md +++ b/.superpowers/sdd/evidence-task-5c-report.md @@ -29,3 +29,27 @@ 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. + +## Review fix wave + +- Publication is now explicit and durable (`PUBLISHED` marker). Retention candidates require a + valid generation manifest and publication marker (ACTIVE remains backward-compatible), so + staged and malformed directories neither consume retention slots nor become deletion targets. +- The policy retains ACTIVE plus exactly `N-1` newest rollback publications, ordered by durable + publication time and generation id. Running and failed-resumable JobRunner checkpoints protect + every referenced plan generation. +- `VectorStore` now exposes exact Evidence generation inventory. Direct pgvector uses a constrained + `SELECT DISTINCT` over `kind='evidence'` and `metadata.vector_generation`; HTTP uses the + allowlisted `list_evidence_generations` RPC and fails closed on legacy 404. The writer RPC SQL, + revokes, and grants are packaged in `create_vector_writer_rpc.sql`. +- Explicit GC reconciles the union of published filesystem generations and vector-only orphans, + preserving vector-before-filesystem deletion and retry semantics. +- `run_as_job` holds the same corpus writer lock across checkpoint recovery, staging, publish, and + retention. Explicit GC already uses this lock, serializing candidate snapshots with publishers. +- Session artifact consumers no longer receive the corpus source path after validation. They get + an owned, read-only copy atomically written from the bytes read and hash-validated on the same + descriptor. + +Fresh verification after the fix wave: full harness `672 passed, 5 deselected`; Docker pgvector, +HTTP parity, and migration suites `43 passed`; exact direct inventory/delete integration `1 passed`; +changed-file Ruff and `git diff --check` clean. diff --git a/harness/scripts/create_vector_writer_rpc.sql b/harness/scripts/create_vector_writer_rpc.sql index 1fe7d11a..d006ac5e 100644 --- a/harness/scripts/create_vector_writer_rpc.sql +++ b/harness/scripts/create_vector_writer_rpc.sql @@ -40,6 +40,25 @@ begin end; $$; +create or replace function public.list_evidence_generations(table_name text, kind text) +returns table(generation text) +language plpgsql +security definer +set search_path = public, vectors, extensions +as $$ +begin + if table_name <> 'evidence' or kind <> 'evidence' then + raise exception 'only exact Evidence generations may be listed'; + end if; + return query + select distinct e.metadata->>'vector_generation' + from vectors.evidence e + where e.kind = 'evidence' + and e.metadata->>'vector_generation' ~ '^gen:[0-9a-f]{32}$' + order by 1; +end; +$$; + create or replace function public.existing_vector_hashes(table_name text, kinds text[]) returns table(record_key text, content_hash text) language plpgsql @@ -125,6 +144,7 @@ revoke all on function public._assert_vector_write_table(text, text[]) from publ revoke all on function public.existing_vector_hashes(text, text[]) from public; revoke all on function public.upsert_vector_records(text, jsonb) from public; revoke all on function public.delete_vector_generation(text, text, text) from public; +revoke all on function public.list_evidence_generations(text, text) from public; -- Su alcuni progetti Supabase le funzioni in `public` ricevono grant automatici: revoca -- esplicitamente dai ruoli client generici, poi abilita solo il writer dedicato. @@ -134,16 +154,19 @@ begin revoke all on function public.existing_vector_hashes(text, text[]) from anon; revoke all on function public.upsert_vector_records(text, jsonb) from anon; revoke all on function public.delete_vector_generation(text, text, text) from anon; + revoke all on function public.list_evidence_generations(text, text) from anon; end if; if exists (select 1 from pg_roles where rolname = 'authenticated') then revoke all on function public.existing_vector_hashes(text, text[]) from authenticated; revoke all on function public.upsert_vector_records(text, jsonb) from authenticated; revoke all on function public.delete_vector_generation(text, text, text) from authenticated; + revoke all on function public.list_evidence_generations(text, text) from authenticated; end if; if exists (select 1 from pg_roles where rolname = 'vector_writer') then grant execute on function public.existing_vector_hashes(text, text[]) to vector_writer; grant execute on function public.upsert_vector_records(text, jsonb) to vector_writer; grant execute on function public.delete_vector_generation(text, text, text) to vector_writer; + grant execute on function public.list_evidence_generations(text, text) to vector_writer; end if; end $$; diff --git a/harness/tests/l0/test_pgvector_store.py b/harness/tests/l0/test_pgvector_store.py index 440d44d2..d01b3a3a 100644 --- a/harness/tests/l0/test_pgvector_store.py +++ b/harness/tests/l0/test_pgvector_store.py @@ -79,6 +79,13 @@ def vector_configs(): f"GRANT INSERT, UPDATE ON vectors.{table} " "TO vector_l0_writer, vector_l0_no_sequence" ) + if table == "evidence": + connection.exec_driver_sql( + "GRANT DELETE ON vectors.evidence TO vector_l0_writer" + ) + connection.exec_driver_sql( + "GRANT SELECT (kind, metadata) ON vectors.evidence TO vector_l0_writer" + ) connection.exec_driver_sql( f"GRANT SELECT (record_key, kind, content_hash) " f"ON vectors.{table} TO vector_l0_writer, vector_l0_no_sequence" @@ -118,6 +125,21 @@ def test_pgvector_round_trip_hash_and_upsert(store): assert store.search(["memory"], [0.8, 0.2], limit=1)[0].id == "record:a" +def test_pgvector_lists_and_deletes_exact_evidence_generation(store): + generation = "gen:" + "a" * 32 + value = VectorWriteRecord( + record=VectorRecord( + id="evidence-generation-a", kind="evidence", ref="doc:a", title="a", + content="content", metadata={"vector_generation": generation}, + ), + embedding=[1.0, 0.0], content_hash="sha256:" + "a" * 64, + ) + store.upsert("evidence", [value]) + assert generation in store.list_evidence_generations("evidence") + assert store.delete_generation("evidence", generation) == 1 + assert generation not in store.list_evidence_generations("evidence") + + def test_pgvector_search_filters_kinds_before_limit(store): store.upsert("memory", [_record("solved", [1.0, 0.0], kind="solved_question")]) hits = store.search("memory".split(), [1.0, 0.0], limit=1, kinds=["memory"]) diff --git a/harness/tests/l0/test_vector_adapter_parity.py b/harness/tests/l0/test_vector_adapter_parity.py index b1655d0f..d5b3435d 100644 --- a/harness/tests/l0/test_vector_adapter_parity.py +++ b/harness/tests/l0/test_vector_adapter_parity.py @@ -250,3 +250,17 @@ def test_http_delete_generation_legacy_404_fails_closed_without_body_leak(monkey 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) + + +def test_http_list_evidence_generations_exact_rpc_and_legacy_fail_closed(monkeypatch): + calls = [] + monkeypatch.setattr( + "tht.vectorstore.rest_client.requests.post", + lambda url, json, **kwargs: calls.append((url, json)) or Response([ + {"generation": "gen:" + "a" * 32} + ]), + ) + client = VectorRestClient(RestConfig(base_url="https://vectors.test", api_key="writer")) + assert client.list_evidence_generations("evidence") == ["gen:" + "a" * 32] + assert calls[0][0].endswith("/rpc/list_evidence_generations") + assert calls[0][1] == {"table_name": "evidence", "kind": "evidence"} diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 373f1738..42d26c31 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -61,6 +61,12 @@ class Vectors: ] return 0 + def list_evidence_generations(self, collection): + return sorted({ + value.record.metadata["vector_generation"] for value in self.records + if value.record.kind == "evidence" + }) + class InterruptingVectors(Vectors): def __init__(self): @@ -142,6 +148,22 @@ def test_retention_keeps_filesystem_when_vector_purge_fails_then_retries(tmp_pat assert len(CorpusStore(tmp_path / "corpus").list_generations()) == 1 +def test_gc_reconciles_vector_only_generation(tmp_path): + vectors = Vectors() + orphan = "gen:" + "f" * 32 + from tht.ports.vector import VectorWriteRecord + from tht.vectorstore.records import VectorRecord + vectors.records.append(VectorWriteRecord( + record=VectorRecord(id="orphan", kind="evidence", ref="doc:x", title="", content="x", + metadata={"vector_generation": orphan}), + embedding=[0.0, 0.0, 0.0], content_hash="sha256:" + "0" * 64, + )) + candidate = pipeline(tmp_path, Source([]), vectors=vectors, retain=1) + report = candidate.gc(workspace_root=tmp_path) + assert report["evicted"] == [orphan] + assert vectors.list_evidence_generations("evidence") == [] + + 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 12c4ad61..97ab41aa 100644 --- a/harness/tests/test_corpus_publish.py +++ b/harness/tests/test_corpus_publish.py @@ -97,3 +97,14 @@ def test_generation_inventory_is_validated_and_sorted(tmp_path): second = store.stage(CorpusManifest(), {}, generation="gen:" + "2" * 32) (store.root / "unrelated").mkdir() assert store.list_generations() == [first, second] + + +def test_published_inventory_excludes_staged_and_invalid_newer_directories(tmp_path): + store = CorpusStore(tmp_path / "corpus") + first = store.stage(CorpusManifest(), {}, generation="gen:" + "1" * 32) + store.publish(first) + store.stage(CorpusManifest(), {}, generation="gen:" + "2" * 32) + invalid = store.generation_path("gen:" + "3" * 32) + invalid.mkdir() + (invalid / "PUBLISHED").write_text("2026-01-01T00:00:00Z\n") + assert store.published_generations() == [first] diff --git a/harness/tht/adapters/vector/pgvector.py b/harness/tht/adapters/vector/pgvector.py index 1b1f9970..6e6021cf 100644 --- a/harness/tht/adapters/vector/pgvector.py +++ b/harness/tht/adapters/vector/pgvector.py @@ -87,6 +87,7 @@ class PgVectorStore: upsert=writable, metadata_filter=self._reader is not None, delete_generation=writable, + list_evidence_generations=writable, ) def _probe( @@ -428,5 +429,26 @@ class PgVectorStore: if raw is not None: raw.close() + def list_evidence_generations(self, collection: str) -> list[str]: + if collection != "evidence": + raise VectorStoreError("Only exact Evidence generations may be listed") + raw = None + try: + raw = self._require_writer().raw_connection() + with raw.cursor() as cursor: + cursor.execute( + sql.SQL( + "SELECT DISTINCT metadata->>'vector_generation' FROM {} " + "WHERE kind = 'evidence' AND metadata->>'vector_generation' " + "~ '^gen:[0-9a-f]{{32}}$' ORDER BY 1" + ).format(_collection(self._schema, collection)) + ) + return [row[0] for row in cursor.fetchall()] + except Exception as exc: + raise VectorWriteUnavailable("Vector generation inventory unavailable") from exc + finally: + if raw is not None: + raw.close() + __all__ = ["ALLOWED_COLLECTIONS", "PgVectorStore"] diff --git a/harness/tht/adapters/vector/thoth_http.py b/harness/tht/adapters/vector/thoth_http.py index 67379600..c8c2a77b 100644 --- a/harness/tht/adapters/vector/thoth_http.py +++ b/harness/tht/adapters/vector/thoth_http.py @@ -42,6 +42,7 @@ class ThothHttpVectorStore: return VectorCapabilities( search=self._reader is not None, existing_hashes=writable, upsert=writable, metadata_filter=self._reader is not None, delete_generation=writable, + list_evidence_generations=writable, ) def health(self) -> VectorHealth: @@ -160,6 +161,14 @@ class ThothHttpVectorStore: except VectorRestError as exc: raise VectorStoreError(str(exc)) from exc + def list_evidence_generations(self, collection: str) -> list[str]: + if collection != "evidence": + raise VectorStoreError("Only exact Evidence generations may be listed") + try: + return self._require_writer().list_evidence_generations(collection) + except VectorRestError as exc: + raise VectorWriteUnavailable("Vector generation inventory unavailable") from exc + @staticmethod def _row(write_record: VectorWriteRecord) -> dict: record = write_record.record diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py index 5afcee21..e0f06beb 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -92,9 +92,16 @@ class CorpusPipeline: return protected def gc(self, *, workspace_root: Path, dry_run: bool = False) -> dict: - generations = self.store.list_generations() + published = self.store.published_generations() + list_vectors = getattr(self.vector_store, "list_evidence_generations", None) + vector_generations = set(list_vectors("evidence")) if list_vectors else set() + generations = sorted(set(published) | vector_generations) protected = self._protected_generations(workspace_root) - keep = set(generations[-self.retain_published_generations:]) | protected + 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 ()) + keep |= protected evicted, failures = [], [] for generation in generations: if generation in keep: @@ -108,7 +115,8 @@ class CorpusPipeline: failures.append({"generation": generation, "error": "vector cleanup failed"}) continue try: - self.store.discard(generation) + if generation in self.store.list_generations(): + self.store.discard(generation) evicted.append(generation) except Exception: failures.append({"generation": generation, "error": "filesystem cleanup failed"}) @@ -131,7 +139,11 @@ class CorpusPipeline: with self.store.writer_lock(): return self._run(dry_run=dry_run, resume=resume) - def run_as_job( + def run_as_job(self, **kwargs) -> PipelineResult: + with self.store.writer_lock(): + return self._run_as_job(**kwargs) + + def _run_as_job( self, *, workspace_id: str, diff --git a/harness/tht/corpus/store.py b/harness/tht/corpus/store.py index 4b2fe831..dfcd65c3 100644 --- a/harness/tht/corpus/store.py +++ b/harness/tht/corpus/store.py @@ -10,6 +10,7 @@ import stat import shutil import uuid import hashlib +from datetime import UTC, datetime from pathlib import Path from contextlib import contextmanager @@ -115,6 +116,7 @@ class CorpusStore: if self.active_generation() == generation: return generation previous = self.active_generation() + published_marker = self.generation_path(generation) / "PUBLISHED" temporary = self.active_path.with_name(f".ACTIVE.{uuid.uuid4().hex}.tmp") replaced = False try: @@ -122,6 +124,10 @@ class CorpusStore: self._replace(temporary, self.active_path) replaced = True self._fsync_directory() + _atomic_write( + published_marker, + (datetime.now(UTC).isoformat().replace("+00:00", "Z") + "\n").encode("ascii"), + ) except BaseException: temporary.unlink(missing_ok=True) if replaced: @@ -179,6 +185,25 @@ class CorpusStore: values.append(f"gen:{match.group(1)}") return sorted(values, key=lambda value: self.generation_path(value).stat().st_mtime_ns) + def published_generations(self) -> list[str]: + active = self.active_generation() + published = [] + for generation in self.list_generations(): + path = self.generation_path(generation) + marker = path / "PUBLISHED" + if generation != active and not marker.is_file(): + continue + try: + manifest = self.manifest(generation) + if manifest.manifest_id != generation: + continue + timestamp = marker.read_text(encoding="ascii").strip() if marker.is_file() else "" + key = (timestamp or manifest.created_at.isoformat(), generation) + published.append((key, generation)) + except (OSError, ValueError): + continue + return [generation for _, generation in sorted(published)] + def resolve_document(self, document_id: str, generation: str | None = None) -> Path | None: generation = generation or self.active_generation() if generation is None: @@ -221,3 +246,15 @@ class CorpusStore: if documents_fd is not None: os.close(documents_fd) os.close(generation_fd) + + def materialize_document( + self, document_id: str, destination: Path, generation: str | None = None, + ) -> Path | None: + content = self.read_document(document_id, generation) + if content is None: + return None + destination = Path(destination) + destination.parent.mkdir(parents=True, exist_ok=True, mode=0o700) + _atomic_write(destination, content.encode("utf-8")) + destination.chmod(0o400) + return destination diff --git a/harness/tht/ports/vector.py b/harness/tht/ports/vector.py index 84f5b6fe..57da32a8 100644 --- a/harness/tht/ports/vector.py +++ b/harness/tht/ports/vector.py @@ -14,6 +14,7 @@ class VectorCapabilities: upsert: bool = False metadata_filter: bool = False delete_generation: bool = False + list_evidence_generations: bool = False @dataclass(frozen=True) @@ -81,6 +82,8 @@ class VectorStore(Protocol): def delete_generation(self, collection: str, generation: str) -> int: ... + def list_evidence_generations(self, collection: str) -> list[str]: ... + __all__ = [ "VectorCapabilities", diff --git a/harness/tht/search/evidence.py b/harness/tht/search/evidence.py index d2529966..ed2a1473 100644 --- a/harness/tht/search/evidence.py +++ b/harness/tht/search/evidence.py @@ -56,7 +56,9 @@ def active_evidence_hits(store: CorpusStore, vector_store, embedding, *, limit: ] -def resolve_evidence_file(store: CorpusStore, evidence_id: str) -> str: +def resolve_evidence_file( + store: CorpusStore, evidence_id: str, *, materialized_root=None, +) -> str: manifest = store.active_manifest() if manifest is None: return "" @@ -64,9 +66,8 @@ 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) + root = materialized_root or (store.root / "runtime") + filename = document.document_id.removeprefix("doc:") + ".md" + path = store.materialize_document(document.document_id, root / filename) return str(path) if path else "" return "" diff --git a/harness/tht/session/artifacts.py b/harness/tht/session/artifacts.py index e0e2c44b..6f4a65cd 100644 --- a/harness/tht/session/artifacts.py +++ b/harness/tht/session/artifacts.py @@ -11,7 +11,10 @@ def _find_evidence_file(evidence_root: Path, evidence_id: str) -> str: if corpus_root.exists(): from tht.corpus.store import CorpusStore from tht.search.evidence import resolve_evidence_file - return resolve_evidence_file(CorpusStore(corpus_root), evidence_id) + return resolve_evidence_file( + CorpusStore(corpus_root), evidence_id, + materialized_root=evidence_root.parent / ".materialized-evidence", + ) for match in evidence_root.rglob(f"{evidence_id}.md"): return str(match) return "" diff --git a/harness/tht/vectorstore/rest_client.py b/harness/tht/vectorstore/rest_client.py index e9c44e77..2fdc2c40 100644 --- a/harness/tht/vectorstore/rest_client.py +++ b/harness/tht/vectorstore/rest_client.py @@ -142,3 +142,17 @@ class VectorRestClient: if isinstance(payload, dict): return int(payload.get("deleted", 0)) return 0 + + def list_evidence_generations(self, table_name: str) -> list[str]: + try: + rows = self._call( + "list_evidence_generations", + {"table_name": table_name, "kind": "evidence"}, + ) or [] + except VectorRestError as error: + if "HTTP 404" in str(error): + raise VectorRestError( + "list_evidence_generations RPC is unavailable; deploy the cleanup migration" + ) from None + raise + return sorted({row["generation"] for row in rows if isinstance(row, dict)})