From 1d5f8c76a799bf8d8d9aecd572dbe71584826573 Mon Sep 17 00:00:00 2001 From: mptyl Date: Sun, 12 Jul 2026 04:39:32 +0200 Subject: [PATCH] feat(preprocess): resume evidence jobs by run id --- .superpowers/sdd/evidence-task-5b-report.md | 43 ++++ harness/tests/test_corpus_pipeline.py | 43 ++++ harness/tests/test_job_runner.py | 21 ++ harness/tests/test_preprocess_cli.py | 24 ++ harness/tht/cli/preprocess_cmd.py | 26 +- harness/tht/corpus/pipeline.py | 257 ++++++++++++++++++++ harness/tht/jobs/runner.py | 4 + 7 files changed, 413 insertions(+), 5 deletions(-) create mode 100644 .superpowers/sdd/evidence-task-5b-report.md diff --git a/.superpowers/sdd/evidence-task-5b-report.md b/.superpowers/sdd/evidence-task-5b-report.md new file mode 100644 index 00000000..0b52316d --- /dev/null +++ b/.superpowers/sdd/evidence-task-5b-report.md @@ -0,0 +1,43 @@ +# Evidence Task 5B implementation report + +## Status + +Integrated Evidence preprocessing with the Task 4 `JobRunner`. The CLI now accepts only a +32-character JobRunner run ID for `--resume`; generation IDs remain outputs. Runs persist the +exact ordered stages `discover`, `acquire_normalize_chunk`, `embed`, `vector_upsert`, +`stage_validate`, `publish`, and `retention_cleanup`. + +Successful-stage artifacts are copied into the new resume run before execution, allowing later +stages to continue without rediscovery, acquisition, normalization, chunking, or embedding. +Job compatibility includes workspace, configuration, discovered-input, pipeline, embedding, and +chunk-policy fingerprints. Generation-specific filesystem/vector compensation is retained, and a +compensated generation is rotated before retry. `ACTIVE` is mutated only by `publish`. + +Dry-run executes discovery/planning and makes every side-effecting stage a no-op. JSON output is +pristine and includes the JobRunner `run_id`, `resumed_from`, generation, plan, and publish status. + +## TDD evidence + +- RED: run-ID rejection and resume-artifact tests failed because generation IDs reached + configuration and resume runs had empty artifact directories. +- GREEN: the two regression tests passed after strict CLI validation and durable artifact carryover. +- Added pipeline job-plan and dry-run counting-fake coverage; both passed. + +## Fresh verification + +- Focused integration/search suite: `62 passed, 4 warnings`. +- Available harness suite excluding sandbox-blocked Docker, loopback HTTP-server, and networked + wheel-build tests: `559 passed, 5 deselected, 18 warnings`. +- Scoped Ruff: `All checks passed!`. +- `git diff --check`: clean. + +## Environment limitations and concerns + +The literal full harness invocation cannot complete in the managed sandbox: Docker socket access, +loopback HTTP test servers, and the `uv build` dependency resolution path are denied. It reached +`575 passed, 5 deselected` before those environment errors. The available-suite rerun above is +green. + +One pre-existing Pydantic serialization warning is exposed by the new end-to-end job test when +canonical metadata contains frozen tuple values; it does not contaminate CLI stdout. Retention is +an explicit stable no-op until a retention policy is configured. diff --git a/harness/tests/test_corpus_pipeline.py b/harness/tests/test_corpus_pipeline.py index 3c41c3ef..0f8fcc71 100644 --- a/harness/tests/test_corpus_pipeline.py +++ b/harness/tests/test_corpus_pipeline.py @@ -49,6 +49,13 @@ class Vectors: raise RuntimeError("partial write") return len(records) + def delete_generation(self, collection, generation): + self.records = [ + value for value in self.records + if value.record.metadata["vector_generation"] != generation + ] + return 0 + def item(name, fingerprint): return SourceObject( @@ -130,3 +137,39 @@ def test_dry_run_and_failed_acquire_never_change_active(tmp_path): with pytest.raises(PipelineError): pipeline(tmp_path, Source([(changed, RuntimeError("boom"))])).run() assert CorpusStore(tmp_path / "corpus").active_generation() == active + + +def test_job_pipeline_uses_ordered_plan_and_returns_run_id(tmp_path): + one = item("one", "a") + candidate = pipeline(tmp_path, Source([(one, "hello")])) + result = candidate.run_as_job( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64, + ) + assert result.status == "succeeded" + assert result.run_id and len(result.run_id) == 32 + checkpoint = tmp_path / ".tht-jobs" / "evidence" / "runs" / result.run_id / "checkpoint.json" + payload = __import__("json").loads(checkpoint.read_text()) + assert [stage["name"] for stage in payload["stages"]] == [ + "discover", "acquire_normalize_chunk", "embed", "vector_upsert", + "stage_validate", "publish", "retention_cleanup", + ] + + +def test_job_pipeline_dry_run_only_discovers_and_reports_changes(tmp_path): + one = item("one", "a") + source = Source([(one, "hello")]) + embedder = Embedder() + vectors = Vectors() + result = pipeline(tmp_path, source, embedder=embedder, vectors=vectors).run_as_job( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint="sha256:" + "1" * 64, + input_fingerprint="sha256:" + "2" * 64, + dry_run=True, + ) + assert result.changed == ("fs:one",) + assert source.acquire_calls == [] + assert embedder.calls == [] + assert vectors.records == [] + assert result.generation is None and result.published is False diff --git a/harness/tests/test_job_runner.py b/harness/tests/test_job_runner.py index 76d6ff69..5553b449 100644 --- a/harness/tests/test_job_runner.py +++ b/harness/tests/test_job_runner.py @@ -64,6 +64,27 @@ def test_failed_stage_is_resumable_and_skips_completed_stage(tmp_path): assert calls == [("discover", False), ("acquire", False), ("recovered", False)] +def test_resume_carries_successful_stage_artifacts_into_new_run(tmp_path): + def discover(context): + artifacts = context.run_dir / "artifacts" + artifacts.mkdir() + (artifacts / "discovery.json").write_text('{"source":"one"}') + + first = run_job( + _spec(tmp_path, stage_ids=("discover", "acquire")), + [discover, lambda _context: (_ for _ in ()).throw(RuntimeError("crash"))], + ) + + def acquire(context): + assert (context.run_dir / "artifacts" / "discovery.json").read_text() == '{"source":"one"}' + + resumed = run_job( + _spec(tmp_path, resume_run_id=first.run_id, stage_ids=("discover", "acquire")), + [discover, acquire], + ) + assert resumed.status == "succeeded" + + def test_successful_job_is_idempotently_resumable(tmp_path): calls = [] diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index bc48c08e..59929cef 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -30,3 +30,27 @@ def test_preprocess_failure_is_structured_and_nonzero(monkeypatch, tmp_path): assert response.exit_code != 0 assert json.loads(response.output) == {"status": "failed", "error": "preprocessing failed"} assert "secret detail" not in response.output + + +def test_preprocess_resume_rejects_generation_id_before_configuration(monkeypatch, tmp_path): + import tht.cli.preprocess_cmd as command + + called = False + + def forbidden(*args, **kwargs): + nonlocal called + called = True + + monkeypatch.setattr(command, "run_from_config", forbidden) + response = CliRunner().invoke( + app, + [ + "preprocess", "evidence", "--resume", "gen:" + "a" * 32, + "--json", "-c", str(tmp_path / "workspace.yaml"), + ], + ) + assert response.exit_code != 0 + assert json.loads(response.output) == { + "status": "failed", "error": "resume requires a preprocessing run id" + } + assert called is False diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index fbbb330e..5465e76c 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -3,6 +3,8 @@ from __future__ import annotations import json +import re +import hashlib from pathlib import Path import typer @@ -25,9 +27,6 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = 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), @@ -36,7 +35,17 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = 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) + def fingerprint(value: str) -> str: + return "sha256:" + hashlib.sha256(value.encode()).hexdigest() + + return pipeline.run_as_job( + workspace_id=config.stem.lower().replace(".", "-").replace("_", "-"), + workspace_root=corpus_root.parent, + config_fingerprint=fingerprint(cfg.model_dump_json()), + input_fingerprint=fingerprint(config.resolve().as_posix()), + dry_run=dry_run, + resume_run_id=resume, + ) @preprocess_app.command("evidence") @@ -46,6 +55,13 @@ def evidence_cmd( resume: str | None = typer.Option(None, "--resume"), json_output: bool = typer.Option(False, "--json"), ) -> None: + 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: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: resume requires a preprocessing run id", fg=typer.colors.RED, err=True) + raise typer.Exit(code=2) try: result = run_from_config(config, dry_run=dry_run, resume=resume) except Exception: @@ -60,6 +76,6 @@ def evidence_cmd( typer.echo(json.dumps(payload, ensure_ascii=False, sort_keys=True)) else: typer.echo( - f"OK: generation={payload['generation']} changed={len(payload['changed'])} " + f"OK: run={payload['run_id']} 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 index 6e0a8ba2..a9ccf7d0 100644 --- a/harness/tht/corpus/pipeline.py +++ b/harness/tht/corpus/pipeline.py @@ -6,6 +6,7 @@ import hashlib import json import uuid from dataclasses import asdict, dataclass +from pathlib import Path from tht.corpus.chunk import ChunkPolicy, chunk from tht.corpus.models import CanonicalChunk, CanonicalDocument, CorpusManifest @@ -14,6 +15,19 @@ 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 +from tht.jobs.models import JobSpec +from tht.jobs.runner import JobContext, run_job + + +EVIDENCE_STAGE_IDS = ( + "discover", + "acquire_normalize_chunk", + "embed", + "vector_upsert", + "stage_validate", + "publish", + "retention_cleanup", +) class PipelineError(RuntimeError): @@ -29,6 +43,8 @@ class PipelineResult: unchanged: tuple[str, ...] removed: tuple[str, ...] manifest: CorpusManifest + run_id: str | None = None + resumed_from: str | None = None def model_dump(self, mode=None): value = asdict(self) @@ -71,6 +87,247 @@ class CorpusPipeline: with self.store.writer_lock(): return self._run(dry_run=dry_run, resume=resume) + def run_as_job( + self, + *, + workspace_id: str, + workspace_root: Path, + config_fingerprint: str, + input_fingerprint: str, + dry_run: bool = False, + resume_run_id: str | None = None, + ) -> PipelineResult: + """Execute preprocessing through the durable shared job envelope.""" + discovered = self._discover() + discovered_fingerprint = _fingerprint( + {item.source_id: item.fingerprint for _, item in discovered} + ) + source_by_id = {item.source_id: (source, item) for source, item in discovered} + compatibility = _fingerprint({ + "pipeline": self.pipeline_version, + "model": self.embedding_model, + "dimensions": self.embedding_dimensions, + "chunk_policy": asdict(self.chunk_policy), + }) + spec = JobSpec( + workspace_id=workspace_id, + job_type="evidence", + workspace_root=workspace_root, + spec_version="jobs-v1", + pipeline_version=self.pipeline_version, + config_fingerprint=config_fingerprint, + input_fingerprint=_fingerprint([input_fingerprint, discovered_fingerprint]), + stage_ids=EVIDENCE_STAGE_IDS, + dry_run=dry_run, + resume_run_id=resume_run_id, + ) + + def artifact(context: JobContext, name: str) -> Path: + root = context.run_dir / "artifacts" + root.mkdir(exist_ok=True) + return root / name + + def write(context: JobContext, name: str, value) -> None: + artifact(context, name).write_text( + json.dumps(value, sort_keys=True, separators=(",", ":")), encoding="utf-8" + ) + + def read(context: JobContext, name: str): + try: + return json.loads(artifact(context, name).read_text(encoding="utf-8")) + except (OSError, ValueError) as error: + raise PipelineError("preprocessing checkpoint artifact is corrupt") from error + + def discover_stage(context: JobContext) -> None: + 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) + 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 + ) + unchanged = sorted(set(fingerprints) - set(changed)) + removed = sorted(set(prior) - set(fingerprints)) + write(context, "plan.json", { + "generation": f"gen:{context.run_id}", + "compatibility": compatibility, + "fingerprints": fingerprints, + "changed": changed, + "unchanged": unchanged, + "removed": removed, + "previous": previous.model_dump(mode="json") if previous else None, + }) + + def acquire_stage(context: JobContext) -> None: + if context.dry_run: + return + plan = read(context, "plan.json") + previous = CorpusManifest.model_validate(plan["previous"]) if plan["previous"] else None + prior = {doc.source_id: doc for doc in previous.documents} if previous else {} + documents = [prior[source_id] for source_id in plan["unchanged"]] + for source_id in plan["changed"]: + source, item = source_by_id[source_id] + documents.append(normalize(source.acquire(item), self.pipeline_version)) + documents.sort(key=lambda value: value.source_id) + chunks = [part for document in documents for part in chunk(document, self.chunk_policy)] + previous_generations = dict(previous.metadata.get("document_generations", {})) if previous else {} + changed = set(plan["changed"]) + generations = { + document.document_id: ( + plan["generation"] if document.source_id in changed + 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=plan["generation"], + documents=tuple(documents), chunks=tuple(chunks), + metadata={ + "compatibility_fingerprint": compatibility, + "fingerprints": plan["fingerprints"], + "removed": plan["removed"], + "document_generations": generations, + }, + ) + write(context, "manifest.json", manifest.model_dump(mode="json")) + + def embed_stage(context: JobContext) -> None: + if context.dry_run: + return + plan = read(context, "plan.json") + manifest = CorpusManifest.model_validate(read(context, "manifest.json")) + changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]} + parts = [part for part in manifest.chunks if part.document_id in changed_docs] + embeddings = self.embedder.embed_documents([part.content for part in parts]) + if len(embeddings) != len(parts) or any( + len(vector) != self.embedding_dimensions for vector in embeddings + ): + raise PipelineError("embedding output is incompatible") + write(context, "embeddings.json", embeddings) + + def records(context: JobContext): + plan = read(context, "plan.json") + manifest = CorpusManifest.model_validate(read(context, "manifest.json")) + changed_docs = {doc.document_id for doc in manifest.documents if doc.source_id in plan["changed"]} + parts = [part for part in manifest.chunks if part.document_id in changed_docs] + embeddings = read(context, "embeddings.json") + return [self._vector_record(part, vector, plan["generation"]) + for part, vector in zip(parts, embeddings, strict=True)] + + def compensate(context: JobContext) -> None: + generation = read(context, "plan.json")["generation"] + self.store.discard(generation) + try: + self.vector_store.delete_generation("evidence", generation) + except Exception: + pass + write(context, "compensated.json", {"generation": generation}) + + def rotate_compensated_generation(context: JobContext) -> None: + marker = artifact(context, "compensated.json") + if not marker.exists(): + return + plan = read(context, "plan.json") + old = plan["generation"] + plan["generation"] = f"gen:{uuid.uuid4().hex}" + write(context, "plan.json", plan) + manifest = CorpusManifest.model_validate(read(context, "manifest.json")) + changed = set(plan["changed"]) + generations = dict(manifest.metadata["document_generations"]) + for document in manifest.documents: + if document.source_id in changed and generations.get(document.document_id) == old: + generations[document.document_id] = plan["generation"] + metadata = dict(manifest.metadata) + metadata["document_generations"] = generations + manifest = manifest.model_copy(update={ + "vector_generation": plan["generation"], "metadata": metadata, + }) + write(context, "manifest.json", manifest.model_dump(mode="json")) + marker.unlink() + + def vector_stage(context: JobContext) -> None: + if context.dry_run: + return + rotate_compensated_generation(context) + values = records(context) + try: + if values and self.vector_store.upsert("evidence", values) != len(values): + raise PipelineError("vector write count mismatch") + except Exception: + compensate(context) + raise + + def stage_stage(context: JobContext) -> None: + if context.dry_run: + return + plan = read(context, "plan.json") + manifest = CorpusManifest.model_validate(read(context, "manifest.json")) + try: + self.store.stage( + manifest, {doc.document_id: doc.content for doc in manifest.documents}, + generation=plan["generation"], + ) + self.store.manifest(plan["generation"]) + except Exception: + compensate(context) + raise + + def publish_stage(context: JobContext) -> None: + if context.dry_run: + return + if artifact(context, "compensated.json").exists(): + rotate_compensated_generation(context) + values = records(context) + if values and self.vector_store.upsert("evidence", values) != len(values): + compensate(context) + raise PipelineError("vector write count mismatch") + manifest = CorpusManifest.model_validate(read(context, "manifest.json")) + generation = read(context, "plan.json")["generation"] + self.store.stage( + manifest, {doc.document_id: doc.content for doc in manifest.documents}, + generation=generation, + ) + generation = read(context, "plan.json")["generation"] + try: + self.store.publish(generation) + except Exception: + compensate(context) + raise + + def retention_stage(context: JobContext) -> None: + # Retention policy is intentionally a stable no-op until configured. + return + + report = run_job(spec, [ + discover_stage, acquire_stage, embed_stage, vector_stage, + stage_stage, publish_stage, retention_stage, + ]) + run_dir = workspace_root / ".tht-jobs" / "evidence" / "runs" / report.run_id + plan = json.loads((run_dir / "artifacts" / "plan.json").read_text()) + if dry_run: + manifest = self.store.active_manifest() or CorpusManifest(pipeline_version=self.pipeline_version) + generation = None + published = False + elif report.status == "succeeded": + generation = plan["generation"] + manifest = self.store.manifest(generation) + published = True + else: + generation = plan["generation"] + manifest_path = run_dir / "artifacts" / "manifest.json" + manifest = (CorpusManifest.model_validate_json(manifest_path.read_text()) + if manifest_path.exists() else CorpusManifest(pipeline_version=self.pipeline_version)) + published = False + return PipelineResult( + report.status, generation, published, tuple(plan["changed"]), + tuple(plan["unchanged"]), tuple(plan["removed"]), manifest, + report.run_id, report.resumed_from, + ) + def _run(self, *, dry_run: bool = False, resume: str | None = None) -> PipelineResult: generation = None vector_written = False diff --git a/harness/tht/jobs/runner.py b/harness/tht/jobs/runner.py index 1e05a869..dc6a23ef 100644 --- a/harness/tht/jobs/runner.py +++ b/harness/tht/jobs/runner.py @@ -7,6 +7,7 @@ import hashlib import os import uuid import stat +import shutil from collections.abc import Callable, Sequence from dataclasses import dataclass from pathlib import Path @@ -155,6 +156,9 @@ def run_job(spec: JobSpec, stages: Sequence[Stage]) -> JobReport: run = _new_run(spec, run_id, stages) else: run = _resume_run(spec, run_id, stages, source) + source_artifacts = jobs_root / source.run_id / "artifacts" + if source_artifacts.exists(): + shutil.copytree(source_artifacts, run_dir / "artifacts") _persist(checkpoint_path, run) context = JobContext(run_id, spec.job_type, spec.dry_run, spec.workspace_root, run_dir)