diff --git a/harness/tests/test_dwh_preprocess_job.py b/harness/tests/test_dwh_preprocess_job.py index 4d8a76db..79d6812f 100644 --- a/harness/tests/test_dwh_preprocess_job.py +++ b/harness/tests/test_dwh_preprocess_job.py @@ -231,6 +231,35 @@ def test_generation_retention_keeps_active_and_one_rollback(tmp_path): assert remaining == set(run_ids[-2:]) +def test_corrupt_newer_directory_does_not_consume_rollback_slot(tmp_path): + run_ids = [] + pipeline = None + for index in range(3): + pipeline = DwhPreprocessPipeline( + workspace_id="demo", workspace_root=tmp_path, + config_fingerprint=FP, input_fingerprint=FP, + introspect=lambda output, i=index: output.write_text(str(i)), + build_lsh=lambda physical, output, i=index: [ + (output / name).write_text(str(i)) + for name in ("demo_lsh.pkl", "demo_minhashes.pkl", "demo_meta.json") + ], + retain_generations=3, + ) + run_ids.append(pipeline.run().run_id) + generations = tmp_path / ".tht-dwh" / "generations" + corrupt = generations / ("f" * 32) + corrupt.mkdir(mode=0o700) + (corrupt / "junk").write_text("not a published generation") + + pipeline.retain_generations = 2 + pipeline._cleanup_generations() + + assert (generations / run_ids[-1]).is_dir() + assert (generations / run_ids[-2]).is_dir() + assert not (generations / run_ids[0]).exists() + assert corrupt.is_dir() + + def test_reader_lease_blocks_retain_one_publisher_until_file_reads_finish(tmp_path): import threading import time diff --git a/harness/tht/jobs/dwh_pipeline.py b/harness/tht/jobs/dwh_pipeline.py index 3581d1ec..e0815f37 100644 --- a/harness/tht/jobs/dwh_pipeline.py +++ b/harness/tht/jobs/dwh_pipeline.py @@ -453,6 +453,10 @@ class DwhPreprocessPipeline: and stat.S_ISDIR(info.st_mode) and info.st_uid == os.getuid() ): + try: + validate_generation(path) + except CorruptCheckpointError: + continue generations.append((info.st_mtime_ns, path.name)) generations.sort() keep_recent = {name for _, name in generations[-self.retain_generations:]}