Compare commits

..

3 Commits

Author SHA1 Message Date
houseme b06cbe326e fix(heal): cleanup consumed MRF replay journals
Do not retain Accepted or Merged replay intents as startup anchors after they have been handed to the heal manager. Only refused or still-pending replay records keep the journal on disk until a successor snapshot can persist them.

This keeps successor snapshots limited to the pending queue, which lets successful replay remove both authoritative and legacy journal paths and restores the crash-boundary tests around successor flush.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
(cherry picked from commit d5b8f49c9d)
2026-09-08 02:12:23 +08:00
houseme 2b4ce0cc1f fix(error): merge equivalent api message branches
Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-08 01:59:25 +08:00
houseme 47e2407362 fix(scanner): retain raw enumeration quantum
Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
2026-09-08 01:03:38 +08:00
5 changed files with 130 additions and 89 deletions
-81
View File
@@ -183,72 +183,6 @@ mod canonical_outcome {
assert_eq!((progress.objects_scanned, progress.objects_healed, progress.objects_failed), (2, 1, 1));
}
#[tokio::test(start_paused = true)]
async fn admin_cluster_lock_timeout_exhaustion_keeps_progress_and_retry_outcome() {
let storage = Arc::new(MockStorage::default());
storage.heal_object_outcomes.lock().expect("outcomes").insert(
"object-a".to_string(),
(0..4).map(|_| MockHealObjectOutcome::RetryableLockTimeout).collect(),
);
let mut request = HealRequest::new(
HealType::Cluster,
HealOptions {
recursive: true,
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
request.source = HealRequestSource::Admin;
let task = HealTask::from_request(request, storage.clone());
let err = task
.execute()
.await
.expect_err("legacy adapter still returns the batch failure detail");
assert!(
err.to_string()
.contains("Lock error: Lock acquisition timeout for resource 'object-a' after 5s"),
"lock timeout must remain actionable in the retained failure detail: {err}"
);
let outcome = task.get_outcome().await;
assert_eq!(outcome.execution, HealExecutionOutcome::CompletedWithErrors);
assert_eq!(outcome.coverage, HealTraversalCoverage::Complete);
assert_eq!(
(
outcome.counters.processed,
outcome.counters.failed,
outcome.counters.unknown,
outcome.counters.attempt_failures
),
(2, 1, 1, 4)
);
let failed = outcome
.objects
.iter()
.find(|item| item.identity.object == "object-a")
.expect("lock-contended object outcome");
assert_eq!(failed.disposition, HealObjectDisposition::Failed(HealFailureClass::RetryExhausted));
assert!(
failed
.detail
.as_deref()
.is_some_and(|detail| detail.contains("Lock error: Lock acquisition timeout for resource 'object-a' after 5s")),
"exhausted lock detail stays observable"
);
let progress = task.get_progress().await;
assert_eq!((progress.objects_scanned, progress.objects_healed, progress.objects_failed), (2, 1, 1));
let (legacy_summary, legacy_detail) = outcome.legacy_status("finished", None);
assert_eq!(legacy_summary, "stopped");
assert_eq!(legacy_detail.as_deref(), Some("heal traversal completed with errors: 1 failed objects"));
assert_eq!(
storage.heal_object_calls.lock().expect("object calls").as_slice(),
["object-a", "object-b", "object-a", "object-a", "object-a"]
);
}
#[tokio::test(start_paused = true)]
async fn retry_success_counts_one_terminal_outcome() {
let storage = Arc::new(MockStorage::default());
@@ -1359,7 +1293,6 @@ fn replacement_identity(
enum MockHealObjectOutcome {
RetryableLock,
RetryableLockTimeout,
OkWithOtherError(&'static str),
ErrOther(&'static str),
DanglingGraceDeferred,
@@ -1503,13 +1436,6 @@ impl HealStorageAPI for MockStorage {
owner: "competing-writer".to_string(),
}))),
)),
MockHealObjectOutcome::RetryableLockTimeout => Ok((
HealResultItem::default(),
Some(Error::Storage(EcstoreError::Lock(rustfs_lock::LockError::Timeout {
resource: object.to_string(),
timeout: Duration::from_secs(5),
}))),
)),
MockHealObjectOutcome::RetryableSlowDown => {
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::SlowDown))))
}
@@ -1542,13 +1468,6 @@ impl HealStorageAPI for MockStorage {
owner: "competing-writer".to_string(),
}))),
)),
MockHealObjectOutcome::RetryableLockTimeout => Ok((
HealResultItem::default(),
Some(Error::Storage(EcstoreError::Lock(rustfs_lock::LockError::Timeout {
resource: object.to_string(),
timeout: Duration::from_secs(5),
}))),
)),
MockHealObjectOutcome::RetryableSlowDown => {
Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::SlowDown))))
}
+31 -8
View File
@@ -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<DataUsageRawEnumerationCursor>, Option<RawEnumerationPageIndex>) {
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) {
@@ -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<_>>(),
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");
@@ -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:
@@ -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()