mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-06 21:33:14 +00:00
feat: preserve request ids across async recovery logs (#3451)
* feat(obs): promote request ids in structured logs * refactor(tracing): propagate spans into request tasks * test(ecstore): baseline recovery monitor log chains * fix(replication): reduce startup resync log noise * chore(docs): stop tracking local recovery baseline * chore(obs): polish request id logging cleanup
This commit is contained in:
@@ -64,9 +64,14 @@ const EVENT_REPLICATION_WORKER_RESIZE_SKIPPED: &str = "replication_worker_resize
|
||||
const EVENT_REPLICATION_WORKER_RESIZED: &str = "replication_worker_resized";
|
||||
const EVENT_REPLICATION_BACKPRESSURE: &str = "replication_backpressure";
|
||||
const EVENT_REPLICATION_RESYNC_LOAD_SKIPPED: &str = "replication_resync_load_skipped";
|
||||
const EVENT_REPLICATION_RESYNC_RECOVERED: &str = "replication_resync_recovered";
|
||||
const EVENT_REPLICATION_CONFIG_LOOKUP_SKIPPED: &str = "replication_config_lookup_skipped";
|
||||
const EVENT_REPLICATION_MRF_QUEUE_OVERFLOW: &str = "replication_mrf_queue_overflow";
|
||||
|
||||
fn should_auto_resume_resync(status: ResyncStatusType) -> bool {
|
||||
matches!(status, ResyncStatusType::ResyncPending | ResyncStatusType::ResyncStarted)
|
||||
}
|
||||
|
||||
// Worker limits
|
||||
pub const WORKER_MAX_LIMIT: usize = 500;
|
||||
pub const WORKER_MIN_LIMIT: usize = 50;
|
||||
@@ -1017,6 +1022,11 @@ impl<S: StorageAPI + NamespaceLocking> ReplicationPool<S> {
|
||||
// Note: Leader lock implementation would be needed here
|
||||
// let _lock_guard = global_leader_lock.get_lock().await?;
|
||||
|
||||
let mut recovered_statuses = Vec::new();
|
||||
let mut restart_opts = Vec::new();
|
||||
let mut recovered_bucket_count = 0usize;
|
||||
let mut skipped_failed_target_count = 0usize;
|
||||
|
||||
for bucket in buckets {
|
||||
let meta = match load_bucket_resync_metadata(bucket, self.storage.clone()).await {
|
||||
Ok(meta) => meta,
|
||||
@@ -1036,39 +1046,53 @@ impl<S: StorageAPI + NamespaceLocking> ReplicationPool<S> {
|
||||
}
|
||||
};
|
||||
|
||||
// Store metadata in resyncer
|
||||
{
|
||||
let mut status_map = self.resyncer.status_map.write().await;
|
||||
status_map.insert(bucket.clone(), meta.clone());
|
||||
if meta.targets_map.is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Process target statistics
|
||||
let target_stats = meta.clone_tgt_stats();
|
||||
for (arn, stats) in target_stats {
|
||||
match stats.resync_status {
|
||||
ResyncStatusType::ResyncFailed | ResyncStatusType::ResyncStarted | ResyncStatusType::ResyncPending => {
|
||||
// Note: This would spawn a resync task in a real implementation
|
||||
// For now, we just log the resync request
|
||||
|
||||
let ctx = CancellationToken::new();
|
||||
let bucket_clone = bucket.clone();
|
||||
let resync = self.resyncer.clone();
|
||||
let storage = self.storage.clone();
|
||||
let opts = ResyncOpts {
|
||||
bucket: bucket_clone,
|
||||
arn,
|
||||
resync_id: stats.resync_id,
|
||||
resync_before: stats.resync_before_date,
|
||||
};
|
||||
tokio::spawn(async move {
|
||||
resync.register_cancel_token(&opts, ctx.clone()).await;
|
||||
Box::pin(resync.clone().resync_bucket(ctx, storage, true, opts.clone())).await;
|
||||
resync.clear_cancel_token(&opts).await;
|
||||
});
|
||||
}
|
||||
_ => {}
|
||||
recovered_bucket_count += 1;
|
||||
for (arn, stats) in &meta.targets_map {
|
||||
if should_auto_resume_resync(stats.resync_status) {
|
||||
restart_opts.push(ResyncOpts {
|
||||
bucket: bucket.clone(),
|
||||
arn: arn.clone(),
|
||||
resync_id: stats.resync_id.clone(),
|
||||
resync_before: stats.resync_before_date,
|
||||
});
|
||||
} else if stats.resync_status == ResyncStatusType::ResyncFailed {
|
||||
skipped_failed_target_count += 1;
|
||||
}
|
||||
}
|
||||
|
||||
recovered_statuses.push((bucket.clone(), meta));
|
||||
}
|
||||
|
||||
if !recovered_statuses.is_empty() {
|
||||
let mut status_map = self.resyncer.status_map.write().await;
|
||||
status_map.extend(recovered_statuses);
|
||||
}
|
||||
|
||||
if !restart_opts.is_empty() || skipped_failed_target_count > 0 {
|
||||
info!(
|
||||
event = EVENT_REPLICATION_RESYNC_RECOVERED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REPLICATION,
|
||||
recovered_buckets = recovered_bucket_count,
|
||||
resumed_targets = restart_opts.len(),
|
||||
skipped_failed_targets = skipped_failed_target_count,
|
||||
"Recovered replication resync state from persisted metadata; failed targets require manual resync restart"
|
||||
);
|
||||
}
|
||||
|
||||
for opts in restart_opts {
|
||||
let ctx = CancellationToken::new();
|
||||
let resync = self.resyncer.clone();
|
||||
let storage = self.storage.clone();
|
||||
tokio::spawn(async move {
|
||||
resync.register_cancel_token(&opts, ctx.clone()).await;
|
||||
Box::pin(resync.clone().resync_bucket(ctx, storage, true, opts.clone())).await;
|
||||
resync.clear_cancel_token(&opts).await;
|
||||
});
|
||||
}
|
||||
|
||||
Ok(())
|
||||
@@ -1498,4 +1522,14 @@ mod tests {
|
||||
admission.merge(ReplicationQueueAdmission::Missed);
|
||||
assert_eq!(admission, ReplicationQueueAdmission::Missed);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn auto_resume_resync_only_for_inflight_states() {
|
||||
assert!(should_auto_resume_resync(ResyncStatusType::ResyncPending));
|
||||
assert!(should_auto_resume_resync(ResyncStatusType::ResyncStarted));
|
||||
assert!(!should_auto_resume_resync(ResyncStatusType::NoResync));
|
||||
assert!(!should_auto_resume_resync(ResyncStatusType::ResyncCanceled));
|
||||
assert!(!should_auto_resume_resync(ResyncStatusType::ResyncCompleted));
|
||||
assert!(!should_auto_resume_resync(ResyncStatusType::ResyncFailed));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -87,7 +87,7 @@ use tokio::task::JoinSet;
|
||||
use tokio::time::Duration as TokioDuration;
|
||||
use tokio_util::io::ReaderStream;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{debug, error, instrument, warn};
|
||||
use tracing::{debug, error, instrument, trace, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
||||
@@ -920,18 +920,32 @@ impl ReplicationResyncer {
|
||||
(roi.size, None)
|
||||
};
|
||||
|
||||
debug!(
|
||||
event = EVENT_RESYNC_OBJECT_PROCESSED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
|
||||
reset_id = %reset_id,
|
||||
bucket = %bucket_name,
|
||||
object = %roi.name,
|
||||
version_id = %roi.version_id.unwrap_or_default(),
|
||||
size,
|
||||
error = ?err,
|
||||
"Processed resync object"
|
||||
);
|
||||
if err.is_some() {
|
||||
debug!(
|
||||
event = EVENT_RESYNC_OBJECT_PROCESSED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
|
||||
reset_id = %reset_id,
|
||||
bucket = %bucket_name,
|
||||
object = %roi.name,
|
||||
version_id = %roi.version_id.unwrap_or_default(),
|
||||
size,
|
||||
error = ?err,
|
||||
"Processed resync object with verification error"
|
||||
);
|
||||
} else {
|
||||
trace!(
|
||||
event = EVENT_RESYNC_OBJECT_PROCESSED,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REPLICATION_RESYNC,
|
||||
reset_id = %reset_id,
|
||||
bucket = %bucket_name,
|
||||
object = %roi.name,
|
||||
version_id = %roi.version_id.unwrap_or_default(),
|
||||
size,
|
||||
"Processed resync object"
|
||||
);
|
||||
}
|
||||
|
||||
if cancel_token.is_cancelled() {
|
||||
return;
|
||||
|
||||
@@ -34,6 +34,7 @@ use std::{
|
||||
};
|
||||
use time::OffsetDateTime;
|
||||
use tokio::sync::mpsc::{self, Receiver, Sender};
|
||||
use tracing::Instrument;
|
||||
|
||||
use crate::client::utils::base64_encode;
|
||||
use crate::client::{
|
||||
@@ -221,11 +222,14 @@ impl TransitionClient {
|
||||
let self_clone = Arc::clone(&self);
|
||||
let bucket_name_owned = bucket_name.to_string();
|
||||
|
||||
tokio::spawn(async move {
|
||||
self_clone
|
||||
.remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts)
|
||||
.await;
|
||||
});
|
||||
tokio::spawn(
|
||||
async move {
|
||||
self_clone
|
||||
.remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts)
|
||||
.await;
|
||||
}
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
result_rx
|
||||
}
|
||||
|
||||
@@ -241,26 +245,32 @@ impl TransitionClient {
|
||||
let bucket_name_owned = bucket_name.to_string();
|
||||
|
||||
let (result_tx, mut result_rx) = mpsc::channel(1);
|
||||
tokio::spawn(async move {
|
||||
self_clone
|
||||
.remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts)
|
||||
.await;
|
||||
});
|
||||
tokio::spawn(async move {
|
||||
while let Some(res) = result_rx.recv().await {
|
||||
if res.err.is_none() {
|
||||
continue;
|
||||
}
|
||||
error_tx
|
||||
.send(RemoveObjectError {
|
||||
object_name: res.object_name,
|
||||
version_id: res.object_version_id,
|
||||
err: res.err,
|
||||
..Default::default()
|
||||
})
|
||||
tokio::spawn(
|
||||
async move {
|
||||
self_clone
|
||||
.remove_objects_inner(&bucket_name_owned, objects_rx, &result_tx, opts)
|
||||
.await;
|
||||
}
|
||||
});
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
tokio::spawn(
|
||||
async move {
|
||||
while let Some(res) = result_rx.recv().await {
|
||||
if res.err.is_none() {
|
||||
continue;
|
||||
}
|
||||
error_tx
|
||||
.send(RemoveObjectError {
|
||||
object_name: res.object_name,
|
||||
version_id: res.object_version_id,
|
||||
err: res.err,
|
||||
..Default::default()
|
||||
})
|
||||
.await;
|
||||
}
|
||||
}
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
|
||||
error_rx
|
||||
}
|
||||
|
||||
@@ -52,7 +52,7 @@ use tokio::{net::TcpStream, time::Duration};
|
||||
use tonic::Request;
|
||||
use tonic::service::interceptor::InterceptedService;
|
||||
use tonic::transport::Channel;
|
||||
use tracing::warn;
|
||||
use tracing::{Instrument, warn};
|
||||
|
||||
pub const PEER_RESTSIGNAL: &str = "signal";
|
||||
pub const PEER_RESTSUB_SYS: &str = "sub-sys";
|
||||
@@ -78,6 +78,16 @@ pub struct PeerRestClient {
|
||||
}
|
||||
|
||||
impl PeerRestClient {
|
||||
fn recovery_monitor_span(grid_host: &str) -> tracing::Span {
|
||||
tracing::info_span!(
|
||||
"recovery-monitor",
|
||||
component = "ecstore",
|
||||
subsystem = "peer_rest_client",
|
||||
kind = "peer_rest",
|
||||
grid_host = %grid_host
|
||||
)
|
||||
}
|
||||
|
||||
pub fn new(host: XHost, grid_host: String) -> Self {
|
||||
Self {
|
||||
host,
|
||||
@@ -172,28 +182,47 @@ impl PeerRestClient {
|
||||
let grid_host = self.grid_host.clone();
|
||||
let offline = Arc::clone(&self.offline);
|
||||
let recovery_running = Arc::clone(&self.recovery_running);
|
||||
tokio::spawn(async move {
|
||||
let mut delay = get_drive_active_check_interval();
|
||||
let connect_timeout = get_drive_active_check_timeout();
|
||||
let span = Self::recovery_monitor_span(&grid_host);
|
||||
tokio::spawn(
|
||||
async move {
|
||||
let mut delay = get_drive_active_check_interval();
|
||||
let connect_timeout = get_drive_active_check_timeout();
|
||||
|
||||
for _ in 0..PEER_REST_RECOVERY_MAX_ATTEMPTS {
|
||||
tokio::time::sleep(delay).await;
|
||||
if Self::perform_connectivity_check(&grid_host, connect_timeout).await.is_ok() {
|
||||
offline.store(false, Ordering::Release);
|
||||
recovery_running.store(false, Ordering::Release);
|
||||
return;
|
||||
for _ in 0..PEER_REST_RECOVERY_MAX_ATTEMPTS {
|
||||
tokio::time::sleep(delay).await;
|
||||
if Self::perform_connectivity_check(&grid_host, connect_timeout).await.is_ok() {
|
||||
offline.store(false, Ordering::Release);
|
||||
recovery_running.store(false, Ordering::Release);
|
||||
return;
|
||||
}
|
||||
|
||||
delay = std::cmp::min(delay.saturating_mul(2), PEER_REST_RECOVERY_MAX_BACKOFF);
|
||||
}
|
||||
|
||||
delay = std::cmp::min(delay.saturating_mul(2), PEER_REST_RECOVERY_MAX_BACKOFF);
|
||||
warn!(
|
||||
grid_host = %grid_host,
|
||||
attempts = PEER_REST_RECOVERY_MAX_ATTEMPTS,
|
||||
"peer recovery monitor reached max attempts; will retry on next request"
|
||||
);
|
||||
recovery_running.store(false, Ordering::Release);
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
}
|
||||
|
||||
warn!(
|
||||
grid_host = %grid_host,
|
||||
attempts = PEER_REST_RECOVERY_MAX_ATTEMPTS,
|
||||
"peer recovery monitor reached max attempts; will retry on next request"
|
||||
);
|
||||
recovery_running.store(false, Ordering::Release);
|
||||
});
|
||||
#[cfg(test)]
|
||||
fn spawn_recovery_monitor_log_probe_for_test(&self) -> tokio::sync::oneshot::Receiver<()> {
|
||||
let (tx, rx) = tokio::sync::oneshot::channel();
|
||||
let grid_host = self.grid_host.clone();
|
||||
let span = Self::recovery_monitor_span(&grid_host);
|
||||
tokio::spawn(
|
||||
async move {
|
||||
warn!(grid_host = %grid_host, "peer recovery monitor log probe");
|
||||
let _ = tx.send(());
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
rx
|
||||
}
|
||||
|
||||
async fn perform_connectivity_check(addr: &str, timeout_duration: Duration) -> Result<()> {
|
||||
@@ -956,6 +985,58 @@ impl PeerRestClient {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use serde_json::Value;
|
||||
use std::io::{self, Write};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use tracing_subscriber::{Registry, fmt::MakeWriter, layer::SubscriberExt};
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct CapturedLogs {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
struct CapturedLogWriter {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl CapturedLogs {
|
||||
fn lines(&self) -> Vec<Value> {
|
||||
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")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<Value>(line).expect("captured log line should be valid JSON"))
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
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_peer_client() -> PeerRestClient {
|
||||
PeerRestClient::new(
|
||||
@@ -1011,4 +1092,40 @@ mod tests {
|
||||
assert!(matches!(err, Error::VolumeNotFound));
|
||||
assert!(!client.offline.load(Ordering::Acquire));
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn peer_rest_recovery_probe_logs_keep_request_id_span_context() {
|
||||
let logs = CapturedLogs::default();
|
||||
let subscriber = Registry::default().with(
|
||||
tracing_subscriber::fmt::layer()
|
||||
.with_writer(logs.clone())
|
||||
.with_ansi(false)
|
||||
.without_time()
|
||||
.json()
|
||||
.flatten_event(true)
|
||||
.with_current_span(true)
|
||||
.with_span_list(true),
|
||||
);
|
||||
let _guard = tracing::subscriber::set_default(subscriber);
|
||||
|
||||
let client = test_peer_client();
|
||||
let span = tracing::info_span!("request-span", request_id = "req-peer-rest");
|
||||
let _entered = span.enter();
|
||||
let done = client.spawn_recovery_monitor_log_probe_for_test();
|
||||
done.await.expect("recovery monitor probe should signal completion");
|
||||
|
||||
let log = logs
|
||||
.lines()
|
||||
.into_iter()
|
||||
.find(|value| value.get("message").and_then(Value::as_str) == Some("peer recovery monitor log probe"))
|
||||
.expect("expected peer recovery monitor probe log");
|
||||
|
||||
assert_eq!(log["span"]["name"], Value::String("recovery-monitor".to_string()));
|
||||
assert_eq!(log["span"]["kind"], Value::String("peer_rest".to_string()));
|
||||
let spans = log["spans"].as_array().expect("spans should be present");
|
||||
assert!(spans.iter().any(|span| {
|
||||
span.get("name").and_then(Value::as_str) == Some("request-span")
|
||||
&& span.get("request_id").and_then(Value::as_str) == Some("req-peer-rest")
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,7 +45,7 @@ use tokio_util::sync::CancellationToken;
|
||||
use tonic::Request;
|
||||
use tonic::service::interceptor::InterceptedService;
|
||||
use tonic::transport::Channel;
|
||||
use tracing::{debug, info, warn};
|
||||
use tracing::{Instrument, debug, info, warn};
|
||||
|
||||
type Client = Arc<Box<dyn PeerS3Client>>;
|
||||
|
||||
@@ -596,6 +596,16 @@ pub struct RemotePeerS3Client {
|
||||
}
|
||||
|
||||
impl RemotePeerS3Client {
|
||||
fn recovery_monitor_span(addr: &str) -> tracing::Span {
|
||||
tracing::info_span!(
|
||||
"recovery-monitor",
|
||||
component = "ecstore",
|
||||
subsystem = "peer_s3_client",
|
||||
kind = "peer_s3",
|
||||
addr = %addr
|
||||
)
|
||||
}
|
||||
|
||||
pub fn new(node: Option<Node>, pools: Option<Vec<usize>>) -> Self {
|
||||
let addr = node.as_ref().map(|v| v.url.to_string()).unwrap_or_default();
|
||||
let client = Self {
|
||||
@@ -661,10 +671,14 @@ impl RemotePeerS3Client {
|
||||
let health_clone = Arc::clone(&health);
|
||||
let addr_clone = addr.clone();
|
||||
let cancel_clone = cancel_token.clone();
|
||||
let span = Self::recovery_monitor_span(&addr_clone);
|
||||
|
||||
tokio::spawn(async move {
|
||||
Self::monitor_remote_peer_recovery(addr_clone, health_clone, cancel_clone).await;
|
||||
});
|
||||
tokio::spawn(
|
||||
async move {
|
||||
Self::monitor_remote_peer_recovery(addr_clone, health_clone, cancel_clone).await;
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -767,9 +781,13 @@ impl RemotePeerS3Client {
|
||||
let health = Arc::clone(&self.health);
|
||||
let cancel_token = self.cancel_token.clone();
|
||||
let addr = self.addr.clone();
|
||||
tokio::spawn(async move {
|
||||
Self::monitor_remote_peer_recovery(addr, health, cancel_token).await;
|
||||
});
|
||||
let span = Self::recovery_monitor_span(&addr);
|
||||
tokio::spawn(
|
||||
async move {
|
||||
Self::monitor_remote_peer_recovery(addr, health, cancel_token).await;
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -65,7 +65,7 @@ use tokio::{
|
||||
};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tonic::{Request, service::interceptor::InterceptedService, transport::Channel};
|
||||
use tracing::{debug, warn};
|
||||
use tracing::{Instrument, debug, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
@@ -117,6 +117,17 @@ pub struct RemoteDisk {
|
||||
}
|
||||
|
||||
impl RemoteDisk {
|
||||
fn recovery_monitor_span(addr: &str, endpoint: &Endpoint) -> tracing::Span {
|
||||
tracing::info_span!(
|
||||
"recovery-monitor",
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REMOTE_DISK,
|
||||
kind = "remote_disk",
|
||||
endpoint = %endpoint,
|
||||
addr = %addr
|
||||
)
|
||||
}
|
||||
|
||||
fn is_retryable_walk_dir_error(err: &DiskError) -> bool {
|
||||
if is_network_like_disk_error(err) {
|
||||
return true;
|
||||
@@ -239,9 +250,37 @@ impl RemoteDisk {
|
||||
let endpoint = self.endpoint.clone();
|
||||
let health = Arc::clone(&self.health);
|
||||
let cancel_token = self.cancel_token.clone();
|
||||
tokio::spawn(async move {
|
||||
Self::monitor_remote_disk_recovery(addr, endpoint, health, cancel_token).await;
|
||||
});
|
||||
let span = Self::recovery_monitor_span(&addr, &endpoint);
|
||||
tokio::spawn(
|
||||
async move {
|
||||
Self::monitor_remote_disk_recovery(addr, endpoint, health, cancel_token).await;
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn spawn_recovery_monitor_log_probe_for_test(&self) -> tokio::sync::oneshot::Receiver<()> {
|
||||
let (tx, rx) = tokio::sync::oneshot::channel();
|
||||
let endpoint = self.endpoint.clone();
|
||||
let addr = self.addr.clone();
|
||||
let span = Self::recovery_monitor_span(&addr, &endpoint);
|
||||
tokio::spawn(
|
||||
async move {
|
||||
warn!(
|
||||
event = EVENT_REMOTE_DISK_HEALTH,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REMOTE_DISK,
|
||||
endpoint = %endpoint,
|
||||
addr,
|
||||
state = "probe",
|
||||
"remote disk recovery monitor log probe"
|
||||
);
|
||||
let _ = tx.send(());
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
rx
|
||||
}
|
||||
|
||||
fn mark_suspect_or_offline(&self, reason: &'static str) -> bool {
|
||||
@@ -295,10 +334,14 @@ impl RemoteDisk {
|
||||
let addr_clone = addr.clone();
|
||||
let endpoint_clone = endpoint.clone();
|
||||
let cancel_clone = cancel_token.clone();
|
||||
let span = Self::recovery_monitor_span(&addr_clone, &endpoint_clone);
|
||||
|
||||
tokio::spawn(async move {
|
||||
Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await;
|
||||
});
|
||||
tokio::spawn(
|
||||
async move {
|
||||
Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await;
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
}
|
||||
|
||||
loop {
|
||||
@@ -357,10 +400,14 @@ impl RemoteDisk {
|
||||
let addr_clone = addr.clone();
|
||||
let endpoint_clone = endpoint.clone();
|
||||
let cancel_clone = cancel_token.clone();
|
||||
let span = Self::recovery_monitor_span(&addr_clone, &endpoint_clone);
|
||||
|
||||
tokio::spawn(async move {
|
||||
Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await;
|
||||
});
|
||||
tokio::spawn(
|
||||
async move {
|
||||
Self::monitor_remote_disk_recovery(addr_clone, endpoint_clone, health_clone, cancel_clone).await;
|
||||
}
|
||||
.instrument(span),
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -375,10 +422,28 @@ impl RemoteDisk {
|
||||
cancel_token: CancellationToken,
|
||||
) {
|
||||
let mut interval = time::interval(get_drive_returning_probe_interval());
|
||||
debug!(
|
||||
event = EVENT_REMOTE_DISK_HEALTH,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REMOTE_DISK,
|
||||
endpoint = %endpoint,
|
||||
addr,
|
||||
state = "recovery_monitor_started",
|
||||
"Remote disk recovery monitor started"
|
||||
);
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = cancel_token.cancelled() => {
|
||||
debug!(
|
||||
event = EVENT_REMOTE_DISK_HEALTH,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
subsystem = LOG_SUBSYSTEM_REMOTE_DISK,
|
||||
endpoint = %endpoint,
|
||||
addr,
|
||||
state = "recovery_monitor_cancelled",
|
||||
"Remote disk recovery monitor cancelled"
|
||||
);
|
||||
return;
|
||||
}
|
||||
_ = interval.tick() => {
|
||||
@@ -2200,17 +2265,68 @@ mod tests {
|
||||
use super::*;
|
||||
use crate::rpc::{InternodeDataTransportCapabilities, TcpHttpInternodeDataTransport};
|
||||
use rustfs_common::GLOBAL_CONN_MAP;
|
||||
use serde_json::Value;
|
||||
use std::io::{self as std_io, Write};
|
||||
use std::pin::Pin;
|
||||
use std::sync::{Mutex as StdMutex, Once};
|
||||
use std::sync::{Arc, Mutex, Mutex as StdMutex, Once};
|
||||
use std::task::{Context, Poll};
|
||||
use tokio::io::{ReadBuf, duplex};
|
||||
use tokio::net::TcpListener;
|
||||
use tonic::transport::Endpoint as TonicEndpoint;
|
||||
use tracing::Level;
|
||||
use tracing_subscriber::{Registry, fmt::MakeWriter, layer::SubscriberExt};
|
||||
use uuid::Uuid;
|
||||
|
||||
static INIT: Once = Once::new();
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct CapturedLogs {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
struct CapturedLogWriter {
|
||||
buffer: Arc<Mutex<Vec<u8>>>,
|
||||
}
|
||||
|
||||
impl CapturedLogs {
|
||||
fn lines(&self) -> Vec<Value> {
|
||||
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")
|
||||
.lines()
|
||||
.map(|line| serde_json::from_str::<Value>(line).expect("captured log line should be valid JSON"))
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
impl Write for CapturedLogWriter {
|
||||
fn write(&mut self, buf: &[u8]) -> std_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) -> std_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),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
enum RecordedTransportCall {
|
||||
Read(ReadStreamRequest),
|
||||
@@ -3500,4 +3616,155 @@ mod tests {
|
||||
assert_eq!(endpoint.set_idx, 2);
|
||||
assert_eq!(endpoint.disk_idx, 3);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn remote_disk_recovery_probe_logs_keep_request_id_span_context() {
|
||||
let logs = CapturedLogs::default();
|
||||
let subscriber = Registry::default().with(
|
||||
tracing_subscriber::fmt::layer()
|
||||
.with_writer(logs.clone())
|
||||
.with_ansi(false)
|
||||
.without_time()
|
||||
.json()
|
||||
.flatten_event(true)
|
||||
.with_current_span(true)
|
||||
.with_span_list(true),
|
||||
);
|
||||
let _guard = tracing::subscriber::set_default(subscriber);
|
||||
|
||||
let endpoint = Endpoint {
|
||||
url: url::Url::parse("http://127.0.0.1:59996/data").expect("endpoint URL should parse"),
|
||||
is_local: false,
|
||||
pool_idx: 0,
|
||||
set_idx: 0,
|
||||
disk_idx: 0,
|
||||
};
|
||||
let remote_disk = RemoteDisk::new(
|
||||
&endpoint,
|
||||
&DiskOption {
|
||||
cleanup: false,
|
||||
health_check: true,
|
||||
},
|
||||
Arc::new(TcpHttpInternodeDataTransport),
|
||||
)
|
||||
.await
|
||||
.expect("remote disk should construct");
|
||||
|
||||
let span = tracing::info_span!("request-span", request_id = "req-remote-disk");
|
||||
let _entered = span.enter();
|
||||
let done = remote_disk.spawn_recovery_monitor_log_probe_for_test();
|
||||
done.await
|
||||
.expect("remote disk recovery monitor probe should signal completion");
|
||||
|
||||
let log = logs
|
||||
.lines()
|
||||
.into_iter()
|
||||
.find(|value| value.get("message").and_then(Value::as_str) == Some("remote disk recovery monitor log probe"))
|
||||
.expect("expected remote disk recovery monitor probe log");
|
||||
|
||||
assert_eq!(log["span"]["name"], Value::String("recovery-monitor".to_string()));
|
||||
assert_eq!(log["span"]["kind"], Value::String("remote_disk".to_string()));
|
||||
let spans = log["spans"].as_array().expect("spans should be present");
|
||||
assert!(spans.iter().any(|span| {
|
||||
span.get("name").and_then(Value::as_str) == Some("request-span")
|
||||
&& span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk")
|
||||
}));
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn remote_disk_network_error_starts_recovery_monitor_with_request_context() {
|
||||
let logs = CapturedLogs::default();
|
||||
let subscriber = Registry::default().with(
|
||||
tracing_subscriber::fmt::layer()
|
||||
.with_writer(logs.clone())
|
||||
.with_ansi(false)
|
||||
.without_time()
|
||||
.json()
|
||||
.flatten_event(true)
|
||||
.with_current_span(true)
|
||||
.with_span_list(true),
|
||||
);
|
||||
let _guard = tracing::subscriber::set_default(subscriber);
|
||||
|
||||
let addr = "http://127.0.0.1:59997".to_string();
|
||||
let endpoint = Endpoint {
|
||||
url: url::Url::parse(&format!("{addr}/data")).expect("endpoint URL should parse"),
|
||||
is_local: false,
|
||||
pool_idx: 0,
|
||||
set_idx: 0,
|
||||
disk_idx: 0,
|
||||
};
|
||||
|
||||
let remote_disk = RemoteDisk::new(
|
||||
&endpoint,
|
||||
&DiskOption {
|
||||
cleanup: false,
|
||||
health_check: true,
|
||||
},
|
||||
Arc::new(TcpHttpInternodeDataTransport),
|
||||
)
|
||||
.await
|
||||
.expect("remote disk should construct");
|
||||
|
||||
let span = tracing::info_span!("request-span", request_id = "req-remote-disk-e2e");
|
||||
let _entered = span.enter();
|
||||
|
||||
let err = remote_disk
|
||||
.execute_with_timeout(
|
||||
|| async {
|
||||
Err::<(), Error>(DiskError::Io(std::io::Error::new(
|
||||
std::io::ErrorKind::ConnectionRefused,
|
||||
"connection refused",
|
||||
)))
|
||||
},
|
||||
Duration::from_secs(1),
|
||||
)
|
||||
.await
|
||||
.expect_err("network-like operation error should fail");
|
||||
assert_eq!(
|
||||
match &err {
|
||||
DiskError::Io(io_err) => io_err.kind(),
|
||||
other => panic!("expected io network error, got {other:?}"),
|
||||
},
|
||||
std::io::ErrorKind::ConnectionRefused
|
||||
);
|
||||
|
||||
tokio::task::yield_now().await;
|
||||
tokio::time::sleep(Duration::from_millis(20)).await;
|
||||
remote_disk.cancel_token.cancel();
|
||||
tokio::task::yield_now().await;
|
||||
|
||||
let lines = logs.lines();
|
||||
let marked_suspect = lines
|
||||
.iter()
|
||||
.find(|value| value.get("state").and_then(Value::as_str) == Some("marked_suspect"))
|
||||
.expect("expected marked_suspect log");
|
||||
assert!(
|
||||
marked_suspect["spans"]
|
||||
.as_array()
|
||||
.expect("spans should be present")
|
||||
.iter()
|
||||
.any(|span| {
|
||||
span.get("name").and_then(Value::as_str) == Some("request-span")
|
||||
&& span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk-e2e")
|
||||
})
|
||||
);
|
||||
|
||||
let recovery_started = lines
|
||||
.iter()
|
||||
.find(|value| value.get("state").and_then(Value::as_str) == Some("recovery_monitor_started"))
|
||||
.expect("expected recovery_monitor_started log");
|
||||
assert_eq!(recovery_started["span"]["name"], Value::String("recovery-monitor".to_string()));
|
||||
assert_eq!(recovery_started["span"]["kind"], Value::String("remote_disk".to_string()));
|
||||
assert!(
|
||||
recovery_started["spans"]
|
||||
.as_array()
|
||||
.expect("spans should be present")
|
||||
.iter()
|
||||
.any(|span| {
|
||||
span.get("name").and_then(Value::as_str) == Some("request-span")
|
||||
&& span.get("request_id").and_then(Value::as_str) == Some("req-remote-disk-e2e")
|
||||
})
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -127,7 +127,7 @@ use tokio::{
|
||||
};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::error;
|
||||
use tracing::{debug, info, warn};
|
||||
use tracing::{Instrument, debug, info, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
|
||||
@@ -1491,55 +1491,58 @@ impl SetDisks {
|
||||
if !pending.is_empty() {
|
||||
let cleanup_requests = requests.clone();
|
||||
let lockers = self.lockers.clone();
|
||||
let handle = tokio::spawn(async move {
|
||||
let mut late_lock_ids_by_client = vec![Vec::new(); lockers.len()];
|
||||
let mut pending = pending;
|
||||
while let Some(join_result) = pending.join_next().await {
|
||||
match join_result {
|
||||
Ok((client_idx, Ok(responses))) => {
|
||||
for (req_idx, request) in cleanup_requests.iter().enumerate() {
|
||||
if let Some(response) = responses.get(req_idx)
|
||||
&& response.success
|
||||
{
|
||||
let lock_id = response
|
||||
.lock_info
|
||||
.as_ref()
|
||||
.map(|lock_info| lock_info.id.clone())
|
||||
.unwrap_or_else(|| request.lock_id.clone());
|
||||
if let Some(client_locks) = late_lock_ids_by_client.get_mut(client_idx) {
|
||||
client_locks.push(lock_id);
|
||||
let handle = tokio::spawn(
|
||||
async move {
|
||||
let mut late_lock_ids_by_client = vec![Vec::new(); lockers.len()];
|
||||
let mut pending = pending;
|
||||
while let Some(join_result) = pending.join_next().await {
|
||||
match join_result {
|
||||
Ok((client_idx, Ok(responses))) => {
|
||||
for (req_idx, request) in cleanup_requests.iter().enumerate() {
|
||||
if let Some(response) = responses.get(req_idx)
|
||||
&& response.success
|
||||
{
|
||||
let lock_id = response
|
||||
.lock_info
|
||||
.as_ref()
|
||||
.map(|lock_info| lock_info.id.clone())
|
||||
.unwrap_or_else(|| request.lock_id.clone());
|
||||
if let Some(client_locks) = late_lock_ids_by_client.get_mut(client_idx) {
|
||||
client_locks.push(lock_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok((_client_idx, Err(err))) => {
|
||||
tracing::warn!("late distributed delete lock batch request failed: {}", err);
|
||||
}
|
||||
Err(err) => {
|
||||
tracing::warn!("late distributed delete lock batch task join failed: {}", err);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
join_all(lockers.iter().cloned().enumerate().filter_map(|(client_idx, client)| {
|
||||
let lock_ids = late_lock_ids_by_client.get(client_idx).cloned().unwrap_or_default();
|
||||
if lock_ids.is_empty() {
|
||||
None
|
||||
} else {
|
||||
Some(async move {
|
||||
if let Err(err) = client.release_locks_batch(&lock_ids).await {
|
||||
tracing::warn!(
|
||||
client_idx,
|
||||
lock_count = lock_ids.len(),
|
||||
"failed to cleanup late distributed delete locks in batch: {}",
|
||||
err
|
||||
);
|
||||
Ok((_client_idx, Err(err))) => {
|
||||
tracing::warn!("late distributed delete lock batch request failed: {}", err);
|
||||
}
|
||||
})
|
||||
Err(err) => {
|
||||
tracing::warn!("late distributed delete lock batch task join failed: {}", err);
|
||||
}
|
||||
}
|
||||
}
|
||||
}))
|
||||
.await;
|
||||
});
|
||||
|
||||
join_all(lockers.iter().cloned().enumerate().filter_map(|(client_idx, client)| {
|
||||
let lock_ids = late_lock_ids_by_client.get(client_idx).cloned().unwrap_or_default();
|
||||
if lock_ids.is_empty() {
|
||||
None
|
||||
} else {
|
||||
Some(async move {
|
||||
if let Err(err) = client.release_locks_batch(&lock_ids).await {
|
||||
tracing::warn!(
|
||||
client_idx,
|
||||
lock_count = lock_ids.len(),
|
||||
"failed to cleanup late distributed delete locks in batch: {}",
|
||||
err
|
||||
);
|
||||
}
|
||||
})
|
||||
}
|
||||
}))
|
||||
.await;
|
||||
}
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
drop(handle);
|
||||
}
|
||||
|
||||
|
||||
@@ -40,7 +40,7 @@ use std::sync::Arc;
|
||||
use tokio::sync::broadcast::{self};
|
||||
use tokio::sync::mpsc::{self, Receiver, Sender};
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::{error, info, warn};
|
||||
use tracing::{Instrument, error, info, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
const MAX_OBJECT_LIST: i32 = 1000;
|
||||
@@ -710,31 +710,37 @@ impl ECStore {
|
||||
let cancel_rx1 = cancel.clone();
|
||||
let cancel_rx1_for_err = cancel_rx1.clone();
|
||||
let err_tx1 = err_tx.clone();
|
||||
let job1 = tokio::spawn(async move {
|
||||
let mut opts = opts;
|
||||
opts.stop_disk_at_limit = true;
|
||||
if let Err(err) = store.list_merged(cancel_rx1, opts, sender).await
|
||||
&& !cancel_rx1_for_err.is_cancelled()
|
||||
{
|
||||
error!("list_merged err {:?}", err);
|
||||
let _ = err_tx1.send(Arc::new(err));
|
||||
let job1 = tokio::spawn(
|
||||
async move {
|
||||
let mut opts = opts;
|
||||
opts.stop_disk_at_limit = true;
|
||||
if let Err(err) = store.list_merged(cancel_rx1, opts, sender).await
|
||||
&& !cancel_rx1_for_err.is_cancelled()
|
||||
{
|
||||
error!("list_merged err {:?}", err);
|
||||
let _ = err_tx1.send(Arc::new(err));
|
||||
}
|
||||
}
|
||||
});
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
|
||||
let cancel_rx2 = cancel.clone();
|
||||
|
||||
let (result_tx, mut result_rx) = mpsc::channel(1);
|
||||
let err_tx2 = err_tx.clone();
|
||||
let opts = o.clone();
|
||||
let job2 = tokio::spawn(async move {
|
||||
if let Err(err) = gather_results(cancel_rx2, opts, recv, result_tx).await {
|
||||
error!("gather_results err {:?}", err);
|
||||
let _ = err_tx2.send(Arc::new(err));
|
||||
}
|
||||
let job2 = tokio::spawn(
|
||||
async move {
|
||||
if let Err(err) = gather_results(cancel_rx2, opts, recv, result_tx).await {
|
||||
error!("gather_results err {:?}", err);
|
||||
let _ = err_tx2.send(Arc::new(err));
|
||||
}
|
||||
|
||||
// cancel call exit spawns
|
||||
cancel.cancel();
|
||||
});
|
||||
// cancel call exit spawns
|
||||
cancel.cancel();
|
||||
}
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
|
||||
let mut result = {
|
||||
// receiver result
|
||||
@@ -815,11 +821,14 @@ impl ECStore {
|
||||
}
|
||||
}
|
||||
|
||||
tokio::spawn(async move {
|
||||
if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await {
|
||||
error!("merge_entry_channels err {:?}", err)
|
||||
tokio::spawn(
|
||||
async move {
|
||||
if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await {
|
||||
error!("merge_entry_channels err {:?}", err)
|
||||
}
|
||||
}
|
||||
});
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
|
||||
// let merge_res = merge_entry_channels(rx, inputs, sender.clone(), 1).await;
|
||||
|
||||
@@ -1007,34 +1016,47 @@ impl ECStore {
|
||||
Err(_) => None,
|
||||
};
|
||||
|
||||
tokio::spawn(async move {
|
||||
let mut sent_err = false;
|
||||
while let Some(entry) = merge_rx.recv().await {
|
||||
if opts.latest_only {
|
||||
let fi = match entry.to_fileinfo(&bucket_clone) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
if !sent_err {
|
||||
tokio::spawn(
|
||||
async move {
|
||||
let mut sent_err = false;
|
||||
while let Some(entry) = merge_rx.recv().await {
|
||||
if opts.latest_only {
|
||||
let fi = match entry.to_fileinfo(&bucket_clone) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
if !sent_err {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: None,
|
||||
err: Some(err.into()),
|
||||
};
|
||||
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
|
||||
sent_err = true;
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(filter) = opts.filter {
|
||||
if filter(&fi) {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: None,
|
||||
err: Some(err.into()),
|
||||
item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, {
|
||||
if let Some(v) = &vcf { v.versioned(&fi.name) } else { false }
|
||||
})),
|
||||
err: None,
|
||||
};
|
||||
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
|
||||
sent_err = true;
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(filter) = opts.filter {
|
||||
if filter(&fi) {
|
||||
} else {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, {
|
||||
if let Some(v) = &vcf { v.versioned(&fi.name) } else { false }
|
||||
@@ -1046,43 +1068,43 @@ impl ECStore {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: Some(ObjectInfo::from_file_info(&fi, &bucket_clone, &fi.name, {
|
||||
if let Some(v) = &vcf { v.versioned(&fi.name) } else { false }
|
||||
})),
|
||||
err: None,
|
||||
};
|
||||
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
let fvs = match entry.file_info_versions(&bucket_clone) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: None,
|
||||
err: Some(err.into()),
|
||||
};
|
||||
let fvs = match entry.file_info_versions(&bucket_clone) {
|
||||
Ok(res) => res,
|
||||
Err(err) => {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: None,
|
||||
err: Some(err.into()),
|
||||
};
|
||||
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
return;
|
||||
}
|
||||
return;
|
||||
};
|
||||
|
||||
if opts.versions_sort == WalkVersionsSortOrder::Ascending {
|
||||
//TODO: SORT
|
||||
}
|
||||
};
|
||||
|
||||
if opts.versions_sort == WalkVersionsSortOrder::Ascending {
|
||||
//TODO: SORT
|
||||
}
|
||||
for fi in fvs.versions.iter() {
|
||||
if let Some(filter) = opts.filter {
|
||||
if filter(fi) {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, {
|
||||
if let Some(v) = &vcf { v.versioned(&fi.name) } else { false }
|
||||
})),
|
||||
err: None,
|
||||
};
|
||||
|
||||
for fi in fvs.versions.iter() {
|
||||
if let Some(filter) = opts.filter {
|
||||
if filter(fi) {
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, {
|
||||
if let Some(v) = &vcf { v.versioned(&fi.name) } else { false }
|
||||
@@ -1094,23 +1116,13 @@ impl ECStore {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
let item = ObjectInfoOrErr {
|
||||
item: Some(ObjectInfo::from_file_info(fi, &bucket_clone, &fi.name, {
|
||||
if let Some(v) = &vcf { v.versioned(&fi.name) } else { false }
|
||||
})),
|
||||
err: None,
|
||||
};
|
||||
|
||||
if let Err(err) = result.send(item).await {
|
||||
error!("walk result send err {:?}", err);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
.instrument(tracing::Span::current()),
|
||||
);
|
||||
|
||||
tokio::spawn(async move { merge_entry_channels(rx, inputs, merge_tx, 1).await });
|
||||
tokio::spawn(async move { merge_entry_channels(rx, inputs, merge_tx, 1).await }.instrument(tracing::Span::current()));
|
||||
|
||||
let walk_results = join_all(futures).await;
|
||||
let mut errs = Vec::new();
|
||||
|
||||
Reference in New Issue
Block a user