diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index c45a9ec32..042a5d403 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -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) -> PathBuf { result } +impl LocalDisk { + async fn reserve_version_delete(&self, volume: &str, object: &str, data_dir: Uuid, rollback_dir: Uuid) -> Result { + 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 { + 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; diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index b8fd83ced..38308cc7a 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -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 = HashSet::new(); - let mut per_disk_dirs: Vec<(usize, Vec)> = 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}"); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 13a791c65..99accf80e 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -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;