diff --git a/.superpowers/sdd/evidence-task-6-report.md b/.superpowers/sdd/evidence-task-6-report.md index 8dcb9d6b..384a8897 100644 --- a/.superpowers/sdd/evidence-task-6-report.md +++ b/.superpowers/sdd/evidence-task-6-report.md @@ -28,3 +28,8 @@ pickle bytes can neither pass validation nor enter the snapshot. Reconciliation generation descriptor in a `finally` block on matches, mismatches, and exceptions. Pipeline-owned snapshot directories are removed and deregistered after `run_job` on both successful and failed runs, preventing repeated pipeline use from accumulating temporary directories or registry entries. + +The cleanup boundary now begins immediately after snapshot materialization. Resume checkpoint +validation and `JobSpec` construction are guarded by the same release routine as `run_job`, so +corrupt/mismatched resume state or constructor failure clears the pipeline holder, removes the +private directory, and restores the snapshot registry to its prior state before propagating. diff --git a/harness/tests/test_dwh_preprocess_job.py b/harness/tests/test_dwh_preprocess_job.py index 8bef80ba..cdc0dbf1 100644 --- a/harness/tests/test_dwh_preprocess_job.py +++ b/harness/tests/test_dwh_preprocess_job.py @@ -608,6 +608,61 @@ def test_pipeline_releases_materialized_snapshot_after_every_run(tmp_path): assert pipeline._snapshot_holder is None +def test_corrupt_resume_checkpoint_releases_materialized_snapshot(tmp_path): + import tht.jobs.dwh_pipeline as module + import pytest + + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: output.write_text("catalog"), + build_lsh=lambda physical, output: _write_lsh([], physical, output), + ) + assert pipeline.run().status == "succeeded" + baseline = set(module._SNAPSHOT_DIRS) + run_id = "e" * 32 + run_dir = tmp_path / ".tht-jobs" / "dwh" / "runs" / run_id + run_dir.mkdir(parents=True) + (run_dir / "checkpoint.json").write_text("not-json") + + with pytest.raises(Exception, match="checkpoint is invalid"): + pipeline.run(resume_run_id=run_id) + assert pipeline._snapshot_holder is None + assert set(module._SNAPSHOT_DIRS) == baseline + + +def test_job_spec_construction_failure_releases_materialized_snapshot(monkeypatch, tmp_path): + import tht.jobs.dwh_pipeline as module + import pytest + + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output: output.write_text("catalog"), + build_lsh=lambda physical, output: _write_lsh([], physical, output), + ) + assert pipeline.run().status == "succeeded" + baseline = set(module._SNAPSHOT_DIRS) + captured = [] + real_materialize = module._materialize_generation_fd + + def capture(*args, **kwargs): + holder, root = real_materialize(*args, **kwargs) + captured.append(holder) + return holder, root + + monkeypatch.setattr(module, "_materialize_generation_fd", capture) + monkeypatch.setattr( + module, "JobSpec", + lambda **kwargs: (_ for _ in ()).throw(RuntimeError("job spec injected")), + ) + with pytest.raises(RuntimeError, match="job spec injected"): + pipeline.run() + assert pipeline._snapshot_holder is None + assert set(module._SNAPSHOT_DIRS) == baseline + assert captured and all(not path.exists() for path in captured) + + def test_publish_root_swap_after_lease_never_writes_replacement(monkeypatch, tmp_path): import tht.jobs.dwh_pipeline as module diff --git a/harness/tht/jobs/dwh_pipeline.py b/harness/tht/jobs/dwh_pipeline.py index 87783fba..a483a404 100644 --- a/harness/tht/jobs/dwh_pipeline.py +++ b/harness/tht/jobs/dwh_pipeline.py @@ -601,19 +601,23 @@ class DwhPreprocessPipeline: self.current_lsh_dir = snapshot_root finally: lease.close() - if resume_run_id is not None: - self._validate_resume_publication(resume_run_id) - spec = JobSpec( - workspace_id=self.workspace_id, - job_type="dwh", - workspace_root=self.workspace_root, - spec_version="jobs-v1", - pipeline_version="dwh-v2", - config_fingerprint=self.config_fingerprint, - input_fingerprint=self.input_fingerprint, - stage_ids=steps, - resume_run_id=resume_run_id, - ) + try: + if resume_run_id is not None: + self._validate_resume_publication(resume_run_id) + spec = JobSpec( + workspace_id=self.workspace_id, + job_type="dwh", + workspace_root=self.workspace_root, + spec_version="jobs-v1", + pipeline_version="dwh-v2", + config_fingerprint=self.config_fingerprint, + input_fingerprint=self.input_fingerprint, + stage_ids=steps, + resume_run_id=resume_run_id, + ) + except BaseException: + self._release_pipeline_snapshot() + raise def introspect_stage(context: JobContext): artifacts = self._artifacts(context) @@ -651,16 +655,19 @@ class DwhPreprocessPipeline: reconcile_effects=self._reconcile_effects, ) finally: - holder, self._snapshot_holder = self._snapshot_holder, None - if isinstance(holder, Path): - shutil.rmtree(holder, ignore_errors=True) - _SNAPSHOT_DIRS.discard(holder) - if self.current_physical is not None and holder in self.current_physical.parents: - self.current_physical = None - if self.current_lsh_dir is not None and ( - self.current_lsh_dir == holder or holder in self.current_lsh_dir.parents - ): - self.current_lsh_dir = None + self._release_pipeline_snapshot() + + def _release_pipeline_snapshot(self) -> None: + holder, self._snapshot_holder = self._snapshot_holder, None + if isinstance(holder, Path): + shutil.rmtree(holder, ignore_errors=True) + _SNAPSHOT_DIRS.discard(holder) + if self.current_physical is not None and holder in self.current_physical.parents: + self.current_physical = None + if self.current_lsh_dir is not None and ( + self.current_lsh_dir == holder or holder in self.current_lsh_dir.parents + ): + self.current_lsh_dir = None def _reconcile_effects(self, source, run_dir: Path) -> set[str]: running = next(