fix(rpc): authenticate internode put file bodies (#5868)

This commit is contained in:
cxymds
2026-08-09 08:05:16 +08:00
committed by GitHub
parent d36166ffb5
commit 3b9c67e79b
8 changed files with 745 additions and 57 deletions
+352 -14
View File
@@ -16,8 +16,9 @@ use crate::server::RPC_PREFIX;
use crate::storage::request_context::spawn_traced;
use crate::storage::storage_api::DiskError;
use crate::storage::storage_api::rpc_consumer::http_service::{
DEFAULT_READ_BUFFER_SIZE, NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, StorageDiskRpcExt as _,
WALK_DIR_STREAM_COMPLETION_V1, WalkDirOptions, find_local_disk_by_ref, sign_ns_scanner_capability, verify_rpc_signature,
DEFAULT_READ_BUFFER_SIZE, NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, PUT_FILE_AUTH_TRAILER_LEN,
PUT_FILE_AUTH_V1, StorageDiskRpcExt as _, WALK_DIR_STREAM_COMPLETION_V1, WalkDirOptions, check_and_record_signed_rpc_nonce,
find_local_disk_by_ref, sign_ns_scanner_capability, verify_put_file_auth_trailer, verify_rpc_signature,
};
#[cfg(test)]
use crate::storage::storage_api::rpc_consumer::http_service::{
@@ -72,6 +73,12 @@ const NS_SCANNER_PATH: &str = "/rustfs/rpc/ns_scanner";
const NS_SCANNER_REQUEST_BODY_TIMEOUT: Duration = Duration::from_secs(15);
const NS_SCANNER_STREAM_BUFFER_SIZE: usize = 64 * 1024;
static NS_SCANNER_SERVER_EPOCH: LazyLock<uuid::Uuid> = LazyLock::new(uuid::Uuid::new_v4);
static PUT_FILE_AUTH_STRICT: LazyLock<bool> = LazyLock::new(|| {
rustfs_utils::get_env_bool(
rustfs_config::ENV_INTERNODE_RPC_BODY_DIGEST_STRICT,
rustfs_config::DEFAULT_INTERNODE_RPC_BODY_DIGEST_STRICT,
)
});
macro_rules! log_internode_rpc_response_failure {
($status:expr, $rpc_path:expr, $method:expr, $operation:expr, $reason:expr, $result:expr, Some(($context_key:expr, $context_value:expr)), Some($error_text:expr)) => {{
@@ -311,13 +318,34 @@ fn validate_walk_dir_completion_request(query: &WalkDirQuery, body: &[u8]) -> Op
Some(propagate_completion_errors)
}
#[derive(Debug, Default, serde::Deserialize)]
#[derive(Clone, Debug, Default, serde::Deserialize)]
struct PutFileQuery {
disk: String,
volume: String,
path: String,
append: bool,
size: i64,
put_file_auth: Option<String>,
put_file_nonce: Option<uuid::Uuid>,
}
fn put_file_auth_nonce(query: &PutFileQuery) -> io::Result<Option<uuid::Uuid>> {
match query.put_file_auth.as_deref() {
None => {
if *PUT_FILE_AUTH_STRICT {
return Err(io::Error::other("put_file auth required"));
}
Ok(None)
}
Some(PUT_FILE_AUTH_V1) => {
let nonce = query
.put_file_nonce
.filter(|nonce| !nonce.is_nil())
.ok_or_else(|| io::Error::other("Invalid RPC nonce"))?;
Ok(Some(nonce))
}
Some(_) => Err(io::Error::other("Unsupported put_file auth version")),
}
}
impl<S> Service<Request<Incoming>> for InternodeRpcService<S>
@@ -1117,10 +1145,48 @@ where
async fn handle_put_file(req: Request<Incoming>) -> Response<Body> {
let method = req.method().clone();
let path = req.uri().path().to_string();
let url = req.uri().to_string();
let query = match parse_query::<PutFileQuery>(&req) {
Ok(query) => query,
Err(response) => return *response,
};
let auth_nonce = match put_file_auth_nonce(&query) {
Ok(nonce) => nonce,
Err(e) => {
log_internode_rpc_response_failure!(
StatusCode::FORBIDDEN,
&path,
&method,
Some(INTERNODE_OPERATION_PUT_FILE_STREAM),
"put_file_auth_invalid",
"rejected",
Some(("disk", query.disk.as_str())),
Some(&e)
);
return response_with_status(StatusCode::FORBIDDEN, format!("invalid put_file auth: {e}"));
}
};
if let Some(nonce) = auth_nonce
&& let Err(e) = check_and_record_signed_rpc_nonce(
req.headers(),
nonce,
PUT_FILE_STREAM_PATH,
INTERNODE_OPERATION_PUT_FILE_STREAM,
INTERNODE_TRANSPORT_BACKEND_TCP_HTTP,
)
{
log_internode_rpc_response_failure!(
StatusCode::FORBIDDEN,
&path,
&method,
Some(INTERNODE_OPERATION_PUT_FILE_STREAM),
"put_file_replay_rejected",
"rejected",
Some(("disk", query.disk.as_str())),
Some(&e)
);
return response_with_status(StatusCode::FORBIDDEN, format!("invalid put_file auth: {e}"));
}
let Some(disk) = find_local_disk_by_ref(&query.disk).await else {
log_internode_rpc_response_failure!(
@@ -1156,14 +1222,16 @@ async fn handle_put_file(req: Request<Incoming>) -> Response<Body> {
}
};
let copied = match write_body_chunks_to_writer(req.into_body().into_data_stream(), &mut file).await {
Ok(copied) => copied,
Err(e) => {
let message = put_file_stage_error_message("write_body", &query, &e);
log_internode_put_file_stage_failure!("write_body", query, e);
return response_with_status(StatusCode::INTERNAL_SERVER_ERROR, message);
}
};
let copied =
match write_put_file_body_chunks_to_writer(req.into_body().into_data_stream(), &mut file, &query, auth_nonce, &url).await
{
Ok(copied) => copied,
Err(e) => {
let message = put_file_stage_error_message("write_body", &query, &e);
log_internode_put_file_stage_failure!("write_body", query, e);
return response_with_status(StatusCode::INTERNAL_SERVER_ERROR, message);
}
};
let metrics = runtime_sources::current_internode_metrics();
metrics.record_incoming_request_for_operation_and_backend(
@@ -1222,6 +1290,112 @@ where
Ok(copied)
}
async fn write_put_file_body_chunks_to_writer<S, E, W>(
body: S,
writer: &mut W,
query: &PutFileQuery,
auth_nonce: Option<uuid::Uuid>,
url: &str,
) -> io::Result<u64>
where
S: futures::TryStream<Ok = Bytes, Error = E> + Unpin,
E: Into<BoxError>,
W: tokio::io::AsyncWrite + Unpin,
{
let Some(nonce) = auth_nonce else {
return write_body_chunks_to_writer(body, writer).await;
};
let expected_size = (!query.append && query.size >= 0)
.then(|| {
u64::try_from(query.size)
.map_err(|_| io::Error::new(io::ErrorKind::InvalidInput, "put_file auth size cannot be represented"))
})
.transpose()?;
let mut body = body;
let mut remaining = expected_size;
let mut copied = 0_u64;
let mut trailer = Vec::with_capacity(PUT_FILE_AUTH_TRAILER_LEN);
let mut hasher = Sha256::new();
while let Some(bytes) = body.try_next().await.map_err(io::Error::other)? {
if let Some(remaining) = remaining.as_mut() {
let chunk_len = u64::try_from(bytes.len()).unwrap_or(u64::MAX);
let data_len = usize::try_from((*remaining).min(chunk_len))
.map_err(|_| io::Error::other("put_file body length cannot be represented"))?;
if data_len > 0 {
hasher.update(&bytes[..data_len]);
copied = copied
.checked_add(
u64::try_from(data_len).map_err(|_| io::Error::other("put_file body length cannot be represented"))?,
)
.ok_or_else(|| io::Error::other("put_file body length overflow"))?;
*remaining -=
u64::try_from(data_len).map_err(|_| io::Error::other("put_file body length cannot be represented"))?;
writer.write_all(&bytes[..data_len]).await?;
}
if data_len < bytes.len() {
trailer.extend_from_slice(&bytes[data_len..]);
if trailer.len() > PUT_FILE_AUTH_TRAILER_LEN {
return Err(io::Error::new(io::ErrorKind::InvalidData, "put_file auth trailer has trailing data"));
}
}
} else {
let write_len = trailer
.len()
.saturating_add(bytes.len())
.saturating_sub(PUT_FILE_AUTH_TRAILER_LEN);
if write_len > 0 {
let buffered_write_len = write_len.min(trailer.len());
if buffered_write_len > 0 {
writer.write_all(&trailer[..buffered_write_len]).await?;
hasher.update(&trailer[..buffered_write_len]);
copied = copied
.checked_add(
u64::try_from(buffered_write_len)
.map_err(|_| io::Error::other("put_file body length cannot be represented"))?,
)
.ok_or_else(|| io::Error::other("put_file body length overflow"))?;
trailer = trailer.split_off(buffered_write_len);
}
let chunk_write_len = write_len - buffered_write_len;
if chunk_write_len > 0 {
writer.write_all(&bytes[..chunk_write_len]).await?;
hasher.update(&bytes[..chunk_write_len]);
}
copied = copied
.checked_add(
u64::try_from(chunk_write_len)
.map_err(|_| io::Error::other("put_file body length cannot be represented"))?,
)
.ok_or_else(|| io::Error::other("put_file body length overflow"))?;
trailer.extend_from_slice(&bytes[chunk_write_len..]);
} else {
trailer.extend_from_slice(&bytes);
}
}
}
if remaining.is_some_and(|remaining| remaining != 0) {
return Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
format!("body size mismatch: expected {} bytes, received {copied}", query.size),
));
}
if trailer.len() != PUT_FILE_AUTH_TRAILER_LEN {
return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "put_file auth trailer is incomplete"));
}
let expected = verify_put_file_auth_trailer(url, &Method::PUT, nonce, &trailer)?;
let actual = hex_simd::encode_to_string(hasher.finalize(), hex_simd::AsciiCase::Lower);
if actual != expected {
return Err(io::Error::new(io::ErrorKind::InvalidData, "put_file body digest mismatch"));
}
Ok(copied)
}
fn parse_query<T>(req: &Request<Incoming>) -> Result<T, RpcErrorResponse>
where
T: DeserializeOwned + Default,
@@ -1308,11 +1482,12 @@ mod tests {
NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, NsScannerQuery, PUT_FILE_STREAM_PATH, PutFileQuery,
READ_FILE_STREAM_PATH, WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_PATH, WalkDirQuery, append_walk_dir_completion,
internode_http_operation, internode_rpc_subsystem, is_internode_rpc_path, ns_scanner_response_body,
ns_scanner_server_epoch_matches, put_body_size_mismatch, put_file_stage_error_message, read_file_body_stream,
remote_scanner_claim_rejection, response_with_disk_error, supports_walk_dir_stream_completion,
ns_scanner_server_epoch_matches, put_body_size_mismatch, put_file_auth_nonce, put_file_stage_error_message,
read_file_body_stream, remote_scanner_claim_rejection, response_with_disk_error, supports_walk_dir_stream_completion,
validate_walk_dir_completion_request, verify_internode_rpc_signature, verify_ns_scanner_body_digest,
verify_walk_dir_body_digest, walk_dir_response_body, write_body_chunks_to_writer,
verify_walk_dir_body_digest, walk_dir_response_body, write_body_chunks_to_writer, write_put_file_body_chunks_to_writer,
};
use crate::storage::storage_api::ecstore_rpc::build_put_file_auth_trailer;
use bytes::Bytes;
use http::{HeaderMap, HeaderValue, Method, StatusCode, Uri};
use http_body_util::BodyExt;
@@ -1433,6 +1608,8 @@ mod tests {
path: "tmp/object/part.1".to_string(),
append: false,
size: 1024,
put_file_auth: None,
put_file_nonce: None,
};
let msg = put_file_stage_error_message("write_body", &query, &"connection reset");
@@ -1452,6 +1629,8 @@ mod tests {
path: "object/part.1".to_string(),
append,
size,
put_file_auth: None,
put_file_nonce: None,
};
// Truncated (or over-long) body on the create path is rejected.
@@ -1579,6 +1758,165 @@ mod tests {
assert_eq!(out, b"hello world");
}
#[test]
fn put_file_auth_nonce_accepts_v1_requests_with_non_nil_nonce() {
let nonce = uuid::Uuid::new_v4();
let query = PutFileQuery {
disk: "disk-a".to_string(),
volume: "bucket".to_string(),
path: "object/part.1".to_string(),
append: false,
size: 11,
put_file_auth: Some("digest-trailer-v1".to_string()),
put_file_nonce: Some(nonce),
};
assert_eq!(put_file_auth_nonce(&query).expect("v1 auth should parse"), Some(nonce));
let mut append = query.clone();
append.append = true;
assert_eq!(put_file_auth_nonce(&append).expect("append auth should parse"), Some(nonce));
let mut unknown_size = query.clone();
unknown_size.size = -1;
assert_eq!(put_file_auth_nonce(&unknown_size).expect("unknown-size auth should parse"), Some(nonce));
let mut nil = query.clone();
nil.put_file_nonce = Some(uuid::Uuid::nil());
assert!(put_file_auth_nonce(&nil).is_err());
let mut unknown = query;
unknown.put_file_auth = Some("digest-trailer-v2".to_string());
assert!(put_file_auth_nonce(&unknown).is_err());
}
#[tokio::test]
async fn put_file_auth_body_writes_only_data_and_verifies_trailer() {
let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string());
let nonce = uuid::Uuid::parse_str("11111111-2222-4333-8444-555555555555").expect("nonce");
let url = concat!(
"/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1",
"&append=false&size=11&put_file_auth=digest-trailer-v1&put_file_nonce=11111111-2222-4333-8444-555555555555"
);
let digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"hello world"), hex_simd::AsciiCase::Lower);
let trailer = build_put_file_auth_trailer(url, &Method::PUT, nonce, &digest).expect("trailer should build");
let query = PutFileQuery {
disk: "disk-a".to_string(),
volume: "bucket".to_string(),
path: "object/part.1".to_string(),
append: false,
size: 11,
put_file_auth: Some("digest-trailer-v1".to_string()),
put_file_nonce: Some(nonce),
};
let mut second = b"world".to_vec();
second.extend_from_slice(&trailer[..7]);
let body = iter(vec![
Ok::<Bytes, io::Error>(Bytes::from_static(b"hello ")),
Ok(Bytes::from(second)),
Ok(Bytes::copy_from_slice(&trailer[7..])),
]);
let mut writer = Vec::new();
let copied = write_put_file_body_chunks_to_writer(body, &mut writer, &query, Some(nonce), url)
.await
.expect("authenticated body should verify");
assert_eq!(copied, 11);
assert_eq!(writer, b"hello world");
}
#[tokio::test]
async fn put_file_auth_body_rejects_tampered_data() {
let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string());
let nonce = uuid::Uuid::parse_str("22222222-3333-4444-8555-666666666666").expect("nonce");
let url = concat!(
"/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1",
"&append=false&size=11&put_file_auth=digest-trailer-v1&put_file_nonce=22222222-3333-4444-8555-666666666666"
);
let signed_digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"hello world"), hex_simd::AsciiCase::Lower);
let trailer = build_put_file_auth_trailer(url, &Method::PUT, nonce, &signed_digest).expect("trailer should build");
let query = PutFileQuery {
disk: "disk-a".to_string(),
volume: "bucket".to_string(),
path: "object/part.1".to_string(),
append: false,
size: 11,
put_file_auth: Some("digest-trailer-v1".to_string()),
put_file_nonce: Some(nonce),
};
let mut payload = b"hello worle".to_vec();
payload.extend_from_slice(&trailer);
let body = iter(vec![Ok::<Bytes, io::Error>(Bytes::from(payload))]);
let mut writer = Vec::new();
let err = write_put_file_body_chunks_to_writer(body, &mut writer, &query, Some(nonce), url)
.await
.expect_err("tampered body must fail digest verification");
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
assert_eq!(err.to_string(), "put_file body digest mismatch");
}
#[tokio::test]
async fn put_file_auth_append_body_uses_trailing_auth_record() {
let _ = rustfs_credentials::set_global_rpc_secret("put-file-auth-body-test-secret".to_string());
let nonce = uuid::Uuid::parse_str("33333333-4444-4555-8666-777777777777").expect("nonce");
let url = concat!(
"/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1",
"&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=33333333-4444-4555-8666-777777777777"
);
let digest = hex_simd::encode_to_string(sha2::Sha256::digest(b"append-data"), hex_simd::AsciiCase::Lower);
let trailer = build_put_file_auth_trailer(url, &Method::PUT, nonce, &digest).expect("trailer should build");
let query = PutFileQuery {
disk: "disk-a".to_string(),
volume: "bucket".to_string(),
path: "object/part.1".to_string(),
append: true,
size: 0,
put_file_auth: Some("digest-trailer-v1".to_string()),
put_file_nonce: Some(nonce),
};
let mut payload = b"append-data".to_vec();
payload.extend_from_slice(&trailer);
let body = iter(vec![Ok::<Bytes, io::Error>(Bytes::from(payload))]);
let mut writer = Vec::new();
let copied = write_put_file_body_chunks_to_writer(body, &mut writer, &query, Some(nonce), url)
.await
.expect("append body should verify");
assert_eq!(copied, 11);
assert_eq!(writer, b"append-data");
}
#[tokio::test]
async fn put_file_auth_append_body_rejects_missing_trailer() {
let nonce = uuid::Uuid::parse_str("44444444-5555-4666-8777-888888888888").expect("nonce");
let url = concat!(
"/rustfs/rpc/put_file_stream?disk=disk-a&volume=bucket&path=object%2Fpart.1",
"&append=true&size=0&put_file_auth=digest-trailer-v1&put_file_nonce=44444444-5555-4666-8777-888888888888"
);
let query = PutFileQuery {
disk: "disk-a".to_string(),
volume: "bucket".to_string(),
path: "object/part.1".to_string(),
append: true,
size: 0,
put_file_auth: Some("digest-trailer-v1".to_string()),
put_file_nonce: Some(nonce),
};
let body = iter(vec![Ok::<Bytes, io::Error>(Bytes::from_static(b"append-data"))]);
let mut writer = Vec::new();
let err = write_put_file_body_chunks_to_writer(body, &mut writer, &query, Some(nonce), url)
.await
.expect_err("missing trailer must fail");
assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
assert_eq!(err.to_string(), "put_file auth trailer is incomplete");
}
#[tokio::test]
async fn walk_dir_body_surfaces_background_failure_after_data() {
let body = walk_dir_response_body(true, |mut writer| async move {
+29 -6
View File
@@ -216,10 +216,12 @@ pub(crate) mod rpc_consumer {
NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, WALK_DIR_BODY_SHA256_QUERY,
};
pub(crate) use super::super::storage_contracts::{
NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, WALK_DIR_STREAM_COMPLETION_V1,
NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, PUT_FILE_AUTH_TRAILER_LEN, PUT_FILE_AUTH_V1,
WALK_DIR_STREAM_COMPLETION_V1,
};
pub(crate) use super::super::{
StorageDiskRpcExt, WalkDirOptions, find_local_disk_by_ref, sign_ns_scanner_capability, verify_rpc_signature,
StorageDiskRpcExt, WalkDirOptions, check_and_record_signed_rpc_nonce, find_local_disk_by_ref,
sign_ns_scanner_capability, verify_put_file_auth_trailer, verify_rpc_signature,
};
}
@@ -498,13 +500,15 @@ pub(crate) mod ecstore_rpc {
pub(crate) use rustfs_ecstore::api::rpc::{
KMS_SIGNAL_SUBSYSTEM, LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient,
PeerS3Client, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, TONIC_RPC_PREFIX,
normalize_tonic_rpc_audience, sign_ns_scanner_capability, sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge,
tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason, verify_rpc_signature,
verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest, verify_tonic_rpc_signature_with_bootstrap,
check_and_record_signed_rpc_nonce, normalize_tonic_rpc_audience, sign_ns_scanner_capability,
sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers,
tonic_rpc_auth_failure_reason, verify_put_file_auth_trailer, verify_rpc_signature, verify_tonic_canonical_body_digest,
verify_tonic_mutation_body_digest, verify_tonic_rpc_signature_with_bootstrap,
};
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::rpc::{
gen_signature_headers, gen_tonic_signature_headers, set_tonic_canonical_body_digest, verify_tonic_rpc_response_proof,
build_put_file_auth_trailer, gen_signature_headers, gen_tonic_signature_headers, set_tonic_canonical_body_digest,
verify_tonic_rpc_response_proof,
};
}
@@ -1655,6 +1659,25 @@ pub(crate) fn verify_rpc_signature(url: &str, method: &http::Method, headers: &h
ecstore_rpc::verify_rpc_signature(url, method, headers)
}
pub(crate) fn check_and_record_signed_rpc_nonce(
headers: &http::HeaderMap,
nonce: uuid::Uuid,
rpc_path: &str,
operation: &'static str,
backend: &'static str,
) -> std::io::Result<()> {
ecstore_rpc::check_and_record_signed_rpc_nonce(headers, nonce, rpc_path, operation, backend)
}
pub(crate) fn verify_put_file_auth_trailer(
url: &str,
method: &http::Method,
nonce: uuid::Uuid,
trailer: &[u8],
) -> std::io::Result<String> {
ecstore_rpc::verify_put_file_auth_trailer(url, method, nonce, trailer)
}
pub(crate) fn sign_ns_scanner_capability(challenge: uuid::Uuid, server_epoch: uuid::Uuid) -> std::io::Result<Vec<u8>> {
ecstore_rpc::sign_ns_scanner_capability(challenge, server_epoch)
}