fix(dwh): release snapshots before job startup

This commit is contained in:
2026-07-12 07:34:56 +02:00
parent 4aae6c433d
commit 09d50ac9bb
3 changed files with 90 additions and 23 deletions
@@ -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 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 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. 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.
+55
View File
@@ -608,6 +608,61 @@ def test_pipeline_releases_materialized_snapshot_after_every_run(tmp_path):
assert pipeline._snapshot_holder is None 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): def test_publish_root_swap_after_lease_never_writes_replacement(monkeypatch, tmp_path):
import tht.jobs.dwh_pipeline as module import tht.jobs.dwh_pipeline as module
+7
View File
@@ -601,6 +601,7 @@ class DwhPreprocessPipeline:
self.current_lsh_dir = snapshot_root self.current_lsh_dir = snapshot_root
finally: finally:
lease.close() lease.close()
try:
if resume_run_id is not None: if resume_run_id is not None:
self._validate_resume_publication(resume_run_id) self._validate_resume_publication(resume_run_id)
spec = JobSpec( spec = JobSpec(
@@ -614,6 +615,9 @@ class DwhPreprocessPipeline:
stage_ids=steps, stage_ids=steps,
resume_run_id=resume_run_id, resume_run_id=resume_run_id,
) )
except BaseException:
self._release_pipeline_snapshot()
raise
def introspect_stage(context: JobContext): def introspect_stage(context: JobContext):
artifacts = self._artifacts(context) artifacts = self._artifacts(context)
@@ -651,6 +655,9 @@ class DwhPreprocessPipeline:
reconcile_effects=self._reconcile_effects, reconcile_effects=self._reconcile_effects,
) )
finally: finally:
self._release_pipeline_snapshot()
def _release_pipeline_snapshot(self) -> None:
holder, self._snapshot_holder = self._snapshot_holder, None holder, self._snapshot_holder = self._snapshot_holder, None
if isinstance(holder, Path): if isinstance(holder, Path):
shutil.rmtree(holder, ignore_errors=True) shutil.rmtree(holder, ignore_errors=True)