diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 30a03ea48..623a08a26 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -846,6 +846,16 @@ impl RawEnumerationProgress { } }) } + + fn has_checkpointable_page_index(&self) -> bool { + self.page_index().is_some() + } + + fn checkpointable_entry_count(&self) -> usize { + self.page_index() + .and_then(|index| index.indexed_entries().ok()) + .map_or(0, |entries| entries.len()) + } } fn update_raw_enumeration_digest(digest: &mut Sha256, label: &[u8], value: &[u8]) { @@ -1165,20 +1175,33 @@ impl FolderScanner { } fn finish_raw_enumeration_parent(&mut self, parent: &str) { + let scan_root = self.old_cache.info.name.as_str(); self.raw_enumeration_progress.retain(|progress| { - progress.parent != parent - && !progress - .parent - .strip_prefix(parent) - .is_some_and(|suffix| suffix.starts_with(SLASH_SEPARATOR)) + if progress.parent == parent { + return parent == scan_root && progress.has_checkpointable_page_index(); + } + + !progress + .parent + .strip_prefix(parent) + .is_some_and(|suffix| suffix.starts_with(SLASH_SEPARATOR)) }); } fn take_raw_enumeration_resume_state(&mut self) -> (Option, Option) { - match self.raw_enumeration_progress.drain(..).next() { - Some(progress) => (progress.cursor(), progress.page_index()), - None => (None, None), + if self.raw_enumeration_progress.is_empty() { + return (None, None); } + let progress_index = self + .raw_enumeration_progress + .iter() + .enumerate() + .max_by_key(|(index, progress)| (progress.checkpointable_entry_count(), std::cmp::Reverse(*index))) + .map(|(index, _)| index) + .unwrap_or(0); + let progress = self.raw_enumeration_progress.swap_remove(progress_index); + self.raw_enumeration_progress.clear(); + (progress.cursor(), progress.page_index()) } fn carry_forward_old_children(&mut self, parent_hash: &DataUsageHash, entry: &mut DataUsageEntry) { diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index 385de36fd..05ba24c72 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -3512,6 +3512,80 @@ fn raw_enumeration_progress_checkpoint_commits_budgeted_page_for_oracle() { ); } +#[tokio::test] +async fn raw_enumeration_root_page_survives_child_partial_boundary() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir), + }; + scanner.old_cache.info.name = "bucket".to_string(); + + let mut root_progress = RawEnumerationProgress::new("bucket", None); + root_progress.record_entry("object-0000"); + root_progress.record_entry("object-0001"); + scanner.raw_enumeration_progress.push(root_progress); + scanner.finish_raw_enumeration_parent("bucket"); + + assert_eq!(scanner.raw_enumeration_progress.len(), 1); + let root_index = scanner.raw_enumeration_progress[0] + .page_index() + .expect("completed scan root should retain its raw-page oracle"); + assert_eq!( + root_index + .committed_entries() + .expect("retained root raw-page oracle should validate"), + vec!["object-0000".to_string(), "object-0001".to_string()] + ); + + let mut child_progress = RawEnumerationProgress::new("bucket/object-0000", None); + child_progress.record_entry("xl.meta"); + scanner.raw_enumeration_progress.push(child_progress); + scanner.finish_raw_enumeration_parent("bucket/object-0000"); + + assert_eq!( + scanner + .raw_enumeration_progress + .iter() + .map(|progress| progress.parent.as_str()) + .collect::>(), + vec!["bucket"] + ); +} + +#[tokio::test] +async fn raw_enumeration_resume_state_keeps_largest_durable_quantum() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir), + }; + + let mut root_progress = RawEnumerationProgress::new("bucket", None); + root_progress.record_entry("object-0000"); + scanner.raw_enumeration_progress.push(root_progress); + + let mut child_progress = RawEnumerationProgress::new("bucket/object-0000", None); + child_progress.record_entry("part-0000"); + child_progress.record_entry("part-0001"); + child_progress.record_entry("part-0002"); + scanner.raw_enumeration_progress.push(child_progress); + + let (cursor, page_index) = scanner.take_raw_enumeration_resume_state(); + assert_eq!( + cursor.as_ref().expect("largest raw quantum should include a cursor").parent, + "bucket/object-0000" + ); + assert_eq!( + page_index + .as_ref() + .expect("largest raw quantum should include a page index") + .indexed_entries() + .expect("selected page index should validate") + .len(), + 3 + ); + assert!(scanner.raw_enumeration_progress.is_empty()); +} + #[test] fn raw_enumeration_progress_retains_resume_index_until_unordered_entries_reappear() { let mut index = RawEnumerationPageIndex::new("bucket", 2).expect("raw page index should initialize"); diff --git a/scripts/diagnose_scanner_enumeration_restart.py b/scripts/diagnose_scanner_enumeration_restart.py index 8acbee1a3..3be70cdfb 100644 --- a/scripts/diagnose_scanner_enumeration_restart.py +++ b/scripts/diagnose_scanner_enumeration_restart.py @@ -85,15 +85,20 @@ def validate_recoverable_quantum(reports, *, objects, budget, require_converged) raise ValueError("no scanner restart reports were produced") previous = None made_enumeration_progress = False + made_raw_page_commit_progress = False made_classification_progress = False made_durable_progress = False for index, report in enumerate(reports): validate_report(report, round_number=index, pid=report["pid"], objects=objects, budget=budget) + if report["raw_page_index_parent"] == "bucket" and report["raw_page_index_committed_entries"] > 0: + made_raw_page_commit_progress = True if previous is not None: if report["objects_before"] != previous["objects_retained"]: raise ValueError("durable retained coverage did not survive process restart") if report["objects_retained"] < previous["objects_retained"]: raise ValueError("durable retained coverage regressed across restart") + if replays_raw_window(previous, report): + raise ValueError("raw enumeration window replayed without durable coverage") if (report["raw_page_index_parent"] == previous["raw_page_index_parent"] and report["raw_page_index_committed_entries"] < previous["raw_page_index_committed_entries"] and not previous["raw_page_index_complete"]): @@ -104,6 +109,8 @@ def validate_recoverable_quantum(reports, *, objects, budget, require_converged) previous = report if not made_enumeration_progress: raise ValueError("restart proof did not exercise raw enumeration") + if not made_raw_page_commit_progress: + raise ValueError("restart proof did not commit a durable raw enumeration page") if not made_classification_progress: raise ValueError("restart proof did not exercise object classification") if not made_durable_progress: diff --git a/scripts/test_diagnose_scanner_enumeration_restart.py b/scripts/test_diagnose_scanner_enumeration_restart.py index 02120f26e..370ebf580 100644 --- a/scripts/test_diagnose_scanner_enumeration_restart.py +++ b/scripts/test_diagnose_scanner_enumeration_restart.py @@ -140,6 +140,15 @@ class ReportTests(unittest.TestCase): advanced = dict(current, objects_retained=1) self.assertFalse(replays_raw_window(previous, advanced)) + def test_recoverable_quantum_rejects_replayed_raw_window(self): + previous = self.report() + previous.update(objects_retained=0, versions_retained=0, bytes_retained=0, + objects_processed=0, snapshot_complete=False, outcome="partial") + current = dict(previous, round=1, pid=124, objects_before=0) + + with self.assertRaisesRegex(ValueError, "raw enumeration window replayed"): + validate_recoverable_quantum([previous, current], objects=4, budget=16, require_converged=False) + def test_recoverable_quantum_requires_three_stage_progress_and_convergence(self): first = self.report() first.update(round=0, pid=123, raw_entries=2, raw_page_index_committed_entries=2, @@ -193,6 +202,15 @@ class ReportTests(unittest.TestCase): with self.assertRaisesRegex(ValueError, "object classification"): validate_recoverable_quantum([report], objects=4, budget=16, require_converged=False) + def test_recoverable_quantum_rejects_missing_raw_page_commit(self): + report = self.report() + report.update(snapshot_complete=False, outcome="partial", + raw_page_index_committed_entries=0, raw_page_index_indexed_entries=1, + objects_retained=1, versions_retained=1, bytes_retained=1) + + with self.assertRaisesRegex(ValueError, "durable raw enumeration page"): + validate_recoverable_quantum([report], objects=4, budget=16, require_converged=False) + if __name__ == "__main__": unittest.main()