diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 71b817879..255bd6a93 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -4755,9 +4755,7 @@ impl LocalDisk { .await); } let rollback_data_path = rollback_path.join(dir.to_string()); - if let Err(err) = rename_all(&dir_path, &rollback_data_path, &rollback_path).await - && !(err == DiskError::FileNotFound || err == DiskError::VolumeNotFound) - { + if 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, @@ -7681,10 +7679,7 @@ impl DiskAPI for LocalDisk { .await); } let rollback_data_path = rollback_path.join(uuid.to_string()); - if let Err(err) = rename_all(&old_path, &rollback_data_path, &rollback_path).await - && err != DiskError::FileNotFound - && err != DiskError::VolumeNotFound - { + if 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, @@ -7962,10 +7957,56 @@ async fn get_disk_info(drive_path: PathBuf) -> Result<(rustfs_utils::os::DiskInf mod test { use super::*; use rustfs_filemeta::ErasureInfo; - use std::io; + use std::io::{self, Write}; use std::pin::Pin; + use std::sync::{Arc, Mutex}; use std::task::{Context, Poll}; use tokio::io::{AsyncReadExt, AsyncWrite, AsyncWriteExt, ReadBuf}; + use tracing_subscriber::fmt::MakeWriter; + + #[derive(Clone, Default)] + struct CapturedLogs { + buffer: Arc>>, + } + + struct CapturedLogWriter { + buffer: Arc>>, + } + + impl CapturedLogs { + fn contents(&self) -> String { + let buffer = self + .buffer + .lock() + .expect("captured logs mutex should not be poisoned") + .clone(); + String::from_utf8(buffer).expect("captured logs should be valid UTF-8") + } + } + + impl Write for CapturedLogWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.buffer + .lock() + .expect("captured logs mutex should not be poisoned") + .extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + impl<'a> MakeWriter<'a> for CapturedLogs { + type Writer = CapturedLogWriter; + + fn make_writer(&'a self) -> Self::Writer { + CapturedLogWriter { + buffer: Arc::clone(&self.buffer), + } + } + } fn test_file_info(name: &str, version_id: Uuid, data_dir: Option, data: Option) -> FileInfo { let size = data @@ -10376,6 +10417,149 @@ mod test { ); } + #[tokio::test(flavor = "current_thread")] + async fn test_delete_version_missing_inline_data_dir_does_not_warn() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "dir/inline-object"; + let version_id = Uuid::parse_str("31313131-1111-2222-3333-444444444444").expect("version id should parse"); + let data_dir = Uuid::parse_str("32323232-1111-2222-3333-444444444444").expect("data dir should parse"); + let rollback_dir = Uuid::parse_str("33333333-1111-2222-3333-444444444444").expect("rollback dir should parse"); + + ensure_test_volume(&disk, bucket).await; + + let object_dir = dir.path().join(bucket).join(object); + fs::create_dir_all(&object_dir).await.expect("object dir should be created"); + let fi = test_file_info(object, version_id, Some(data_dir), Some(Bytes::from_static(b"inline"))); + fs::write(object_dir.join(STORAGE_FORMAT_FILE), test_meta(fi.clone())) + .await + .expect("inline metadata should be written"); + assert!(!object_dir.join(data_dir.to_string()).exists()); + + let logs = CapturedLogs::default(); + let subscriber = tracing_subscriber::fmt() + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .finish(); + let _guard = tracing::subscriber::set_default(subscriber); + + disk.delete_version( + bucket, + object, + fi, + false, + DeleteOptions { + old_data_dir: Some(rollback_dir), + ..Default::default() + }, + ) + .await + .expect("missing inline data dir should not fail deletion"); + + assert!(!object_dir.join(STORAGE_FORMAT_FILE).exists()); + assert!( + object_dir + .join(rollback_dir.to_string()) + .join(STORAGE_FORMAT_FILE_BACKUP) + .exists() + ); + assert!(!logs.contents().contains("reliable_rename failed")); + } + + #[tokio::test(flavor = "current_thread")] + async fn test_delete_versions_missing_inline_data_dir_does_not_warn() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + + let bucket = "bucket"; + let object = "dir/inline-object-batch"; + let version_id = Uuid::parse_str("34343434-1111-2222-3333-444444444444").expect("version id should parse"); + let data_dir = Uuid::parse_str("35353535-1111-2222-3333-444444444444").expect("data dir should parse"); + let rollback_dir = Uuid::parse_str("36363636-1111-2222-3333-444444444444").expect("rollback dir should parse"); + + ensure_test_volume(&disk, bucket).await; + + let object_dir = dir.path().join(bucket).join(object); + fs::create_dir_all(&object_dir).await.expect("object dir should be created"); + let fi = test_file_info(object, version_id, Some(data_dir), Some(Bytes::from_static(b"inline"))); + fs::write(object_dir.join(STORAGE_FORMAT_FILE), test_meta(fi.clone())) + .await + .expect("inline metadata should be written"); + assert!(!object_dir.join(data_dir.to_string()).exists()); + + let logs = CapturedLogs::default(); + let subscriber = tracing_subscriber::fmt() + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .finish(); + let _guard = tracing::subscriber::set_default(subscriber); + + disk.delete_versions_internal( + bucket, + object, + &[fi], + &DeleteOptions { + old_data_dir: Some(rollback_dir), + ..Default::default() + }, + ) + .await + .expect("missing inline data dir should not fail batch deletion"); + + assert!(!object_dir.join(STORAGE_FORMAT_FILE).exists()); + assert!( + object_dir + .join(rollback_dir.to_string()) + .join(STORAGE_FORMAT_FILE_BACKUP) + .exists() + ); + assert!(!logs.contents().contains("reliable_rename failed")); + } + + #[tokio::test(flavor = "current_thread")] + async fn test_ignore_missing_source_helper_still_warns_on_real_error() { + use tempfile::tempdir; + + let dir = tempdir().expect("temp dir should be created"); + let src = dir.path().join("src"); + let dst = dir.path().join("dst"); + fs::create_dir_all(&src).await.expect("source dir should be created"); + fs::write(src.join("part.1"), b"live-data") + .await + .expect("source data should be written"); + fs::create_dir_all(&dst).await.expect("destination dir should be created"); + fs::write(dst.join("sentinel"), b"conflict") + .await + .expect("destination conflict should be written"); + + let logs = CapturedLogs::default(); + let subscriber = tracing_subscriber::fmt() + .with_writer(logs.clone()) + .with_ansi(false) + .without_time() + .finish(); + let _guard = tracing::subscriber::set_default(subscriber); + + let err = rename_all_ignore_missing_source(&src, &dst, dir.path()) + .await + .expect_err("a non-empty rollback destination must reject rename"); + + assert_ne!(err, DiskError::FileNotFound); + assert!(src.exists()); + assert!(dst.join("sentinel").exists()); + assert!(logs.contents().contains("reliable_rename failed")); + } + #[tokio::test] async fn test_delete_version_rollback_restores_staged_data_dir() { use tempfile::tempdir; diff --git a/crates/ecstore/src/disk/os.rs b/crates/ecstore/src/disk/os.rs index 7acb84e9d..bc5f430fa 100644 --- a/crates/ecstore/src/disk/os.rs +++ b/crates/ecstore/src/disk/os.rs @@ -26,7 +26,7 @@ use std::{ }; use tokio::fs; use tokio::sync::{OwnedSemaphorePermit, Semaphore, SemaphorePermit}; -use tracing::{debug, warn}; +use tracing::warn; /// Check path length according to OS limits. pub fn check_path_length(path_name: &str) -> Result<()> { @@ -527,7 +527,7 @@ async fn reliable_rename_inner( src_file_path: impl AsRef, dst_file_path: impl AsRef, base_dir: impl AsRef, - warn_on_failure: bool, + warn_on_missing_source: bool, ) -> io::Result<()> { if let Some(parent) = dst_file_path.as_ref().parent() && !file_exists(parent) @@ -542,32 +542,8 @@ async fn reliable_rename_inner( i += 1; continue; } - if warn_on_failure { - // NotFound is an expected outcome for speculative renames (staging - // a dataDir that inline objects never had, trashing an already- - // removed tmp path, restoring an absent rollback backup): those - // callers treat FileNotFound as acceptable, and on per-request - // delete paths a WARN becomes a log storm (one line per disk per - // DELETE). Keep it at debug so quorum-masked cases (e.g. a tmp - // file lost before a commit rename) stay diagnosable; genuine - // failures keep the WARN. The error is returned either way. - if e.kind() == io::ErrorKind::NotFound { - debug!( - "reliable_rename failed. src_file_path: {:?}, dst_file_path: {:?}, base_dir: {:?}, err: {:?}", - src_file_path.as_ref(), - dst_file_path.as_ref(), - base_dir.as_ref(), - e - ); - } else { - warn!( - "reliable_rename failed. src_file_path: {:?}, dst_file_path: {:?}, base_dir: {:?}, err: {:?}", - src_file_path.as_ref(), - dst_file_path.as_ref(), - base_dir.as_ref(), - e - ); - } + if warn_on_missing_source || e.kind() != io::ErrorKind::NotFound { + warn_reliable_rename_failure(src_file_path.as_ref(), dst_file_path.as_ref(), base_dir.as_ref(), &e); } return Err(e); } @@ -578,6 +554,13 @@ async fn reliable_rename_inner( Ok(()) } +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: {:?}", + src_file_path, dst_file_path, base_dir, err + ); +} + /// Whether a failed `rename` in [`reliable_rename_inner`] should be retried. /// /// Only the first failure is retried, and `NotFound` is never retried: the @@ -755,24 +738,22 @@ mod tests { } #[tokio::test] - async fn rename_all_ignore_missing_source_returns_ok() { + async fn rename_all_ignore_missing_source_returns_ok_without_warn() { let temp_dir = tempdir().expect("create temp dir"); let src = temp_dir.path().join("missing"); let dst = temp_dir.path().join("dst"); + let (logs, _guard) = warn_capture(); rename_all_ignore_missing_source(&src, &dst, temp_dir.path()) .await .expect("missing cleanup source must be ignored"); assert!(!dst.exists()); + assert!(!logs.contents().contains("reliable_rename failed")); } #[tokio::test] - async fn rename_all_missing_source_skips_rename_warn() { - // Callers like delete_version's rollback staging accept FileNotFound - // (inline objects have no dataDir), so a missing source must not emit - // the per-disk "reliable_rename failed" WARN — under delete storms it - // was pure log noise. The error itself is still returned. + async fn rename_all_missing_source_still_warns() { let temp_dir = tempdir().expect("create temp dir"); let src = temp_dir.path().join("missing"); let dst = temp_dir.path().join("dst"); @@ -785,8 +766,8 @@ mod tests { assert!(matches!(err, DiskError::FileNotFound)); let captured = logs.contents(); assert!( - !captured.contains("reliable_rename failed"), - "NotFound must not emit the reliable_rename WARN, got: {captured}" + captured.contains("reliable_rename failed"), + "ordinary missing-source failures must keep the WARN, got: {captured}" ); }