mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-28 00:58:59 +00:00
fix(ecstore): suppress missing rollback rename warnings (#4792)
This commit is contained in:
@@ -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<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
struct CapturedLogWriter {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
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<usize> {
|
||||
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<Uuid>, data: Option<Bytes>) -> 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;
|
||||
|
||||
@@ -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<Path>,
|
||||
dst_file_path: impl AsRef<Path>,
|
||||
base_dir: impl AsRef<Path>,
|
||||
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}"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user