Compare commits

...

1 Commits

Author SHA1 Message Date
唐小鸭 741a1cac59 fix(filemeta): keep data dir of a version awaiting purge replication 2026-09-07 03:29:15 +08:00
2 changed files with 139 additions and 2 deletions
@@ -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();
+59 -2
View File
@@ -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();