fix(rpc): preserve walk_dir missing-path errors (#8226)

## Related Issues

rustfs/rustfs#8217 and rustfs/backlog#2683.

## Summary of Changes

Reviewed 539ed275db against d2175d1e1e. No findings. Capable callers requesting missing-path reporting recover explicit pre-output FileNotFound or VolumeNotFound errors; partial-output and other failures continue to fail the stream.

## Verification

Two independent source reviews completed across all lenses below; the frozen diff matches Git and passes diff whitespace checks. Traced the authenticated handler, bounded first-chunk preflight, cancellation ownership, HTTP status/token classification and quorum error handling. The new tests cover ordinary streaming, typed missing errors, byte preservation and mixed missing/I/O failures.

Author-reported focused tests and Docker reproduction were not rerun or independently audited in this review. The reported Docker comparison predates the final report_notfound-only gate; its scope is not distributed acceptance of the final head.

## Impact

Correctness: no findings; missing errors remain typed only before output. Security/trust: no findings; authentication, body digest and operation/status/token restrictions remain. Compatibility: no findings; legacy clean EOF remains, with the documented old-server limitation. Concurrency/durability: no findings; receiver ownership cancels an abandoned producer. Simplicity: no findings. Coverage: no blocking gap identified. Performance: no findings; preflight retains one bounded chunk and ordinary report_notfound=false streams start immediately.

## Additional Notes

This review does not establish that the patch fixes any current main CI failure. Full main CI and release acceptance remain separate gates.
This commit is contained in:
Hauser
2026-09-29 15:29:59 +08:00
committed by GitHub
parent 0a61f4a81e
commit fcb60e7322
4 changed files with 205 additions and 32 deletions
@@ -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));
@@ -901,6 +901,7 @@ pub fn build_internode_data_transport_from_env() -> Result<Arc<dyn InternodeData
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use tokio::io::AsyncWriteExt;
use tokio::sync::{Barrier, Notify};
async fn wait_for_capability_flight_waiters(entry: &PutFileCapabilityCacheEntry, waiters: usize) {
@@ -1026,6 +1027,45 @@ mod tests {
assert!(matches!(open_err, Error::MethodNotAllowed));
}
#[tokio::test]
async fn walk_dir_restores_explicit_remote_missing_error() {
let _ = rustfs_credentials::set_global_rpc_secret("walk-dir-error-test-secret".to_string());
let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => 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;
+20 -7
View File
@@ -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)
+134 -25
View File
@@ -812,7 +812,12 @@ async fn handle_walk_dir(req: Request<Incoming>) -> Response<Body> {
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<Incoming>) -> Response<Body> {
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<F, Fut>(propagate_completion_errors: bool, producer: F) -> Body
async fn walk_dir_response_body<F, Fut>(
propagate_completion_errors: bool,
preflight_missing_path_error: bool,
producer: F,
) -> Result<Body, DiskError>
where
F: FnOnce(tokio::io::DuplexStream) -> Fut + Send + 'static,
Fut: Future<Output = io::Result<()>> + Send + 'static,
Fut: Future<Output = Result<(), DiskError>> + 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<S>(
stream: S,
completion_rx: oneshot::Receiver<io::Result<()>>,
completion_rx: oneshot::Receiver<Result<(), DiskError>>,
propagate_completion_errors: bool,
) -> impl Stream<Item = io::Result<Bytes>>
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::<io::Result<()>>().await
});
std::future::pending::<Result<(), DiskError>>().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