From 14c99a994c7f2fdbbfc4f0c7a5cd4d952a7fc856 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=94=90=E5=B0=8F=E9=B8=AD?= Date: Mon, 7 Sep 2026 05:30:44 +0800 Subject: [PATCH] fix(filemeta): keep data dir of a version awaiting purge replication (#7307) --- .../src/replication_extension_test.rs | 80 +++++++++++++++++++ crates/filemeta/src/filemeta.rs | 61 +++++++++++++- 2 files changed, 139 insertions(+), 2 deletions(-) diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index 001ff2bb3..f28f9c153 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -3908,6 +3908,86 @@ async fn test_bucket_replication_replicates_directory_marker_in_versioned_bucket Ok(()) } +/// Regression for rustfs/backlog#2340 (not Wasabi specific): permanently +/// deleting a version whose payload lives in a data dir must leave the source +/// clean once the purge replicates. Managed-SSE objects are never inlined and a +/// plain object above the inline threshold takes the same layout. The version +/// retained with a pending purge used to lose its data dir, so the purge state +/// could never be applied (`VersionNotFound` on every retry) and the bucket +/// stayed `BucketNotEmpty` while `ListObjectVersions` was already empty. +#[tokio::test] +async fn test_bucket_replication_version_purge_of_non_inline_object_releases_source_bucket() -> TestResult { + init_logging(); + + let (source_env, target_env, source_bucket, target_bucket) = build_sse_replication_pair("purge-datadir", true, true).await?; + let target_arn = wait_for_remote_target_arn(&source_env, &source_bucket).await?; + put_bucket_replication_with_delete_statuses(&source_env, &source_bucket, &target_arn, "Enabled", Some("Enabled")).await?; + let source_client = source_env.create_s3_client(); + let target_client = target_env.create_s3_client(); + + let sse_key = "sse-object.bin"; + let large_key = "large-object.bin"; + let sse_put = source_client + .put_object() + .bucket(&source_bucket) + .key(sse_key) + .body(ByteStream::from_static(b"encrypted source payload")) + .server_side_encryption(ServerSideEncryption::Aes256) + .send() + .await?; + let large_put = source_client + .put_object() + .bucket(&source_bucket) + .key(large_key) + .body(ByteStream::from(vec![0x5a; 2 * 1024 * 1024])) + .send() + .await?; + let purged = [ + (sse_key, sse_put.version_id().ok_or("SSE PUT omitted version ID")?.to_string()), + (large_key, large_put.version_id().ok_or("large PUT omitted version ID")?.to_string()), + ]; + assert_replication_converged(&source_client, &source_bucket, &target_client, &target_bucket).await?; + + for (key, version_id) in &purged { + source_client + .delete_object() + .bucket(&source_bucket) + .key(*key) + .version_id(version_id) + .send() + .await?; + } + assert_replication_converged(&source_client, &source_bucket, &target_client, &target_bucket).await?; + let target_state = list_replication_state(&target_client, &target_bucket).await?; + assert!(target_state.is_empty(), "target retained an explicitly purged version: {target_state:?}"); + + // The purge state is applied on the source asynchronously after the target + // acknowledges the delete; only then does the retained version go away and + // the bucket become deletable. A listing that is empty while DeleteBucket + // keeps answering BucketNotEmpty is exactly the regression. + let deadline = tokio::time::Instant::now() + Duration::from_secs(60); + loop { + let listing = source_client.list_object_versions().bucket(&source_bucket).send().await?; + let listed = listing.versions().len() + listing.delete_markers().len(); + match source_client.delete_bucket().bucket(&source_bucket).send().await { + Ok(_) => break, + Err(err) if err.code() == Some("BucketNotEmpty") => { + if tokio::time::Instant::now() >= deadline { + return Err(format!( + "source bucket stayed BucketNotEmpty after the version purge replicated; \ + ListObjectVersions shows {listed} entries" + ) + .into()); + } + sleep(Duration::from_millis(500)).await; + } + Err(err) => return Err(err.into()), + } + } + + Ok(()) +} + #[tokio::test] async fn test_bucket_replication_disabled_delete_marker_does_not_propagate() -> TestResult { init_logging(); diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index 70bd66c13..2fc3f4320 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -692,10 +692,15 @@ impl FileMeta { } } - let old_dir = v.object.as_ref().map(|v| v.data_dir).unwrap_or_default(); + // The version stays on disk while the purge replicates + // (status PENDING/FAILED); its data dir must stay with + // it. Returning the dir here made the disk layer delete + // it, which turned every non-inline retained version + // into an unreadable zombie: the purge state could never + // be applied and the bucket could never be deleted. self.set_idx(i, v)?; - return Ok(old_dir); + return Ok(None); } found_index = Some(i); } @@ -2702,6 +2707,58 @@ mod test { ); } + /// Regression for rustfs/backlog#2340: a version purge that still awaits + /// the replication target keeps the object version on disk with a pending + /// purge status. Its data dir must be retained with it; handing the dir + /// back here made the disk layer delete it, leaving every non-inline + /// retained version unreadable. The dir is released only once the purge + /// completes and the version itself goes away. + #[test] + fn delete_version_pending_version_purge_retains_object_data_dir() { + let version_id = Uuid::new_v4(); + let data_dir = Uuid::new_v4(); + let mut fm = FileMeta::new(); + let mut fi = FileInfo::new("object", 2, 2); + fi.version_id = Some(version_id); + fi.data_dir = Some(data_dir); + fi.mod_time = Some(OffsetDateTime::now_utc()); + fm.add_version(fi).unwrap(); + + let pending_purge = FileInfo { + name: "object".to_string(), + version_id: Some(version_id), + mark_deleted: true, + replication_state_internal: Some(ReplicationState { + version_purge_status_internal: Some("target=PENDING;".to_string()), + purge_targets: version_purge_statuses_map("target=PENDING;"), + ..Default::default() + }), + ..Default::default() + }; + let freed = fm.delete_version(&pending_purge).unwrap(); + assert_eq!(freed, None, "a pending purge must not release the retained version's data dir"); + assert_eq!(fm.versions.len(), 1, "the version must stay until the purge replicates"); + let retained = fm + .into_fileinfo("vol", "object", &version_id.to_string(), false, false, true) + .unwrap(); + assert_eq!(retained.data_dir, Some(data_dir)); + assert_eq!(retained.version_purge_status(), VersionPurgeStatusType::Pending); + + let completed_purge = FileInfo { + name: "object".to_string(), + version_id: Some(version_id), + replication_state_internal: Some(ReplicationState { + version_purge_status_internal: Some("target=COMPLETE;".to_string()), + purge_targets: version_purge_statuses_map("target=COMPLETE;"), + ..Default::default() + }), + ..Default::default() + }; + let freed = fm.delete_version(&completed_purge).unwrap(); + assert_eq!(freed, Some(data_dir), "a completed purge removes the version and releases its data dir"); + assert!(fm.versions.is_empty()); + } + #[test] fn delete_version_accepts_delete_only_marker_and_free_version_paths() { let marker_version_id = Uuid::new_v4();