fix(ecstore): defer delete cleanup for snapshot reads (#5408)

* feat(ecstore): add local snapshot leases

* feat(ecstore): add remote snapshot lease RPCs

* feat(ecstore): protect streaming GETs with snapshot leases

* fix(ecstore): defer version cleanup for snapshot reads

* fix(ecstore): cover batch snapshot cleanup safely

* fix(e2e): stub snapshot lease RPCs in lock mock

* fix(e2e): stub snapshot lease RPCs in lock mock

* fix(ecstore): bind deferred delete cleanup intents

* fix(rpc): keep snapshot lease checks CI-compatible
This commit is contained in:
cxymds
2026-07-29 22:23:01 +08:00
committed by GitHub
parent c195b18fb8
commit cc24ef173c
3 changed files with 833 additions and 25 deletions
+605 -12
View File
@@ -77,6 +77,8 @@ const STALE_TMP_OBJECT_EXPIRY: Duration = Duration::from_secs(24 * 60 * 60);
const RUSTFS_META_TMP_OLD_BUCKET: &str = ".rustfs.sys/tmp-old";
const INLINE_METADATA_ROLLBACK_DIR_XOR: u128 = 0x7275737466735f696e6c696e655f7262;
const DELETE_MARKER_ROLLBACK_FILE: &str = "xl.meta.delete-marker.rollback";
pub(crate) const DELETE_DATA_DIR_MARKER_PREFIX: &str = "delete-data.";
pub(crate) const RESERVED_DELETE_DATA_DIR_MARKER_PREFIX: &str = "reserve-delete-data.";
const STARTUP_CLEANUP_WAIT_TIMEOUT: Duration = Duration::from_secs(2);
const ENV_BITROT_SIZE_MISMATCH_RETRY_COUNT: &str = "RUSTFS_BITROT_SIZE_MISMATCH_RETRY_COUNT";
const ENV_BITROT_SIZE_MISMATCH_RETRY_DELAY_MS: &str = "RUSTFS_BITROT_SIZE_MISMATCH_RETRY_DELAY_MS";
@@ -243,6 +245,7 @@ async fn restore_metadata_backup(object_dir: &Path, xl_path: &Path, rollback_dir
}
async fn restore_delete_rollback(object_dir: &Path, xl_path: &Path, rollback_dir: Uuid) -> Result<()> {
remove_version_delete_markers(object_dir, rollback_dir).await?;
let rollback_path = object_dir.join(rollback_dir.to_string());
let mut staged_paths = Vec::new();
let mut remove_new_metadata = false;
@@ -300,6 +303,31 @@ async fn restore_delete_rollback(object_dir: &Path, xl_path: &Path, rollback_dir
}
}
async fn remove_version_delete_markers(object_dir: &Path, rollback_dir: Uuid) -> Result<()> {
let reserved_name = format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}");
let committed_name = format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}");
let mut entries = match fs::read_dir(object_dir).await {
Ok(entries) => entries,
Err(err) if err.kind() == ErrorKind::NotFound => return Ok(()),
Err(err) => return Err(to_file_error(err).into()),
};
while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? {
if !entry.file_type().await.map_err(to_file_error)?.is_dir()
|| !entry.file_name().to_str().is_some_and(|name| Uuid::parse_str(name).is_ok())
{
continue;
}
for marker_name in [&reserved_name, &committed_name] {
match fs::remove_file(entry.path().join(marker_name)).await {
Ok(()) => {}
Err(err) if err.kind() == ErrorKind::NotFound => {}
Err(err) => return Err(to_file_error(err).into()),
}
}
}
Ok(())
}
async fn restore_delete_rollback_after_error(
object_dir: &Path,
xl_path: &Path,
@@ -4917,6 +4945,7 @@ impl LocalDisk {
fm.unmarshal_msg(&data)?;
let rollback_dir = opts.old_data_dir;
let mut reserved_version_delete = false;
if let Some(rollback_dir) = rollback_dir {
write_metadata_rollback_backup(object_dir, rollback_dir, &data).await?;
}
@@ -4930,6 +4959,18 @@ impl LocalDisk {
continue;
}
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_metadata_update",
err,
)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -4950,6 +4991,18 @@ impl LocalDisk {
let dir_path = match self.get_object_path(volume, format!("{path}/{dir}").as_str()) {
Ok(dir_path) => dir_path,
Err(err) => {
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_data_path",
err,
)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -4966,6 +5019,18 @@ impl LocalDisk {
let rollback_path = object_dir.join(rollback_dir.to_string());
if let Err(err) = fs::create_dir_all(&rollback_path).await {
let err: DiskError = to_file_error(err).into();
if reserved_version_delete {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_rollback_dir",
err,
)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -4977,8 +5042,26 @@ impl LocalDisk {
)
.await);
}
let reserved = match self.reserve_version_delete(volume, path, dir, rollback_dir).await {
Ok(reserved) => reserved,
Err(err) => {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_reserve_data",
err,
)
.await);
}
};
reserved_version_delete |= reserved;
let rollback_data_path = rollback_path.join(dir.to_string());
if let Err(err) = rename_all_ignore_missing_source(&dir_path, &rollback_data_path, &rollback_path).await {
if !reserved
&& let Err(err) = rename_all_ignore_missing_source(&dir_path, &rollback_data_path, &rollback_path).await
{
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -4991,6 +5074,18 @@ impl LocalDisk {
.await);
}
if should_fail_after_delete_data_staged(path) {
if reserved_version_delete {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_test_after_stage",
DiskError::Unexpected,
)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -5018,6 +5113,18 @@ impl LocalDisk {
// Remove xl.meta when no versions remain
if fm.versions.is_empty() {
if let Err(err) = self.delete_file(&volume_dir, &xlpath, true, false).await {
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_commit_delete",
err,
)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -5029,6 +5136,14 @@ impl LocalDisk {
)
.await);
}
if reserved_version_delete
&& let Some(rollback_dir) = rollback_dir
&& let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await
{
return Err(self
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err)
.await);
}
return Ok(());
}
@@ -5038,6 +5153,18 @@ impl LocalDisk {
Ok(buf) => buf,
Err(err) => {
let err: DiskError = err.into();
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(
object_dir,
rollback_dir,
volume,
path,
"delete_versions_metadata_encode",
err,
)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -5055,6 +5182,11 @@ impl LocalDisk {
.write_all_meta(volume, format!("{path}/{STORAGE_FORMAT_FILE}").as_str(), &buf, true)
.await
{
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_write", err)
.await);
}
return Err(restore_delete_rollback_after_error(
object_dir,
&xlpath,
@@ -5067,6 +5199,15 @@ impl LocalDisk {
.await);
}
if reserved_version_delete
&& let Some(rollback_dir) = rollback_dir
&& let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await
{
return Err(self
.abort_reserved_version_delete(object_dir, rollback_dir, volume, path, "delete_versions_commit_intent", err)
.await);
}
Ok(())
}
@@ -6086,6 +6227,108 @@ fn normalize_path_components(path: impl AsRef<Path>) -> PathBuf {
result
}
impl LocalDisk {
async fn reserve_version_delete(&self, volume: &str, object: &str, data_dir: Uuid, rollback_dir: Uuid) -> Result<bool> {
let path = format!("{object}/{data_dir}");
let data_path = self.get_object_path(volume, &path)?;
match fs::metadata(&data_path).await {
Ok(metadata) if metadata.is_dir() => {}
Ok(_) => return Ok(false),
Err(err) if err.kind() == ErrorKind::NotFound => return Ok(false),
Err(err) => return Err(to_file_error(err).into()),
}
let marker_path = data_path.join(format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"));
let marker = File::create(marker_path).await.map_err(to_file_error)?;
if effective_durability(volume).syncs_commit_metadata() {
marker.sync_all().await.map_err(to_file_error)?;
os::fsync_dir(&data_path).await.map_err(to_file_error)?;
}
Ok(true)
}
async fn commit_reserved_version_delete(&self, volume: &str, object: &str, rollback_dir: Uuid) -> Result<()> {
let object_path = self.get_object_path(volume, object)?;
let mut entries = match fs::read_dir(object_path).await {
Ok(entries) => entries,
Err(err) if err.kind() == ErrorKind::NotFound => return Ok(()),
Err(err) => return Err(to_file_error(err).into()),
};
let reserved_name = format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}");
let committed_name = format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}");
while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? {
if !entry.file_type().await.map_err(to_file_error)?.is_dir()
|| !entry.file_name().to_str().is_some_and(|name| Uuid::parse_str(name).is_ok())
{
continue;
}
let reserved_path = entry.path().join(&reserved_name);
match fs::rename(&reserved_path, entry.path().join(&committed_name)).await {
Ok(()) => {
if effective_durability(volume).syncs_commit_metadata() {
os::fsync_dir(&entry.path()).await.map_err(to_file_error)?;
}
}
Err(err) if err.kind() == ErrorKind::NotFound => {}
Err(err) => return Err(to_file_error(err).into()),
}
}
Ok(())
}
async fn finish_version_delete(&self, volume: &str, object: &str, rollback_dir: Uuid) -> Result<bool> {
let object_path = self.get_object_path(volume, object)?;
let mut entries = match fs::read_dir(object_path).await {
Ok(entries) => entries,
Err(err) if err.kind() == ErrorKind::NotFound => return Ok(false),
Err(err) => return Err(to_file_error(err).into()),
};
let marker_name = format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}");
let mut first_err = None;
let mut found = false;
while let Some(entry) = entries.next_entry().await.map_err(to_file_error)? {
let Some(data_dir) = entry.file_name().to_str().and_then(|data_dir| Uuid::parse_str(data_dir).ok()) else {
continue;
};
match fs::metadata(entry.path().join(&marker_name)).await {
Ok(metadata) if metadata.is_file() => found = true,
Ok(_) => continue,
Err(err) if err.kind() == ErrorKind::NotFound => continue,
Err(err) => return Err(to_file_error(err).into()),
}
if let Err(err) = self
.delete_data_dir(
volume,
&format!("{object}/{data_dir}"),
DeleteOptions {
recursive: true,
..Default::default()
},
)
.await
&& first_err.is_none()
&& err != DiskError::FileNotFound
&& err != DiskError::VolumeNotFound
{
first_err = Some(err);
}
}
first_err.map_or(Ok(found), Err)
}
async fn abort_reserved_version_delete(
&self,
object_dir: &Path,
rollback_dir: Uuid,
volume: &str,
object: &str,
stage: &'static str,
err: DiskError,
) -> DiskError {
let xl_path = object_dir.join(STORAGE_FORMAT_FILE);
restore_delete_rollback_after_error(object_dir, &xl_path, Some(rollback_dir), volume, object, stage, err).await
}
}
#[async_trait::async_trait]
impl DiskAPI for LocalDisk {
fn to_string(&self) -> String {
@@ -6254,7 +6497,19 @@ impl DiskAPI for LocalDisk {
#[tracing::instrument(level = "trace", skip_all)]
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> {
crate::hp_guard!("LocalDisk::delete");
self.delete_unleased(volume, path, &opt).await
let handled_version_delete = if opt.recursive
&& opt.immediate
&& let Some((object, transaction_id)) = path.rsplit_once('/')
&& let Ok(transaction_id) = Uuid::parse_str(transaction_id)
{
self.finish_version_delete(volume, object, transaction_id).await?
} else {
false
};
match self.delete_unleased(volume, path, &opt).await {
Err(DiskError::FileNotFound) if handled_version_delete => Ok(()),
result => result,
}
}
#[tracing::instrument(level = "trace", skip_all)]
@@ -8252,6 +8507,7 @@ impl DiskAPI for LocalDisk {
let mut meta = FileMeta::load(&buf)?;
let old_dir = meta.delete_version(&fi)?;
let mut reserved_version_delete = false;
if let Some(rollback_dir) = rollback_dir {
write_metadata_rollback_backup(file_path.as_path(), rollback_dir, &buf).await?;
}
@@ -8301,8 +8557,25 @@ impl DiskAPI for LocalDisk {
)
.await);
}
reserved_version_delete = match self.reserve_version_delete(volume, path, uuid, rollback_dir).await {
Ok(reserved) => reserved,
Err(err) => {
return Err(restore_delete_rollback_after_error(
file_path.as_path(),
&xl_path,
Some(rollback_dir),
volume,
path,
"delete_version_reserve_data",
err,
)
.await);
}
};
let rollback_data_path = rollback_path.join(uuid.to_string());
if let Err(err) = rename_all_ignore_missing_source(&old_path, &rollback_data_path, &rollback_path).await {
if !reserved_version_delete
&& let Err(err) = rename_all_ignore_missing_source(&old_path, &rollback_data_path, &rollback_path).await
{
return Err(restore_delete_rollback_after_error(
file_path.as_path(),
&xl_path,
@@ -8315,6 +8588,18 @@ impl DiskAPI for LocalDisk {
.await);
}
if should_fail_after_delete_data_staged(path) {
if reserved_version_delete {
return Err(self
.abort_reserved_version_delete(
file_path.as_path(),
rollback_dir,
volume,
path,
"delete_version_test_after_stage",
DiskError::Unexpected,
)
.await);
}
return Err(restore_delete_rollback_after_error(
file_path.as_path(),
&xl_path,
@@ -8346,6 +8631,18 @@ impl DiskAPI for LocalDisk {
Ok(buf) => buf,
Err(err) => {
let err: DiskError = err.into();
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(
file_path.as_path(),
rollback_dir,
volume,
path,
"delete_version_metadata_encode",
err,
)
.await);
}
return Err(restore_delete_rollback_after_error(
file_path.as_path(),
&xl_path,
@@ -8365,6 +8662,11 @@ impl DiskAPI for LocalDisk {
};
if let Err(err) = commit_result {
if reserved_version_delete && let Some(rollback_dir) = rollback_dir {
return Err(self
.abort_reserved_version_delete(file_path.as_path(), rollback_dir, volume, path, "delete_version_commit", err)
.await);
}
return Err(restore_delete_rollback_after_error(
file_path.as_path(),
&xl_path,
@@ -8377,6 +8679,22 @@ impl DiskAPI for LocalDisk {
.await);
}
if reserved_version_delete
&& let Some(rollback_dir) = rollback_dir
&& let Err(err) = self.commit_reserved_version_delete(volume, path, rollback_dir).await
{
return Err(self
.abort_reserved_version_delete(
file_path.as_path(),
rollback_dir,
volume,
path,
"delete_version_commit_intent",
err,
)
.await);
}
if should_fail_after_delete_commit(self.root.as_path(), path) {
return Err(DiskError::Unexpected);
}
@@ -11784,7 +12102,7 @@ mod test {
}
#[tokio::test]
async fn test_delete_version_rollback_restores_staged_data_dir() {
async fn test_delete_version_rollback_releases_reserved_data_dir() {
use tempfile::tempdir;
let dir = tempdir().expect("temp dir should be created");
@@ -11826,20 +12144,17 @@ mod test {
.expect("delete should stage rollback state");
assert!(!object_dir.join(STORAGE_FORMAT_FILE).exists());
assert!(!data_path.exists());
assert!(
data_path.exists(),
"the delete transaction must reserve the original data dir instead of moving it"
);
assert!(
object_dir
.join(rollback_dir.to_string())
.join(STORAGE_FORMAT_FILE_BACKUP)
.exists()
);
assert!(
object_dir
.join(rollback_dir.to_string())
.join(data_dir.to_string())
.join("part.1")
.exists()
);
assert!(!object_dir.join(rollback_dir.to_string()).join(data_dir.to_string()).exists());
disk.delete_version(
bucket,
@@ -15016,6 +15331,284 @@ mod test {
assert!(matches!(disk.read_all(volume, &first_part).await, Err(DiskError::FileNotFound)));
}
#[tokio::test]
async fn delete_version_keeps_later_part_until_snapshot_release() {
use tempfile::tempdir;
let root_dir = tempdir().expect("temp dir should be created");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
let volume = "snapshot-version-delete";
let object = "object";
let version_id = Uuid::new_v4();
let data_dir = Uuid::new_v4();
let rollback_dir = Uuid::new_v4();
let data_path = path_join_buf(&[object, &data_dir.to_string()]);
let first_part = path_join_buf(&[&data_path, "part.1"]);
let later_part = path_join_buf(&[&data_path, "part.2"]);
ensure_test_volume(&disk, volume).await;
disk.write_all(volume, &first_part, Bytes::from_static(b"first"))
.await
.expect("first shard should be written");
disk.write_all(volume, &later_part, Bytes::from_static(b"later"))
.await
.expect("later shard should be written");
let fi = test_file_info(object, version_id, Some(data_dir), None);
disk.write_all(volume, &path_join_buf(&[object, STORAGE_FORMAT_FILE]), test_meta(fi.clone()).into())
.await
.expect("metadata should be written");
let snapshot = disk
.acquire_snapshot_lease(volume, &data_path)
.await
.expect("snapshot lease should be acquired");
disk.delete_version(
volume,
object,
fi.clone(),
false,
DeleteOptions {
old_data_dir: Some(rollback_dir),
..Default::default()
},
)
.await
.expect("version delete should commit metadata");
disk.delete(
volume,
&format!("{object}/{rollback_dir}"),
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("version delete should schedule physical cleanup");
assert_eq!(
disk.read_all(volume, &later_part)
.await
.expect("a later multipart shard must remain openable while leased"),
Bytes::from_static(b"later")
);
disk.release_snapshot_lease(volume, &data_path, snapshot)
.await
.expect("snapshot release should run deferred cleanup");
assert!(matches!(disk.read_all(volume, &first_part).await, Err(DiskError::FileNotFound)));
}
#[tokio::test]
async fn version_delete_cleanup_intent_survives_local_disk_restart() {
use tempfile::tempdir;
let root_dir = tempdir().expect("temp dir should be created");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
let volume = "snapshot-version-delete-restart";
let object = "object";
let data_dir = Uuid::new_v4();
let rollback_dir = Uuid::new_v4();
let data_path = path_join_buf(&[object, &data_dir.to_string()]);
let part = path_join_buf(&[&data_path, "part.1"]);
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
ensure_test_volume(&disk, volume).await;
disk.write_all(volume, &part, Bytes::from_static(b"part"))
.await
.expect("shard should be written");
fs::create_dir_all(root_dir.path().join(volume).join(object).join(rollback_dir.to_string()))
.await
.expect("rollback directory should be created");
assert!(
disk.reserve_version_delete(volume, object, data_dir, rollback_dir)
.await
.expect("cleanup intent should be persisted")
);
disk.commit_reserved_version_delete(volume, object, rollback_dir)
.await
.expect("cleanup intent should be committed");
drop(disk);
let restarted = LocalDisk::new(&endpoint, false).await.expect("local disk should restart");
restarted
.delete(
volume,
&format!("{object}/{rollback_dir}"),
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("rollback cleanup should recover persisted intent");
assert!(matches!(restarted.read_all(volume, &part).await, Err(DiskError::FileNotFound)));
}
#[tokio::test]
async fn uuid_suffix_delete_does_not_run_version_cleanup_without_bound_marker() {
use tempfile::tempdir;
let root_dir = tempdir().expect("temp dir should be created");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
let volume = "snapshot-non-transaction-delete";
let object = "object";
let requested_dir = Uuid::new_v4();
let victim_dir = Uuid::new_v4();
ensure_test_volume(&disk, volume).await;
disk.write_all(
volume,
&format!("{object}/{requested_dir}/{DELETE_DATA_DIR_MARKER_PREFIX}{victim_dir}"),
Bytes::new(),
)
.await
.expect("legacy-shaped marker should be written");
disk.write_all(volume, &format!("{object}/{victim_dir}/part.1"), Bytes::from_static(b"live"))
.await
.expect("victim shard should be written");
disk.delete(
volume,
&format!("{object}/{requested_dir}"),
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("ordinary UUID directory delete should succeed");
assert_eq!(
disk.read_all(volume, &format!("{object}/{victim_dir}/part.1"))
.await
.expect("unbound sibling must not be deleted"),
Bytes::from_static(b"live")
);
}
#[tokio::test]
async fn version_delete_marker_is_durable_and_marker_errors_propagate() {
use tempfile::tempdir;
let _mode = durability_mode_override::set(DurabilityMode::Strict);
let root_dir = tempdir().expect("temp dir should be created");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
let volume = "snapshot-marker-durability";
let object = "object";
let data_dir = Uuid::new_v4();
let rollback_dir = Uuid::new_v4();
ensure_test_volume(&disk, volume).await;
let data_path = disk
.get_object_path(volume, &format!("{object}/{data_dir}"))
.expect("data path should resolve");
fs::create_dir_all(&data_path).await.expect("data dir should be created");
assert!(
disk.reserve_version_delete(volume, object, data_dir, rollback_dir)
.await
.expect("reserved marker should be written")
);
assert!(
os::fsync_dir_recorder::was_fsynced(&data_path),
"strict durability must fsync the data directory after marker creation"
);
let committed_path = data_path.join(format!("{DELETE_DATA_DIR_MARKER_PREFIX}{rollback_dir}"));
fs::create_dir_all(&committed_path)
.await
.expect("conflicting committed marker directory should be created");
fs::write(committed_path.join("entry"), b"conflict")
.await
.expect("conflicting marker directory should be non-empty");
disk.commit_reserved_version_delete(volume, object, rollback_dir)
.await
.expect_err("marker rename failure must propagate");
let second_data_dir = Uuid::new_v4();
let second_rollback = Uuid::new_v4();
let second_path = disk
.get_object_path(volume, &format!("{object}/{second_data_dir}"))
.expect("second data path should resolve");
fs::create_dir_all(second_path.join(format!("{RESERVED_DELETE_DATA_DIR_MARKER_PREFIX}{second_rollback}")))
.await
.expect("reserved marker conflict directory should be created");
assert!(
disk.reserve_version_delete(volume, object, second_data_dir, second_rollback)
.await
.is_err(),
"marker creation failure must propagate"
);
}
#[tokio::test]
async fn deferred_version_delete_replays_after_restart_without_rollback_dir() {
use tempfile::tempdir;
let root_dir = tempdir().expect("temp dir should be created");
let endpoint = Endpoint::try_from(root_dir.path().to_string_lossy().as_ref()).expect("endpoint should parse");
let volume = "snapshot-deferred-delete-restart";
let object = "object";
let version_id = Uuid::new_v4();
let data_dir = Uuid::new_v4();
let rollback_dir = Uuid::new_v4();
let data_path = path_join_buf(&[object, &data_dir.to_string()]);
let part = path_join_buf(&[&data_path, "part.1"]);
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
ensure_test_volume(&disk, volume).await;
disk.write_all(volume, &part, Bytes::from_static(b"part"))
.await
.expect("shard should be written");
let fi = test_file_info(object, version_id, Some(data_dir), None);
disk.write_all(volume, &path_join_buf(&[object, STORAGE_FORMAT_FILE]), test_meta(fi.clone()).into())
.await
.expect("metadata should be written");
let _lease = disk
.acquire_snapshot_lease(volume, &data_path)
.await
.expect("snapshot lease should be acquired");
disk.delete_version(
volume,
object,
fi,
false,
DeleteOptions {
old_data_dir: Some(rollback_dir),
..Default::default()
},
)
.await
.expect("version delete should commit");
disk.delete(
volume,
&format!("{object}/{rollback_dir}"),
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("physical cleanup should be deferred");
assert!(disk.read_all(volume, &part).await.is_ok(), "leased data must remain");
drop(disk);
let restarted = LocalDisk::new(&endpoint, false).await.expect("local disk should restart");
restarted
.delete(
volume,
&format!("{object}/{rollback_dir}"),
DeleteOptions {
recursive: true,
immediate: true,
..Default::default()
},
)
.await
.expect("committed marker should replay without rollback directory");
assert!(matches!(restarted.read_all(volume, &part).await, Err(DiskError::FileNotFound)));
}
#[tokio::test]
async fn data_dir_cleanup_without_a_lease_keeps_existing_behavior() {
use tempfile::tempdir;
@@ -46,6 +46,7 @@ use crate::diagnostics::get::{
GetObjectFailureReason, classify_disk_error, get_stage_timer_if_enabled, record_get_object_pipeline_failure,
record_get_object_pipeline_failure_for_path, record_get_stage_duration_if_enabled,
};
use crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX;
use crate::disk::{
DataDirDeleteStatus, OldCurrentSize, PART_TRANSACTION_NEW_META, PART_TRANSACTION_OLD_META, PART_TRANSACTION_ROLLBACK,
PartTransactionAction, part_transaction_path,
@@ -3952,9 +3953,9 @@ impl SetDisks {
/// * The set of referenced data dirs is the UNION of `get_data_dirs()` across
/// every online disk's `xl.meta`, so a dir named by *any* replica is kept.
/// * If a disk holds the object directory but its `xl.meta` is missing or
/// unparsable, the object is treated as degraded and NOTHING is removed —
/// the unreadable copy could be the only one naming a live data dir, and a
/// heal must run first.
/// unparsable, the object is treated as degraded and unmarked data dirs are
/// never removed. A data dir carrying a committed delete-transaction marker
/// remains reclaimable after a downgrade/re-upgrade cleanup interruption.
/// * Only subdirectories whose names parse as a UUID are ever considered;
/// removal is non-recursive-safe via a recursive delete of the full stray
/// data-dir path only.
@@ -3969,7 +3970,7 @@ impl SetDisks {
// physical UUID subdirectories present on each disk. Abort on any degraded
// copy so a healable object is never stripped of a referenced data dir.
let mut referenced: HashSet<Uuid> = HashSet::new();
let mut per_disk_dirs: Vec<(usize, Vec<Uuid>)> = Vec::new();
let mut per_disk_dirs: Vec<(usize, Vec<(Uuid, bool)>)> = Vec::new();
let mut healthy_metas = 0usize;
for (i, disk) in disks.iter().enumerate() {
@@ -4005,6 +4006,22 @@ impl SetDisks {
// to the orphan-dir / dangling-object heal paths.
continue;
}
let mut committed = Vec::with_capacity(physical.len());
for dir in physical {
let data_dir = format!("{object}/{dir}");
let committed_delete = disk.list_dir("", bucket, &data_dir, 0).await.is_ok_and(|entries| {
entries.iter().any(|entry| {
entry
.strip_prefix(DELETE_DATA_DIR_MARKER_PREFIX)
.is_some_and(|transaction| Uuid::parse_str(transaction).is_ok())
})
});
committed.push((dir, committed_delete));
}
if committed.iter().all(|(_, committed_delete)| *committed_delete) {
per_disk_dirs.push((i, committed));
continue;
}
warn!(
target: "rustfs_ecstore::set_disk",
bucket, object,
@@ -4041,22 +4058,16 @@ impl SetDisks {
healthy_metas += 1;
if !physical.is_empty() {
per_disk_dirs.push((i, physical));
per_disk_dirs.push((i, physical.into_iter().map(|dir| (dir, false)).collect()));
}
}
// No healthy metadata anywhere: this is not a live object, so surplus dirs
// (if any) belong to the dangling-object heal path, not here.
if healthy_metas == 0 {
return Ok(0);
}
// Phase 2: delete every physical data dir not referenced by the union.
let mut removed = 0usize;
for (i, physical) in per_disk_dirs {
let Some(disk) = disks[i].as_ref() else { continue };
for dir in physical {
if referenced.contains(&dir) {
for (dir, committed_delete) in physical {
if referenced.contains(&dir) || (healthy_metas == 0 && !committed_delete) {
continue;
}
let stray = format!("{object}/{dir}");
+204
View File
@@ -4698,6 +4698,7 @@ mod tests {
use crate::bucket::replication::{replication_statuses_map, version_purge_statuses_map};
use crate::disk::CHECK_PART_UNKNOWN;
use crate::disk::CHECK_PART_VOLUME_NOT_FOUND;
use crate::disk::DataDirDeleteStatus;
use crate::disk::DiskOption;
use crate::disk::RUSTFS_META_BUCKET;
use crate::disk::RUSTFS_META_TMP_BUCKET;
@@ -6403,6 +6404,62 @@ mod tests {
assert!(object_dir.join(STORAGE_FORMAT_FILE).exists(), "metadata must be preserved");
}
#[tokio::test]
async fn reclaim_orphan_data_dirs_recovers_deferred_cleanup_after_restart() {
let (dir, disk) = make_single_local_disk().await;
let live = Uuid::new_v4();
let orphan = Uuid::new_v4();
let object_dir = dir.path().join("bucket").join("obj");
write_object_meta_with_data_dirs(&object_dir, "bucket", "obj", &[live]).await;
fs::create_dir_all(object_dir.join(live.to_string()))
.await
.expect("live data dir should be created");
let orphan_path = format!("obj/{orphan}");
disk.write_all("bucket", &format!("{orphan_path}/part.1"), Bytes::from_static(b"stale"))
.await
.expect("orphan part should be written");
let _token = disk
.acquire_snapshot_lease("bucket", &orphan_path)
.await
.expect("snapshot lease should be acquired");
assert_eq!(
disk.delete_data_dir(
"bucket",
&orphan_path,
DeleteOptions {
recursive: true,
..Default::default()
},
)
.await
.expect("cleanup should be deferred"),
DataDirDeleteStatus::Deferred
);
drop(disk);
let endpoint =
Endpoint::try_from(dir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
let restarted = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("disk should restart");
let set = make_set_disks_with(vec![Some(restarted)]).await;
let removed = set
.reclaim_orphan_data_dirs("bucket", "obj")
.await
.expect("restart reclaim should succeed");
assert_eq!(removed, 1, "the deferred orphan should be reclaimed after restart");
assert!(object_dir.join(live.to_string()).exists(), "referenced data dir must be preserved");
assert!(!object_dir.join(orphan.to_string()).exists(), "deferred orphan must be removed");
}
// Nothing to reclaim when every physical data dir is still referenced.
#[tokio::test]
async fn reclaim_orphan_data_dirs_keeps_referenced_dir() {
@@ -6441,6 +6498,16 @@ mod tests {
fs::write(object_dir.join(stray.to_string()).join("part.1"), b"data")
.await
.expect("part should be written");
fs::write(
object_dir.join(stray.to_string()).join(format!(
"{}{}",
crate::disk::local::RESERVED_DELETE_DATA_DIR_MARKER_PREFIX,
Uuid::new_v4()
)),
[],
)
.await
.expect("pre-commit delete reservation should be written");
let set = make_set_disks_with(vec![Some(disk)]).await;
let removed = set
@@ -6455,6 +6522,36 @@ mod tests {
);
}
#[tokio::test]
async fn reclaim_orphan_data_dirs_recovers_committed_delete_marker_without_meta() {
let (dir, disk) = make_single_local_disk().await;
let stale = Uuid::new_v4();
let transaction = Uuid::new_v4();
let object_dir = dir.path().join("bucket").join("obj");
let stale_dir = object_dir.join(stale.to_string());
fs::create_dir_all(&stale_dir)
.await
.expect("committed stale data dir should be created");
fs::write(stale_dir.join("part.1"), b"stale")
.await
.expect("stale part should be written");
fs::write(
stale_dir.join(format!("{}{}", crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX, transaction)),
[],
)
.await
.expect("committed delete marker should be written");
let set = make_set_disks_with(vec![Some(disk)]).await;
let removed = set
.reclaim_orphan_data_dirs("bucket", "obj")
.await
.expect("upgrade reclaim should succeed");
assert_eq!(removed, 1, "the committed delete residue should be reclaimed");
assert!(!stale_dir.exists(), "the committed stale data dir should be removed");
}
// Cross-replica union: a data dir referenced by ANOTHER disk's xl.meta must be
// kept even where the local replica does not name it.
#[tokio::test]
@@ -9774,6 +9871,113 @@ mod tests {
.await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn streaming_get_snapshot_survives_concurrent_delete() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "snapshot-streaming-delete";
let object = "object";
let body = vec![0x41; 2 * 1024 * 1024];
let opts = ObjectOptions::default();
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(body.clone());
set_disks
.put_object(bucket, object, &mut reader, &opts)
.await
.expect("object should be written");
let mut snapshot = set_disks
.get_object_reader(bucket, object, None, HeaderMap::new(), &opts)
.await
.expect("snapshot reader should open");
let delete_set = Arc::clone(&set_disks);
let delete_opts = opts.clone();
let delete = tokio::spawn(async move { delete_set.delete_object(bucket, object, delete_opts).await });
tokio::time::timeout(Duration::from_secs(30), delete)
.await
.expect("delete should not wait for the response body")
.expect("delete task should join")
.expect("delete should succeed");
let mut restored = Vec::new();
snapshot
.stream
.read_to_end(&mut restored)
.await
.expect("leased snapshot should remain readable after delete");
assert_eq!(restored, body);
let err = match set_disks
.get_object_reader(bucket, object, None, HeaderMap::new(), &opts)
.await
{
Ok(_) => panic!("a new read must not observe the deleted object"),
Err(err) => err,
};
assert!(is_err_object_not_found(&err));
})
.await;
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn streaming_get_snapshot_survives_concurrent_delete_objects() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_OPTIMIZATION_ENABLE, Some("true"))], async {
let set_disks = make_local_bucket_test_set_disks().await;
let bucket = "snapshot-streaming-delete-objects";
let object = "object";
let body = vec![0x41; 2 * 1024 * 1024];
let opts = ObjectOptions::default();
set_disks
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
let mut reader = PutObjReader::from_vec(body.clone());
set_disks
.put_object(bucket, object, &mut reader, &opts)
.await
.expect("object should be written");
let mut snapshot = set_disks
.get_object_reader(bucket, object, None, HeaderMap::new(), &opts)
.await
.expect("snapshot reader should open");
let delete_set = Arc::clone(&set_disks);
let delete_opts = opts.clone();
let delete = tokio::spawn(async move {
delete_set
.delete_objects(
bucket,
vec![ObjectToDelete {
object_name: object.to_string(),
..Default::default()
}],
delete_opts,
)
.await
});
let (_, errors) = tokio::time::timeout(Duration::from_secs(30), delete)
.await
.expect("batch delete should not wait for the response body")
.expect("batch delete task should join");
assert!(errors.iter().all(Option::is_none));
let mut restored = Vec::new();
snapshot
.stream
.read_to_end(&mut restored)
.await
.expect("leased snapshot should remain readable after batch delete");
assert_eq!(restored, body);
})
.await;
}
#[tokio::test]
async fn set_level_batched_large_put_get_restores_body() {
const BATCHED_LARGE_SIZE: usize = 64 * 1024 * 1024;