mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 04:21:35 +00:00
fix(rpc): preserve walk_dir missing-path errors
Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -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;
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user