From 4d7d17f456ef719c0a1af483aa3c28d5ae7f5481 Mon Sep 17 00:00:00 2001 From: Xuban <59646791+EHxuban11@users.noreply.github.com> Date: Mon, 10 Aug 2026 21:19:21 +0200 Subject: [PATCH] Reclaim a dataset's cache and resume checkpoint once it finishes Nothing removed either, so both grew for the length of a campaign. On a 250 GB box the post-resize cache reached 157 GB, filled the disk, and deadlocked every worker for six hours; with 16 concurrent lanes it reaches 71 GB within two epochs because longest-first schedules the biggest datasets together. last.pt is 635 MB for a 51M-parameter model, 63 GB across 100 datasets. Neither is read again once a dataset is done. The cache is consumed only by the epochs of the dataset that wrote it, and last.pt exists to resume an interrupted run. best.pt and the copy at the weights root are what the uploader ships and what the skip logic reads, so those stay. Reclaim runs in the worker after stats and the checkpoint are written, never before, so a dataset that fails to record itself as done keeps its resume point. Bytes freed are reported in the worker result. --keep-cache opts out. Operators have been working around this with an external janitor. That is fragile: the cache filename is a LibreYOLO implementation detail, and when a version bump changed it from "*.r640x640.npy" to ".jpg.npy" the janitor silently matched nothing, kept reporting success, and the disk filled anyway. The code that writes the cache should own deleting it. --- tests/test_rf100vl_train.py | 48 ++++++++++++++++++++++++++++++ va_bench/cli.py | 9 ++++++ va_bench/rf100vl_train.py | 59 +++++++++++++++++++++++++++++++++++++ 3 files changed, 116 insertions(+) diff --git a/tests/test_rf100vl_train.py b/tests/test_rf100vl_train.py index f0c3a54..07c12ce 100644 --- a/tests/test_rf100vl_train.py +++ b/tests/test_rf100vl_train.py @@ -592,3 +592,51 @@ def test_child_registry_terminates_live_children_and_blocks_new_ones(): if sleeper.poll() is None: sleeper.kill() sleeper.wait() + + +def test_finished_dataset_reclaims_cache_and_resume_checkpoint(tmp_path): + """A finished dataset keeps what is read again and drops what is not. + + The post-resize cache and last.pt are the two largest consumers on a + campaign box. A campaign has already filled a 250 GB disk and deadlocked + every worker because nothing removed them. + """ + dataset = tmp_path / "some-dataset" + (dataset / "train").mkdir(parents=True) + image = dataset / "train" / "a.jpg" + image.write_bytes(b"jpeg") + cached = dataset / "train" / "a.jpg.npy" + cached.write_bytes(b"x" * 4096) + (dataset / "train" / "_annotations.coco.json").write_text("{}") + + weights = tmp_path / "run" / "weights" + weights.mkdir(parents=True) + (weights / "last.pt").write_bytes(b"y" * 2048) + (weights / "best.pt").write_bytes(b"z" * 1024) + + freed = rf100vl_train.reclaim_finished_dataset(dataset, tmp_path / "run") + + assert freed == 4096 + 2048 + assert not cached.exists() + assert not (weights / "last.pt").exists() + # what the uploader ships and what the images are must survive + assert (weights / "best.pt").read_bytes() == b"z" * 1024 + assert image.read_bytes() == b"jpeg" + assert (dataset / "train" / "_annotations.coco.json").exists() + + +def test_keep_cache_opts_out_of_reclaiming_the_cache(tmp_path): + dataset = tmp_path / "some-dataset" + (dataset / "train").mkdir(parents=True) + cached = dataset / "train" / "a.jpg.npy" + cached.write_bytes(b"x" * 4096) + weights = tmp_path / "run" / "weights" + weights.mkdir(parents=True) + (weights / "last.pt").write_bytes(b"y" * 2048) + + freed = rf100vl_train.reclaim_finished_dataset( + dataset, tmp_path / "run", keep_cache=True + ) + + assert cached.exists() + assert freed == 2048 diff --git a/va_bench/cli.py b/va_bench/cli.py index e22a205..dca96ed 100644 --- a/va_bench/cli.py +++ b/va_bench/cli.py @@ -266,6 +266,7 @@ def cmd_rf100vl_train(args: argparse.Namespace) -> None: state_root=args.state_root, smoke_epochs=args.smoke_epochs, force=args.force, + keep_cache=args.keep_cache, ) _stop_syncer(syncer) print( @@ -1213,6 +1214,14 @@ def main(argv: list[str] | None = None) -> None: action="store_true", help="Disable GPU telemetry capture (on by default; ~4.4 MB per campaign)", ) + rc.add_argument( + "--keep-cache", + action="store_true", + help="Keep each dataset's post-resize .npy cache and its last.pt once " + "the dataset finishes. Off by default: across 100 datasets those are " + "the two largest consumers on a campaign box and neither is read again " + "after a dataset is done.", + ) rc.add_argument( "--dollars-per-hour", type=float, diff --git a/va_bench/rf100vl_train.py b/va_bench/rf100vl_train.py index 551ae35..ed13b86 100644 --- a/va_bench/rf100vl_train.py +++ b/va_bench/rf100vl_train.py @@ -570,6 +570,50 @@ def _atomic_copy(source: Path, target: Path) -> None: temporary.replace(target) +def reclaim_finished_dataset( + dataset_dir: Path, + run_dir: Path, + *, + keep_cache: bool = False, +) -> int: + """Drop what a finished dataset no longer needs, and return bytes freed. + + Two things pile up per dataset and neither is needed once it is done: + + The post-resize image cache, one ``.npy`` beside each image. It is read + only by the epochs of the dataset that wrote it. Left behind it is the + single largest consumer on a campaign box: measured at 157 GB across the + 100 datasets, and 71 GB for the 16 concurrently active ones, against a + 250 GB disk. A campaign has filled its disk and deadlocked every worker on + this alone. + + ``last.pt``, which exists to resume an interrupted dataset. ``best.pt`` + and the copy at the weights root are what the uploader ships and what the + skip logic reads, so those stay. At 51M parameters ``last.pt`` is 635 MB, + or 63 GB across a campaign. + + Deleting the cache of a dataset another campaign is training on the same + box costs that campaign a re-cache, not correctness. + """ + freed = 0 + if not keep_cache: + for cached in dataset_dir.rglob("*.npy"): + try: + size = cached.stat().st_size + cached.unlink() + freed += size + except OSError: + continue + resume_checkpoint = run_dir / "weights" / "last.pt" + try: + if resume_checkpoint.is_file(): + freed += resume_checkpoint.stat().st_size + resume_checkpoint.unlink() + except OSError: + pass + return freed + + def run_dataset_worker(config_path: str | Path) -> int: """Train one dataset inside an isolated child process.""" config = load_json(config_path) @@ -685,6 +729,14 @@ def run_dataset_worker(config_path: str | Path) -> int: } stats_path = target_checkpoint.parent / "stats.json" atomic_write_json(stats_path, stats) + # Only after stats and the checkpoint are safely written: reclaiming + # first would risk deleting a resume point for a dataset that then + # failed to record itself as done. + freed = reclaim_finished_dataset( + dataset_dir, + run_dir, + keep_cache=bool(config.get("keep_cache", False)), + ) atomic_write_json( result_path, { @@ -692,6 +744,7 @@ def run_dataset_worker(config_path: str | Path) -> int: "stats_path": str(stats_path), "target_checkpoint": str(target_checkpoint), "wall_seconds": wall_seconds, + "reclaimed_bytes": freed, }, ) return 0 @@ -1147,6 +1200,7 @@ def _run_attempt( restart_reason: str | None, children: _ChildProcesses | None = None, disable_cuda_graph: bool = False, + keep_cache: bool = False, ) -> tuple[dict[str, Any], dict[str, Any], Path]: spec = get_spec(model_key) dataset_dir = data_dir / dataset_name @@ -1230,6 +1284,7 @@ def _run_attempt( "smoke_epochs": smoke_epochs, "restart_reason": restart_reason, "disable_cuda_graph": disable_cuda_graph, + "keep_cache": keep_cache, } worker_config_path = _worker_config_path(state_root, dataset_name) atomic_write_json(worker_config_path, worker_config) @@ -1323,6 +1378,7 @@ def orchestrate_training( state_root: str | Path | None = None, smoke_epochs: int | None = None, force: bool = False, + keep_cache: bool = False, ) -> dict[str, Any]: """Run a name-addressed dataset queue with ``jobs_per_gpu`` lanes per GPU.""" # (see order_longest_first for why the queue is not alphabetical) @@ -1450,6 +1506,7 @@ def consume(gpu: str, work_queue: "queue.Queue[str]", solo: bool = False) -> Non recipe_path=recipe_path, recipe=recipe, version_lock=version_lock, + keep_cache=keep_cache, versions_sha256=versions_sha256, gpu=gpu, timeout_seconds=timeout_hours * 3600, @@ -1502,6 +1559,7 @@ def consume(gpu: str, work_queue: "queue.Queue[str]", solo: bool = False) -> Non recipe_path=recipe_path, recipe=recipe, version_lock=version_lock, + keep_cache=keep_cache, versions_sha256=versions_sha256, gpu=gpu, timeout_seconds=timeout_hours * 3600, @@ -1544,6 +1602,7 @@ def consume(gpu: str, work_queue: "queue.Queue[str]", solo: bool = False) -> Non recipe_path=recipe_path, recipe=recipe, version_lock=version_lock, + keep_cache=keep_cache, versions_sha256=versions_sha256, gpu=gpu, timeout_seconds=timeout_hours * 3600,