diff --git a/crates/ecstore/src/cache_value/metacache_set.rs b/crates/ecstore/src/cache_value/metacache_set.rs index 203bcfc42..c0c97231e 100644 --- a/crates/ecstore/src/cache_value/metacache_set.rs +++ b/crates/ecstore/src/cache_value/metacache_set.rs @@ -1206,6 +1206,17 @@ mod tests { assert_eq!(mixed_timeout, DiskError::Timeout); } + #[test] + fn mixed_missing_and_body_io_errors_remain_an_actionable_quorum_failure() { + let failures = [ + DiskError::FileNotFound, + DiskError::Io(std::io::Error::other("remote body stream aborted")), + ]; + + assert!(!is_benign_not_found_listing_failure(&failures)); + assert_eq!(classify_listing_quorum_failure(&failures), DiskError::ErasureReadQuorum); + } + #[test] fn missing_path_error_classification_excludes_actionable_failures() { assert!(is_missing_path_error(&DiskError::FileNotFound)); diff --git a/crates/ecstore/src/cluster/rpc/internode_data_transport.rs b/crates/ecstore/src/cluster/rpc/internode_data_transport.rs index e1081ffd3..5df74f19a 100644 --- a/crates/ecstore/src/cluster/rpc/internode_data_transport.rs +++ b/crates/ecstore/src/cluster/rpc/internode_data_transport.rs @@ -901,6 +901,7 @@ pub fn build_internode_data_transport_from_env() -> Result listener, + Err(err) if err.kind() == io::ErrorKind::PermissionDenied => return, + Err(err) => panic!("test listener should bind: {err}"), + }; + let addr = listener.local_addr().expect("test listener address should be available"); + let server = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("test server should accept the request"); + let mut request = [0_u8; 2048]; + let _ = socket.read(&mut request).await.expect("request should be readable"); + socket + .write_all( + b"HTTP/1.1 500 Internal Server Error\r\nx-rustfs-disk-error: file-not-found\r\ncontent-length: 0\r\nconnection: close\r\n\r\n", + ) + .await + .expect("typed disk error response should be writable"); + }); + + let endpoint = format!("http://{addr}"); + let error = match TcpHttpInternodeDataTransport + .open_walk_dir(WalkDirStreamRequest { + endpoint: endpoint.clone(), + disk: "disk-a".to_string(), + body: b"{}".to_vec(), + stall_timeout: Some(Duration::from_secs(2)), + }) + .await + { + Ok(_) => panic!("an explicit remote missing-path response must fail the open"), + Err(error) => error, + }; + + assert!(matches!(error, Error::FileNotFound)); + server.await.expect("test server task should finish"); + } + #[test] fn tcp_http_capabilities_are_behavior_preserving() { let transport = TcpHttpInternodeDataTransport; diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index 7fd230caf..3a1e58fe8 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -861,16 +861,19 @@ fn classify_http_response( operation: Option<&'static str>, ) -> ClassifiedHttpResponse { let kind = classify_http_status(status); - if status != reqwest::StatusCode::INTERNAL_SERVER_ERROR || operation != Some(INTERNODE_OPERATION_READ_FILE_STREAM) { + if status != reqwest::StatusCode::INTERNAL_SERVER_ERROR { return ClassifiedHttpResponse { kind, remote_disk_error: None, }; } - let remote_disk_error = match headers.get(INTERNODE_DISK_ERROR_HEADER).and_then(|value| value.to_str().ok()) { - Some(INTERNODE_FILE_NOT_FOUND) => Some(RemoteDiskErrorKind::FileNotFound), - Some(INTERNODE_VOLUME_NOT_FOUND) => Some(RemoteDiskErrorKind::VolumeNotFound), - Some(INTERNODE_FILE_CORRUPT) => Some(RemoteDiskErrorKind::FileCorrupt), + let token = headers.get(INTERNODE_DISK_ERROR_HEADER).and_then(|value| value.to_str().ok()); + let remote_disk_error = match (operation, token) { + (Some(INTERNODE_OPERATION_READ_FILE_STREAM), Some(INTERNODE_FILE_NOT_FOUND)) + | (Some(INTERNODE_OPERATION_WALK_DIR), Some(INTERNODE_FILE_NOT_FOUND)) => Some(RemoteDiskErrorKind::FileNotFound), + (Some(INTERNODE_OPERATION_READ_FILE_STREAM), Some(INTERNODE_VOLUME_NOT_FOUND)) + | (Some(INTERNODE_OPERATION_WALK_DIR), Some(INTERNODE_VOLUME_NOT_FOUND)) => Some(RemoteDiskErrorKind::VolumeNotFound), + (Some(INTERNODE_OPERATION_READ_FILE_STREAM), Some(INTERNODE_FILE_CORRUPT)) => Some(RemoteDiskErrorKind::FileCorrupt), _ => None, }; ClassifiedHttpResponse { kind, remote_disk_error } @@ -2126,8 +2129,18 @@ mod tests { let wrong_status = classify_http_response(reqwest::StatusCode::NOT_FOUND, &headers, Some(INTERNODE_OPERATION_READ_FILE_STREAM)); assert!(wrong_status.remote_disk_error.is_none()); - let wrong_operation = + let walk_dir_missing = classify_http_response(reqwest::StatusCode::INTERNAL_SERVER_ERROR, &headers, Some(INTERNODE_OPERATION_WALK_DIR)); + assert_eq!(walk_dir_missing.remote_disk_error, Some(RemoteDiskErrorKind::FileNotFound)); + headers.insert(INTERNODE_DISK_ERROR_HEADER, INTERNODE_VOLUME_NOT_FOUND.parse().unwrap()); + let walk_dir_volume_missing = + classify_http_response(reqwest::StatusCode::INTERNAL_SERVER_ERROR, &headers, Some(INTERNODE_OPERATION_WALK_DIR)); + assert_eq!(walk_dir_volume_missing.remote_disk_error, Some(RemoteDiskErrorKind::VolumeNotFound)); + let wrong_operation = classify_http_response( + reqwest::StatusCode::INTERNAL_SERVER_ERROR, + &headers, + Some(INTERNODE_OPERATION_PUT_FILE_STREAM), + ); assert!(wrong_operation.remote_disk_error.is_none()); } @@ -2150,8 +2163,8 @@ mod tests { } for operation in [ None, - Some(INTERNODE_OPERATION_WALK_DIR), Some(INTERNODE_OPERATION_PUT_FILE_STREAM), + Some(INTERNODE_OPERATION_WALK_DIR), ] { assert!( classify_http_response(reqwest::StatusCode::INTERNAL_SERVER_ERROR, &headers, operation) diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index 302ed12d9..f7231268d 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -812,7 +812,12 @@ async fn handle_walk_dir(req: Request) -> Response { let log_limit = args.limit; let log_disk_id = args.disk_id.clone(); let log_skip_total_timeout = args.skip_total_timeout; - let body = walk_dir_response_body(propagate_completion_errors, move |mut writer| async move { + // Only callers requesting missing-path reporting need a typed pre-stream status. + let preflight_missing_path_error = propagate_completion_errors && args.report_notfound; + runtime_sources::current_internode_metrics() + .record_incoming_request_for_operation_and_backend(INTERNODE_OPERATION_WALK_DIR, INTERNODE_TRANSPORT_BACKEND_TCP_HTTP); + + let body = walk_dir_response_body(propagate_completion_errors, preflight_missing_path_error, move |mut writer| async move { walk_dir_with_legacy_meta_fallback(&disk, args, &mut writer) .await .map_err(|e| { @@ -835,12 +840,15 @@ async fn handle_walk_dir(req: Request) -> Response { error = %e, "internode rpc background task failed" ); - io::Error::other("remote walk_dir failed") + e }) - }); + }) + .await; - runtime_sources::current_internode_metrics() - .record_incoming_request_for_operation_and_backend(INTERNODE_OPERATION_WALK_DIR, INTERNODE_TRANSPORT_BACKEND_TCP_HTTP); + let body = match body { + Ok(body) => body, + Err(error) => return response_with_disk_error(&error, error.to_string()), + }; Response::builder() .status(StatusCode::OK) @@ -1198,10 +1206,14 @@ fn remote_scanner_claim_rejection(error: &rustfs_scanner::ScannerError) -> (Stat } } -fn walk_dir_response_body(propagate_completion_errors: bool, producer: F) -> Body +async fn walk_dir_response_body( + propagate_completion_errors: bool, + preflight_missing_path_error: bool, + producer: F, +) -> Result where F: FnOnce(tokio::io::DuplexStream) -> Fut + Send + 'static, - Fut: Future> + Send + 'static, + Fut: Future> + Send + 'static, { let (reader, writer) = tokio::io::duplex(DEFAULT_READ_BUFFER_SIZE); let (mut completion_tx, completion_rx) = oneshot::channel(); @@ -1224,14 +1236,45 @@ where ); bytes }); - let stream = append_walk_dir_completion(stream, completion_rx, propagate_completion_errors); + let mut stream = Box::pin(stream); + if !preflight_missing_path_error { + let stream = append_walk_dir_completion(stream, completion_rx, propagate_completion_errors); + return Ok(Body::from(StreamingBlob::wrap(stream))); + } + // Keep the first chunk bounded in memory so a missing-path error can use + // the typed HTTP error headers before the response status is committed. + match stream.next().await { + Some(Ok(first_bytes)) => { + let stream = stream::once(async move { Ok(first_bytes) }).chain(stream); + let stream = append_walk_dir_completion(stream, completion_rx, propagate_completion_errors); + Ok(Body::from(StreamingBlob::wrap(stream))) + } + Some(Err(first_error)) => { + let stream = stream::once(async move { Err(first_error) }).chain(stream); + let stream = append_walk_dir_completion(stream, completion_rx, propagate_completion_errors); + Ok(Body::from(StreamingBlob::wrap(stream))) + } + None => match completion_rx.await { + Ok(Ok(())) => Ok(Body::empty()), + Ok(Err(error)) if propagate_completion_errors => match error { + DiskError::FileNotFound | DiskError::VolumeNotFound => Err(error), + _ => Ok(walk_dir_error_body("remote walk_dir failed")), + }, + Err(_) if propagate_completion_errors => Ok(walk_dir_error_body("remote walk_dir task ended without a result")), + Ok(Err(_)) | Err(_) => Ok(Body::empty()), + }, + } +} + +fn walk_dir_error_body(message: &'static str) -> Body { + let stream = stream::once(async move { Err(io::Error::other(message)) }); Body::from(StreamingBlob::wrap(stream)) } fn append_walk_dir_completion( stream: S, - completion_rx: oneshot::Receiver>, + completion_rx: oneshot::Receiver>, propagate_completion_errors: bool, ) -> impl Stream> where @@ -1241,9 +1284,9 @@ where stream::once(async move { match completion_rx.await { Ok(Ok(())) => None, - Ok(Err(err)) if propagate_completion_errors => Some(Err(err)), - Err(err) if propagate_completion_errors => { - Some(Err(io::Error::other(format!("remote walk_dir task ended without a result: {err}")))) + Ok(Err(_)) if propagate_completion_errors => Some(Err(io::Error::other("remote walk_dir failed"))), + Err(_) if propagate_completion_errors => { + Some(Err(io::Error::other("remote walk_dir task ended without a result"))) } Ok(Err(_)) | Err(_) => None, } @@ -2905,10 +2948,12 @@ mod tests { #[tokio::test] async fn walk_dir_body_surfaces_background_failure_after_data() { - let body = walk_dir_response_body(true, |mut writer| async move { + let body = walk_dir_response_body(true, true, |mut writer| async move { writer.write_all(b"partial walk data").await?; - Err(io::Error::other("remote walk_dir failed")) - }); + Err(DiskError::Io(io::Error::other("remote walk_dir failed"))) + }) + .await + .expect("partial output should retain the streaming response"); let err = BodyExt::collect(body) .await .expect_err("failed completion must fail body collection"); @@ -2918,10 +2963,12 @@ mod tests { #[tokio::test] async fn walk_dir_body_preserves_data_after_success() { - let body = walk_dir_response_body(true, |mut writer| async move { + let body = walk_dir_response_body(true, true, |mut writer| async move { writer.write_all(b"complete walk data").await?; Ok(()) - }); + }) + .await + .expect("successful response should be prepared"); let bytes = BodyExt::collect(body) .await .expect("successful completion should preserve the body") @@ -2930,16 +2977,74 @@ mod tests { assert_eq!(bytes, Bytes::from_static(b"complete walk data")); } + #[tokio::test] + async fn walk_dir_pre_stream_missing_errors_remain_typed_for_capable_peers() { + let error = match walk_dir_response_body(true, true, |_writer| async { Err(DiskError::FileNotFound) }).await { + Ok(_) => panic!("capable peers should receive a pre-stream missing-path result"), + Err(error) => error, + }; + assert!(matches!(error, DiskError::FileNotFound)); + + let body = walk_dir_response_body(false, false, |_writer| async { Err(DiskError::FileNotFound) }) + .await + .expect("legacy peers should retain their clean-EOF behavior"); + assert!( + BodyExt::collect(body) + .await + .expect("legacy missing-path response should remain clean EOF") + .to_bytes() + .is_empty() + ); + + let body = walk_dir_response_body(true, true, |_writer| async { Ok(()) }) + .await + .expect("an empty successful walk should complete during preflight"); + assert!( + BodyExt::collect(body) + .await + .expect("empty successful walk should remain a complete body") + .to_bytes() + .is_empty() + ); + } + + #[tokio::test] + async fn capable_walk_without_missing_report_streams_before_producer_finishes() { + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let (release_tx, release_rx) = tokio::sync::oneshot::channel(); + let body = tokio::time::timeout( + std::time::Duration::from_secs(1), + walk_dir_response_body(true, false, move |mut writer| async move { + let _ = started_tx.send(()); + release_rx.await.expect("test should release the walk producer"); + writer.write_all(b"partial walk data").await?; + Err(DiskError::Io(io::Error::other("remote walk_dir failed"))) + }), + ) + .await + .expect("ordinary v1 streams must return before the producer's first output") + .expect("ordinary v1 stream response should be prepared"); + + started_rx.await.expect("walk producer should start"); + release_tx.send(()).expect("walk producer should be released"); + let error = BodyExt::collect(body) + .await + .expect_err("terminal producer errors must still fail ordinary v1 streams"); + assert!(error.to_string().contains("remote walk_dir failed")); + } + #[tokio::test] async fn walk_dir_body_records_operation_sent_bytes() { let metrics = global_internode_metrics(); let before = metrics.snapshot().sent_bytes_total; let payload = Bytes::from_static(b"metered walk data"); let expected_len = u64::try_from(payload.len()).expect("test payload length should fit u64"); - let body = walk_dir_response_body(true, move |mut writer| async move { + let body = walk_dir_response_body(true, true, move |mut writer| async move { writer.write_all(&payload).await?; Ok(()) - }); + }) + .await + .expect("successful response should be prepared"); let bytes = BodyExt::collect(body) .await @@ -2972,11 +3077,13 @@ mod tests { async fn dropping_walk_dir_body_cancels_blocked_producer() { let (started_tx, started_rx) = tokio::sync::oneshot::channel(); let (dropped_tx, dropped_rx) = tokio::sync::oneshot::channel(); - let body = walk_dir_response_body(true, move |_writer| async move { + let body = walk_dir_response_body(true, false, move |_writer| async move { let _drop_notifier = DropNotifier(Some(dropped_tx)); let _ = started_tx.send(()); - std::future::pending::>().await - }); + std::future::pending::>().await + }) + .await + .expect("ordinary streams should start without a first-chunk preflight"); started_rx.await.expect("walk producer should start"); drop(body); @@ -2989,10 +3096,12 @@ mod tests { #[tokio::test] async fn legacy_walk_dir_client_keeps_clean_eof_compatibility() { - let body = walk_dir_response_body(false, |mut writer| async move { + let body = walk_dir_response_body(false, false, |mut writer| async move { writer.write_all(b"legacy partial data").await?; - Err(io::Error::other("remote walk_dir failed")) - }); + Err(DiskError::Io(io::Error::other("remote walk_dir failed"))) + }) + .await + .expect("legacy response should preserve clean EOF compatibility"); let bytes = BodyExt::collect(body) .await