mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-30 10:08:58 +00:00
fix(ecstore): prevent recursive delete after empty bucket scan (#5453)
This commit is contained in:
Generated
+2
@@ -9218,6 +9218,8 @@ dependencies = [
|
||||
"url",
|
||||
"urlencoding",
|
||||
"uuid",
|
||||
"winapi-util",
|
||||
"windows-sys 0.61.2",
|
||||
"xxhash-rust",
|
||||
]
|
||||
|
||||
|
||||
@@ -314,7 +314,9 @@ uuid = { version = "1.24.0" }
|
||||
vaultrs = { version = "0.8.0" }
|
||||
tar = "0.4.46"
|
||||
walkdir = "2.5.0"
|
||||
winapi-util = "0.1.11"
|
||||
windows = { version = "0.62.2" }
|
||||
windows-sys = "0.61.2"
|
||||
xxhash-rust = { version = "0.8.18" }
|
||||
zip = "8.6.0"
|
||||
zstd = "0.13.3"
|
||||
|
||||
@@ -153,6 +153,10 @@ metrics = { workspace = true }
|
||||
[target.'cfg(target_os = "linux")'.dependencies]
|
||||
rustfs-uring = "0.2.1"
|
||||
|
||||
[target.'cfg(windows)'.dependencies]
|
||||
winapi-util.workspace = true
|
||||
windows-sys = { workspace = true, features = ["Win32_Foundation", "Win32_Storage_FileSystem"] }
|
||||
|
||||
[dev-dependencies]
|
||||
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "test-util", "fs"] }
|
||||
criterion = { workspace = true, features = ["html_reports"] }
|
||||
|
||||
@@ -40,7 +40,14 @@ use rustfs_protos::proto_gen::node_service::node_service_client::NodeServiceClie
|
||||
use rustfs_protos::proto_gen::node_service::{
|
||||
DeleteBucketRequest, GetBucketInfoRequest, HealBucketRequest, ListBucketRequest, MakeBucketRequest,
|
||||
};
|
||||
#[cfg(test)]
|
||||
use std::sync::{
|
||||
Mutex as StdMutex,
|
||||
atomic::{AtomicBool, Ordering},
|
||||
};
|
||||
use std::{collections::HashMap, fmt::Debug, sync::Arc, time::Duration};
|
||||
#[cfg(test)]
|
||||
use tokio::sync::Notify;
|
||||
use tokio::{net::TcpStream, sync::RwLock, time};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tonic::Request;
|
||||
@@ -50,6 +57,68 @@ use tracing::{debug, info, warn};
|
||||
|
||||
type Client = Arc<Box<dyn PeerS3Client>>;
|
||||
|
||||
#[cfg(test)]
|
||||
#[derive(Default)]
|
||||
pub(crate) struct DeleteBucketEmptyScanBarrier {
|
||||
arrived: AtomicBool,
|
||||
arrived_notify: Notify,
|
||||
released: AtomicBool,
|
||||
release_notify: Notify,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl DeleteBucketEmptyScanBarrier {
|
||||
pub(crate) async fn wait_until_paused(&self) {
|
||||
loop {
|
||||
let notified = self.arrived_notify.notified();
|
||||
if self.arrived.load(Ordering::Acquire) {
|
||||
return;
|
||||
}
|
||||
notified.await;
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn release(&self) {
|
||||
self.released.store(true, Ordering::Release);
|
||||
self.release_notify.notify_waiters();
|
||||
}
|
||||
|
||||
async fn pause(&self) {
|
||||
self.arrived.store(true, Ordering::Release);
|
||||
self.arrived_notify.notify_waiters();
|
||||
loop {
|
||||
let notified = self.release_notify.notified();
|
||||
if self.released.load(Ordering::Acquire) {
|
||||
return;
|
||||
}
|
||||
notified.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
static DELETE_BUCKET_EMPTY_SCAN_BARRIER: StdMutex<Option<Arc<DeleteBucketEmptyScanBarrier>>> = StdMutex::new(None);
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn install_delete_bucket_empty_scan_barrier() -> Arc<DeleteBucketEmptyScanBarrier> {
|
||||
let barrier = Arc::new(DeleteBucketEmptyScanBarrier::default());
|
||||
*DELETE_BUCKET_EMPTY_SCAN_BARRIER
|
||||
.lock()
|
||||
.expect("empty scan barrier lock should not be poisoned") = Some(barrier.clone());
|
||||
barrier
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn pause_after_delete_bucket_empty_scan() {
|
||||
let barrier = DELETE_BUCKET_EMPTY_SCAN_BARRIER
|
||||
.lock()
|
||||
.expect("empty scan barrier lock should not be poisoned")
|
||||
.take();
|
||||
if let Some(barrier) = barrier {
|
||||
barrier.pause().await;
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct ScannerBucketListing {
|
||||
pub buckets: Vec<BucketInfo>,
|
||||
@@ -650,24 +719,22 @@ impl PeerS3Client for LocalPeerS3Client {
|
||||
return Err(Error::ErasureWriteQuorum);
|
||||
}
|
||||
|
||||
let force = if opts.force_if_empty && !opts.force {
|
||||
if opts.force_if_empty && !opts.force {
|
||||
for disk in local_disks.iter() {
|
||||
if has_xlmeta_files(&disk.path().join(bucket)).await.map_err(Error::Io)? {
|
||||
return Err(Error::VolumeNotEmpty);
|
||||
}
|
||||
}
|
||||
true
|
||||
} else {
|
||||
opts.force
|
||||
};
|
||||
#[cfg(test)]
|
||||
pause_after_delete_bucket_empty_scan().await;
|
||||
}
|
||||
|
||||
let mut futures = Vec::with_capacity(local_disks.len());
|
||||
|
||||
for disk in local_disks.iter() {
|
||||
// Non-force delete refuses a non-empty bucket (VolumeNotEmpty), which
|
||||
// the recreate loop below turns into BucketNotEmpty; only an explicit
|
||||
// force delete removes recursively (backlog#799 B1).
|
||||
futures.push(disk.delete_volume(bucket, force));
|
||||
// `force_if_empty` is validation-only. Passing it as force would let
|
||||
// a PutObject committed after the scan be removed recursively.
|
||||
futures.push(disk.delete_volume(bucket, opts.force));
|
||||
}
|
||||
|
||||
let results = join_all(futures).await;
|
||||
@@ -730,6 +797,15 @@ pub struct RemotePeerS3Client {
|
||||
}
|
||||
|
||||
impl RemotePeerS3Client {
|
||||
fn encode_delete_bucket_options(opts: &DeleteBucketOptions) -> Result<String> {
|
||||
let mut remote_opts = opts.clone();
|
||||
// Older peers promote `force_if_empty` to recursive force after their
|
||||
// metadata scan. Keep this coordinator-only hint off the wire so a
|
||||
// mixed-version delete fails closed on non-empty directory remnants.
|
||||
remote_opts.force_if_empty = false;
|
||||
serde_json::to_string(&remote_opts).map_err(Into::into)
|
||||
}
|
||||
|
||||
fn recovery_monitor_span(addr: &str) -> tracing::Span {
|
||||
tracing::info_span!(
|
||||
"recovery-monitor",
|
||||
@@ -1040,7 +1116,7 @@ impl PeerS3Client for RemotePeerS3Client {
|
||||
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
|
||||
self.execute_with_timeout(
|
||||
|| async {
|
||||
let options = serde_json::to_string(opts)?;
|
||||
let options = Self::encode_delete_bucket_options(opts)?;
|
||||
let mut client = self.get_client().await?;
|
||||
|
||||
let mut request = Request::new(DeleteBucketRequest {
|
||||
@@ -1414,6 +1490,49 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_delete_bucket_options_fail_closed_for_legacy_peers() {
|
||||
let encoded = RemotePeerS3Client::encode_delete_bucket_options(&DeleteBucketOptions {
|
||||
no_lock: true,
|
||||
no_recreate: true,
|
||||
force_if_empty: true,
|
||||
..Default::default()
|
||||
})
|
||||
.expect("remote delete options should serialize");
|
||||
let legacy_opts: DeleteBucketOptions =
|
||||
serde_json::from_str(&encoded).expect("legacy peer should decode remote delete options");
|
||||
|
||||
assert!(legacy_opts.no_lock);
|
||||
assert!(legacy_opts.no_recreate);
|
||||
assert!(!legacy_opts.force);
|
||||
assert!(!legacy_opts.force_if_empty);
|
||||
|
||||
let legacy_recursive_force = if legacy_opts.force_if_empty && !legacy_opts.force {
|
||||
true
|
||||
} else {
|
||||
legacy_opts.force
|
||||
};
|
||||
assert!(
|
||||
!legacy_recursive_force,
|
||||
"legacy peer must not upgrade empty-only delete to recursive force"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_delete_bucket_options_preserve_explicit_force() {
|
||||
let encoded = RemotePeerS3Client::encode_delete_bucket_options(&DeleteBucketOptions {
|
||||
force: true,
|
||||
force_if_empty: true,
|
||||
..Default::default()
|
||||
})
|
||||
.expect("remote force-delete options should serialize");
|
||||
let remote_opts: DeleteBucketOptions =
|
||||
serde_json::from_str(&encoded).expect("remote peer should decode force-delete options");
|
||||
|
||||
assert!(remote_opts.force);
|
||||
assert!(!remote_opts.force_if_empty);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_execute_with_timeout_marks_remote_peer_faulty_on_network_like_error() {
|
||||
let client = test_remote_peer("http://peer-network-error:9000");
|
||||
@@ -1571,6 +1690,54 @@ mod tests {
|
||||
reset_local_disk_test_state().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn local_peer_force_if_empty_preserves_unclassified_file_in_selected_pool() {
|
||||
reset_local_disk_test_state().await;
|
||||
|
||||
let temp_dir = TempDir::new().expect("create temp dir for empty-only delete regression");
|
||||
let disks = init_test_local_disks_for_pools(
|
||||
&temp_dir,
|
||||
&[(0, 1), (1, 1)],
|
||||
"local-peer-force-if-empty-preserves-unclassified-file",
|
||||
)
|
||||
.await;
|
||||
let bucket = "empty-only-delete-bucket";
|
||||
let marker = "object/commit-marker";
|
||||
let data = bytes::Bytes::from_static(b"committed object data");
|
||||
|
||||
disks[1]
|
||||
.make_volume(bucket)
|
||||
.await
|
||||
.expect("bucket should be created in the selected pool");
|
||||
disks[1]
|
||||
.write_all(bucket, marker, data.clone())
|
||||
.await
|
||||
.expect("unclassified committed file should be written");
|
||||
|
||||
let err = LocalPeerS3Client::new_with_local_disks(None, Some(vec![1]), disks.clone())
|
||||
.delete_bucket(
|
||||
bucket,
|
||||
&DeleteBucketOptions {
|
||||
force_if_empty: true,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("empty-only delete must not recursively remove an unclassified file");
|
||||
|
||||
assert_eq!(err, Error::VolumeNotEmpty);
|
||||
assert_eq!(
|
||||
disks[1]
|
||||
.read_all(bucket, marker)
|
||||
.await
|
||||
.expect("unclassified committed file should be preserved"),
|
||||
data
|
||||
);
|
||||
|
||||
reset_local_disk_test_state().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn heal_bucket_local_recreates_missing_bucket_volumes() {
|
||||
|
||||
@@ -381,6 +381,291 @@ fn classify_delete_volume_error(err: std::io::Error) -> DiskError {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
struct EmptyDirectoryFrame {
|
||||
path: PathBuf,
|
||||
name_in_parent: std::ffi::CString,
|
||||
entries: rustix::fs::Dir,
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn empty_tree_io_error(err: rustix::io::Errno) -> std::io::Error {
|
||||
match err {
|
||||
rustix::io::Errno::NOTDIR | rustix::io::Errno::LOOP => std::io::Error::from(ErrorKind::DirectoryNotEmpty),
|
||||
_ => err.into(),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn remove_empty_directory_tree_unix_with(
|
||||
root: &Path,
|
||||
mut before_descend: impl FnMut(&Path) -> std::io::Result<()>,
|
||||
mut before_remove: impl FnMut(&Path) -> std::io::Result<()>,
|
||||
) -> std::io::Result<()> {
|
||||
use rustix::{
|
||||
fs::{AtFlags, Dir, Mode, OFlags, fstat, open, openat, statat, unlinkat},
|
||||
io::Errno,
|
||||
};
|
||||
use std::os::fd::AsFd;
|
||||
use std::os::unix::ffi::OsStrExt;
|
||||
|
||||
let flags = OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC;
|
||||
let root_parent_path = root
|
||||
.parent()
|
||||
.ok_or_else(|| std::io::Error::from(ErrorKind::DirectoryNotEmpty))?;
|
||||
let root_name = root
|
||||
.file_name()
|
||||
.ok_or_else(|| std::io::Error::from(ErrorKind::DirectoryNotEmpty))?
|
||||
.as_bytes();
|
||||
let root_name = std::ffi::CString::new(root_name).map_err(|_| std::io::Error::from(ErrorKind::DirectoryNotEmpty))?;
|
||||
let root_parent = open(root_parent_path, flags, Mode::empty()).map_err(empty_tree_io_error)?;
|
||||
let root_fd = openat(&root_parent, root_name.as_c_str(), flags, Mode::empty()).map_err(empty_tree_io_error)?;
|
||||
// Each frame owns one directory iterator/FD, so memory and descriptors are
|
||||
// bounded by path depth rather than by the number of empty remnants.
|
||||
let mut stack = vec![EmptyDirectoryFrame {
|
||||
path: root.to_path_buf(),
|
||||
name_in_parent: root_name,
|
||||
entries: Dir::new(root_fd).map_err(empty_tree_io_error)?,
|
||||
}];
|
||||
|
||||
while let Some(mut frame) = stack.pop() {
|
||||
let next_child = loop {
|
||||
let Some(entry) = frame.entries.next() else {
|
||||
break None;
|
||||
};
|
||||
let entry = entry.map_err(std::io::Error::from)?;
|
||||
let name = entry.file_name();
|
||||
if name.to_bytes() == b"." || name.to_bytes() == b".." {
|
||||
continue;
|
||||
}
|
||||
|
||||
let name = name.to_owned();
|
||||
let child_path = frame.path.join(std::ffi::OsStr::from_bytes(name.as_bytes()));
|
||||
before_descend(&child_path)?;
|
||||
break Some((child_path, name));
|
||||
};
|
||||
|
||||
if let Some((child_path, name)) = next_child {
|
||||
let parent = frame.entries.fd().map_err(std::io::Error::from)?;
|
||||
let child = match openat(parent, name.as_c_str(), flags, Mode::empty()) {
|
||||
Ok(child) => child,
|
||||
// A concurrent cleanup may remove an empty child after readdir
|
||||
// returns it. Resume the parent instead of treating the whole
|
||||
// bucket as missing and leaving its root behind.
|
||||
Err(Errno::NOENT) => {
|
||||
stack.push(frame);
|
||||
continue;
|
||||
}
|
||||
Err(err) => return Err(empty_tree_io_error(err)),
|
||||
};
|
||||
stack.push(frame);
|
||||
stack.push(EmptyDirectoryFrame {
|
||||
path: child_path,
|
||||
name_in_parent: name,
|
||||
entries: Dir::new(child).map_err(empty_tree_io_error)?,
|
||||
});
|
||||
continue;
|
||||
}
|
||||
before_remove(&frame.path)?;
|
||||
let parent = if let Some(parent) = stack.last() {
|
||||
parent.entries.fd().map_err(std::io::Error::from)?
|
||||
} else {
|
||||
root_parent.as_fd()
|
||||
};
|
||||
let expected = fstat(frame.entries.fd().map_err(std::io::Error::from)?).map_err(empty_tree_io_error)?;
|
||||
let current = statat(parent, frame.name_in_parent.as_c_str(), AtFlags::SYMLINK_NOFOLLOW).map_err(empty_tree_io_error)?;
|
||||
if current.st_dev != expected.st_dev || current.st_ino != expected.st_ino {
|
||||
return Err(std::io::Error::from(ErrorKind::DirectoryNotEmpty));
|
||||
}
|
||||
match unlinkat(parent, frame.name_in_parent.as_c_str(), AtFlags::REMOVEDIR).map_err(empty_tree_io_error) {
|
||||
Ok(()) => {}
|
||||
Err(err) if err.kind() == ErrorKind::NotFound => {}
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
async fn remove_empty_directory_tree_with(
|
||||
root: &Path,
|
||||
before_descend: impl FnMut(&Path) -> std::io::Result<()>,
|
||||
before_remove: impl FnMut(&Path) -> std::io::Result<()>,
|
||||
) -> std::io::Result<()> {
|
||||
remove_empty_directory_tree_unix_with(root, before_descend, before_remove)
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
async fn remove_empty_directory_tree(root: &Path) -> std::io::Result<()> {
|
||||
let root = root.to_path_buf();
|
||||
tokio::task::spawn_blocking(move || remove_empty_directory_tree_unix_with(&root, |_| Ok(()), |_| Ok(()))).await?
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[derive(Debug)]
|
||||
struct LockedEmptyDirectory {
|
||||
handle: winapi_util::Handle,
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
fn validate_windows_empty_directory(file_attributes: u64) -> std::io::Result<()> {
|
||||
const FILE_ATTRIBUTE_DIRECTORY: u64 = 0x10;
|
||||
const FILE_ATTRIBUTE_REPARSE_POINT: u64 = 0x400;
|
||||
|
||||
if file_attributes & FILE_ATTRIBUTE_DIRECTORY == 0 || file_attributes & FILE_ATTRIBUTE_REPARSE_POINT != 0 {
|
||||
return Err(std::io::Error::from(ErrorKind::DirectoryNotEmpty));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
async fn lock_windows_empty_directory(path: &Path, canonical_root: Option<&Path>) -> std::io::Result<LockedEmptyDirectory> {
|
||||
use std::os::windows::fs::OpenOptionsExt;
|
||||
use windows_sys::Win32::{
|
||||
Foundation::GENERIC_READ,
|
||||
Storage::FileSystem::{DELETE, FILE_SHARE_READ},
|
||||
};
|
||||
|
||||
const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
|
||||
const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
|
||||
|
||||
let path = path.to_path_buf();
|
||||
let canonical_root = canonical_root.map(Path::to_path_buf);
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let file = std::fs::OpenOptions::new()
|
||||
.access_mode(GENERIC_READ | DELETE)
|
||||
.share_mode(FILE_SHARE_READ)
|
||||
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
|
||||
.open(&path)?;
|
||||
let handle = winapi_util::Handle::from_file(file);
|
||||
let info = winapi_util::file::information(&handle)?;
|
||||
validate_windows_empty_directory(info.file_attributes())?;
|
||||
if let Some(canonical_root) = canonical_root {
|
||||
let canonical_path = std::fs::canonicalize(path)?;
|
||||
if !canonical_path.starts_with(canonical_root) {
|
||||
return Err(std::io::Error::from(ErrorKind::DirectoryNotEmpty));
|
||||
}
|
||||
}
|
||||
Ok::<_, std::io::Error>(LockedEmptyDirectory { handle })
|
||||
})
|
||||
.await?
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
// SAFETY: This helper only passes an owned live handle and one initialized
|
||||
// FILE_DISPOSITION_INFO to the synchronous Windows deletion API.
|
||||
#[allow(unsafe_code)]
|
||||
async fn remove_windows_empty_directory(directory: LockedEmptyDirectory) -> std::io::Result<()> {
|
||||
tokio::task::spawn_blocking(move || {
|
||||
use std::os::windows::io::AsRawHandle;
|
||||
use windows_sys::Win32::Storage::FileSystem::{FILE_DISPOSITION_INFO, FileDispositionInfo, SetFileInformationByHandle};
|
||||
|
||||
let disposition = FILE_DISPOSITION_INFO { DeleteFile: true };
|
||||
let disposition_size = u32::try_from(std::mem::size_of_val(&disposition))
|
||||
.map_err(|_| std::io::Error::other("FILE_DISPOSITION_INFO size exceeds the Win32 API limit"))?;
|
||||
let handle = directory.handle.as_raw_handle();
|
||||
// SAFETY: `handle` is owned by `directory` and stays live for this synchronous
|
||||
// call. `disposition` is initialized with the exact structure and byte size
|
||||
// required by `FileDispositionInfo`; Windows does not retain the pointer.
|
||||
let deleted = unsafe {
|
||||
SetFileInformationByHandle(handle, FileDispositionInfo, std::ptr::from_ref(&disposition).cast(), disposition_size)
|
||||
};
|
||||
if deleted == 0 {
|
||||
Err(std::io::Error::last_os_error())
|
||||
} else {
|
||||
Ok(())
|
||||
}
|
||||
})
|
||||
.await?
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
struct WindowsEmptyDirectoryFrame {
|
||||
path: PathBuf,
|
||||
directory: LockedEmptyDirectory,
|
||||
entries: fs::ReadDir,
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
async fn remove_empty_directory_tree_with(
|
||||
root: &Path,
|
||||
mut before_descend: impl FnMut(&Path) -> std::io::Result<()>,
|
||||
mut before_remove: impl FnMut(&Path) -> std::io::Result<()>,
|
||||
) -> std::io::Result<()> {
|
||||
let root_directory = lock_windows_empty_directory(root, None).await?;
|
||||
let canonical_root = fs::canonicalize(root).await?;
|
||||
let root_entries = fs::read_dir(root).await?;
|
||||
|
||||
// Holding each validated directory without delete sharing keeps its path
|
||||
// generation stable until handle-relative deletion. State is O(depth).
|
||||
let mut stack = vec![WindowsEmptyDirectoryFrame {
|
||||
path: root.to_path_buf(),
|
||||
directory: root_directory,
|
||||
entries: root_entries,
|
||||
}];
|
||||
|
||||
while let Some(mut frame) = stack.pop() {
|
||||
match frame.entries.next_entry().await {
|
||||
Ok(Some(entry)) => {
|
||||
let child = entry.path();
|
||||
before_descend(&child)?;
|
||||
let child_directory = match lock_windows_empty_directory(&child, Some(&canonical_root)).await {
|
||||
Ok(directory) => directory,
|
||||
Err(err) if err.kind() == ErrorKind::NotFound => {
|
||||
stack.push(frame);
|
||||
continue;
|
||||
}
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
let child_entries = match fs::read_dir(&child).await {
|
||||
Ok(entries) => entries,
|
||||
Err(err) if err.kind() == ErrorKind::NotFound => {
|
||||
stack.push(frame);
|
||||
continue;
|
||||
}
|
||||
Err(err) if err.kind() == ErrorKind::NotADirectory => {
|
||||
return Err(std::io::Error::from(ErrorKind::DirectoryNotEmpty));
|
||||
}
|
||||
Err(err) => return Err(err),
|
||||
};
|
||||
stack.push(frame);
|
||||
stack.push(WindowsEmptyDirectoryFrame {
|
||||
path: child,
|
||||
directory: child_directory,
|
||||
entries: child_entries,
|
||||
});
|
||||
}
|
||||
Ok(None) => {
|
||||
before_remove(&frame.path)?;
|
||||
drop(frame.entries);
|
||||
match remove_windows_empty_directory(frame.directory).await {
|
||||
Ok(()) => {}
|
||||
Err(err) if err.kind() == ErrorKind::NotFound => {}
|
||||
Err(err) if err.kind() == ErrorKind::NotADirectory => {
|
||||
return Err(std::io::Error::from(ErrorKind::DirectoryNotEmpty));
|
||||
}
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
}
|
||||
Err(err) if err.kind() == ErrorKind::NotFound => {}
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
async fn remove_empty_directory_tree(root: &Path) -> std::io::Result<()> {
|
||||
remove_empty_directory_tree_with(root, |_| Ok(()), |_| Ok(())).await
|
||||
}
|
||||
|
||||
#[cfg(all(not(unix), not(windows)))]
|
||||
async fn remove_empty_directory_tree(root: &Path) -> std::io::Result<()> {
|
||||
fs::remove_dir(root).await
|
||||
}
|
||||
|
||||
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
||||
const LOG_SUBSYSTEM_DISK_LOCAL: &str = "disk_local";
|
||||
const EVENT_DISK_LOCAL_STARTUP_CLEANUP: &str = "disk_local_startup_cleanup";
|
||||
@@ -1556,6 +1841,8 @@ static RENAME_DATA_FAIL_COMMIT_RENAME: std::sync::Mutex<Option<String>> = std::s
|
||||
#[cfg(test)]
|
||||
static LOCAL_INLINE_ROLLBACK_HARDLINK_FAILURE: std::sync::Mutex<Option<PathBuf>> = std::sync::Mutex::new(None);
|
||||
#[cfg(test)]
|
||||
static RENAME_DATA_REMOVE_DST_BASE_BEFORE_COMMIT: std::sync::Mutex<Option<(String, PathBuf)>> = std::sync::Mutex::new(None);
|
||||
#[cfg(test)]
|
||||
static DELETE_VERSION_FAIL_AFTER_DATA_STAGED: std::sync::Mutex<Vec<String>> = std::sync::Mutex::new(Vec::new());
|
||||
#[cfg(test)]
|
||||
static DELETE_VERSION_FAIL_AFTER_COMMIT: std::sync::Mutex<Vec<(PathBuf, String)>> = std::sync::Mutex::new(Vec::new());
|
||||
@@ -1588,6 +1875,13 @@ fn set_local_inline_rollback_hardlink_failure(dst_path: &Path) {
|
||||
.expect("test failpoint lock should not be poisoned") = Some(dst_path.to_path_buf());
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn set_rename_data_remove_dst_base_before_commit(dst_path: &str, dst_base: &Path) {
|
||||
*RENAME_DATA_REMOVE_DST_BASE_BEFORE_COMMIT
|
||||
.lock()
|
||||
.expect("test failpoint lock should not be poisoned") = Some((dst_path.to_string(), dst_base.to_path_buf()));
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn set_delete_version_fail_after_data_staged(path: &str) {
|
||||
DELETE_VERSION_FAIL_AFTER_DATA_STAGED
|
||||
@@ -1656,6 +1950,19 @@ fn should_fail_local_inline_rollback_hardlink(dst_path: &Path) -> bool {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn remove_dst_base_before_commit(dst_path: &str) -> std::io::Result<()> {
|
||||
let mut target = RENAME_DATA_REMOVE_DST_BASE_BEFORE_COMMIT
|
||||
.lock()
|
||||
.expect("test failpoint lock should not be poisoned");
|
||||
let Some((_, base)) = target.as_ref().filter(|(target_path, _)| target_path == dst_path) else {
|
||||
return Ok(());
|
||||
};
|
||||
std::fs::remove_dir_all(base)?;
|
||||
target.take();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn should_fail_after_delete_data_staged(path: &str) -> bool {
|
||||
let mut targets = DELETE_VERSION_FAIL_AFTER_DATA_STAGED
|
||||
@@ -1705,6 +2012,11 @@ fn should_fail_local_inline_rollback_hardlink(_dst_path: &Path) -> bool {
|
||||
false
|
||||
}
|
||||
|
||||
#[cfg(not(test))]
|
||||
fn remove_dst_base_before_commit(_dst_path: &str) -> std::io::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(not(test))]
|
||||
fn should_fail_after_delete_data_staged(_path: &str) -> bool {
|
||||
false
|
||||
@@ -7359,6 +7671,10 @@ impl DiskAPI for LocalDisk {
|
||||
dst_path: &str,
|
||||
) -> Result<RenameDataResp> {
|
||||
crate::hp_guard!("LocalDisk::rename_data");
|
||||
// A non-force DeleteBucket must not remove a directory while a local
|
||||
// object commit is publishing into it. The peer's empty scan remains
|
||||
// optimistic; this guard establishes the local commit/delete order.
|
||||
let _volume_mutation_guard = os::disk_volume_mutation_lock(&self.root, dst_volume).read_owned().await;
|
||||
if fi.is_legacy_indexed_delete_marker() {
|
||||
fi.erasure.index = 0;
|
||||
}
|
||||
@@ -7555,6 +7871,7 @@ impl DiskAPI for LocalDisk {
|
||||
// sequential version did.
|
||||
tmp_meta_res?;
|
||||
shard_sync_res?;
|
||||
remove_dst_base_before_commit(dst_path).map_err(to_file_error)?;
|
||||
|
||||
if let Some((src_data_path, dst_data_path)) = has_data_dir_path.as_ref()
|
||||
&& let Err(err) = rename_all(src_data_path, dst_data_path, &skip_parent).await
|
||||
@@ -7814,6 +8131,7 @@ impl DiskAPI for LocalDisk {
|
||||
xlmeta.add_version(fi)?;
|
||||
let version_signature = rename_data_versions_signature(&xlmeta);
|
||||
let new_buf = xlmeta.marshal_msg()?;
|
||||
remove_dst_base_before_commit(&dst_path_for_failpoint).map_err(to_file_error)?;
|
||||
|
||||
// Write new xl.meta + rename. Inline objects carry their data
|
||||
// inside xl.meta, so this whole sequence is a metadata commit:
|
||||
@@ -7838,9 +8156,11 @@ impl DiskAPI for LocalDisk {
|
||||
{
|
||||
let old_path = dst_parent.join(old_dir.to_string()).join(STORAGE_FORMAT_FILE_BACKUP);
|
||||
let old_parent = old_path.parent().map(|p| p.to_path_buf());
|
||||
if let Some(ref old_parent) = old_parent {
|
||||
std::fs::create_dir_all(old_parent)?;
|
||||
}
|
||||
let _old_parent_guard = old_parent
|
||||
.as_deref()
|
||||
.map(|parent| os::mkdir_all_below_existing_base_std(parent, &bucket_dir))
|
||||
.transpose()
|
||||
.map_err(to_file_error)?;
|
||||
// This rollback backup is the sole restore source for a later
|
||||
// undo_write when the set-level write quorum fails. Persist it as
|
||||
// durably as the new xl.meta written above (and as the non-inline
|
||||
@@ -7866,6 +8186,12 @@ impl DiskAPI for LocalDisk {
|
||||
local_rollback_path = Some(create_local_inline_rollback_backup(&dst, &src, old_metadata)?);
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
let _commit_parent_guard = if let Some(parent) = dst.parent() {
|
||||
Some(os::mkdir_all_below_existing_base_std(parent, &bucket_dir).map_err(to_file_error)?)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let commit_result = if should_fail_commit_rename(&dst_path_for_failpoint) {
|
||||
Err(std::io::Error::other("test fail during metadata commit rename"))
|
||||
} else {
|
||||
@@ -7874,7 +8200,8 @@ impl DiskAPI for LocalDisk {
|
||||
Err(err) if err.kind() == ErrorKind::NotFound && !src.exists() => Ok(()),
|
||||
Err(err) if err.kind() == ErrorKind::NotFound => {
|
||||
if let Some(parent) = dst.parent() {
|
||||
std::fs::create_dir_all(parent)?;
|
||||
let _parent_guard =
|
||||
os::mkdir_all_below_existing_base_std(parent, &bucket_dir).map_err(to_file_error)?;
|
||||
}
|
||||
std::fs::rename(&src, &dst).map_err(to_file_error)?;
|
||||
Ok(())
|
||||
@@ -8791,18 +9118,16 @@ impl DiskAPI for LocalDisk {
|
||||
#[tracing::instrument(level = "trace", skip_all)]
|
||||
async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()> {
|
||||
let p = self.get_bucket_path(volume)?;
|
||||
let _volume_mutation_guard = os::disk_volume_mutation_lock(&self.root, volume).write_owned().await;
|
||||
|
||||
// Non-force is non-recursive: `remove_dir` (rmdir) fails atomically with
|
||||
// `DirectoryNotEmpty` -> VolumeNotEmpty if the bucket still holds any
|
||||
// object data, so a misclassified "dangling" bucket on the heal path
|
||||
// (or a non-force S3 DeleteBucket on a populated bucket) can never be
|
||||
// recursively wiped. Only an explicit `force_delete` (e.g. S3 force
|
||||
// bucket delete) removes recursively. Mirrors MinIO's
|
||||
// xlStorage.DeleteVol (Remove vs RemoveAll). (backlog#799 B1)
|
||||
// Non-force removes empty directory remnants children-first with
|
||||
// non-recursive rmdir calls. A file that exists during the scan, or
|
||||
// appears before its parent is removed, fails closed with
|
||||
// VolumeNotEmpty. Only an explicit force delete removes recursively.
|
||||
let res = if force_delete {
|
||||
fs::remove_dir_all(&p).await
|
||||
} else {
|
||||
fs::remove_dir(&p).await
|
||||
remove_empty_directory_tree(&p).await
|
||||
};
|
||||
|
||||
if let Err(err) = res {
|
||||
@@ -9061,6 +9386,35 @@ mod test {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[test]
|
||||
fn windows_empty_tree_requires_non_reparse_directory() {
|
||||
validate_windows_empty_directory(0x10).expect("ordinary directories should be accepted");
|
||||
assert!(validate_windows_empty_directory(0).is_err());
|
||||
assert!(validate_windows_empty_directory(0x10 | 0x400).is_err());
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[tokio::test]
|
||||
async fn windows_empty_tree_blocks_replacement_at_final_delete_boundary() {
|
||||
let root = tempfile::tempdir().expect("temporary root should be created");
|
||||
let bucket_path = root.path().join("bucket");
|
||||
let child_path = bucket_path.join("child");
|
||||
fs::create_dir_all(&child_path).await.expect("bucket child should be created");
|
||||
let canonical_root = fs::canonicalize(&bucket_path).await.expect("bucket path should canonicalize");
|
||||
let directory = lock_windows_empty_directory(&child_path, Some(&canonical_root))
|
||||
.await
|
||||
.expect("child directory should be locked");
|
||||
|
||||
std::fs::rename(&child_path, bucket_path.join("replacement"))
|
||||
.expect_err("the locked directory must not be replaceable at the final deletion boundary");
|
||||
remove_windows_empty_directory(directory)
|
||||
.await
|
||||
.expect("handle-relative deletion should remove the locked directory");
|
||||
|
||||
assert!(!child_path.exists(), "the exact locked directory should be removed");
|
||||
}
|
||||
|
||||
async fn ensure_test_volume(disk: &LocalDisk, volume: &str) {
|
||||
match disk.make_volume(volume).await {
|
||||
Ok(()) | Err(DiskError::VolumeExists) => {}
|
||||
@@ -10787,6 +11141,72 @@ mod test {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(rename_data_deleted_bucket)]
|
||||
async fn rename_data_non_inline_does_not_recreate_bucket_deleted_before_commit() {
|
||||
let dir = tempfile::tempdir().expect("temp dir should be created");
|
||||
let endpoint = Endpoint::try_from(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 bucket = "deleted-before-non-inline-commit";
|
||||
let object = "prefix/object";
|
||||
let tmp_object = "tmp-non-inline-delete-race";
|
||||
let data_dir = Uuid::parse_str("aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee").expect("data dir should parse");
|
||||
ensure_test_volume(&disk, bucket).await;
|
||||
ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await;
|
||||
|
||||
let staged_data = disk
|
||||
.get_object_path(RUSTFS_META_TMP_BUCKET, &format!("{tmp_object}/{data_dir}/part.1"))
|
||||
.expect("staged data path should resolve");
|
||||
fs::create_dir_all(staged_data.parent().expect("staged data should have a parent"))
|
||||
.await
|
||||
.expect("staged data parent should be created");
|
||||
fs::write(&staged_data, b"payload")
|
||||
.await
|
||||
.expect("staged shard should be written");
|
||||
|
||||
let bucket_path = disk.get_bucket_path(bucket).expect("bucket path should resolve");
|
||||
set_rename_data_remove_dst_base_before_commit(object, &bucket_path);
|
||||
let fi = test_file_info(object, Uuid::new_v4(), Some(data_dir), None);
|
||||
let err = disk
|
||||
.rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, fi, bucket, object)
|
||||
.await
|
||||
.expect_err("commit must fail after the destination bucket is deleted");
|
||||
|
||||
assert_eq!(err, DiskError::FileNotFound);
|
||||
assert!(!bucket_path.exists(), "rename_data must not recreate the deleted bucket");
|
||||
assert!(staged_data.exists(), "failed commit must preserve the staged shard");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(rename_data_deleted_bucket)]
|
||||
async fn rename_data_inline_does_not_recreate_bucket_deleted_before_commit() {
|
||||
let dir = tempfile::tempdir().expect("temp dir should be created");
|
||||
let endpoint = Endpoint::try_from(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 bucket = "deleted-before-inline-commit";
|
||||
let object = "prefix/object";
|
||||
let tmp_object = "tmp-inline-delete-race";
|
||||
ensure_test_volume(&disk, bucket).await;
|
||||
ensure_test_volume(&disk, RUSTFS_META_TMP_BUCKET).await;
|
||||
|
||||
let bucket_path = disk.get_bucket_path(bucket).expect("bucket path should resolve");
|
||||
set_rename_data_remove_dst_base_before_commit(object, &bucket_path);
|
||||
let fi = test_file_info(object, Uuid::new_v4(), None, Some(Bytes::from_static(b"inline payload")));
|
||||
let err = disk
|
||||
.rename_data(RUSTFS_META_TMP_BUCKET, tmp_object, fi, bucket, object)
|
||||
.await
|
||||
.expect_err("inline commit must fail after the destination bucket is deleted");
|
||||
|
||||
assert_eq!(err, DiskError::FileNotFound);
|
||||
assert!(!bucket_path.exists(), "inline rename_data must not recreate the deleted bucket");
|
||||
assert!(
|
||||
disk.get_object_path(RUSTFS_META_TMP_BUCKET, &format!("{tmp_object}/{STORAGE_FORMAT_FILE}"))
|
||||
.expect("staged metadata path should resolve")
|
||||
.exists(),
|
||||
"failed inline commit must preserve staged metadata"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_resolve_durability_mode_mapping() {
|
||||
// Default: nothing set -> strict (current main behavior).
|
||||
@@ -14517,6 +14937,154 @@ mod test {
|
||||
let _ = fs::remove_dir_all(&test_dir).await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_volume_non_force_removes_nested_empty_directories() {
|
||||
let root = tempfile::tempdir().expect("temporary disk root should be created");
|
||||
let endpoint = Endpoint::try_from(root.path().to_string_lossy().as_ref()).expect("endpoint should parse");
|
||||
let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created");
|
||||
let bucket = "nested-empty-bucket";
|
||||
|
||||
disk.make_volume(bucket).await.expect("bucket should be created");
|
||||
fs::create_dir_all(disk.path().join(bucket).join("a/b/c"))
|
||||
.await
|
||||
.expect("nested empty directories should be created");
|
||||
|
||||
disk.delete_volume(bucket, false)
|
||||
.await
|
||||
.expect("non-force delete should remove an empty directory tree");
|
||||
|
||||
assert!(matches!(disk.stat_volume(bucket).await, Err(DiskError::VolumeNotFound)));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn empty_tree_delete_preserves_xlmeta_published_after_scan() {
|
||||
let root = tempfile::tempdir().expect("temporary disk root should be created");
|
||||
let bucket_path = root.path().join("bucket");
|
||||
let object_path = bucket_path.join("object").join(STORAGE_FORMAT_FILE);
|
||||
fs::create_dir_all(object_path.parent().expect("object path should have a parent"))
|
||||
.await
|
||||
.expect("empty object directory should be created");
|
||||
|
||||
let data = b"committed object metadata";
|
||||
let mut published = false;
|
||||
let err = remove_empty_directory_tree_with(
|
||||
&bucket_path,
|
||||
|_| Ok(()),
|
||||
|directory| {
|
||||
if !published && directory == object_path.parent().expect("object path should have a parent") {
|
||||
std::fs::write(&object_path, data)?;
|
||||
published = true;
|
||||
}
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("rmdir should refuse metadata published after the directory scan");
|
||||
|
||||
assert!(published, "test barrier should publish metadata before rmdir");
|
||||
assert!(matches!(classify_delete_volume_error(err), DiskError::VolumeNotEmpty));
|
||||
assert_eq!(std::fs::read(&object_path).expect("object metadata should be preserved"), data);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn empty_tree_delete_rejects_child_replaced_with_external_symlink() {
|
||||
use std::os::unix::fs::symlink;
|
||||
|
||||
let root = tempfile::tempdir().expect("temporary root should be created");
|
||||
let bucket_path = root.path().join("bucket");
|
||||
let child_path = bucket_path.join("child");
|
||||
let outside_path = root.path().join("outside");
|
||||
let outside_empty = outside_path.join("must-remain");
|
||||
fs::create_dir_all(&child_path).await.expect("bucket child should be created");
|
||||
fs::create_dir_all(&outside_empty)
|
||||
.await
|
||||
.expect("outside directory should be created");
|
||||
|
||||
let mut replaced = false;
|
||||
let err = remove_empty_directory_tree_with(
|
||||
&bucket_path,
|
||||
|child| {
|
||||
if !replaced && child == child_path {
|
||||
std::fs::remove_dir(&child_path)?;
|
||||
symlink(&outside_path, &child_path)?;
|
||||
replaced = true;
|
||||
}
|
||||
Ok(())
|
||||
},
|
||||
|_| Ok(()),
|
||||
)
|
||||
.await
|
||||
.expect_err("a replaced child must fail closed");
|
||||
|
||||
assert!(replaced, "test barrier should replace the child before it is opened");
|
||||
assert!(matches!(classify_delete_volume_error(err), DiskError::VolumeNotEmpty));
|
||||
assert!(outside_empty.exists(), "bucket deletion must not remove directories outside the bucket");
|
||||
assert!(
|
||||
std::fs::symlink_metadata(&child_path)
|
||||
.expect("replacement symlink should remain")
|
||||
.file_type()
|
||||
.is_symlink()
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn empty_tree_delete_rechecks_parent_after_child_disappears() {
|
||||
let root = tempfile::tempdir().expect("temporary root should be created");
|
||||
let bucket_path = root.path().join("bucket");
|
||||
let child_path = bucket_path.join("child");
|
||||
fs::create_dir_all(&child_path).await.expect("bucket child should be created");
|
||||
|
||||
let mut removed = false;
|
||||
remove_empty_directory_tree_with(
|
||||
&bucket_path,
|
||||
|child| {
|
||||
if !removed && child == child_path {
|
||||
std::fs::remove_dir(&child_path)?;
|
||||
removed = true;
|
||||
}
|
||||
Ok(())
|
||||
},
|
||||
|_| Ok(()),
|
||||
)
|
||||
.await
|
||||
.expect("a vanished empty child should not leave the bucket root behind");
|
||||
|
||||
assert!(removed, "test hook should remove the child before openat");
|
||||
assert!(!bucket_path.exists(), "parent should be rechecked and removed after the child disappears");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn empty_tree_delete_rejects_root_generation_replacement() {
|
||||
let root = tempfile::tempdir().expect("temporary root should be created");
|
||||
let bucket_path = root.path().join("bucket");
|
||||
let moved_path = root.path().join("bucket-before-replacement");
|
||||
fs::create_dir(&bucket_path).await.expect("bucket should be created");
|
||||
|
||||
let mut replaced = false;
|
||||
let err = remove_empty_directory_tree_with(
|
||||
&bucket_path,
|
||||
|_| Ok(()),
|
||||
|directory| {
|
||||
if !replaced && directory == bucket_path {
|
||||
std::fs::rename(&bucket_path, &moved_path)?;
|
||||
std::fs::create_dir(&bucket_path)?;
|
||||
replaced = true;
|
||||
}
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("a replacement root directory must fail closed");
|
||||
|
||||
assert!(replaced, "test hook should replace the root generation before rmdir");
|
||||
assert!(matches!(classify_delete_volume_error(err), DiskError::VolumeNotEmpty));
|
||||
assert!(moved_path.exists(), "the originally scanned root should not be removed by its old name");
|
||||
assert!(bucket_path.exists(), "the replacement root should remain after identity validation fails");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_local_disk_volume_operations() {
|
||||
let test_dir = "./test_local_disk_volumes";
|
||||
|
||||
@@ -25,7 +25,7 @@ use std::{
|
||||
sync::{Arc, LazyLock, Weak},
|
||||
};
|
||||
use tokio::fs;
|
||||
use tokio::sync::{OwnedSemaphorePermit, Semaphore, SemaphorePermit};
|
||||
use tokio::sync::{OwnedSemaphorePermit, RwLock, Semaphore, SemaphorePermit};
|
||||
use tracing::warn;
|
||||
|
||||
/// Check path length according to OS limits.
|
||||
@@ -123,6 +123,8 @@ const TEST_GLOBAL_FILE_SYNCS: usize = 64;
|
||||
|
||||
static FILE_SYNC_PERMITS: LazyLock<Semaphore> = LazyLock::new(|| Semaphore::new(global_file_sync_limit()));
|
||||
static DISK_FILE_SYNC_LIMITERS: LazyLock<Mutex<HashMap<PathBuf, Weak<Semaphore>>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
static DISK_VOLUME_MUTATION_LOCKS: LazyLock<Mutex<HashMap<PathBuf, Weak<RwLock<()>>>>> =
|
||||
LazyLock::new(|| Mutex::new(HashMap::new()));
|
||||
|
||||
fn default_global_file_sync_limit(cpu_count: usize, max_blocking_threads: usize) -> usize {
|
||||
let cpu_scaled = cpu_count
|
||||
@@ -157,6 +159,24 @@ pub(crate) fn disk_file_sync_limiter(root: &Path) -> Arc<Semaphore> {
|
||||
limiter
|
||||
}
|
||||
|
||||
/// Serialize a bucket's local metadata commits with physical bucket removal.
|
||||
///
|
||||
/// The key includes the canonical disk root, so independently reconnected
|
||||
/// [`LocalDisk`](super::local::LocalDisk) instances share the same lock while
|
||||
/// disconnected disks do not keep the registry alive.
|
||||
pub(crate) fn disk_volume_mutation_lock(root: &Path, volume: &str) -> Arc<RwLock<()>> {
|
||||
let key = root.join(volume);
|
||||
let mut locks = DISK_VOLUME_MUTATION_LOCKS.lock();
|
||||
locks.retain(|_, lock| lock.strong_count() > 0);
|
||||
if let Some(lock) = locks.get(&key).and_then(Weak::upgrade) {
|
||||
return lock;
|
||||
}
|
||||
|
||||
let lock = Arc::new(RwLock::new(()));
|
||||
locks.insert(key, Arc::downgrade(&lock));
|
||||
lock
|
||||
}
|
||||
|
||||
/// Always acquire the per-disk permit before the process-wide permit. Keeping
|
||||
/// this order uniform prevents one slow disk from reserving global capacity
|
||||
/// while it waits for its own concurrency slot.
|
||||
@@ -572,15 +592,14 @@ async fn reliable_rename_inner(
|
||||
base_dir: impl AsRef<Path>,
|
||||
warn_on_missing_source: bool,
|
||||
) -> io::Result<()> {
|
||||
if let Some(parent) = dst_file_path.as_ref().parent()
|
||||
&& !file_exists(parent)
|
||||
{
|
||||
reliable_mkdir_all(parent, base_dir.as_ref()).await?;
|
||||
}
|
||||
let parent_guard = match dst_file_path.as_ref().parent() {
|
||||
Some(parent) => Some(mkdir_all_below_existing_base(parent, base_dir.as_ref()).await?),
|
||||
None => None,
|
||||
};
|
||||
|
||||
let mut i = 0;
|
||||
loop {
|
||||
if let Err(e) = super::fs::rename_std(src_file_path.as_ref(), dst_file_path.as_ref()) {
|
||||
if let Err(e) = rename_into_existing_parent(src_file_path.as_ref(), dst_file_path.as_ref(), parent_guard.as_ref()) {
|
||||
if should_retry_rename(&e, i) {
|
||||
i += 1;
|
||||
continue;
|
||||
@@ -597,6 +616,159 @@ async fn reliable_rename_inner(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
fn rename_into_existing_parent(
|
||||
src_file_path: &Path,
|
||||
dst_file_path: &Path,
|
||||
parent_guard: Option<&ExistingBaseDirectoryGuard>,
|
||||
) -> io::Result<()> {
|
||||
use rustix::fs::{Mode, OFlags, open, renameat};
|
||||
|
||||
let Some(parent_guard) = parent_guard else {
|
||||
return super::fs::rename_std(src_file_path, dst_file_path);
|
||||
};
|
||||
let src_parent = src_file_path
|
||||
.parent()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "rename source must have a parent directory"))?;
|
||||
let src_name = src_file_path
|
||||
.file_name()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "rename source must have a file name"))?;
|
||||
let dst_name = dst_file_path
|
||||
.file_name()
|
||||
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "rename destination must have a file name"))?;
|
||||
let src_parent = open(
|
||||
src_parent,
|
||||
OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
|
||||
Mode::empty(),
|
||||
)
|
||||
.map_err(io::Error::from)?;
|
||||
let dst_parent = parent_guard
|
||||
.last()
|
||||
.ok_or_else(|| io::Error::other("rename destination parent guard is empty"))?;
|
||||
|
||||
renameat(&src_parent, src_name, dst_parent, dst_name).map_err(io::Error::from)
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
fn rename_into_existing_parent(
|
||||
src_file_path: &Path,
|
||||
dst_file_path: &Path,
|
||||
_parent_guard: Option<&ExistingBaseDirectoryGuard>,
|
||||
) -> io::Result<()> {
|
||||
super::fs::rename_std(src_file_path, dst_file_path)
|
||||
}
|
||||
|
||||
async fn mkdir_all_below_existing_base(dir_path: &Path, base_dir: &Path) -> io::Result<ExistingBaseDirectoryGuard> {
|
||||
let dir_path = dir_path.to_path_buf();
|
||||
let base_dir = base_dir.to_path_buf();
|
||||
|
||||
tokio::task::spawn_blocking(move || mkdir_all_below_existing_base_std(&dir_path, &base_dir)).await?
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
pub(crate) type ExistingBaseDirectoryGuard = Vec<winapi_util::Handle>;
|
||||
|
||||
#[cfg(unix)]
|
||||
pub(crate) type ExistingBaseDirectoryGuard = Vec<std::os::fd::OwnedFd>;
|
||||
|
||||
#[cfg(all(not(unix), not(windows)))]
|
||||
pub(crate) type ExistingBaseDirectoryGuard = ();
|
||||
|
||||
#[cfg(windows)]
|
||||
fn lock_windows_directory(path: &Path) -> io::Result<winapi_util::Handle> {
|
||||
use std::os::windows::fs::OpenOptionsExt;
|
||||
|
||||
const FILE_ATTRIBUTE_DIRECTORY: u64 = 0x10;
|
||||
const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x400;
|
||||
const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
|
||||
const FILE_FLAG_OPEN_REPARSE_POINT: u32 = 0x0020_0000;
|
||||
const FILE_SHARE_READ: u32 = 0x1;
|
||||
|
||||
let file = std::fs::OpenOptions::new()
|
||||
.read(true)
|
||||
.share_mode(FILE_SHARE_READ)
|
||||
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OPEN_REPARSE_POINT)
|
||||
.open(path)?;
|
||||
let handle = winapi_util::Handle::from_file(file);
|
||||
let info = winapi_util::file::information(&handle)?;
|
||||
if info.file_attributes() & FILE_ATTRIBUTE_DIRECTORY == 0
|
||||
|| info.file_attributes() & u64::from(FILE_ATTRIBUTE_REPARSE_POINT) != 0
|
||||
{
|
||||
return Err(io::Error::from(io::ErrorKind::NotADirectory));
|
||||
}
|
||||
Ok(handle)
|
||||
}
|
||||
|
||||
pub(crate) fn mkdir_all_below_existing_base_std(dir_path: &Path, base_dir: &Path) -> io::Result<ExistingBaseDirectoryGuard> {
|
||||
let relative = dir_path
|
||||
.strip_prefix(base_dir)
|
||||
.map_err(|_| io::Error::new(io::ErrorKind::InvalidInput, "rename destination must remain below its base directory"))?;
|
||||
for component in relative.components() {
|
||||
if !matches!(component, Component::Normal(_) | Component::CurDir) {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidInput,
|
||||
"rename destination contains an invalid path component",
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use rustix::fs::{Mode, OFlags, mkdirat, open, openat};
|
||||
use rustix::io::Errno;
|
||||
|
||||
let flags = OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC;
|
||||
let mode = Mode::RWXU | Mode::RWXG | Mode::RWXO;
|
||||
let mut parents = vec![open(base_dir, flags, Mode::empty()).map_err(io::Error::from)?];
|
||||
|
||||
for component in relative.components() {
|
||||
let Component::Normal(component) = component else {
|
||||
continue;
|
||||
};
|
||||
let parent = parents
|
||||
.last()
|
||||
.expect("base directory guard should contain the base directory");
|
||||
match mkdirat(parent, component, mode) {
|
||||
Ok(()) => {}
|
||||
Err(Errno::EXIST) => {}
|
||||
Err(err) => return Err(err.into()),
|
||||
}
|
||||
parents.push(openat(parent, component, flags, Mode::empty()).map_err(io::Error::from)?);
|
||||
}
|
||||
|
||||
Ok(parents)
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
{
|
||||
let mut handles = vec![lock_windows_directory(base_dir)?];
|
||||
let mut current = base_dir.to_path_buf();
|
||||
for component in relative.components() {
|
||||
let Component::Normal(component) = component else {
|
||||
continue;
|
||||
};
|
||||
current.push(component);
|
||||
match std::fs::create_dir(¤t) {
|
||||
Ok(()) => {}
|
||||
Err(err) if err.kind() == io::ErrorKind::AlreadyExists => {}
|
||||
Err(err) => return Err(err),
|
||||
}
|
||||
handles.push(lock_windows_directory(¤t)?);
|
||||
}
|
||||
|
||||
Ok(handles)
|
||||
}
|
||||
|
||||
#[cfg(all(not(unix), not(windows)))]
|
||||
{
|
||||
let _ = relative;
|
||||
Err(io::Error::new(
|
||||
io::ErrorKind::Unsupported,
|
||||
"safe recursive directory creation is unavailable on this platform",
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
fn warn_reliable_rename_failure(src_file_path: &Path, dst_file_path: &Path, base_dir: &Path, err: &io::Error) {
|
||||
warn!(
|
||||
"reliable_rename failed. src_file_path: {:?}, dst_file_path: {:?}, base_dir: {:?}, err: {:?}",
|
||||
@@ -727,6 +899,20 @@ mod tests {
|
||||
Arc::new(Semaphore::new(MAX_PARALLEL_FILE_SYNCS))
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn disk_volume_mutation_lock_is_shared_per_root_and_volume() {
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
let first = disk_volume_mutation_lock(temp_dir.path(), "bucket");
|
||||
let second = disk_volume_mutation_lock(temp_dir.path(), "bucket");
|
||||
let other = disk_volume_mutation_lock(temp_dir.path(), "other-bucket");
|
||||
|
||||
assert!(Arc::ptr_eq(&first, &second), "reconnected disks must share a bucket mutation lock");
|
||||
assert!(!Arc::ptr_eq(&first, &other), "different buckets must not serialize each other");
|
||||
|
||||
let _write_guard = first.write().await;
|
||||
assert!(second.try_read().is_err(), "a bucket delete lock must exclude local commits");
|
||||
}
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct CapturedLogs {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
@@ -968,6 +1154,112 @@ mod tests {
|
||||
assert_eq!(std::fs::read(dst.join("nested").join("part.1")).expect("read moved part"), b"payload");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rename_all_does_not_recreate_missing_base_directory() {
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
let base = temp_dir.path().join("bucket");
|
||||
std::fs::create_dir(&base).expect("create destination base");
|
||||
let src = temp_dir.path().join("staged-object");
|
||||
std::fs::write(&src, b"payload").expect("write staged object");
|
||||
let dst = base.join("object").join("xl.meta");
|
||||
std::fs::remove_dir(&base).expect("delete destination base before commit");
|
||||
|
||||
let err = rename_all(&src, &dst, &base)
|
||||
.await
|
||||
.expect_err("rename must not recreate a deleted destination base");
|
||||
|
||||
assert!(matches!(err, DiskError::FileNotFound));
|
||||
assert!(src.exists(), "failed commit must preserve the staged source");
|
||||
assert!(!base.exists(), "failed commit must not recreate the deleted bucket");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn rename_all_rejects_a_replaced_base_with_an_existing_parent() {
|
||||
use std::os::unix::fs::symlink;
|
||||
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
let base = temp_dir.path().join("bucket");
|
||||
let outside = temp_dir.path().join("outside");
|
||||
std::fs::create_dir(&base).expect("create destination base");
|
||||
std::fs::create_dir_all(outside.join("object")).expect("create outside destination parent");
|
||||
let src = temp_dir.path().join("staged-object");
|
||||
std::fs::write(&src, b"payload").expect("write staged object");
|
||||
let dst = base.join("object").join("xl.meta");
|
||||
|
||||
std::fs::remove_dir(&base).expect("remove destination base before replacement");
|
||||
symlink(&outside, &base).expect("replace destination base with a symlink");
|
||||
|
||||
rename_all(&src, &dst, &base)
|
||||
.await
|
||||
.expect_err("rename must reject an existing destination parent below a replaced base");
|
||||
|
||||
assert!(src.exists(), "rejected rename must preserve the staged source");
|
||||
assert!(
|
||||
!outside.join("object/xl.meta").exists(),
|
||||
"rename must not publish through the replacement symlink"
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[test]
|
||||
fn windows_parent_guard_blocks_base_and_intermediate_replacement() {
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
let base = temp_dir.path().join("bucket");
|
||||
std::fs::create_dir(&base).expect("create destination base");
|
||||
let parent = base.join("object").join("nested");
|
||||
let guard = mkdir_all_below_existing_base_std(&parent, &base).expect("create and lock destination parents");
|
||||
|
||||
std::fs::rename(&base, temp_dir.path().join("replacement-base"))
|
||||
.expect_err("the locked base must not be replaceable before commit");
|
||||
std::fs::rename(base.join("object"), base.join("replacement-object"))
|
||||
.expect_err("a locked intermediate directory must not be replaceable before commit");
|
||||
|
||||
drop(guard);
|
||||
std::fs::rename(base.join("object"), base.join("replacement-object"))
|
||||
.expect("replacement should succeed after the commit guard is released");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn rename_parent_creation_rejects_symlinked_base() {
|
||||
use std::os::unix::fs::symlink;
|
||||
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
let outside = temp_dir.path().join("outside");
|
||||
std::fs::create_dir(&outside).expect("create outside directory");
|
||||
let base = temp_dir.path().join("bucket");
|
||||
symlink(&outside, &base).expect("create symlinked base");
|
||||
|
||||
mkdir_all_below_existing_base(&base.join("object"), &base)
|
||||
.await
|
||||
.expect_err("symlinked base must be rejected");
|
||||
|
||||
assert!(!outside.join("object").exists(), "parent creation must remain confined to the base");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn rename_parent_creation_rejects_symlink_below_base() {
|
||||
use std::os::unix::fs::symlink;
|
||||
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
let base = temp_dir.path().join("bucket");
|
||||
let outside = temp_dir.path().join("outside");
|
||||
std::fs::create_dir(&base).expect("create destination base");
|
||||
std::fs::create_dir(&outside).expect("create outside directory");
|
||||
symlink(&outside, base.join("linked")).expect("create symlink below base");
|
||||
|
||||
mkdir_all_below_existing_base(&base.join("linked/object"), &base)
|
||||
.await
|
||||
.expect_err("symlink below base must be rejected");
|
||||
|
||||
assert!(
|
||||
!outside.join("object").exists(),
|
||||
"parent creation must not follow a symlink outside the base"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn fsync_dir_succeeds_on_directory() {
|
||||
let temp_dir = tempdir().expect("create temp dir");
|
||||
|
||||
@@ -568,6 +568,7 @@ mod tests {
|
||||
};
|
||||
use crate::bucket::metadata::table_bucket_catalog_metadata_prefix;
|
||||
use crate::bucket::metadata_sys;
|
||||
use crate::cluster::rpc::peer_s3_client::install_delete_bucket_empty_scan_barrier;
|
||||
use crate::disk::{BUCKET_META_PREFIX, RUSTFS_META_BUCKET};
|
||||
use crate::error::StorageError;
|
||||
use crate::object_api::{ObjectOptions, PutObjReader};
|
||||
@@ -589,6 +590,7 @@ mod tests {
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::time::{Duration, SystemTime};
|
||||
use time::OffsetDateTime;
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::{Notify, OnceCell};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use uuid::Uuid;
|
||||
@@ -1048,7 +1050,7 @@ mod tests {
|
||||
.await
|
||||
.expect("bucket should be created");
|
||||
for disk_path in &disk_paths {
|
||||
tokio::fs::create_dir_all(disk_path.join(&bucket).join("empty-directory"))
|
||||
tokio::fs::create_dir_all(disk_path.join(&bucket).join("empty-directory/nested/leaf"))
|
||||
.await
|
||||
.expect("empty directory remnant should be created");
|
||||
}
|
||||
@@ -1257,6 +1259,55 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn bucket_delete_preserves_put_committed_after_empty_scan() {
|
||||
let (_, ecstore) = setup_bucket_delete_test_env().await;
|
||||
let bucket = format!("bucket-delete-empty-scan-race-{}", Uuid::new_v4().simple());
|
||||
let object = "committed-after-empty-scan";
|
||||
let payload = b"object committed after DeleteBucket empty scan".to_vec();
|
||||
|
||||
ecstore
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("bucket should be created");
|
||||
|
||||
let barrier = install_delete_bucket_empty_scan_barrier();
|
||||
let delete_store = ecstore.clone();
|
||||
let delete_bucket = bucket.clone();
|
||||
let delete = tokio::spawn(async move {
|
||||
delete_store
|
||||
.delete_bucket(&delete_bucket, &DeleteBucketOptions::default())
|
||||
.await
|
||||
});
|
||||
barrier.wait_until_paused().await;
|
||||
|
||||
let mut put_reader = PutObjReader::from_vec(payload.clone());
|
||||
ecstore
|
||||
.put_object(&bucket, object, &mut put_reader, &ObjectOptions::default())
|
||||
.await
|
||||
.expect("PUT should commit while DeleteBucket is paused after its empty scan");
|
||||
|
||||
barrier.release();
|
||||
let err = delete
|
||||
.await
|
||||
.expect("DeleteBucket task should join")
|
||||
.expect_err("DeleteBucket must reject a PUT committed after its empty scan");
|
||||
assert!(matches!(err, StorageError::BucketNotEmpty(name) if name == bucket));
|
||||
|
||||
let mut reader = ecstore
|
||||
.get_object_reader(&bucket, object, None, http::HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("committed object should remain readable after DeleteBucket fails");
|
||||
let mut restored = Vec::new();
|
||||
reader
|
||||
.stream
|
||||
.read_to_end(&mut restored)
|
||||
.await
|
||||
.expect("object body should remain readable");
|
||||
assert_eq!(restored, payload);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn bucket_recreation_does_not_publish_unverified_usage() {
|
||||
|
||||
@@ -6895,6 +6895,7 @@ mod test {
|
||||
)
|
||||
.await
|
||||
.expect("local disk should be created");
|
||||
disk.make_volume("bucket").await.expect("test bucket should be created");
|
||||
let mut fi = FileInfo::new("object", 1, 1);
|
||||
fi.volume = "bucket".to_owned();
|
||||
fi.name = "object".to_owned();
|
||||
|
||||
@@ -53,7 +53,7 @@ pub struct DeleteBucketOptions {
|
||||
pub no_recreate: bool,
|
||||
/// Force deletion even if bucket is not empty.
|
||||
pub force: bool,
|
||||
/// Force deletion only after the local peer verifies the bucket is empty.
|
||||
/// Delete empty directory remnants after the local peer verifies no object metadata exists.
|
||||
pub force_if_empty: bool,
|
||||
/// Site replication delete operation.
|
||||
pub srdelete_op: SRBucketDeleteOp,
|
||||
|
||||
Reference in New Issue
Block a user