diff --git a/.superpowers/sdd/evidence-task-7-report.md b/.superpowers/sdd/evidence-task-7-report.md index d2e2ab89..237d4c79 100644 --- a/.superpowers/sdd/evidence-task-7-report.md +++ b/.superpowers/sdd/evidence-task-7-report.md @@ -62,3 +62,18 @@ Failure injection runs a real exit-97 command after resources exist and reaches Cleanup aggregates Compose-down, residual container/volume/network, and temp-directory failures while preserving the original failure status. S3 prefixes are validated before any client request for leading slash, UTF-8 byte length, controls, and DEL. + +## Canonical unchanged-run correction + +The durable job now persists a deterministic source snapshot keyed by source identity. Each entry +binds canonical URI, exact source fingerprint, UTC modification time, canonical immutable metadata, +and explicit media type and size contract fields. The manifest also binds document-to-source +provenance, supplied config/input fingerprints, compatibility, embedding settings, and pipeline and +chunk-policy versions. + +An unchanged run reuses ACTIVE only when ownership, bindings, the complete snapshot, document +provenance, materialized document hashes, and every required vector ID/content hash match exactly. +Snapshot changes rebuild only the affected sources; job input/config changes publish a new manifest +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`. diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 511c6fc0..c61dc117 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -1,3 +1,5 @@ +from datetime import UTC, datetime, timedelta + import pytest from tht.corpus.chunk import ChunkPolicy @@ -432,12 +434,27 @@ def test_unchanged_documents_skip_acquire_normalize_chunk_and_embed(tmp_path): def test_unchanged_job_reuses_active_generation_without_new_directory(tmp_path): - source = Source([(item("one", "a"), "hello")]) + source = Source([( + 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": 5, "nested": {"b": 2, "a": 1}}, + ), + "hello", + )]) candidate = pipeline(tmp_path, source) args = dict(workspace_id="demo", workspace_root=tmp_path, config_fingerprint="sha256:" + "1" * 64, input_fingerprint="sha256:" + "2" * 64) first = candidate.run_as_job(**args) + snapshot = first.manifest.metadata["source_snapshot"]["fs:one"] + assert snapshot == { + "source_id": "fs:one", "uri": "file:///safe/one.md", "fingerprint": "sha256:a", + "modified_at": "2026-01-01T00:00:00Z", + "metadata": {"media_type": "text/markdown", "size": 5, + "nested": {"a": 1, "b": 2}}, + "media_type": "text/markdown", "size": 5, + } count = len(candidate.store.list_generations()) second = candidate.run_as_job(**args) assert second.generation == first.generation @@ -445,6 +462,82 @@ def test_unchanged_job_reuses_active_generation_without_new_directory(tmp_path): assert len(candidate.store.list_generations()) == count +@pytest.mark.parametrize("field", ["uri", "modified_at", "metadata"]) +def test_job_source_snapshot_change_forces_publish_with_same_fingerprint(tmp_path, field): + original = 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": 5, "label": "original"}, + ) + vectors = Vectors() + args = dict(workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64) + first = pipeline(tmp_path, Source([(original, "hello")]), vectors=vectors).run_as_job(**args) + updates = { + "uri": "file:///safe/renamed.md", + "modified_at": original.modified_at + timedelta(seconds=1), + "metadata": {"media_type": "text/markdown", "size": 5, "label": "changed"}, + } + changed = original.model_copy(update={field: updates[field]}) + source = Source([(changed, "hello")]) + result = pipeline(tmp_path, source, vectors=vectors).run_as_job(**args) + assert result.published is True + assert result.generation != first.generation + assert source.acquire_calls == ["fs:one"] + + +@pytest.mark.parametrize("fingerprint_name", ["config_fingerprint", "input_fingerprint"]) +def test_job_binding_change_forces_publish(tmp_path, fingerprint_name): + vectors = Vectors() + source_object = item("one", "a") + args = dict(workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64) + first = pipeline(tmp_path, Source([(source_object, "hello")]), vectors=vectors).run_as_job(**args) + args[fingerprint_name] = "sha256:" + "3" * 64 + source = Source([(source_object, "hello")]) + result = pipeline(tmp_path, source, vectors=vectors).run_as_job(**args) + assert result.published is True + assert result.generation != first.generation + assert source.acquire_calls == [] + + +@pytest.mark.parametrize( + "damage", ["legacy_metadata", "corrupt_document_sources", "document", "vector"] +) +def test_job_incomplete_active_contract_never_noops(tmp_path, damage): + import json + + vectors = Vectors() + source_object = item("one", "a") + 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")]), vectors=vectors) + first = candidate.run_as_job(**args) + if damage == "legacy_metadata": + manifest_path = candidate.store.generation_path(first.generation) / "manifest.json" + payload = json.loads(manifest_path.read_text()) + payload["metadata"].pop("source_snapshot") + manifest_path.write_text(json.dumps(payload)) + elif damage == "corrupt_document_sources": + manifest_path = candidate.store.generation_path(first.generation) / "manifest.json" + payload = json.loads(manifest_path.read_text()) + payload["metadata"]["document_sources"] = {} + manifest_path.write_text(json.dumps(payload)) + elif damage == "document": + path = candidate.store.resolve_document(first.manifest.documents[0].document_id) + path.unlink() + else: + vectors.records.clear() + source = Source([(source_object, "hello")]) + result = pipeline(tmp_path, source, vectors=vectors).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 2b860ff7..0f1c2b15 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -6,7 +6,9 @@ import hashlib import json import re import uuid +from collections.abc import Mapping, Sequence from dataclasses import asdict, dataclass +from datetime import UTC from pathlib import Path from tht.corpus.chunk import ChunkPolicy, chunk @@ -66,6 +68,31 @@ def _fingerprint(value) -> str: return "sha256:" + hashlib.sha256(payload.encode()).hexdigest() +def _canonical_json(value): + if isinstance(value, Mapping): + return {str(key): _canonical_json(value[key]) for key in sorted(value)} + if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): + return [_canonical_json(child) for child in value] + return value + + +def _source_snapshot(discovered) -> dict[str, dict]: + snapshot = {} + for _, item in discovered: + modified_at = item.modified_at.astimezone(UTC) if item.modified_at else None + metadata = _canonical_json(item.metadata) + snapshot[item.source_id] = { + "source_id": item.source_id, + "uri": item.uri, + "fingerprint": item.fingerprint, + "modified_at": modified_at.isoformat().replace("+00:00", "Z") if modified_at else None, + "metadata": metadata, + "media_type": metadata.get("media_type"), + "size": metadata.get("size"), + } + return snapshot + + class CorpusPipeline: def __init__( self, *, store: CorpusStore, sources: list[EvidenceSource], embedder, @@ -213,6 +240,7 @@ class CorpusPipeline: discovered_fingerprint = _fingerprint( {item.source_id: item.fingerprint for _, item in discovered} ) + source_snapshot = _source_snapshot(discovered) source_by_id = {item.source_id: (source, item) for source, item in discovered} compatibility = _fingerprint({ "pipeline": self.pipeline_version, @@ -220,12 +248,72 @@ class CorpusPipeline: "dimensions": self.embedding_dimensions, "chunk_policy": asdict(self.chunk_policy), }) + job_binding = { + "config_fingerprint": config_fingerprint, + "input_fingerprint": input_fingerprint, + "compatibility_fingerprint": compatibility, + "pipeline_version": self.pipeline_version, + "chunk_policy_version": self.chunk_policy.version, + "embedding_model": self.embedding_model, + "embedding_dimensions": self.embedding_dimensions, + } previous = self.store.active_manifest() - if (not dry_run and resume_run_id is None and previous is not None - and previous.metadata.get("compatibility_fingerprint") == compatibility - and previous.metadata.get("fingerprints") == { - item.source_id: item.fingerprint for _, item in discovered - }): + + def document_sources(manifest: CorpusManifest) -> dict[str, dict[str, str]]: + return { + document.document_id: { + "source_id": document.source_id, + "source_uri": document.source_uri, + "source_fingerprint": document.source_fingerprint, + } + for document in manifest.documents + } + + def active_assets_are_valid(manifest: CorpusManifest | None) -> bool: + if manifest is None or manifest.metadata.get("workspace_id") != workspace_id: + return False + actual_documents = {document.source_id: document for document in manifest.documents} + persisted_snapshot = _canonical_json(manifest.metadata.get("source_snapshot")) + if ( + not isinstance(persisted_snapshot, dict) + or set(persisted_snapshot) != set(actual_documents) + or manifest.metadata.get("compatibility_fingerprint") != compatibility + or _canonical_json(manifest.metadata.get("document_sources")) + != document_sources(manifest) + ): + return False + for source_id, document in actual_documents.items(): + source = persisted_snapshot[source_id] + 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 + ): + return False + generations = manifest.metadata.get("document_generations") + if not isinstance(generations, Mapping): + return False + existing = self.vector_store.existing_hashes("evidence", ["evidence"]) + for part in manifest.chunks: + generation = generations.get(part.document_id) + if not isinstance(generation, str): + return False + record_id = f"{workspace_id}:{generation}:{part.chunk_id}" + if existing.get(record_id) != part.content_hash: + return False + return True + + try: + active_assets_valid = active_assets_are_valid(previous) + except Exception: + active_assets_valid = False + reusable = ( + active_assets_valid + and _canonical_json(previous.metadata.get("source_snapshot")) == source_snapshot + and _canonical_json(previous.metadata.get("job_binding")) == job_binding + ) + if not dry_run and resume_run_id is None and reusable: return PipelineResult( "succeeded", previous.manifest_id, False, (), tuple(sorted(item.source_id for _, item in discovered)), (), previous, @@ -263,17 +351,22 @@ class CorpusPipeline: previous = self.store.active_manifest() prior = {doc.source_id: doc for doc in previous.documents} if previous else {} fingerprints = {item.source_id: item.fingerprint for _, item in discovered} - rebuild = bool(previous and previous.metadata.get("compatibility_fingerprint") != compatibility) + previous_snapshot = ( + _canonical_json(previous.metadata.get("source_snapshot")) if previous else {} + ) + rebuild = bool(previous and not active_assets_valid) changed = sorted( item.source_id for _, item in discovered if rebuild or item.source_id not in prior - or prior[item.source_id].source_fingerprint != item.fingerprint + or previous_snapshot.get(item.source_id) != source_snapshot[item.source_id] ) unchanged = sorted(set(fingerprints) - set(changed)) removed = sorted(set(prior) - set(fingerprints)) write(context, "plan.json", { "generation": f"gen:{context.run_id}", "compatibility": compatibility, + "job_binding": job_binding, + "source_snapshot": source_snapshot, "fingerprints": fingerprints, "changed": changed, "unchanged": unchanged, @@ -311,6 +404,16 @@ class CorpusPipeline: metadata={ "workspace_id": self.workspace_id, "compatibility_fingerprint": compatibility, + "job_binding": plan["job_binding"], + "source_snapshot": plan["source_snapshot"], + "document_sources": { + document.document_id: { + "source_id": document.source_id, + "source_uri": document.source_uri, + "source_fingerprint": document.source_fingerprint, + } + for document in documents + }, "fingerprints": plan["fingerprints"], "removed": plan["removed"], "document_generations": generations, diff --git a/scripts/preprocess-smoke.sh b/scripts/preprocess-smoke.sh index ca58403a..77d4dc40 100755 --- a/scripts/preprocess-smoke.sh +++ b/scripts/preprocess-smoke.sh @@ -114,7 +114,7 @@ before, first_count, second_count, third_count = map(int, sys.argv[4:]) assert len(a["changed"]) == 1 and not a["unchanged"] assert len(b["unchanged"]) == 1 and not b["changed"] assert len(c["changed"]) == 1 and c["generation"] != a["generation"] -assert all(row["published"] for row in (a, b, c)) +assert a["published"] and not b["published"] and c["published"] assert first_count == before + 1 assert second_count == first_count assert third_count == second_count + 1