fix(ecstore): fail listing on walk_dir producer errors (#2937)

This commit is contained in:
Henry Guo
2026-05-13 19:08:05 +08:00
committed by GitHub
parent fd3fb77ed5
commit 26341742e0
3 changed files with 96 additions and 57 deletions
+94 -55
View File
@@ -18,7 +18,12 @@ use crate::disk::{self, DiskAPI, DiskStore, WalkDirOptions};
use futures::future::join_all;
use metrics::counter;
use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetacacheReader, is_io_eof};
use std::{future::Future, pin::Pin, time::Duration};
use std::{
future::Future,
pin::Pin,
sync::{Arc, Mutex},
time::Duration,
};
use tokio::io::AsyncRead;
use tokio::spawn;
use tokio::time::timeout;
@@ -46,10 +51,12 @@ async fn peek_with_timeout<R: AsyncRead + Unpin>(reader: &mut MetacacheReader<R>
}
#[cfg(test)]
#[derive(Clone, Copy)]
#[derive(Clone)]
pub(crate) enum TestReaderBehavior {
Eof,
Stall,
ProducerError(DiskError),
PartialThenTimeout(Vec<MetaCacheEntry>),
}
#[derive(Default)]
@@ -107,21 +114,21 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
let mut readers = Vec::with_capacity(opts.disks.len());
let fds = opts.fallback_disks.iter().flatten().cloned().collect::<Vec<_>>();
let max_disk_failures = opts.disks.len().saturating_sub(opts.min_disks);
let producer_errs = Arc::new(Mutex::new(vec![None; opts.disks.len()]));
let cancel_rx = CancellationToken::new();
for (disk_idx, disk) in opts.disks.iter().enumerate() {
#[cfg(not(test))]
let _ = disk_idx;
let opdisk = disk.clone();
let opts_clone = opts.clone();
let mut fds_clone = fds.clone();
let cancel_rx_clone = cancel_rx.clone();
let producer_errs_clone = producer_errs.clone();
let (rd, wr) = tokio::io::duplex(64);
readers.push(MetacacheReader::new(rd));
jobs.push(spawn(async move {
#[cfg(test)]
if let Some(behavior) = opts_clone.test_reader_behaviors.get(disk_idx).copied() {
if let Some(behavior) = opts_clone.test_reader_behaviors.get(disk_idx).cloned() {
match behavior {
TestReaderBehavior::Eof => return Ok(()),
TestReaderBehavior::Stall => {
@@ -129,6 +136,19 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
cancel_rx_clone.cancelled().await;
return Ok(());
}
TestReaderBehavior::ProducerError(err) => {
producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone());
return Err(err);
}
TestReaderBehavior::PartialThenTimeout(entries) => {
let mut wr = wr;
let mut out = rustfs_filemeta::MetacacheWriter::new(&mut wr);
let err = DiskError::Timeout;
producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone());
let _ = out.write(&entries).await;
drop(out);
return Err(err);
}
}
}
@@ -166,18 +186,20 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
}
while need_fallback {
let disk_op = {
if fds_clone.is_empty() {
None
} else {
let disk = fds_clone.remove(0);
if disk.is_online().await { Some(disk.clone()) } else { None }
let mut disk_op = None;
while !fds_clone.is_empty() {
let disk = fds_clone.remove(0);
if disk.is_online().await {
disk_op = Some(disk);
break;
}
};
}
let Some(disk) = disk_op else {
warn!("list_path_raw: fallback disk is none");
break;
let err = last_err.unwrap_or(DiskError::DiskNotFound);
producer_errs_clone.lock().expect("producer error mutex poisoned")[disk_idx] = Some(err.clone());
return Err(err);
};
match disk
@@ -204,7 +226,6 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
Err(err) => {
error!("walk dir2 err {:?}", &err);
last_err = Some(err);
break;
}
}
}
@@ -260,6 +281,11 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
// info!("read entry disk: {}, name: {}", i, entry.name);
entry
} else {
if let Some(err) = producer_errs.lock().expect("producer error mutex poisoned")[i].clone() {
has_err += 1;
errs[i] = Some(err);
continue;
}
// eof
at_eof += 1;
// warn!("list_path_raw: peek eof, disk: {}", i);
@@ -267,6 +293,12 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
}
}
PeekOutcome::Error(err) => {
if let Some(err) = producer_errs.lock().expect("producer error mutex poisoned")[i].clone() {
has_err += 1;
errs[i] = Some(err);
continue;
}
if err == rustfs_filemeta::Error::Unexpected {
at_eof += 1;
// warn!("list_path_raw: peek err eof, disk: {}", i);
@@ -376,6 +408,15 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
if let Some(finished_fn) = opts.finished.as_ref() {
finished_fn(&errs).await;
}
if errs.iter().flatten().any(|err| *err == DiskError::Timeout) {
return Err(DiskError::Timeout);
}
let mut err_iter = errs.iter().flatten();
if let Some(err) = err_iter.next()
&& err_iter.next().is_none()
{
return Err(err.clone());
}
let mut combined_err = Vec::new();
errs.iter().zip(opts.disks.iter()).for_each(|(err, disk)| match (err, disk) {
(Some(err), Some(disk)) => {
@@ -451,7 +492,7 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
match result {
Ok(Ok(())) => {}
Ok(Err(err)) => {
error!("list_path_raw err {:?}", err);
error!("list_path_raw producer err {:?}", err);
job_errs.push(err);
}
Err(err) => {
@@ -472,8 +513,6 @@ pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> d
#[cfg(test)]
mod tests {
use super::*;
use crate::disk::endpoint::Endpoint;
use crate::disk::{DiskOption, new_disk};
use rustfs_filemeta::MetacacheWriter;
#[tokio::test]
@@ -503,6 +542,38 @@ mod tests {
assert_eq!(err, DiskError::Timeout);
}
#[tokio::test]
async fn list_path_raw_returns_timeout_when_producer_fails_after_partial_entry() {
let seen = Arc::new(Mutex::new(Vec::new()));
let seen_clone = seen.clone();
let err = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None],
min_disks: 1,
test_reader_behaviors: vec![TestReaderBehavior::PartialThenTimeout(vec![MetaCacheEntry {
name: "bucket/object".to_string(),
metadata: vec![1, 2, 3],
cached: None,
reusable: false,
}])],
agreed: Some(Box::new(move |entry: MetaCacheEntry| {
let seen = seen_clone.clone();
Box::pin(async move {
seen.lock().expect("seen mutex poisoned").push(entry.name);
})
})),
..Default::default()
},
)
.await
.expect_err("producer timeout after partial output must fail the listing");
assert_eq!(err, DiskError::Timeout);
assert_eq!(seen.lock().expect("seen mutex poisoned").as_slice(), &["bucket/object".to_string()]);
}
#[tokio::test]
async fn peek_with_timeout_times_out_on_silent_reader() {
let (_writer, reader) = tokio::io::duplex(64);
@@ -536,52 +607,20 @@ mod tests {
}
}
#[cfg(unix)]
#[tokio::test]
async fn list_path_raw_propagates_parent_scan_access_denied() {
use std::fs::Permissions;
use std::os::unix::fs::PermissionsExt;
use tempfile::tempdir;
use tokio::fs;
let dir = tempdir().expect("tempdir should be created");
let bucket = "test-bucket";
let parent = dir.path().join(bucket).join("shiplog");
fs::create_dir_all(parent.join("nano/a.txt"))
.await
.expect("test object directory should be created");
fs::write(parent.join("nano/a.txt/xl.meta"), b"meta")
.await
.expect("test metadata should be written");
std::fs::set_permissions(&parent, Permissions::from_mode(0o111)).expect("parent permissions should be changed");
if fs::read_dir(&parent).await.is_ok() {
std::fs::set_permissions(&parent, Permissions::from_mode(0o755)).expect("parent permissions should be restored");
return;
}
let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp path should be valid utf-8"))
.expect("endpoint should be created");
let disk = new_disk(&endpoint, &DiskOption::default())
.await
.expect("local disk should be created");
let result = list_path_raw(
async fn list_path_raw_propagates_producer_access_denied() {
let err = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![Some(disk)],
bucket: bucket.to_string(),
path: "shiplog/".to_string(),
disks: vec![None],
min_disks: 1,
test_reader_behaviors: vec![TestReaderBehavior::ProducerError(DiskError::FileAccessDenied)],
..Default::default()
},
)
.await;
.await
.expect_err("producer access failure must not be treated as an empty listing");
std::fs::set_permissions(&parent, Permissions::from_mode(0o755)).expect("parent permissions should be restored");
let err = result.expect_err("parent directory access failure must not be treated as an empty listing");
assert_eq!(err, DiskError::FileAccessDenied);
}
}
+1 -1
View File
@@ -863,7 +863,7 @@ impl LocalDiskWrapper {
timeout_ms = timeout_duration.as_millis(),
"Local disk operation timed out"
);
Err(DiskError::other(format!("disk operation timeout after {timeout_duration:?}")))
Err(DiskError::Timeout)
}
}
}
+1 -1
View File
@@ -385,7 +385,7 @@ impl RemoteDisk {
timeout_ms = timeout_duration.as_millis(),
"Remote disk operation timed out"
);
Err(Error::other(format!("Remote disk operation timeout after {timeout_duration:?}")))
Err(DiskError::Timeout)
}
}
}