feat(internode): harden p0 transport boundary and baseline tooling (#3017)

* feat(internode): p0 transport baseline and ci hardening

* fix(internode): avoid double wrapping transport errors

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Henry Guo
2026-05-20 19:06:26 +08:00
committed by GitHub
parent c684438625
commit 19b69abe5c
4 changed files with 366 additions and 82 deletions
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::disk::error::Result;
use crate::disk::error::{Error, Result};
use crate::disk::{FileReader, FileWriter};
use crate::rpc::build_auth_headers;
use async_trait::async_trait;
@@ -21,9 +21,43 @@ use rustfs_config::{DEFAULT_INTERNODE_DATA_TRANSPORT, ENV_RUSTFS_INTERNODE_DATA_
use rustfs_rio::{HttpReader, HttpWriter};
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use tracing::warn;
static INTERNODE_DATA_TRANSPORT: OnceLock<Arc<dyn InternodeDataTransport>> = OnceLock::new();
pub const INTERNODE_DATA_TRANSPORT_TCP: &str = "tcp";
static INTERNODE_DATA_TRANSPORT: OnceLock<std::result::Result<Arc<dyn InternodeDataTransport>, String>> = OnceLock::new();
const READ_FILE_STREAM_PATH: &str = "/rustfs/rpc/read_file_stream";
const PUT_FILE_STREAM_PATH: &str = "/rustfs/rpc/put_file_stream";
const WALK_DIR_PATH: &str = "/rustfs/rpc/walk_dir";
const CONTENT_TYPE_JSON: &str = "application/json";
fn unsupported_transport_message(transport: &str) -> String {
format!(
"invalid {ENV_RUSTFS_INTERNODE_DATA_TRANSPORT}={transport:?}; supported values: {DEFAULT_INTERNODE_DATA_TRANSPORT}, {INTERNODE_DATA_TRANSPORT_TCP}"
)
}
#[derive(Debug, Clone, Copy, Eq, PartialEq)]
pub struct InternodeDataTransportCapabilities {
pub stream_read: bool,
pub stream_write: bool,
pub walk_dir: bool,
pub registered_memory: bool,
pub scatter_gather: bool,
pub zero_copy_receive: bool,
}
impl InternodeDataTransportCapabilities {
pub const fn tcp_http() -> Self {
Self {
stream_read: true,
stream_write: true,
walk_dir: true,
registered_memory: false,
scatter_gather: false,
zero_copy_receive: false,
}
}
}
#[derive(Debug, Clone)]
pub struct ReadStreamRequest {
@@ -59,6 +93,7 @@ pub trait InternodeDataTransport: Send + Sync + std::fmt::Debug {
async fn open_write(&self, request: WriteStreamRequest) -> Result<FileWriter>;
async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result<FileReader>;
fn name(&self) -> &'static str;
fn capabilities(&self) -> InternodeDataTransportCapabilities;
}
#[derive(Debug, Default)]
@@ -67,41 +102,22 @@ pub struct TcpHttpInternodeDataTransport;
#[async_trait]
impl InternodeDataTransport for TcpHttpInternodeDataTransport {
async fn open_read(&self, request: ReadStreamRequest) -> Result<FileReader> {
let url = format!(
"{}/rustfs/rpc/read_file_stream?disk={}&volume={}&path={}&offset={}&length={}",
request.endpoint,
urlencoding::encode(&request.disk),
urlencoding::encode(&request.volume),
urlencoding::encode(&request.path),
request.offset,
request.length
);
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
let url = build_read_file_stream_url(&request);
let mut headers = json_headers();
build_auth_headers(&url, &Method::GET, &mut headers)?;
Ok(Box::new(HttpReader::new(url, Method::GET, headers, None).await?))
}
async fn open_write(&self, request: WriteStreamRequest) -> Result<FileWriter> {
let url = format!(
"{}/rustfs/rpc/put_file_stream?disk={}&volume={}&path={}&append={}&size={}",
request.endpoint,
urlencoding::encode(&request.disk),
urlencoding::encode(&request.volume),
urlencoding::encode(&request.path),
request.append,
request.size
);
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
let url = build_put_file_stream_url(&request);
let mut headers = json_headers();
build_auth_headers(&url, &Method::PUT, &mut headers)?;
Ok(Box::new(HttpWriter::new(url, Method::PUT, headers).await?))
}
async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result<FileReader> {
let url = format!("{}/rustfs/rpc/walk_dir?disk={}", request.endpoint, urlencoding::encode(&request.disk));
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
let url = build_walk_dir_url(&request);
let mut headers = json_headers();
build_auth_headers(&url, &Method::GET, &mut headers)?;
Ok(Box::new(
HttpReader::new_with_stall_timeout(url, Method::GET, headers, Some(request.body), request.stall_timeout).await?,
@@ -111,55 +127,195 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport {
fn name(&self) -> &'static str {
DEFAULT_INTERNODE_DATA_TRANSPORT
}
}
fn build_internode_data_transport(configured_transport: Option<&str>) -> Arc<dyn InternodeDataTransport> {
match configured_transport.map(str::trim).filter(|transport| !transport.is_empty()) {
Some(transport) if transport.eq_ignore_ascii_case(DEFAULT_INTERNODE_DATA_TRANSPORT) => {
Arc::new(TcpHttpInternodeDataTransport)
}
Some(transport) => {
warn!(
env = ENV_RUSTFS_INTERNODE_DATA_TRANSPORT,
requested = %transport,
fallback = DEFAULT_INTERNODE_DATA_TRANSPORT,
"unknown internode data transport, using default backend"
);
Arc::new(TcpHttpInternodeDataTransport)
}
None => Arc::new(TcpHttpInternodeDataTransport),
fn capabilities(&self) -> InternodeDataTransportCapabilities {
InternodeDataTransportCapabilities::tcp_http()
}
}
pub fn build_internode_data_transport_from_env() -> Arc<dyn InternodeDataTransport> {
Arc::clone(
INTERNODE_DATA_TRANSPORT
.get_or_init(|| build_internode_data_transport(std::env::var(ENV_RUSTFS_INTERNODE_DATA_TRANSPORT).ok().as_deref())),
fn build_read_file_stream_url(request: &ReadStreamRequest) -> String {
format!(
"{}{}?disk={}&volume={}&path={}&offset={}&length={}",
request.endpoint,
READ_FILE_STREAM_PATH,
urlencoding::encode(&request.disk),
urlencoding::encode(&request.volume),
urlencoding::encode(&request.path),
request.offset,
request.length
)
}
fn build_put_file_stream_url(request: &WriteStreamRequest) -> String {
format!(
"{}{}?disk={}&volume={}&path={}&append={}&size={}",
request.endpoint,
PUT_FILE_STREAM_PATH,
urlencoding::encode(&request.disk),
urlencoding::encode(&request.volume),
urlencoding::encode(&request.path),
request.append,
request.size
)
}
fn build_walk_dir_url(request: &WalkDirStreamRequest) -> String {
format!("{}{}?disk={}", request.endpoint, WALK_DIR_PATH, urlencoding::encode(&request.disk))
}
fn json_headers() -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert(CONTENT_TYPE, HeaderValue::from_static(CONTENT_TYPE_JSON));
headers
}
fn build_internode_data_transport_result(
configured_transport: Option<&str>,
) -> std::result::Result<Arc<dyn InternodeDataTransport>, String> {
match configured_transport.map(str::trim).filter(|transport| !transport.is_empty()) {
None => Ok(Arc::new(TcpHttpInternodeDataTransport)),
Some(transport)
if transport.eq_ignore_ascii_case(DEFAULT_INTERNODE_DATA_TRANSPORT)
|| transport.eq_ignore_ascii_case(INTERNODE_DATA_TRANSPORT_TCP) =>
{
Ok(Arc::new(TcpHttpInternodeDataTransport))
}
Some(transport) => Err(unsupported_transport_message(transport)),
}
}
pub fn build_internode_data_transport(configured_transport: Option<&str>) -> Result<Arc<dyn InternodeDataTransport>> {
build_internode_data_transport_result(configured_transport).map_err(Error::other)
}
pub fn build_internode_data_transport_from_env() -> Result<Arc<dyn InternodeDataTransport>> {
let configured_transport = std::env::var(ENV_RUSTFS_INTERNODE_DATA_TRANSPORT).ok();
#[cfg(test)]
{
build_internode_data_transport(configured_transport.as_deref())
}
#[cfg(not(test))]
INTERNODE_DATA_TRANSPORT
.get_or_init(|| build_internode_data_transport_result(configured_transport.as_deref()))
.as_ref()
.map(Arc::clone)
.map_err(|err| Error::other(err.clone()))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn tcp_http_capabilities_are_behavior_preserving() {
let transport = TcpHttpInternodeDataTransport;
assert_eq!(transport.name(), DEFAULT_INTERNODE_DATA_TRANSPORT);
assert_eq!(
transport.capabilities(),
InternodeDataTransportCapabilities {
stream_read: true,
stream_write: true,
walk_dir: true,
registered_memory: false,
scatter_gather: false,
zero_copy_receive: false,
}
);
}
#[test]
fn read_file_stream_url_encodes_query_values() {
let url = build_read_file_stream_url(&ReadStreamRequest {
endpoint: "http://node1:9000".to_string(),
disk: "http://node1:9000/data/rustfs0".to_string(),
volume: ".rustfs.sys".to_string(),
path: "pool.bin/../part.1".to_string(),
offset: 7,
length: 11,
});
assert_eq!(
url,
"http://node1:9000/rustfs/rpc/read_file_stream?disk=http%3A%2F%2Fnode1%3A9000%2Fdata%2Frustfs0&volume=.rustfs.sys&path=pool.bin%2F..%2Fpart.1&offset=7&length=11"
);
}
#[test]
fn put_file_stream_url_encodes_query_values() {
let url = build_put_file_stream_url(&WriteStreamRequest {
endpoint: "http://node1:9000".to_string(),
disk: "http://node1:9000/data/rustfs0".to_string(),
volume: "bucket".to_string(),
path: "object/part.1".to_string(),
append: false,
size: 4096,
});
assert_eq!(
url,
"http://node1:9000/rustfs/rpc/put_file_stream?disk=http%3A%2F%2Fnode1%3A9000%2Fdata%2Frustfs0&volume=bucket&path=object%2Fpart.1&append=false&size=4096"
);
}
#[test]
fn walk_dir_url_encodes_disk_ref() {
let url = build_walk_dir_url(&WalkDirStreamRequest {
endpoint: "http://node1:9000".to_string(),
disk: "http://node1:9000/data/rustfs0".to_string(),
body: Vec::new(),
stall_timeout: None,
});
assert_eq!(
url,
"http://node1:9000/rustfs/rpc/walk_dir?disk=http%3A%2F%2Fnode1%3A9000%2Fdata%2Frustfs0"
);
}
#[test]
fn transport_config_defaults_to_tcp_http() {
let transport = build_internode_data_transport(None);
let transport = build_internode_data_transport(None).unwrap();
assert_eq!(transport.name(), DEFAULT_INTERNODE_DATA_TRANSPORT);
}
#[test]
fn transport_config_accepts_tcp_http() {
let transport = build_internode_data_transport(Some("TCP-HTTP"));
fn transport_config_blank_value_falls_back_to_default() {
let transport = build_internode_data_transport(Some(" ")).unwrap();
assert_eq!(transport.name(), DEFAULT_INTERNODE_DATA_TRANSPORT);
}
#[test]
fn transport_config_falls_back_to_tcp_http_for_unknown_backend() {
let transport = build_internode_data_transport(Some("rdma"));
fn transport_config_accepts_tcp_aliases() {
for configured in [
DEFAULT_INTERNODE_DATA_TRANSPORT,
INTERNODE_DATA_TRANSPORT_TCP,
"TCP-HTTP",
"TCP",
] {
let transport = build_internode_data_transport(Some(configured)).unwrap();
assert_eq!(transport.name(), DEFAULT_INTERNODE_DATA_TRANSPORT);
assert_eq!(transport.name(), DEFAULT_INTERNODE_DATA_TRANSPORT);
}
}
#[test]
fn transport_config_rejects_unknown_backend() {
let err = build_internode_data_transport(Some("rdma")).expect_err("unknown backend should fail closed");
assert!(err.to_string().contains(ENV_RUSTFS_INTERNODE_DATA_TRANSPORT));
assert!(err.to_string().contains("rdma"));
}
#[test]
fn cached_transport_config_error_uses_raw_message() {
let err = build_internode_data_transport_result(Some("rdma")).expect_err("unknown backend should fail closed");
assert!(!err.starts_with("io error "));
assert!(err.contains(ENV_RUSTFS_INTERNODE_DATA_TRANSPORT));
assert!(err.contains("rdma"));
}
}