Files
rustfs/crates/ecstore/src/cache_value/metacache_set.rs
T
2026-07-09 04:10:57 +08:00

1374 lines
55 KiB
Rust

// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::disk::disk_store::get_drive_walkdir_stall_timeout;
use crate::disk::error::DiskError;
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::{
collections::{HashSet, VecDeque},
future::Future,
io::ErrorKind,
pin::Pin,
sync::{Arc, OnceLock},
time::Duration,
};
use tokio::io::AsyncRead;
use tokio::spawn;
use tokio::sync::Mutex as TokioMutex;
use tokio::time::timeout;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, warn};
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
const LOG_SUBSYSTEM_METACACHE: &str = "metacache";
const EVENT_METACACHE_LISTING: &str = "metacache_listing";
pub type AgreedFn = Box<dyn Fn(MetaCacheEntry) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
pub type PartialFn =
Box<dyn Fn(MetaCacheEntries, &[Option<DiskError>]) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
type FinishedFn = Box<dyn Fn(&[Option<DiskError>]) -> Pin<Box<dyn Future<Output = ()> + Send>> + Send + 'static>;
#[derive(Clone, Default)]
pub(crate) struct FallbackClaimTracker {
claimed: Arc<TokioMutex<HashSet<String>>>,
}
impl FallbackClaimTracker {
pub(crate) async fn claim_disk(&self, disk: &DiskStore) {
self.claimed.lock().await.insert(disk.endpoint().to_string());
}
pub(crate) async fn claimed_keys(&self) -> HashSet<String> {
self.claimed.lock().await.clone()
}
#[cfg(test)]
pub(crate) async fn claim_test_fallback(&self) {
let mut claimed = self.claimed.lock().await;
let key = format!("test-fallback-{}", claimed.len());
claimed.insert(key);
}
#[cfg(test)]
pub(crate) async fn contains_key(&self, key: &str) -> bool {
self.claimed.lock().await.contains(key)
}
}
#[derive(Debug)]
enum PeekOutcome {
Ready(Option<MetaCacheEntry>),
Error(rustfs_filemeta::Error),
TimedOut,
}
async fn peek_with_timeout<R: AsyncRead + Unpin>(reader: &mut MetacacheReader<R>, timeout_duration: Duration) -> PeekOutcome {
match timeout(timeout_duration, reader.peek()).await {
Ok(Ok(entry)) => PeekOutcome::Ready(entry),
Ok(Err(err)) => PeekOutcome::Error(err),
Err(_) => PeekOutcome::TimedOut,
}
}
fn is_missing_path_error(err: &DiskError) -> bool {
matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound)
}
fn is_tolerated_producer_completion_error(err: &DiskError) -> bool {
matches!(err, DiskError::DiskOngoingReq)
|| err.is_metacache_output_stream_closed()
|| err.contains_io_error_kind(ErrorKind::BrokenPipe)
}
async fn take_fallback_candidate<T>(fallback_items: &Arc<TokioMutex<VecDeque<T>>>) -> Option<T> {
fallback_items.lock().await.pop_front()
}
fn duration_millis(duration: Duration) -> u64 {
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
}
#[cfg(test)]
#[derive(Clone)]
pub(crate) enum TestReaderBehavior {
Eof,
Entries(Vec<MetaCacheEntry>),
Stall,
IgnoreCancel,
ProducerError(DiskError),
PrimaryErrorThenFallback(DiskError),
PartialThenTimeout(Vec<MetaCacheEntry>),
}
#[derive(Default)]
pub struct ListPathRawOptions {
pub disks: Vec<Option<DiskStore>>,
pub fallback_disks: Vec<Option<DiskStore>>,
pub bucket: String,
pub path: String,
pub recursive: bool,
pub incl_deleted: bool,
pub filter_prefix: Option<String>,
pub forward_to: Option<String>,
pub min_disks: usize,
pub report_not_found: bool,
pub per_disk_limit: i32,
pub skip_walkdir_total_timeout: bool,
pub walkdir_timeout: Option<Duration>,
pub walkdir_stall_timeout: Option<Duration>,
pub agreed: Option<AgreedFn>,
pub partial: Option<PartialFn>,
pub finished: Option<FinishedFn>,
#[cfg(test)]
pub(crate) test_reader_behaviors: Vec<TestReaderBehavior>,
#[cfg(test)]
pub(crate) test_fallback_reader_behaviors: Vec<TestReaderBehavior>,
#[cfg(test)]
pub(crate) peek_timeout: Option<Duration>,
// pub agreed: Option<Arc<dyn Fn(MetaCacheEntry) + Send + Sync>>,
// pub partial: Option<Arc<dyn Fn(MetaCacheEntries, &[Option<Error>]) + Send + Sync>>,
// pub finished: Option<Arc<dyn Fn(&[Option<Error>]) + Send + Sync>>,
}
impl Clone for ListPathRawOptions {
fn clone(&self) -> Self {
Self {
disks: self.disks.clone(),
fallback_disks: self.fallback_disks.clone(),
bucket: self.bucket.clone(),
path: self.path.clone(),
recursive: self.recursive,
incl_deleted: self.incl_deleted,
filter_prefix: self.filter_prefix.clone(),
forward_to: self.forward_to.clone(),
min_disks: self.min_disks,
report_not_found: self.report_not_found,
per_disk_limit: self.per_disk_limit,
skip_walkdir_total_timeout: self.skip_walkdir_total_timeout,
walkdir_timeout: self.walkdir_timeout,
walkdir_stall_timeout: self.walkdir_stall_timeout,
#[cfg(test)]
test_reader_behaviors: self.test_reader_behaviors.clone(),
#[cfg(test)]
test_fallback_reader_behaviors: self.test_fallback_reader_behaviors.clone(),
#[cfg(test)]
peek_timeout: self.peek_timeout,
..Default::default()
}
}
}
pub async fn list_path_raw(rx: CancellationToken, opts: ListPathRawOptions) -> disk::error::Result<()> {
list_path_raw_inner(rx, opts, None).await
}
pub(crate) async fn list_path_raw_with_claim_tracker(
rx: CancellationToken,
opts: ListPathRawOptions,
claim_tracker: FallbackClaimTracker,
) -> disk::error::Result<()> {
list_path_raw_inner(rx, opts, Some(claim_tracker)).await
}
async fn list_path_raw_inner(
rx: CancellationToken,
opts: ListPathRawOptions,
fallback_claim_tracker: Option<FallbackClaimTracker>,
) -> disk::error::Result<()> {
if opts.disks.is_empty() {
return Err(DiskError::ErasureReadQuorum);
}
if opts.min_disks > opts.disks.len() {
return Err(DiskError::ErasureReadQuorum);
}
let log_bucket = opts.bucket.clone();
let log_path = opts.path.clone();
let mut jobs: Vec<tokio::task::JoinHandle<std::result::Result<(), DiskError>>> = Vec::new();
let mut readers = Vec::with_capacity(opts.disks.len());
let fds = Arc::new(TokioMutex::new(opts.fallback_disks.iter().flatten().cloned().collect::<VecDeque<_>>()));
#[cfg(test)]
let test_fallbacks = Arc::new(TokioMutex::new(
opts.test_fallback_reader_behaviors.iter().cloned().collect::<VecDeque<_>>(),
));
let max_disk_failures = opts.disks.len().saturating_sub(opts.min_disks);
let producer_errs: Arc<[OnceLock<DiskError>]> = (0..opts.disks.len()).map(|_| OnceLock::new()).collect::<Vec<_>>().into();
let cancel_rx = CancellationToken::new();
for (disk_idx, disk) in opts.disks.iter().enumerate() {
let opdisk = disk.clone();
let opts_clone = opts.clone();
let fallback_claim_tracker = fallback_claim_tracker.clone();
let fds_clone = fds.clone();
#[cfg(test)]
let test_fallbacks_clone = test_fallbacks.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)]
let test_primary_error = if let Some(behavior) = opts_clone.test_reader_behaviors.get(disk_idx).cloned() {
match behavior {
TestReaderBehavior::Eof => return Ok(()),
TestReaderBehavior::Entries(entries) => {
let mut wr = wr;
let mut out = rustfs_filemeta::MetacacheWriter::new(&mut wr);
out.write(&entries).await.expect("test entries should be written");
out.close().await.expect("test entries should close");
return Ok(());
}
TestReaderBehavior::Stall => {
let _held_writer = wr;
cancel_rx_clone.cancelled().await;
return Ok(());
}
TestReaderBehavior::IgnoreCancel => {
let _held_writer = wr;
std::future::pending::<()>().await;
return Ok(());
}
TestReaderBehavior::ProducerError(err) => {
record_producer_error(&producer_errs_clone, disk_idx, &err);
return Err(err);
}
TestReaderBehavior::PrimaryErrorThenFallback(err) => Some(err),
TestReaderBehavior::PartialThenTimeout(entries) => {
let mut wr = wr;
let mut out = rustfs_filemeta::MetacacheWriter::new(&mut wr);
let err = DiskError::Timeout;
record_producer_error(&producer_errs_clone, disk_idx, &err);
let _ = out.write(&entries).await;
drop(out);
return Err(err);
}
}
} else {
None
};
let mut wr = wr;
let wakl_opts = WalkDirOptions {
bucket: opts_clone.bucket.clone(),
base_dir: opts_clone.path.clone(),
recursive: opts_clone.recursive,
incl_deleted: opts_clone.incl_deleted,
report_notfound: opts_clone.report_not_found,
filter_prefix: opts_clone.filter_prefix.clone(),
forward_to: opts_clone.forward_to.clone(),
limit: opts_clone.per_disk_limit,
skip_total_timeout: opts_clone.skip_walkdir_total_timeout,
timeout_ms: opts_clone.walkdir_timeout.map(duration_millis),
stall_timeout_ms: opts_clone.walkdir_stall_timeout.map(duration_millis),
..Default::default()
};
let mut need_fallback = false;
let mut last_err = None;
#[cfg(test)]
if let Some(err) = test_primary_error {
last_err = Some(err);
need_fallback = true;
}
if !need_fallback && let Some(disk) = opdisk {
let primary_walk_started = std::time::Instant::now();
match disk.walk_dir(wakl_opts, &mut wr).await {
Ok(_res) => {
rustfs_io_metrics::record_stage_duration(
"metacache_walk_dir_primary",
primary_walk_started.elapsed().as_secs_f64() * 1000.0,
);
}
Err(err) => {
if err.is_metacache_output_stream_closed() {
rustfs_io_metrics::record_stage_duration(
"metacache_walk_dir_primary",
primary_walk_started.elapsed().as_secs_f64() * 1000.0,
);
return Ok(());
}
rustfs_io_metrics::record_stage_duration(
"metacache_walk_dir_primary_failed",
primary_walk_started.elapsed().as_secs_f64() * 1000.0,
);
if is_missing_path_error(&err) {
debug!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts_clone.bucket,
path = %opts_clone.path,
disk_index = disk_idx,
state = "walk_dir_missing_path",
error = ?err,
"Metacache walk_dir missing path skipped"
);
} else {
warn!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts_clone.bucket,
path = %opts_clone.path,
disk_index = disk_idx,
state = "walk_dir_failed",
error = ?err,
"Metacache walk_dir failed"
);
}
last_err = Some(err);
need_fallback = true;
}
}
} else if !need_fallback {
last_err = Some(DiskError::DiskNotFound);
need_fallback = true;
}
if cancel_rx_clone.is_cancelled() {
// warn!("list_path_raw: cancel_rx_clone.is_cancelled()");
return Ok(());
}
while need_fallback {
#[cfg(test)]
if let Some(behavior) = take_fallback_candidate(&test_fallbacks_clone).await {
if let Some(claim_tracker) = fallback_claim_tracker.as_ref() {
claim_tracker.claim_test_fallback().await;
}
match behavior {
TestReaderBehavior::Eof => {
need_fallback = false;
last_err = None;
continue;
}
TestReaderBehavior::Entries(entries) => {
let mut out = rustfs_filemeta::MetacacheWriter::new(&mut wr);
out.write(&entries).await.expect("test fallback entries should be written");
out.close().await.expect("test fallback entries should close");
need_fallback = false;
last_err = None;
continue;
}
TestReaderBehavior::ProducerError(err) | TestReaderBehavior::PrimaryErrorThenFallback(err) => {
last_err = Some(err);
continue;
}
TestReaderBehavior::Stall
| TestReaderBehavior::IgnoreCancel
| TestReaderBehavior::PartialThenTimeout(_) => {
last_err = Some(DiskError::Timeout);
continue;
}
}
}
let mut disk_op = None;
while let Some(disk) = take_fallback_candidate(&fds_clone).await {
if disk.is_online().await {
disk_op = Some(disk);
break;
}
}
let Some(disk) = disk_op else {
let err = last_err.unwrap_or(DiskError::DiskNotFound);
if is_missing_path_error(&err) {
debug!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts_clone.bucket,
path = %opts_clone.path,
disk_index = disk_idx,
state = "fallback_disk_missing_for_path",
error = ?err,
"Metacache fallback disk unavailable for missing path"
);
} else {
warn!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts_clone.bucket,
path = %opts_clone.path,
disk_index = disk_idx,
state = "fallback_disk_missing",
error = ?err,
"Metacache fallback disk missing"
);
}
record_producer_error(&producer_errs_clone, disk_idx, &err);
return Err(err);
};
if let Some(claim_tracker) = fallback_claim_tracker.as_ref() {
claim_tracker.claim_disk(&disk).await;
}
let fallback_walk_started = std::time::Instant::now();
match disk
.as_ref()
.walk_dir(
WalkDirOptions {
bucket: opts_clone.bucket.clone(),
base_dir: opts_clone.path.clone(),
recursive: opts_clone.recursive,
incl_deleted: opts_clone.incl_deleted,
report_notfound: opts_clone.report_not_found,
filter_prefix: opts_clone.filter_prefix.clone(),
forward_to: opts_clone.forward_to.clone(),
limit: opts_clone.per_disk_limit,
skip_total_timeout: opts_clone.skip_walkdir_total_timeout,
timeout_ms: opts_clone.walkdir_timeout.map(duration_millis),
stall_timeout_ms: opts_clone.walkdir_stall_timeout.map(duration_millis),
..Default::default()
},
&mut wr,
)
.await
{
Ok(_r) => {
rustfs_io_metrics::record_stage_duration(
"metacache_walk_dir_fallback",
fallback_walk_started.elapsed().as_secs_f64() * 1000.0,
);
need_fallback = false;
last_err = None;
}
Err(err) => {
if err.is_metacache_output_stream_closed() {
rustfs_io_metrics::record_stage_duration(
"metacache_walk_dir_fallback",
fallback_walk_started.elapsed().as_secs_f64() * 1000.0,
);
return Ok(());
}
rustfs_io_metrics::record_stage_duration(
"metacache_walk_dir_fallback_failed",
fallback_walk_started.elapsed().as_secs_f64() * 1000.0,
);
if is_missing_path_error(&err) {
debug!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts_clone.bucket,
path = %opts_clone.path,
disk_index = disk_idx,
state = "fallback_walk_dir_missing_path",
error = ?err,
"Metacache fallback walk_dir missing path skipped"
);
} else {
error!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts_clone.bucket,
path = %opts_clone.path,
disk_index = disk_idx,
state = "fallback_walk_dir_failed",
error = ?err,
"Metacache fallback walk_dir failed"
);
}
last_err = Some(err);
}
}
}
if need_fallback {
return Err(last_err.unwrap_or(DiskError::DiskNotFound));
}
// warn!("list_path_raw: while need_fallback done");
Ok(())
}));
}
let revjob = spawn(async move {
let peek_timeout = opts
.walkdir_stall_timeout
.or({
#[cfg(test)]
{
opts.peek_timeout
}
#[cfg(not(test))]
{
None
}
})
.unwrap_or_else(get_drive_walkdir_stall_timeout);
let mut errs: Vec<Option<DiskError>> = Vec::with_capacity(readers.len());
for _ in 0..readers.len() {
errs.push(None);
}
let mut pending_entries: Vec<Option<MetaCacheEntry>> = vec![None; readers.len()];
loop {
let mut current = MetaCacheEntry::default();
// warn!(
// "list_path_raw: loop start, bucket: {}, path: {}, current: {:?}",
// opts.bucket, opts.path, &current.name
// );
if rx.is_cancelled() {
return Ok(());
}
let mut top_entries: Vec<Option<MetaCacheEntry>> = vec![None; readers.len()];
let mut at_eof = 0;
let mut fnf = 0;
let mut vnf = 0;
let mut has_err = 0;
let mut agree = 0;
for (i, r) in readers.iter_mut().enumerate() {
if errs[i].is_some() {
has_err += 1;
continue;
}
let entry = if let Some(entry) = pending_entries[i].take() {
entry
} else {
match peek_with_timeout(r, peek_timeout).await {
PeekOutcome::Ready(res) => {
if let Some(entry) = res {
// info!("read entry disk: {}, name: {}", i, entry.name);
entry
} else {
if let Some(err) = producer_error(&producer_errs, i) {
has_err += 1;
errs[i] = Some(err);
continue;
}
// eof
at_eof += 1;
// warn!("list_path_raw: peek eof, disk: {}", i);
continue;
}
}
PeekOutcome::Error(err) => {
if let Some(err) = producer_error(&producer_errs, i) {
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);
continue;
}
// warn!("list_path_raw: peek err00, err: {:?}", err);
if is_io_eof(&err) {
at_eof += 1;
// warn!("list_path_raw: peek eof, disk: {}", i);
continue;
}
if err == rustfs_filemeta::Error::FileNotFound {
at_eof += 1;
fnf += 1;
// warn!("list_path_raw: peek fnf, disk: {}", i);
continue;
} else if err == rustfs_filemeta::Error::VolumeNotFound {
at_eof += 1;
fnf += 1;
vnf += 1;
// warn!("list_path_raw: peek vnf, disk: {}", i);
continue;
} else {
has_err += 1;
errs[i] = Some(err.into());
// warn!("list_path_raw: peek err, disk: {}", i);
continue;
}
}
PeekOutcome::TimedOut => {
has_err += 1;
errs[i] = Some(DiskError::Timeout);
let endpoint = opts
.disks
.get(i)
.and_then(|disk| disk.as_ref().map(|disk| disk.endpoint().to_string()))
.unwrap_or_else(|| "missing".to_string());
counter!(
"rustfs_list_path_raw_stall_total",
"drive" => endpoint.clone()
)
.increment(1);
warn!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
drive = %endpoint,
bucket = %opts.bucket,
path = %opts.path,
timeout_ms = peek_timeout.as_millis(),
state = "peek_timed_out",
"Metacache reader peek timed out"
);
let (detached_rd, write_half) = tokio::io::duplex(1);
drop(write_half);
*r = MetacacheReader::new(detached_rd);
continue;
}
}
};
// warn!("list_path_raw: loop entry: {:?}, disk: {}", &entry.name, i);
// If no current, add it.
if current.name.is_empty() {
top_entries[i] = Some(entry.clone());
current = entry;
agree += 1;
continue;
}
// If exact match, we agree.
if let (_, true) = current.matches(Some(&entry), true) {
top_entries[i] = Some(entry);
agree += 1;
continue;
}
// If only the name matches we didn't agree, but add it for resolution.
if entry.name == current.name {
top_entries[i] = Some(entry);
continue;
}
// We got different entries
if entry.name > current.name {
pending_entries[i] = Some(entry);
continue;
}
for (idx, item) in top_entries.iter_mut().enumerate().take(i) {
if let Some(entry) = item.take() {
pending_entries[idx] = Some(entry);
}
}
agree = 1;
top_entries[i] = Some(entry.clone());
current = entry;
}
if vnf > 0 && vnf >= (readers.len() - opts.min_disks) {
// warn!("list_path_raw: vnf > 0 && vnf >= (readers.len() - opts.min_disks) break");
return Err(DiskError::VolumeNotFound);
}
if fnf > 0 && fnf >= (readers.len() - opts.min_disks) {
// warn!("list_path_raw: fnf > 0 && fnf >= (readers.len() - opts.min_disks) break");
return Err(DiskError::FileNotFound);
}
if has_err > 0 && has_err > opts.disks.len() - opts.min_disks {
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)) => {
combined_err.push(format!("drive {} returned: {}", disk.to_string(), err));
}
(Some(err), None) => {
combined_err.push(err.to_string());
}
_ => {}
});
error!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %opts.bucket,
path = %opts.path,
state = "quorum_failed",
error = %combined_err.join(", "),
"Metacache listing quorum failed"
);
return Err(DiskError::other(combined_err.join(", ")));
}
// Break if all at EOF or error.
if at_eof + has_err == readers.len() {
if has_err > 0
&& let Some(finished_fn) = opts.finished.as_ref()
{
finished_fn(&errs).await;
}
// Tolerated reader failures, including timeouts, must not turn a
// quorum EOF into a failed listing. The quorum failure branch above
// is the single place that escalates too many reader errors.
// error!("list_path_raw: at_eof + has_err == readers.len() break {:?}", &errs);
break;
}
if agree == readers.len() {
for r in readers.iter_mut() {
let _ = r.skip(1).await;
}
if let Some(agreed_fn) = opts.agreed.as_ref() {
// warn!("list_path_raw: agreed_fn start, current: {:?}", &current.name);
agreed_fn(current).await;
// warn!("list_path_raw: agreed_fn done");
}
continue;
}
// warn!("list_path_raw: skip start, current: {:?}", &current.name);
for (i, r) in readers.iter_mut().enumerate() {
if top_entries[i].is_some() {
let _ = r.skip(1).await;
}
}
if let Some(partial_fn) = opts.partial.as_ref() {
partial_fn(MetaCacheEntries(top_entries), &errs).await;
}
}
Ok(())
});
let merge_started = std::time::Instant::now();
if let Err(err) = revjob.await.map_err(std::io::Error::other)? {
rustfs_io_metrics::record_stage_duration("metacache_merge_failed", merge_started.elapsed().as_secs_f64() * 1000.0);
if is_missing_path_error(&err) {
debug!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %log_bucket,
path = %log_path,
state = "merge_job_missing_path",
error = ?err,
"Metacache merge job missing path skipped"
);
} else {
error!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
bucket = %log_bucket,
path = %log_path,
state = "merge_job_failed",
error = ?err,
"Metacache merge job failed"
);
}
cancel_rx.cancel();
for job in jobs {
job.abort();
}
return Err(err);
}
rustfs_io_metrics::record_stage_duration("metacache_merge", merge_started.elapsed().as_secs_f64() * 1000.0);
// The merge consumer can finish successfully before every producer finishes
// (for example after reaching EOF quorum while a tolerated drive is stalled,
// or after the requested listing limit is satisfied). Cancel remaining walk
// jobs before aborting them so list calls do not wait for slow remote streams.
cancel_rx.cancel();
for job in jobs.iter() {
if !job.is_finished() {
job.abort();
}
}
let results = join_all(jobs).await;
let mut job_errs = Vec::new();
for result in results {
match result {
Ok(Ok(())) => {}
Ok(Err(err)) => {
if is_tolerated_producer_completion_error(&err) {
debug!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
state = "producer_tolerated_after_merge",
error = ?err,
"Metacache producer stopped after merge completed"
);
continue;
}
if is_missing_path_error(&err) {
debug!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
state = "producer_missing_path",
error = ?err,
"Metacache producer missing path"
);
} else {
error!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
state = "producer_failed",
error = ?err,
"Metacache producer failed"
);
}
job_errs.push(err);
}
Err(err) => {
if err.is_cancelled() {
continue;
}
error!(
event = EVENT_METACACHE_LISTING,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_METACACHE,
state = "producer_join_failed",
error = ?err,
"Metacache producer join failed"
);
job_errs.push(err.into());
}
}
}
if job_errs.len() > max_disk_failures {
return Err(job_errs.remove(0));
}
// warn!("list_path_raw: done");
Ok(())
}
#[inline]
fn record_producer_error(producer_errs: &[OnceLock<DiskError>], idx: usize, err: &DiskError) {
let _ = producer_errs[idx].set(err.clone());
}
#[inline]
fn producer_error(producer_errs: &[OnceLock<DiskError>], idx: usize) -> Option<DiskError> {
producer_errs[idx].get().cloned()
}
#[cfg(test)]
mod tests {
use super::*;
use rustfs_filemeta::MetacacheWriter;
use rustfs_filemeta::{FileInfo, FileMeta, MetadataResolutionParams};
use std::sync::Mutex;
use time::OffsetDateTime;
use uuid::Uuid;
#[tokio::test]
async fn list_path_raw_empty_disks_returns_read_quorum() {
let err = list_path_raw(CancellationToken::new(), ListPathRawOptions::default())
.await
.expect_err("empty drive list should fail");
assert_eq!(err, DiskError::ErasureReadQuorum);
}
#[tokio::test]
async fn list_path_raw_rejects_impossible_min_disks() {
let err = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None],
min_disks: 3,
..Default::default()
},
)
.await
.expect_err("impossible listing quorum should fail before producing partial results");
assert_eq!(err, DiskError::ErasureReadQuorum);
}
#[test]
fn missing_path_error_classification_excludes_actionable_failures() {
assert!(is_missing_path_error(&DiskError::FileNotFound));
assert!(is_missing_path_error(&DiskError::FileVersionNotFound));
assert!(is_missing_path_error(&DiskError::VolumeNotFound));
assert!(!is_missing_path_error(&DiskError::Timeout));
assert!(!is_missing_path_error(&DiskError::DiskNotFound));
assert!(!is_missing_path_error(&DiskError::FileAccessDenied));
}
#[test]
fn tolerated_producer_completion_error_classification_excludes_actionable_failures() {
assert!(is_tolerated_producer_completion_error(&DiskError::DiskOngoingReq));
assert!(is_tolerated_producer_completion_error(&DiskError::metacache_output_stream_closed()));
assert!(is_tolerated_producer_completion_error(&DiskError::Io(std::io::Error::new(
ErrorKind::BrokenPipe,
"reader closed after merge",
))));
assert!(!is_tolerated_producer_completion_error(&DiskError::Timeout));
assert!(!is_tolerated_producer_completion_error(&DiskError::DiskNotFound));
assert!(!is_tolerated_producer_completion_error(&DiskError::FileAccessDenied));
}
#[tokio::test]
async fn fallback_candidates_are_claimed_once_across_producers() {
let mut queue = VecDeque::new();
queue.push_back(1usize);
let candidates = Arc::new(TokioMutex::new(queue));
let first = take_fallback_candidate(&candidates).await;
let second = take_fallback_candidate(&candidates).await;
assert_eq!(first, Some(1));
assert_eq!(second, None);
}
fn fallback_test_entry() -> MetaCacheEntry {
MetaCacheEntry {
name: "bucket/object".to_string(),
metadata: vec![1, 2, 3],
cached: None,
reusable: false,
}
}
#[tokio::test]
async fn list_path_raw_does_not_reuse_one_fallback_for_multiple_failed_producers() {
let err = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::PrimaryErrorThenFallback(DiskError::DiskNotFound),
TestReaderBehavior::PrimaryErrorThenFallback(DiskError::DiskNotFound),
],
test_fallback_reader_behaviors: vec![TestReaderBehavior::Entries(vec![fallback_test_entry()])],
..Default::default()
},
)
.await
.expect_err("one fallback must not be counted as two failed primary producers");
assert_eq!(err, DiskError::DiskNotFound);
}
#[tokio::test]
async fn list_path_raw_uses_distinct_fallbacks_to_restore_quorum() {
let seen = Arc::new(Mutex::new(Vec::new()));
let seen_clone = seen.clone();
list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::PrimaryErrorThenFallback(DiskError::DiskNotFound),
TestReaderBehavior::PrimaryErrorThenFallback(DiskError::DiskNotFound),
],
test_fallback_reader_behaviors: vec![
TestReaderBehavior::Entries(vec![fallback_test_entry()]),
TestReaderBehavior::Entries(vec![fallback_test_entry()]),
],
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
let seen = seen_clone.clone();
Box::pin(async move {
let matches = entries
.0
.iter()
.flatten()
.filter(|entry| entry.name == "bucket/object")
.count();
seen.lock().expect("seen mutex poisoned").push(matches);
})
})),
..Default::default()
},
)
.await
.expect("distinct fallback producers should restore listing quorum");
assert_eq!(seen.lock().expect("seen mutex poisoned").as_slice(), &[2]);
}
#[tokio::test]
async fn list_path_raw_records_claimed_fallback_candidates() {
let claim_tracker = FallbackClaimTracker::default();
list_path_raw_with_claim_tracker(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None],
min_disks: 1,
test_reader_behaviors: vec![TestReaderBehavior::PrimaryErrorThenFallback(DiskError::DiskNotFound)],
test_fallback_reader_behaviors: vec![TestReaderBehavior::Entries(vec![fallback_test_entry()])],
..Default::default()
},
claim_tracker.clone(),
)
.await
.expect("fallback producer should restore the single logical reader");
assert!(claim_tracker.contains_key("test-fallback-0").await);
}
#[tokio::test]
async fn list_path_raw_returns_timeout_when_reader_stalls_before_completion() {
let err = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None],
min_disks: 2,
test_reader_behaviors: vec![TestReaderBehavior::Stall, TestReaderBehavior::Eof],
peek_timeout: Some(Duration::from_millis(20)),
..Default::default()
},
)
.await
.expect_err("stalled reader should fail when read quorum cannot be met");
assert_eq!(err, DiskError::Timeout);
}
#[tokio::test]
async fn list_path_raw_tolerates_stalled_reader_after_quorum_eof() {
list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None, None],
min_disks: 2,
test_reader_behaviors: vec![TestReaderBehavior::Eof, TestReaderBehavior::Eof, TestReaderBehavior::Stall],
peek_timeout: Some(Duration::from_millis(20)),
..Default::default()
},
)
.await
.expect("stalled reader within the tolerated failure budget should not fail quorum EOF");
}
#[tokio::test]
async fn list_path_raw_completes_after_partial_quorum_when_reader_stalls() {
let entry = MetaCacheEntry {
name: "bucket/object".to_string(),
metadata: vec![1, 2, 3],
cached: None,
reusable: false,
};
let seen = Arc::new(Mutex::new(Vec::new()));
let seen_clone = seen.clone();
list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::Entries(vec![entry.clone()]),
TestReaderBehavior::Entries(vec![entry]),
TestReaderBehavior::Stall,
],
peek_timeout: Some(Duration::from_millis(20)),
partial: Some(Box::new(move |entries: MetaCacheEntries, _: &[Option<DiskError>]| {
let seen = seen_clone.clone();
Box::pin(async move {
let mut names = entries.0.iter().flatten().map(|entry| entry.name.clone());
if let Some(name) = names.next()
&& names.any(|next| next == name)
{
seen.lock().expect("seen mutex poisoned").push(name);
}
})
})),
..Default::default()
},
)
.await
.expect("stalled reader within failure budget should not fail a quorum listing");
assert_eq!(seen.lock().expect("seen mutex poisoned").as_slice(), &["bucket/object".to_string()]);
}
#[tokio::test]
async fn list_path_raw_aborts_unresponsive_producer_after_quorum_eof() {
let result = timeout(
Duration::from_millis(200),
list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::Eof,
TestReaderBehavior::Eof,
TestReaderBehavior::IgnoreCancel,
],
peek_timeout: Some(Duration::from_millis(20)),
..Default::default()
},
),
)
.await;
let listing = result.expect("list_path_raw should abort unresponsive producer instead of hanging");
listing.expect("unresponsive producer within failure budget should not fail quorum EOF");
}
#[tokio::test]
async fn list_path_raw_treats_external_cancel_as_successful_shutdown() {
let cancel = CancellationToken::new();
cancel.cancel();
list_path_raw(
cancel,
ListPathRawOptions {
disks: vec![None],
min_disks: 1,
test_reader_behaviors: vec![TestReaderBehavior::IgnoreCancel],
..Default::default()
},
)
.await
.expect("external cancellation should stop listing without a synthetic canceled error");
}
#[tokio::test]
async fn list_path_raw_ignores_closed_output_stream_after_quorum_eof() {
list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::Eof,
TestReaderBehavior::Eof,
TestReaderBehavior::ProducerError(DiskError::metacache_output_stream_closed()),
],
..Default::default()
},
)
.await
.expect("closed metacache output after quorum EOF should not fail listing");
}
#[tokio::test]
async fn list_path_raw_tolerates_concurrent_scan_after_quorum_eof() {
let result = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::Eof,
TestReaderBehavior::Eof,
TestReaderBehavior::ProducerError(DiskError::DiskOngoingReq),
],
..Default::default()
},
)
.await;
assert!(result.is_ok(), "concurrent scan after quorum EOF should not fail listing");
}
#[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 list_path_raw_continues_after_single_disk_delete_marker_gap() {
fn object_entry(name: &str, version_id: &str) -> MetaCacheEntry {
let mut meta = FileMeta::default();
let mut fi = FileInfo::new(name, 1, 1);
fi.version_id = Some(Uuid::parse_str(version_id).expect("test version id should parse"));
fi.mod_time = Some(OffsetDateTime::now_utc());
meta.add_version(fi).expect("object metadata should be valid");
MetaCacheEntry {
name: name.to_string(),
metadata: meta.marshal_msg().expect("object metadata should encode"),
cached: Some(meta),
reusable: false,
}
}
fn delete_marker_entry(name: &str, version_id: &str) -> MetaCacheEntry {
let mut meta = FileMeta::default();
meta.add_version(FileInfo {
deleted: true,
version_id: Some(Uuid::parse_str(version_id).expect("test version id should parse")),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
})
.expect("delete marker metadata should be valid");
MetaCacheEntry {
name: name.to_string(),
metadata: meta.marshal_msg().expect("delete marker metadata should encode"),
cached: Some(meta),
reusable: false,
}
}
let visible_before = object_entry("object-000003", "11111111-1111-1111-1111-111111111111");
let hidden_delete = delete_marker_entry("object-000004", "22222222-2222-2222-2222-222222222222");
let visible_after = object_entry("object-000005", "33333333-3333-3333-3333-333333333333");
let seen = Arc::new(Mutex::new(Vec::new()));
let agreed_seen = seen.clone();
let partial_seen = seen.clone();
list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None, None, None, None],
min_disks: 2,
test_reader_behaviors: vec![
TestReaderBehavior::Entries(vec![visible_before.clone(), hidden_delete, visible_after.clone()]),
TestReaderBehavior::Entries(vec![visible_before.clone(), visible_after.clone()]),
TestReaderBehavior::Entries(vec![visible_before.clone(), visible_after.clone()]),
TestReaderBehavior::Entries(vec![visible_before, visible_after]),
],
agreed: Some(Box::new(move |entry| {
let seen = agreed_seen.clone();
Box::pin(async move {
seen.lock().expect("seen mutex poisoned").push(entry.name);
})
})),
partial: Some(Box::new(move |entries, _| {
let seen = partial_seen.clone();
Box::pin(async move {
if let Some(entry) = entries.resolve(MetadataResolutionParams {
obj_quorum: 2,
requested_versions: 1,
bucket: "bucket".to_string(),
..Default::default()
}) {
seen.lock().expect("seen mutex poisoned").push(entry.name);
}
})
})),
..Default::default()
},
)
.await
.expect("single-disk delete marker gap should not end listing");
assert_eq!(
seen.lock().expect("seen mutex poisoned").as_slice(),
&["object-000003".to_string(), "object-000005".to_string()]
);
}
#[tokio::test]
async fn peek_with_timeout_times_out_on_silent_reader() {
let (_writer, reader) = tokio::io::duplex(64);
let mut reader = MetacacheReader::new(reader);
let outcome = peek_with_timeout(&mut reader, Duration::from_millis(20)).await;
assert!(matches!(outcome, PeekOutcome::TimedOut));
}
#[tokio::test]
async fn peek_with_timeout_reads_entry_before_deadline() {
let (reader, writer) = tokio::io::duplex(256);
let mut metacache_reader = MetacacheReader::new(reader);
tokio::spawn(async move {
let mut writer = MetacacheWriter::new(writer);
let entry = MetaCacheEntry {
name: "bucket/object".to_string(),
metadata: vec![1, 2, 3],
cached: None,
reusable: false,
};
writer.write(&[entry]).await.expect("entry should be written");
writer.close().await.expect("writer should close");
});
let outcome = peek_with_timeout(&mut metacache_reader, Duration::from_secs(1)).await;
match outcome {
PeekOutcome::Ready(Some(entry)) => assert_eq!(entry.name, "bucket/object"),
other => panic!("expected ready entry, got {other:?}"),
}
}
#[tokio::test]
async fn list_path_raw_propagates_producer_access_denied() {
let err = list_path_raw(
CancellationToken::new(),
ListPathRawOptions {
disks: vec![None],
min_disks: 1,
test_reader_behaviors: vec![TestReaderBehavior::ProducerError(DiskError::FileAccessDenied)],
..Default::default()
},
)
.await
.expect_err("producer access failure must not be treated as an empty listing");
assert_eq!(err, DiskError::FileAccessDenied);
}
}