fix(preprocess): ignore corrupt DWH retention entries
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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:]}
|
||||
|
||||
Reference in New Issue
Block a user