From 03ee3ffda879c20159e2a0a5b0f2d189cc6dde30 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sun, 12 Jul 2026 06:35:27 +0200 Subject: [PATCH] fix(evidence): validate active corpus artifacts --- .superpowers/sdd/evidence-task-7-report.md | 14 ++++ harness/tests/test_corpus_pipeline.py | 74 ++++++++++++++++++++- harness/tht/corpus/pipeline.py | 76 +++++++++++++++++++--- 3 files changed, 153 insertions(+), 11 deletions(-) diff --git a/.superpowers/sdd/evidence-task-7-report.md b/.superpowers/sdd/evidence-task-7-report.md index 237d4c79..2c217cd2 100644 --- a/.superpowers/sdd/evidence-task-7-report.md +++ b/.superpowers/sdd/evidence-task-7-report.md @@ -77,3 +77,17 @@ Snapshot changes rebuild only the affected sources; job input/config changes pub while retaining valid stable vector-generation dependencies. Missing or corrupt legacy contract metadata, documents, or vectors fails closed and rebuilds. The Compose smoke now explicitly expects the unchanged no-op to report `published=false` while proving generation deltas `+1`, `+0`, `+1`. + +## Corrupt ACTIVE reconstruction correction + +ACTIVE reuse now reconstructs each source contract from the persisted discovery snapshot and checks +the deterministic document identity, canonical URI, source fingerprint, UTC modification time, +source metadata, applicable media type, content hash, and pipeline identity against the owned +materialized document. The persisted document-source map carries the same exact binding. + +Chunks are recomputed under the current chunk policy and must match the manifest exactly in count, +order, IDs, ordinals, content, hashes, linkage, provenance, and policy metadata. Vector health must +report the configured dimension, and every recomputed chunk must have its generation-scoped vector +ID with the exact content hash. Missing, altered, or extra chunks and corrupt document or vector +contracts therefore disable the no-op and rebuild, while a valid unchanged run still performs no +source acquisition. diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index c61dc117..6bf3f812 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -5,9 +5,9 @@ import pytest from tht.corpus.chunk import ChunkPolicy from tht.corpus.pipeline import CorpusPipeline, PipelineError, PipelineResult from tht.corpus.store import CorpusStore -from tht.corpus.models import CorpusManifest +from tht.corpus.models import CanonicalChunk, CorpusManifest from tht.ports.evidence import AcquiredDocument, SourceObject -from tht.ports.vector import VectorCapabilities +from tht.ports.vector import VectorCapabilities, VectorHealth class Source: @@ -23,7 +23,9 @@ class Source: 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()) + return AcquiredDocument( + source=item, content=payload.encode(), media_type=item.metadata.get("media_type") + ) class Embedder: @@ -45,6 +47,7 @@ class Vectors: def __init__(self, fail=False): self.fail = fail self.records = [] + self.dimension = 3 def upsert(self, collection, records): self.records.extend(records[:1] if self.fail else records) @@ -57,6 +60,12 @@ class Vectors: value.record.id: value.content_hash for value in self.records } + def health(self): + return VectorHealth( + ok=True, expected_dimension=3, observed_dimensions=(self.dimension,), + dimension_compatible=self.dimension == 3, + ) + def delete_generation(self, collection, generation, workspace_id): self.records = [ value for value in self.records @@ -538,6 +547,65 @@ def test_job_incomplete_active_contract_never_noops(tmp_path, damage): assert source.acquire_calls == ["fs:one"] +@pytest.mark.parametrize( + "damage", ["modified_at", "source_metadata", "media_type", "missing_chunk", + "altered_chunk", "extra_chunk", "vector_dimension"] +) +def test_job_corrupt_canonical_document_or_chunk_never_noops(tmp_path, damage): + import hashlib + import json + + vectors = Vectors() + source_object = SourceObject( + source_id="fs:one", uri="file:///safe/one.md", fingerprint="sha256:a", + modified_at=datetime(2026, 1, 1, tzinfo=UTC), + metadata={"media_type": "text/markdown", "size": 11, "owner": "docs"}, + ) + args = dict(workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64) + candidate = pipeline( + tmp_path, Source([(source_object, "hello world")]), vectors=vectors, + policy=ChunkPolicy(version="chunk-v1", max_chars=6), + ) + first = candidate.run_as_job(**args) + manifest_path = candidate.store.generation_path(first.generation) / "manifest.json" + payload = json.loads(manifest_path.read_text()) + document = payload["documents"][0] + chunks = payload["chunks"] + if damage == "modified_at": + document["modified_at"] = "2026-01-01T00:00:01Z" + elif damage == "source_metadata": + document["metadata"]["source"]["owner"] = "attacker" + elif damage == "media_type": + document["media_type"] = "text/plain" + elif damage == "missing_chunk": + payload["chunks"] = chunks[:-1] + elif damage == "altered_chunk": + chunks[0]["content"] = "HELLO " + chunks[0]["content_hash"] = "sha256:" + hashlib.sha256(b"HELLO ").hexdigest() + chunks[0]["chunk_id"] = "chunk:" + "a" * 64 + elif damage == "extra_chunk": + extra = CanonicalChunk( + chunk_id="chunk:" + "b" * 64, document_id=document["document_id"], + ordinal=len(chunks), content="", content_hash="sha256:" + hashlib.sha256(b"").hexdigest(), + source_uri=document["source_uri"], pipeline_version=document["pipeline_version"], + ) + chunks.append(extra.model_dump(mode="json")) + else: + vectors.dimension = 4 + manifest_path.write_text(json.dumps(payload)) + + source = Source([(source_object, "hello world")]) + result = pipeline( + tmp_path, source, vectors=vectors, + policy=ChunkPolicy(version="chunk-v1", max_chars=6), + ).run_as_job(**args) + assert result.published is True + assert result.generation != first.generation + assert source.acquire_calls == ["fs:one"] + + 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() diff --git a/harness/tht/corpus/pipeline.py b/harness/tht/corpus/pipeline.py index 0f1c2b15..19945268 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -15,7 +15,7 @@ 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.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 @@ -259,12 +259,21 @@ class CorpusPipeline: } previous = self.store.active_manifest() - def document_sources(manifest: CorpusManifest) -> dict[str, dict[str, str]]: + 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 } @@ -283,19 +292,59 @@ class CorpusPipeline: ): return False for source_id, document in actual_documents.items(): - source = persisted_snapshot[source_id] + 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 ( - document.source_uri != source["uri"] - or document.source_fingerprint != source["fingerprint"] - or self.store.read_document(document.document_id, manifest.manifest_id) - != document.content + 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 manifest.chunks: + for part in expected_chunks: generation = generations.get(part.document_id) if not isinstance(generation, str): return False @@ -408,9 +457,20 @@ class CorpusPipeline: "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 },