@@ -410,20 +410,20 @@ def prune_run_bundles(
410410 protected .add (_safe_resolved_path (current_run_root ))
411411 clock = time .time () if now is None else now
412412
413- # Filesystem discovery and recursive size accounting are deliberately
414- # outside the lock. The destructive phase revalidates each candidate
415- # under the lock so another invocation can never turn a live bundle into a
416- # deletion candidate while discovery is in progress.
417- bundles = _discover_run_bundles (
418- runs_root ,
419- protected = protected ,
420- max_age_seconds = effective .max_age_seconds ,
421- now = clock ,
422- measure_sizes = effective .max_total_bytes is not None ,
423- size_budget = _RETENTION_SIZE_MEASUREMENT_BUDGET ,
424- )
425413 try :
426414 with _retention_lock (runs_root ):
415+ # Keep cursor read, size walk, and index update in one critical
416+ # section so concurrent pruners cannot overwrite scan progress.
417+ size_scan_cursor = _read_size_scan_cursor (runs_root )
418+ bundles , size_scan_cursor = _discover_run_bundles (
419+ runs_root ,
420+ protected = protected ,
421+ max_age_seconds = effective .max_age_seconds ,
422+ now = clock ,
423+ measure_sizes = effective .max_total_bytes is not None ,
424+ size_budget = _RETENTION_SIZE_MEASUREMENT_BUDGET ,
425+ size_scan_cursor = size_scan_cursor ,
426+ )
427427 _apply_bundle_retention (
428428 runs_root ,
429429 bundles ,
@@ -439,6 +439,7 @@ def prune_run_bundles(
439439 log ,
440440 current_run_root = current_run_root ,
441441 now = clock ,
442+ size_scan_cursor = size_scan_cursor ,
442443 )
443444 except (OSError , RuntimeError ) as exc :
444445 # Retention is maintenance. An unavailable lock or a transient
@@ -460,7 +461,7 @@ def refresh_run_bundle_index(
460461 if not runs_root .exists () or runs_root .is_symlink ():
461462 return
462463 try :
463- bundles = _discover_run_bundles (
464+ bundles , _size_scan_cursor = _discover_run_bundles (
464465 runs_root ,
465466 protected = set (),
466467 max_age_seconds = None ,
@@ -482,13 +483,14 @@ def _discover_run_bundles(
482483 now : float ,
483484 measure_sizes : bool ,
484485 size_budget : int ,
485- ) -> list [dict [str , Any ]]:
486+ size_scan_cursor : str | None = None ,
487+ ) -> tuple [list [dict [str , Any ]], str | None ]:
486488 bundles : list [dict [str , Any ]] = []
487- measured_sizes = 0
489+ scan_order : list [ dict [ str , Any ]] = []
488490 try :
489491 children = sorted (runs_root .iterdir (), key = lambda path : path .name )
490492 except OSError :
491- return bundles
493+ return bundles , size_scan_cursor
492494 for child in children :
493495 if child .name .startswith ("." ) or child .is_symlink () or not child .is_dir ():
494496 continue
@@ -524,18 +526,6 @@ def _discover_run_bundles(
524526 if status not in {"running" , "ok" , "aborted" , "error" }:
525527 continue
526528 resolved = _safe_resolved_path (child )
527- size = 0
528- size_known = False
529- if measure_sizes and measured_sizes < size_budget :
530- try :
531- size = _bundle_size (child )
532- size_known = True
533- measured_sizes += 1
534- except OSError :
535- # A file that disappears or becomes unreadable remains a
536- # retention candidate for count/age policy, but its byte
537- # contribution is unknown and must be reported below.
538- pass
539529 retention_metadata = metadata .get ("retention" )
540530 preserve = bool (metadata .get ("preserve" )) or (
541531 isinstance (retention_metadata , dict ) and retention_metadata .get ("preserve" ) is True
@@ -548,14 +538,60 @@ def _discover_run_bundles(
548538 "status" : status ,
549539 "started_at" : started_at ,
550540 "age" : age ,
551- "size" : size ,
552- "size_known" : size_known ,
541+ "size" : 0 ,
542+ "size_known" : False ,
553543 "preserve" : preserve ,
554544 "protected" : resolved in protected ,
555545 }
556546 )
547+ if measure_sizes :
548+ scan_order .append (bundles [- 1 ])
557549 bundles .sort (key = lambda bundle : (float (bundle ["started_at" ]), str (bundle ["path" ])))
558- return bundles
550+ if measure_sizes and bundles and size_budget > 0 :
551+ # The run index's cursor affects only which discovered bundles receive
552+ # an expensive size walk. It never authorizes deletion; every candidate
553+ # is re-read and revalidated before the destructive phase.
554+ if size_scan_cursor is not None :
555+ start_index = next (
556+ (index for index , bundle in enumerate (scan_order ) if bundle ["path" ].name > size_scan_cursor ),
557+ 0 ,
558+ )
559+ scan_order = scan_order [start_index :] + scan_order [:start_index ]
560+ attempted = 0
561+ for bundle in scan_order :
562+ if attempted >= size_budget :
563+ break
564+ attempted += 1
565+ path = bundle ["path" ]
566+ size_scan_cursor = path .name
567+ try :
568+ bundle ["size" ] = _bundle_size (path )
569+ bundle ["size_known" ] = True
570+ except OSError :
571+ # A file that disappears or becomes unreadable remains a
572+ # retention candidate for count/age policy, but its byte
573+ # contribution is unknown and reported below. Advancing the
574+ # cursor prevents one unreadable entry from starving others.
575+ pass
576+ return bundles , size_scan_cursor
577+
578+
579+ def _read_size_scan_cursor (runs_root : Path ) -> str | None :
580+ """Read the advisory byte-scan cursor; never use it to select deletions."""
581+
582+ index_path = runs_root / _RUN_INDEX_NAME
583+ try :
584+ if index_path .is_symlink () or not index_path .is_file () or index_path .stat ().st_size > 1_048_576 :
585+ return None
586+ payload = json .loads (index_path .read_text (encoding = "utf-8" ))
587+ except (OSError , UnicodeDecodeError , json .JSONDecodeError ):
588+ return None
589+ if not isinstance (payload , dict ):
590+ return None
591+ cursor = payload .get ("byte_scan_cursor" )
592+ if not isinstance (cursor , str ) or not cursor or len (cursor ) > 1024 or "/" in cursor or "\\ " in cursor :
593+ return None
594+ return cursor
559595
560596
561597def _apply_bundle_retention (
@@ -705,6 +741,7 @@ def _write_run_index(
705741 * ,
706742 current_run_root : Path | None = None ,
707743 now : float | None = None ,
744+ size_scan_cursor : str | None = None ,
708745) -> None :
709746 indexed = list (bundles )
710747 if current_run_root is not None and current_run_root .exists ():
@@ -734,6 +771,7 @@ def _write_run_index(
734771 "version" : 1 ,
735772 "complete" : omitted_bundles == 0 ,
736773 "omitted_bundles" : omitted_bundles ,
774+ "byte_scan_cursor" : size_scan_cursor if size_scan_cursor is not None else _read_size_scan_cursor (runs_root ),
737775 "bundles" : [
738776 {
739777 "path" : str (bundle ["path" ]),
0 commit comments