From a28cb343816e253be9922a1a51884e23d6bc9864 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Sat, 23 May 2026 17:41:25 +0800 Subject: [PATCH] test(internode): cover RemoteDisk adapter routing (#3070) Co-authored-by: Henry Guo --- .../internode-transport-adapter-rfc.md | 29 ++- crates/ecstore/src/rpc/remote_disk.rs | 194 +++++++++++++++++- 2 files changed, 204 insertions(+), 19 deletions(-) diff --git a/crates/ecstore/docs/internode-transport/internode-transport-adapter-rfc.md b/crates/ecstore/docs/internode-transport/internode-transport-adapter-rfc.md index 94544624a..50083af2f 100644 --- a/crates/ecstore/docs/internode-transport/internode-transport-adapter-rfc.md +++ b/crates/ecstore/docs/internode-transport/internode-transport-adapter-rfc.md @@ -381,31 +381,27 @@ The initial TCP backend can keep the current signed HTTP URLs internally. `RemoteDisk` delegates only these methods to the data transport: - `read_file_stream` -- `read_file_zero_copy` as a wrapper over `read_file_stream` unless the backend - supports a stronger zero-copy API +- `read_file_zero_copy` as the current wrapper over `read_file_stream` - `append_file` - `create_file` - `walk_dir` All other `RemoteDisk` methods continue using the current gRPC client -until measurements prove otherwise. +in this adapter scope. ### Capability model Avoid hard-coding transport-specific assumptions into the generic interface. -Use capabilities: +The current conservative capability fields are: -- stream read -- stream write -- bounded range read -- bidirectional streaming -- backend-specific buffer registration or staging requirements -- stable buffer ownership support -- copy-reduced receive into caller-owned or backend-owned buffers -- authenticated out-of-band transfer -- transport fallback support +- streaming read +- streaming write +- streaming walk-dir +- ordered delivery +- maximum transfer size +- fallback support -The first TCP backend should report only capabilities that it actually provides. +The TCP/HTTP backend should report only capabilities that it actually provides. ## TCP Fallback Requirements @@ -474,8 +470,9 @@ The current adapter boundary has these constraints: - Metrics must identify the selected backend and operation without high-cardinality labels. -`walk_dir`, metadata RPCs, locks, admin RPCs, and bucket coordination remain -outside the current data-plane boundary. +The current `RemoteDisk::walk_dir` stream is routed through the adapter. +Metadata RPCs, locks, admin RPCs, bucket coordination, and the legacy gRPC +`WalkDir` handler remain outside the current data-plane boundary. ## Out of Scope diff --git a/crates/ecstore/src/rpc/remote_disk.rs b/crates/ecstore/src/rpc/remote_disk.rs index 5a8844bf1..a47dc6350 100644 --- a/crates/ecstore/src/rpc/remote_disk.rs +++ b/crates/ecstore/src/rpc/remote_disk.rs @@ -1703,10 +1703,12 @@ impl DiskAPI for RemoteDisk { #[cfg(test)] mod tests { use super::*; - use crate::rpc::TcpHttpInternodeDataTransport; + use crate::rpc::{InternodeDataTransportCapabilities, TcpHttpInternodeDataTransport}; use rustfs_common::GLOBAL_CONN_MAP; - use std::sync::Once; - use tokio::io::duplex; + use std::pin::Pin; + use std::sync::{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; @@ -1714,6 +1716,96 @@ mod tests { static INIT: Once = Once::new(); + #[derive(Debug, Clone)] + enum RecordedTransportCall { + Read(ReadStreamRequest), + Write(WriteStreamRequest), + WalkDir(WalkDirStreamRequest), + } + + #[derive(Debug, Clone, Default)] + struct RecordingInternodeDataTransport { + calls: Arc>>, + } + + impl RecordingInternodeDataTransport { + 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; + + impl AsyncRead for EmptyTestReader { + fn poll_read(self: Pin<&mut Self>, _cx: &mut Context<'_>, _buf: &mut ReadBuf<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + } + + #[derive(Debug, Default)] + struct SinkTestWriter; + + impl AsyncWrite for SinkTestWriter { + fn poll_write(self: Pin<&mut Self>, _cx: &mut Context<'_>, buf: &[u8]) -> Poll> { + Poll::Ready(Ok(buf.len())) + } + + fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + } + + #[async_trait::async_trait] + impl InternodeDataTransport for RecordingInternodeDataTransport { + async fn open_read(&self, request: ReadStreamRequest) -> Result { + self.record(RecordedTransportCall::Read(request)); + Ok(Box::new(EmptyTestReader)) + } + + async fn open_write(&self, request: WriteStreamRequest) -> Result { + self.record(RecordedTransportCall::Write(request)); + Ok(Box::new(SinkTestWriter)) + } + + async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result { + self.record(RecordedTransportCall::WalkDir(request)); + Ok(Box::new(EmptyTestReader)) + } + + fn name(&self) -> &'static str { + "recording" + } + + 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(), + is_local: false, + pool_idx: 0, + set_idx: 0, + disk_idx: 0, + }; + let disk_option = DiskOption { + cleanup: false, + health_check: false, + }; + + RemoteDisk::new(&endpoint, &disk_option, data_transport).await.unwrap() + } + fn init_tracing(filter_level: Level) { INIT.call_once(|| { let _ = tracing_subscriber::fmt() @@ -1957,6 +2049,102 @@ mod tests { assert_eq!(remote_disk.disk_ref().await, disk_id.to_string()); } + #[tokio::test] + async fn test_remote_disk_read_file_stream_uses_configured_data_transport() { + let transport = RecordingInternodeDataTransport::default(); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + let expected_disk = remote_disk.disk_ref().await; + + let _reader = remote_disk.read_file_stream("bucket", "object/part.1", 7, 11).await.unwrap(); + + let calls = transport.calls(); + assert_eq!(calls.len(), 1); + match &calls[0] { + RecordedTransportCall::Read(request) => { + assert_eq!(request.endpoint, "http://remote-node:9000"); + assert_eq!(request.disk, expected_disk); + assert_eq!(request.volume, "bucket"); + assert_eq!(request.path, "object/part.1"); + assert_eq!(request.offset, 7); + assert_eq!(request.length, 11); + } + other => panic!("expected read transport call, got {other:?}"), + } + } + + #[tokio::test] + async fn test_remote_disk_create_and_append_file_use_configured_data_transport() { + let transport = RecordingInternodeDataTransport::default(); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + let expected_disk = remote_disk.disk_ref().await; + + let _created = remote_disk + .create_file("orig-bucket", "bucket", "object/part.1", 4096) + .await + .unwrap(); + let _appended = remote_disk.append_file("bucket", "object/part.2").await.unwrap(); + + let calls = transport.calls(); + assert_eq!(calls.len(), 2); + + match &calls[0] { + RecordedTransportCall::Write(request) => { + assert_eq!(request.endpoint, "http://remote-node:9000"); + assert_eq!(request.disk, expected_disk); + assert_eq!(request.volume, "bucket"); + assert_eq!(request.path, "object/part.1"); + assert!(!request.append); + assert_eq!(request.size, 4096); + } + other => panic!("expected create write transport call, got {other:?}"), + } + + match &calls[1] { + RecordedTransportCall::Write(request) => { + assert_eq!(request.endpoint, "http://remote-node:9000"); + assert_eq!(request.disk, expected_disk); + assert_eq!(request.volume, "bucket"); + assert_eq!(request.path, "object/part.2"); + assert!(request.append); + assert_eq!(request.size, 0); + } + other => panic!("expected append write transport call, got {other:?}"), + } + } + + #[tokio::test] + async fn test_remote_disk_walk_dir_uses_configured_data_transport() { + let transport = RecordingInternodeDataTransport::default(); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await; + let expected_disk = remote_disk.disk_ref().await; + let opts = WalkDirOptions { + bucket: "bucket".to_string(), + base_dir: "prefix".to_string(), + recursive: true, + report_notfound: false, + filter_prefix: Some("part".to_string()), + forward_to: None, + limit: 10, + disk_id: String::new(), + }; + let expected_body = serde_json::to_vec(&opts).unwrap(); + let mut writer = Vec::new(); + + remote_disk.walk_dir(opts, &mut writer).await.unwrap(); + + let calls = transport.calls(); + assert_eq!(calls.len(), 1); + match &calls[0] { + RecordedTransportCall::WalkDir(request) => { + assert_eq!(request.endpoint, "http://remote-node:9000"); + assert_eq!(request.disk, expected_disk); + assert_eq!(request.body, expected_body); + assert_eq!(request.stall_timeout, Some(get_drive_walkdir_stall_timeout())); + } + other => panic!("expected walk-dir transport call, got {other:?}"), + } + } + #[tokio::test] async fn test_remote_disk_endpoints_with_different_schemes() { let test_cases = vec![