From b5ecfcd28c3fe9a50fa57d33770ca718ffa22f84 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Wed, 17 Jun 2026 00:56:13 -0700 Subject: [PATCH] volume: validate remote S3 endpoints in FetchAndWriteNeedle (Rust) (#10001) * volume: validate remote S3 endpoints in FetchAndWriteNeedle (Rust) Port the Go volume server's SSRF guard to the Rust volume server. The gRPC FetchAndWriteNeedle reads from a caller-supplied S3 endpoint, so an unguarded server can be pointed at loopback / link-local / RFC1918 / CGNAT / cloud-metadata hosts to read internal services. Resolve and reject those endpoints unless -volume.allowUntrustedRemoteEndpoints is set (default off), mirroring weed/server/volume_grpc_remote.go. Connect-time re-validation against DNS rebinding (Go's guardedDialer) is not yet ported: the aws-sdk-s3 client builds its own connector, so the up-front resolve-and-check leaves a narrow TOCTOU window. Left as a follow-up. * volume: harden remote endpoint guard (IPv4-mapped IPv6, all S3 types) Address SSRF review feedback: - Normalize IPv4-mapped IPv6 (::ffff:a.b.c.d) to IPv4 before the deny checks, so ::ffff:127.0.0.1 / ::ffff:169.254.169.254 no longer slip past the IPv4 rules. - Validate the endpoint for every S3-compatible backend, not just type "s3"; wasabi/backblaze/aliyun/... all dial a caller-supplied endpoint through the same client. Skip validation when the endpoint is empty (the provider default, e.g. real AWS S3, which cannot reach internal hosts). * volume: set allow_untrusted_remote_endpoints in integration-test state The tests/http_integration.rs VolumeServerState literal was missed, which broke cargo test compilation (it builds the integration tests, unlike the cargo test --lib used locally). --- seaweed-volume/src/config.rs | 10 + seaweed-volume/src/main.rs | 1 + .../src/remote_storage/endpoint_guard.rs | 359 ++++++++++++++++++ seaweed-volume/src/remote_storage/mod.rs | 57 +++ seaweed-volume/src/server/grpc_client.rs | 1 + seaweed-volume/src/server/grpc_server.rs | 22 ++ seaweed-volume/src/server/heartbeat.rs | 1 + seaweed-volume/src/server/profiling.rs | 1 + seaweed-volume/src/server/volume_server.rs | 3 + seaweed-volume/src/server/write_queue.rs | 1 + seaweed-volume/tests/http_integration.rs | 1 + 11 files changed, 457 insertions(+) create mode 100644 seaweed-volume/src/remote_storage/endpoint_guard.rs diff --git a/seaweed-volume/src/config.rs b/seaweed-volume/src/config.rs index cfcfe71cd..d5962d2bb 100644 --- a/seaweed-volume/src/config.rs +++ b/seaweed-volume/src/config.rs @@ -188,6 +188,12 @@ pub struct Cli { #[arg(long = "securityFile", default_value = "")] pub security_file: String, + /// If true, FetchAndWriteNeedle accepts arbitrary remote S3 endpoints + /// including loopback / link-local hosts. Default rejects internal / + /// metadata endpoints. + #[arg(long = "volume.allowUntrustedRemoteEndpoints", default_value_t = false)] + pub allow_untrusted_remote_endpoints: bool, + /// A file of command line options, each line in optionName=optionValue format. #[arg(long = "options", default_value = "")] pub options: String, @@ -258,6 +264,9 @@ pub struct VolumeServerConfig { pub enable_write_queue: bool, /// Path to security.toml — stored for SIGHUP reload. pub security_file: String, + /// If true, FetchAndWriteNeedle skips remote S3 endpoint validation + /// (allows loopback / link-local / metadata hosts). + pub allow_untrusted_remote_endpoints: bool, } pub use crate::storage::needle_map::NeedleMapKind; @@ -804,6 +813,7 @@ fn resolve_config(cli: Cli) -> VolumeServerConfig { .map(|v| v == "1" || v == "true") .unwrap_or(false), security_file: cli.security_file, + allow_untrusted_remote_endpoints: cli.allow_untrusted_remote_endpoints, } } diff --git a/seaweed-volume/src/main.rs b/seaweed-volume/src/main.rs index 56819f5a8..fa7a7d083 100644 --- a/seaweed-volume/src/main.rs +++ b/seaweed-volume/src/main.rs @@ -327,6 +327,7 @@ async fn run( seaweed_volume::remote_storage::s3_tier::S3TierRegistry::new(), ), read_mode: config.read_mode, + allow_untrusted_remote_endpoints: config.allow_untrusted_remote_endpoints, master_url, master_urls, seed_master_set, diff --git a/seaweed-volume/src/remote_storage/endpoint_guard.rs b/seaweed-volume/src/remote_storage/endpoint_guard.rs new file mode 100644 index 000000000..90e051ff6 --- /dev/null +++ b/seaweed-volume/src/remote_storage/endpoint_guard.rs @@ -0,0 +1,359 @@ +//! SSRF guard for caller-supplied remote storage endpoints. +//! +//! `FetchAndWriteNeedle` accepts an S3 endpoint URL chosen by the caller and +//! dials it directly from the volume server, which typically has network +//! access to cluster-internal hosts. Without validation a caller could point +//! the server at loopback / link-local / RFC 1918 / cloud-metadata addresses +//! and read internal services (SSRF). This mirrors the Go volume server's +//! `validateRemoteEndpoint` (`weed/server/volume_grpc_remote.go`); operators +//! that legitimately fetch from private hosts opt out with +//! `-volume.allowUntrustedRemoteEndpoints`. +//! +//! NOTE: the Go server additionally pins the validated endpoint against DNS +//! rebinding with a custom dialer that re-checks the resolved IP at TCP connect +//! time (`guardedDialer`). The `aws-sdk-s3` client used here builds its own +//! connector, so that connect-time re-validation is not yet ported; the +//! up-front resolve-and-check below still blocks the common SSRF vectors. A +//! narrow TOCTOU window (a hostname that resolves to a public IP here and then +//! flips to a blocked one when the SDK dials) remains as a follow-up. + +use std::net::{IpAddr, Ipv4Addr}; + +/// AWS/Azure/GCP IPv4 instance-metadata-service (IMDS) address. It is +/// link-local and thus already covered by [`is_link_local`], but is named +/// explicitly so the rejection reason is unambiguous in logs. +const IMDS_IPV4: IpAddr = IpAddr::V4(Ipv4Addr::new(169, 254, 169, 254)); + +/// Hostnames that target cloud instance metadata services, blocked regardless +/// of how they resolve because some environments alias the IMDS address under +/// a name. +fn is_blocked_imds_host(host: &str) -> bool { + matches!(host, "metadata.google.internal" | "metadata") +} + +/// Whether `ip` is in a link-local range: IPv4 169.254.0.0/16, IPv6 fe80::/10 +/// (unicast) and interface-local / link-local multicast (scope 1 and 2). +fn is_link_local(ip: IpAddr) -> bool { + match ip { + IpAddr::V4(v4) => v4.is_link_local(), + IpAddr::V6(v6) => { + let seg0 = v6.segments()[0]; + if (seg0 & 0xffc0) == 0xfe80 { + return true; // fe80::/10 link-local unicast + } + if (seg0 & 0xff00) == 0xff00 { + // ffXS:: multicast; reject interface-local (scope 1) and + // link-local (scope 2) the way Go's IsInterfaceLocalMulticast / + // IsLinkLocalMulticast do. + let scope = seg0 & 0x000f; + return scope == 1 || scope == 2; + } + false + } + } +} + +/// Whether `ip` is in a private range: IPv4 RFC 1918, IPv6 fc00::/7 unique-local +/// (matching Go's `net.IP.IsPrivate`). +fn is_private(ip: IpAddr) -> bool { + match ip { + IpAddr::V4(v4) => v4.is_private(), + IpAddr::V6(v6) => (v6.segments()[0] & 0xfe00) == 0xfc00, + } +} + +/// Whether `ip` is in the RFC 6598 carrier-grade NAT range (100.64.0.0/10). +/// The IPv4 private check does not cover CGNAT, so check it explicitly. +fn is_cgnat(ip: IpAddr) -> bool { + match ip { + IpAddr::V4(v4) => { + let o = v4.octets(); + o[0] == 100 && (o[1] & 0xc0) == 0x40 // second octet 64..=127 + } + IpAddr::V6(_) => false, + } +} + +/// Returns an error if `ip` is not safe to dial from a server that can reach +/// cluster-internal hosts. Mirrors Go's `checkBlockedIP`. +pub fn check_blocked_ip(endpoint: &str, ip: IpAddr) -> Result<(), String> { + // Normalize IPv4-mapped IPv6 (`::ffff:a.b.c.d`) to its IPv4 form so the + // IPv4 deny rules apply. The OS routes these to the embedded IPv4 address, + // so without this `::ffff:127.0.0.1` / `::ffff:169.254.169.254` would slip + // past the IPv6-only checks. Mirrors Go's reliance on net.IP.To4(). + let ip = match ip { + IpAddr::V6(v6) => match v6.to_ipv4_mapped() { + Some(v4) => IpAddr::V4(v4), + None => IpAddr::V6(v6), + }, + other => other, + }; + if ip == IMDS_IPV4 { + return Err(format!( + "remote endpoint {:?} targets instance metadata service {}", + endpoint, ip + )); + } + if ip.is_loopback() { + return Err(format!( + "remote endpoint {:?} resolves to loopback address {}", + endpoint, ip + )); + } + if ip.is_unspecified() { + return Err(format!( + "remote endpoint {:?} resolves to unspecified address {}", + endpoint, ip + )); + } + if is_link_local(ip) { + return Err(format!( + "remote endpoint {:?} resolves to link-local address {}", + endpoint, ip + )); + } + if is_private(ip) { + return Err(format!( + "remote endpoint {:?} resolves to private address {}", + endpoint, ip + )); + } + if is_cgnat(ip) { + return Err(format!( + "remote endpoint {:?} resolves to CGNAT address {}", + endpoint, ip + )); + } + Ok(()) +} + +/// Outcome of the synchronous, DNS-free portion of endpoint validation. +#[derive(Debug)] +enum HostCheck { + /// Host was a literal IP and already passed [`check_blocked_ip`]. + Validated, + /// Host is a name that must be resolved and each address checked. + NeedsResolution(String), +} + +/// Validate the scheme and host without touching DNS. Pure and unit-testable. +fn precheck_endpoint(endpoint: &str) -> Result { + let trimmed = endpoint.trim(); + if trimmed.is_empty() { + return Err("remote endpoint is empty".to_string()); + } + + let Some(scheme_end) = trimmed.find("://") else { + return Err(format!( + "remote endpoint {:?} must use http or https", + endpoint + )); + }; + let scheme = trimmed[..scheme_end].to_ascii_lowercase(); + if scheme != "http" && scheme != "https" { + return Err(format!( + "remote endpoint {:?} must use http or https, got {:?}", + endpoint, + &trimmed[..scheme_end] + )); + } + + // Authority is everything up to the first '/', '?', or '#'. + let after = &trimmed[scheme_end + 3..]; + let authority_end = after + .find(|c| c == '/' || c == '?' || c == '#') + .unwrap_or(after.len()); + let authority = &after[..authority_end]; + + // Strip optional userinfo ("user:pass@"). + let host_port = match authority.rfind('@') { + Some(at) => &authority[at + 1..], + None => authority, + }; + + // Extract the host, handling bracketed IPv6 literals and host:port. + let host = if let Some(rest) = host_port.strip_prefix('[') { + match rest.find(']') { + Some(end) => &rest[..end], + None => { + return Err(format!( + "remote endpoint {:?} has a malformed IPv6 host", + endpoint + )) + } + } + } else { + match host_port.find(':') { + Some(colon) => &host_port[..colon], + None => host_port, + } + }; + + if host.is_empty() { + return Err(format!("remote endpoint {:?} has no host", endpoint)); + } + + if is_blocked_imds_host(&host.to_ascii_lowercase()) { + return Err(format!( + "remote endpoint {:?} targets instance metadata service", + endpoint + )); + } + + if let Ok(ip) = host.parse::() { + check_blocked_ip(endpoint, ip)?; + return Ok(HostCheck::Validated); + } + + Ok(HostCheck::NeedsResolution(host.to_string())) +} + +/// Resolve `host` to its addresses with a short timeout, matching Go's 2s +/// resolver deadline. +async fn resolve_host(host: &str) -> Result, String> { + let lookup = tokio::net::lookup_host((host, 0u16)); + let addrs = tokio::time::timeout(std::time::Duration::from_secs(2), lookup) + .await + .map_err(|_| format!("resolve remote endpoint host {:?}: timed out", host))? + .map_err(|e| format!("resolve remote endpoint host {:?}: {}", host, e))?; + Ok(addrs.map(|sock| sock.ip()).collect()) +} + +/// Returns an error if `endpoint` is not safe to dial from a server with access +/// to cluster-internal hosts. Rejects empty / non-http(s) endpoints, +/// loopback / unspecified / link-local / RFC 1918 / CGNAT addresses, and +/// well-known IMDS hostnames. Hostnames are resolved and every returned address +/// is checked. Mirrors Go's `validateRemoteEndpoint`. +pub async fn validate_remote_endpoint(endpoint: &str) -> Result<(), String> { + match precheck_endpoint(endpoint)? { + HostCheck::Validated => Ok(()), + HostCheck::NeedsResolution(host) => { + let addrs = resolve_host(&host).await?; + if addrs.is_empty() { + return Err(format!( + "resolve remote endpoint host {:?}: no addresses", + host + )); + } + for ip in addrs { + check_blocked_ip(endpoint, ip)?; + } + Ok(()) + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn ip(s: &str) -> IpAddr { + s.parse().unwrap() + } + + #[test] + fn rejects_blocked_literal_endpoints() { + let cases = [ + ("http://127.0.0.1:8080", "loopback"), + ("http://[::1]:8080", "loopback"), + ("http://169.254.169.254/", "metadata"), + ("http://0.0.0.0/", "unspecified"), + ("http://[fe80::1]/", "link-local"), + ("http://10.0.0.1/", "private"), + ("http://172.16.5.5/", "private"), + ("http://192.168.0.1/", "private"), + ("http://100.64.0.1/", "CGNAT"), + ]; + for (endpoint, want) in cases { + let err = precheck_endpoint(endpoint).expect_err(endpoint); + assert!(err.contains(want), "{endpoint}: {err:?} missing {want:?}"); + } + } + + #[test] + fn rejects_empty_and_bad_scheme() { + assert!(precheck_endpoint("").unwrap_err().contains("empty")); + assert!(precheck_endpoint("ftp://example.com/") + .unwrap_err() + .contains("http or https")); + assert!(precheck_endpoint("example.com/") + .unwrap_err() + .contains("http or https")); + } + + #[test] + fn rejects_imds_hostnames() { + assert!(precheck_endpoint("http://metadata.google.internal/") + .unwrap_err() + .contains("metadata service")); + assert!(precheck_endpoint("http://metadata/") + .unwrap_err() + .contains("metadata service")); + } + + #[test] + fn accepts_public_literal_and_defers_hostname() { + assert!(matches!( + precheck_endpoint("https://52.216.10.10/"), + Ok(HostCheck::Validated) + )); + match precheck_endpoint("https://s3.us-east-1.amazonaws.com/") { + Ok(HostCheck::NeedsResolution(host)) => { + assert_eq!(host, "s3.us-east-1.amazonaws.com") + } + other => panic!("expected resolution, got {other:?}"), + } + } + + #[test] + fn check_blocked_ip_matches_resolved_categories() { + // Mirror Go's "host resolves to X" cases at the address level. + assert!(check_blocked_ip("e", ip("127.0.0.1")) + .unwrap_err() + .contains("loopback")); + assert!(check_blocked_ip("e", ip("169.254.10.20")) + .unwrap_err() + .contains("link-local")); + assert!(check_blocked_ip("e", ip("10.1.2.3")) + .unwrap_err() + .contains("private")); + assert!(check_blocked_ip("e", ip("172.20.0.5")) + .unwrap_err() + .contains("private")); + assert!(check_blocked_ip("e", ip("192.168.1.1")) + .unwrap_err() + .contains("private")); + assert!(check_blocked_ip("e", ip("100.64.0.42")) + .unwrap_err() + .contains("CGNAT")); + assert!(check_blocked_ip("e", ip("fc00::1")) + .unwrap_err() + .contains("private")); + assert!(check_blocked_ip("e", ip("52.216.10.10")).is_ok()); + assert!(check_blocked_ip("e", ip("2606:4700:4700::1111")).is_ok()); + } + + #[test] + fn rejects_ipv4_mapped_ipv6() { + // IPv4-mapped IPv6 must be unmapped so the IPv4 rules catch it. + assert!(check_blocked_ip("e", ip("::ffff:127.0.0.1")) + .unwrap_err() + .contains("loopback")); + assert!(check_blocked_ip("e", ip("::ffff:169.254.169.254")) + .unwrap_err() + .contains("metadata")); + assert!(check_blocked_ip("e", ip("::ffff:10.0.0.1")) + .unwrap_err() + .contains("private")); + // A mapped public address still passes, and genuine IPv6 loopback is + // still caught by the V6 path. + assert!(check_blocked_ip("e", ip("::ffff:52.216.10.10")).is_ok()); + assert!(check_blocked_ip("e", ip("::1")) + .unwrap_err() + .contains("loopback")); + // Bracketed mapped literal via the full endpoint path. + assert!(precheck_endpoint("http://[::ffff:127.0.0.1]/") + .unwrap_err() + .contains("loopback")); + } +} diff --git a/seaweed-volume/src/remote_storage/mod.rs b/seaweed-volume/src/remote_storage/mod.rs index 599333ede..49be142d6 100644 --- a/seaweed-volume/src/remote_storage/mod.rs +++ b/seaweed-volume/src/remote_storage/mod.rs @@ -3,9 +3,12 @@ //! Provides a trait-based abstraction over cloud storage providers (S3, GCS, Azure, etc.) //! and a registry to create clients from protobuf RemoteConf messages. +pub mod endpoint_guard; pub mod s3; pub mod s3_tier; +pub use endpoint_guard::validate_remote_endpoint; + use crate::pb::remote_pb::{RemoteConf, RemoteStorageLocation}; /// Error type for remote storage operations. @@ -86,6 +89,26 @@ pub fn make_remote_storage_client( } } +/// Endpoint URL that the volume server would dial directly for `conf`, or +/// `None` for non-S3-compatible backends. Every type handled here routes +/// through [`s3::S3RemoteStorageClient`] with a caller-supplied endpoint, so +/// the SSRF guard must validate each one. Keep this match in sync with +/// [`make_remote_storage_client`]. +pub fn s3_compatible_endpoint(conf: &RemoteConf) -> Option<&str> { + match conf.r#type.as_str() { + "s3" => Some(&conf.s3_endpoint), + "wasabi" => Some(&conf.wasabi_endpoint), + "backblaze" => Some(&conf.backblaze_endpoint), + "aliyun" => Some(&conf.aliyun_endpoint), + "tencent" => Some(&conf.tencent_endpoint), + "baidu" => Some(&conf.baidu_endpoint), + "filebase" => Some(&conf.filebase_endpoint), + "storj" => Some(&conf.storj_endpoint), + "contabo" => Some(&conf.contabo_endpoint), + _ => None, + } +} + /// Extract S3-compatible credentials from a RemoteConf based on its type. fn extract_s3_credentials(conf: &RemoteConf) -> (String, String, String, String) { match conf.r#type.as_str() { @@ -155,3 +178,37 @@ fn extract_s3_credentials(conf: &RemoteConf) -> (String, String, String, String) ), } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn s3_compatible_endpoint_covers_all_s3_backends() { + let s3 = RemoteConf { + r#type: "s3".to_string(), + s3_endpoint: "http://s3.internal".to_string(), + ..Default::default() + }; + assert_eq!(s3_compatible_endpoint(&s3), Some("http://s3.internal")); + + // A non-"s3" S3-compatible type still surfaces its own endpoint, so the + // SSRF guard cannot be bypassed by picking a different alias. + let wasabi = RemoteConf { + r#type: "wasabi".to_string(), + wasabi_endpoint: "http://wasabi.internal".to_string(), + ..Default::default() + }; + assert_eq!( + s3_compatible_endpoint(&wasabi), + Some("http://wasabi.internal") + ); + + // Non-S3 backends do not dial a caller-supplied URL directly. + let gcs = RemoteConf { + r#type: "gcs".to_string(), + ..Default::default() + }; + assert_eq!(s3_compatible_endpoint(&gcs), None); + } +} diff --git a/seaweed-volume/src/server/grpc_client.rs b/seaweed-volume/src/server/grpc_client.rs index e222478fd..40caf8b59 100644 --- a/seaweed-volume/src/server/grpc_client.rs +++ b/seaweed-volume/src/server/grpc_client.rs @@ -189,6 +189,7 @@ mod tests { white_list: vec![], fix_jpg_orientation: false, read_mode: ReadMode::Local, + allow_untrusted_remote_endpoints: false, cpu_profile: String::new(), mem_profile: String::new(), compaction_byte_per_second: 0, diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 6a0f80801..036d1e2bc 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -3490,6 +3490,25 @@ impl VolumeServer for VolumeGrpcService { .as_ref() .ok_or_else(|| Status::invalid_argument("remote storage configuration is required"))?; + // Reject SSRF-prone endpoints unless the operator opted out. Every + // S3-compatible backend (s3, wasabi, backblaze, ...) dials a + // caller-supplied endpoint URL directly via make_remote_storage_client, + // so validate whichever endpoint applies. An empty endpoint means the + // provider default (e.g. real AWS S3) and cannot target an internal + // host, so skip it. Extends the Go volume server's validateRemoteEndpoint + // gate, which only covered type "s3". + if !self.state.allow_untrusted_remote_endpoints { + if let Some(endpoint) = crate::remote_storage::s3_compatible_endpoint(remote_conf) { + if !endpoint.trim().is_empty() { + crate::remote_storage::validate_remote_endpoint(endpoint) + .await + .map_err(|e| { + Status::invalid_argument(format!("reject remote endpoint: {}", e)) + })?; + } + } + } + // Create remote storage client let client = crate::remote_storage::make_remote_storage_client(remote_conf).map_err(|e| { @@ -4728,6 +4747,7 @@ mod tests { crate::remote_storage::s3_tier::S3TierRegistry::new(), ), read_mode: crate::config::ReadMode::Local, + allow_untrusted_remote_endpoints: false, master_url: String::new(), master_urls: Vec::new(), seed_master_set: std::collections::HashSet::new(), @@ -4832,6 +4852,7 @@ mod tests { crate::remote_storage::s3_tier::S3TierRegistry::new(), ), read_mode: crate::config::ReadMode::Local, + allow_untrusted_remote_endpoints: false, master_url: String::new(), master_urls: Vec::new(), seed_master_set: std::collections::HashSet::new(), @@ -4940,6 +4961,7 @@ mod tests { crate::remote_storage::s3_tier::S3TierRegistry::new(), ), read_mode: crate::config::ReadMode::Local, + allow_untrusted_remote_endpoints: false, master_url: master_urls.first().cloned().unwrap_or_default(), master_urls, seed_master_set, diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index 0906cc257..2d668c67d 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -1047,6 +1047,7 @@ mod tests { write_queue: std::sync::OnceLock::new(), s3_tier_registry: std::sync::RwLock::new(S3TierRegistry::new()), read_mode: ReadMode::Local, + allow_untrusted_remote_endpoints: false, master_url: String::new(), master_urls: Vec::new(), seed_master_set: std::collections::HashSet::new(), diff --git a/seaweed-volume/src/server/profiling.rs b/seaweed-volume/src/server/profiling.rs index 1965d227f..51e4f174c 100644 --- a/seaweed-volume/src/server/profiling.rs +++ b/seaweed-volume/src/server/profiling.rs @@ -131,6 +131,7 @@ mod tests { white_list: vec![], fix_jpg_orientation: false, read_mode: ReadMode::Local, + allow_untrusted_remote_endpoints: false, cpu_profile: String::new(), mem_profile: String::new(), compaction_byte_per_second: 0, diff --git a/seaweed-volume/src/server/volume_server.rs b/seaweed-volume/src/server/volume_server.rs index 513ef92aa..54928f61d 100644 --- a/seaweed-volume/src/server/volume_server.rs +++ b/seaweed-volume/src/server/volume_server.rs @@ -77,6 +77,9 @@ pub struct VolumeServerState { pub s3_tier_registry: std::sync::RwLock, /// Read mode: local, proxy, or redirect for non-local volumes. pub read_mode: ReadMode, + /// If true, FetchAndWriteNeedle skips remote S3 endpoint validation, + /// allowing arbitrary (incl. loopback / link-local / metadata) hosts. + pub allow_untrusted_remote_endpoints: bool, /// First master address for volume lookups (e.g., "localhost:9333"). pub master_url: String, /// Seed master addresses for UI rendering. diff --git a/seaweed-volume/src/server/write_queue.rs b/seaweed-volume/src/server/write_queue.rs index 3f99d517b..df809db38 100644 --- a/seaweed-volume/src/server/write_queue.rs +++ b/seaweed-volume/src/server/write_queue.rs @@ -205,6 +205,7 @@ mod tests { crate::remote_storage::s3_tier::S3TierRegistry::new(), ), read_mode: crate::config::ReadMode::Local, + allow_untrusted_remote_endpoints: false, master_url: String::new(), master_urls: Vec::new(), seed_master_set: std::collections::HashSet::new(), diff --git a/seaweed-volume/tests/http_integration.rs b/seaweed-volume/tests/http_integration.rs index d8fe04c80..13d39000c 100644 --- a/seaweed-volume/tests/http_integration.rs +++ b/seaweed-volume/tests/http_integration.rs @@ -99,6 +99,7 @@ fn test_state_with_guard( seaweed_volume::remote_storage::s3_tier::S3TierRegistry::new(), ), read_mode: seaweed_volume::config::ReadMode::Local, + allow_untrusted_remote_endpoints: false, master_url: String::new(), master_urls: Vec::new(), seed_master_set: std::collections::HashSet::new(),