diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 9ce5ec1fd..f9c91cd91 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -110,6 +110,15 @@ pub struct RemoteDisk { } impl RemoteDisk { + fn is_retryable_walk_dir_error(err: &DiskError) -> bool { + if is_network_like_disk_error(err) { + return true; + } + + let err_text = err.to_string().to_ascii_lowercase(); + err_text.contains("httpreader stream error") || err_text.contains("error decoding response body") + } + pub(crate) async fn new(ep: &Endpoint, opt: &DiskOption, data_transport: Arc) -> Result { let addr = if let Some(port) = ep.url.port() { format!("{}://{}:{}", ep.url.scheme(), ep.url.host_str().unwrap(), port) @@ -1216,24 +1225,58 @@ impl DiskAPI for RemoteDisk { async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()> { info!("walk_dir {}", self.endpoint.to_string()); + let disk = self.disk_ref().await; + let body = serde_json::to_vec(&opts)?; + let stall_timeout = get_drive_walkdir_stall_timeout(); + let bucket = opts.bucket.clone(); + let base_dir = opts.base_dir.clone(); + let disk_for_log = disk.clone(); + self.execute_with_timeout_for_op_and_health_action( "walk_dir", || async { - let disk = self.disk_ref().await; - let opts = serde_json::to_vec(&opts)?; - let mut reader = self - .data_transport - .open_walk_dir(WalkDirStreamRequest { - endpoint: self.endpoint.grid_host(), - disk, - body: opts, - stall_timeout: Some(get_drive_walkdir_stall_timeout()), - }) - .await?; + let mut last_err = None; - copy_stream_with_buffer(&mut reader, wr, DEFAULT_READ_BUFFER_SIZE).await?; + for attempt in 1..=2 { + let mut reader = match self + .data_transport + .open_walk_dir(WalkDirStreamRequest { + endpoint: self.endpoint.grid_host(), + disk: disk.clone(), + body: body.clone(), + stall_timeout: Some(stall_timeout), + }) + .await + { + Ok(reader) => reader, + Err(err) => { + if attempt == 1 && Self::is_retryable_walk_dir_error(&err) { + warn!( + endpoint = %self.endpoint, + addr = %self.addr, + disk = %disk_for_log, + bucket = %bucket, + base_dir = %base_dir, + attempt, + stall_timeout_ms = stall_timeout.as_millis(), + error = %err, + "remote walk_dir returned retryable transport error; retrying" + ); + last_err = Some(err); + continue; + } - Ok(()) + return Err(err); + } + }; + + match copy_stream_with_buffer(&mut reader, wr, DEFAULT_READ_BUFFER_SIZE).await { + Ok(_) => return Ok(()), + Err(io_err) => return Err(DiskError::Io(io_err)), + } + } + + Err(last_err.unwrap_or_else(|| DiskError::other("walk_dir retry exhausted without captured error"))) }, get_drive_walkdir_timeout(), FailureHealthAction::IgnoreFailure, @@ -1759,6 +1802,36 @@ mod tests { } } + #[derive(Debug)] + enum WalkDirTestStep { + Error(DiskError), + Data(Vec), + PartialDataThenError { data: Vec, error: io::Error }, + } + + #[derive(Debug, Clone, Default)] + struct RetryingWalkDirInternodeDataTransport { + calls: Arc>>, + steps: Arc>>, + } + + impl RetryingWalkDirInternodeDataTransport { + fn with_steps(steps: Vec) -> Self { + Self { + calls: Arc::new(StdMutex::new(Vec::new())), + steps: Arc::new(StdMutex::new(steps)), + } + } + + fn calls(&self) -> Vec { + self.calls.lock().expect("recorded transport calls lock poisoned").clone() + } + + fn record(&self, call: RecordedTransportCall) { + self.calls.lock().expect("recorded transport calls lock poisoned").push(call); + } + } + #[derive(Debug, Default)] struct EmptyTestReader; @@ -1811,6 +1884,38 @@ mod tests { } } + #[async_trait::async_trait] + impl InternodeDataTransport for RetryingWalkDirInternodeDataTransport { + async fn open_read(&self, _request: ReadStreamRequest) -> Result { + panic!("open_read should not be used in walk_dir retry test"); + } + + async fn open_write(&self, _request: WriteStreamRequest) -> Result { + panic!("open_write should not be used in walk_dir retry test"); + } + + async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result { + self.record(RecordedTransportCall::WalkDir(request)); + let step = self.steps.lock().expect("walk_dir retry steps lock poisoned").remove(0); + match step { + WalkDirTestStep::Error(err) => Err(err), + WalkDirTestStep::Data(data) => Ok(Box::new(Cursor::new(data))), + WalkDirTestStep::PartialDataThenError { data, error } => Ok(Box::new(PartialThenErrorReader { + cursor: Cursor::new(data), + error: Some(error), + })), + } + } + + fn name(&self) -> &'static str { + "retrying-walk-dir" + } + + fn capabilities(&self) -> InternodeDataTransportCapabilities { + InternodeDataTransportCapabilities::tcp_http() + } + } + async fn new_remote_disk_with_transport(data_transport: Arc) -> RemoteDisk { let endpoint = Endpoint { url: url::Url::parse("http://remote-node:9000/data/rustfs0").unwrap(), @@ -1827,6 +1932,32 @@ mod tests { RemoteDisk::new(&endpoint, &disk_option, data_transport).await.unwrap() } + #[derive(Debug)] + struct PartialThenErrorReader { + cursor: Cursor>, + error: Option, + } + + impl AsyncRead for PartialThenErrorReader { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + let filled_before = buf.filled().len(); + match Pin::new(&mut self.cursor).poll_read(cx, buf) { + Poll::Ready(Ok(())) => { + if buf.filled().len() > filled_before { + return Poll::Ready(Ok(())); + } + + if let Some(err) = self.error.take() { + return Poll::Ready(Err(err)); + } + + Poll::Ready(Ok(())) + } + other => other, + } + } + } + fn init_tracing(filter_level: Level) { INIT.call_once(|| { let _ = tracing_subscriber::fmt() @@ -2185,6 +2316,66 @@ mod tests { } } + #[tokio::test] + async fn test_remote_disk_walk_dir_retries_once_on_retryable_transport_error() { + let transport = RetryingWalkDirInternodeDataTransport::with_steps(vec![ + WalkDirTestStep::Error(DiskError::other("HttpReader stream error: error decoding response body")), + WalkDirTestStep::Data(b"walk-dir-retry-ok".to_vec()), + ]); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + let opts = WalkDirOptions { + bucket: "bucket".to_string(), + base_dir: "config/iam".to_string(), + recursive: true, + report_notfound: false, + filter_prefix: None, + forward_to: None, + limit: 10, + disk_id: String::new(), + }; + let mut writer = Vec::new(); + + remote_disk + .walk_dir(opts, &mut writer) + .await + .expect("retryable walk_dir error should recover"); + + assert_eq!(writer, b"walk-dir-retry-ok"); + assert_eq!(transport.calls().len(), 2, "walk_dir should retry exactly once"); + } + + #[tokio::test] + async fn test_remote_disk_walk_dir_does_not_retry_after_partial_stream_failure() { + let transport = RetryingWalkDirInternodeDataTransport::with_steps(vec![ + WalkDirTestStep::PartialDataThenError { + data: b"partial-walk-dir".to_vec(), + error: io::Error::new(io::ErrorKind::ConnectionReset, "connection reset"), + }, + WalkDirTestStep::Data(b"walk-dir-retry-ok".to_vec()), + ]); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + let opts = WalkDirOptions { + bucket: "bucket".to_string(), + base_dir: "config/iam".to_string(), + recursive: true, + report_notfound: false, + filter_prefix: None, + forward_to: None, + limit: 10, + disk_id: String::new(), + }; + let mut writer = Vec::new(); + + let err = remote_disk + .walk_dir(opts, &mut writer) + .await + .expect_err("partial stream failure should be returned without retry"); + + assert!(matches!(err, DiskError::Io(ref io_err) if io_err.kind() == io::ErrorKind::ConnectionReset)); + assert_eq!(writer, b"partial-walk-dir"); + assert_eq!(transport.calls().len(), 1, "walk_dir should not retry after writing partial bytes"); + } + #[tokio::test] async fn test_remote_disk_endpoints_with_different_schemes() { let test_cases = vec![ diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index a3930330a..956b21dc7 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -248,22 +248,24 @@ impl HttpReader { let resp = request.send().await.map_err(|e| { record_internode_error(track_internode_metrics, internode_operation); - Error::other(format!("HttpReader HTTP request error: {e}")) + Error::other(format!("HttpReader HTTP request error for {method} {url}: {e}")) })?; if resp.status().is_success().not() { record_internode_error(track_internode_metrics, internode_operation); return Err(Error::other(format!( - "HttpReader HTTP request failed with non-200 status {}", - resp.status() + "HttpReader HTTP request failed for {method} {url} with non-200 status {}", + resp.status(), ))); } record_internode_outgoing_request(track_internode_metrics, internode_operation); + let stream_error_url = url.clone(); + let stream_error_method = method.clone(); let stream = resp.bytes_stream().map_err(move |e| { record_internode_error(track_internode_metrics, internode_operation); - Error::other(format!("HttpReader stream error: {e}")) + Error::other(format!("HttpReader stream error for {stream_error_method} {stream_error_url}: {e}")) }); Ok(Self { @@ -837,10 +839,13 @@ mod tests { assert_eq!(&first, b"hello"); let mut next = [0u8; 1]; - let err = tokio::time::timeout(Duration::from_secs(1), reader.read(&mut next)) + let read_result = tokio::time::timeout(Duration::from_secs(1), reader.read(&mut next)) .await - .expect("stall timeout should wake reader") - .expect_err("reader should return a timeout error"); + .expect("stall timeout should wake reader"); + let err = match read_result { + Ok(_) => panic!("reader should return a timeout error"), + Err(err) => err, + }; assert_eq!(err.kind(), io::ErrorKind::TimedOut); handle.abort(); @@ -898,6 +903,23 @@ mod tests { handle.abort(); } + #[tokio::test] + async fn http_reader_request_error_includes_method_and_url() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + drop(listener); + + let url = format!("http://{addr}/stream"); + let err = match HttpReader::new(url.clone(), Method::GET, HeaderMap::new(), None).await { + Ok(_) => panic!("closed listener should trigger request error"), + Err(err) => err, + }; + + let err_text = err.to_string(); + assert!(err_text.contains("HttpReader HTTP request error for GET")); + assert!(err_text.contains(&url)); + } + #[test] fn loopback_urls_bypass_proxy_selection() { assert!(should_bypass_proxy_for_url("http://127.0.0.1:9000/stream"));