diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 6a0532f9c..32bb3e794 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -265,6 +265,18 @@ pub struct VolumeGrpcService { } impl VolumeGrpcService { + fn store_read(&self) -> std::sync::RwLockReadGuard<'_, crate::storage::store::Store> { + self.state.store.read().unwrap_or_else(|e| { + tracing::warn!("recovering poisoned store read lock"); + e.into_inner() + }) + } + fn store_write(&self) -> std::sync::RwLockWriteGuard<'_, crate::storage::store::Store> { + self.state.store.write().unwrap_or_else(|e| { + tracing::warn!("recovering poisoned store write lock"); + e.into_inner() + }) + } /// Verifies the gRPC caller is allowed to invoke a destructive admin /// operation. Mirrors the Go side's checkGrpcAdminAuth: an empty /// whitelist accepts everyone (insecure-by-default for tests and @@ -2595,12 +2607,32 @@ impl VolumeServer for VolumeGrpcService { let vid = VolumeId(req.volume_id); let offset = req.offset; let size = Size(req.size); + if let Err(msg) = crate::storage::needle::needle::validate_wire_size(size) { + return Err(Status::invalid_argument(msg)); + } - let store = self.state.store.read().unwrap(); + let store = self.store_read(); let (_, vol) = store .find_volume(vid) .ok_or_else(|| Status::not_found(format!("not found volume id {}", vid)))?; + // Framing: request `size` is body size, but read_needle_blob returns + // GetActualSize bytes + protobuf tag/len, so body==1<<30 always exceeds + // GRPC_MAX_MESSAGE_SIZE on the wire after allocation+disk read. + // Check actual encoded size post-lock (have vol.version() here). + // The 6-byte headroom is the exact protobuf overhead for this + // response at the boundary: 1 tag byte (`bytes needle_blob = 1`) + // plus a 5-byte varint length for sizes near 1 GiB. Sizes whose + // encoded form still fits are therefore accepted; anything larger + // is rejected before paying for the disk read. + let actual = crate::storage::needle::needle::get_actual_size(size, vol.version()) as u64; + if actual + 6 > crate::server::grpc_client::GRPC_MAX_MESSAGE_SIZE as u64 { + return Err(Status::invalid_argument(format!( + "needle blob size {} exceeds transport limit", + actual + ))); + } + let blob = vol.read_needle_blob(offset, size).map_err(|e| { Status::internal(format!( "read needle blob offset {} size {}: {}", @@ -2659,8 +2691,11 @@ impl VolumeServer for VolumeGrpcService { let vid = VolumeId(req.volume_id); let needle_id = NeedleId(req.needle_id); let size = Size(req.size); + if let Err(msg) = crate::storage::needle::needle::validate_wire_size(size) { + return Err(Status::invalid_argument(msg)); + } - let mut store = self.state.store.write().unwrap(); + let mut store = self.store_write(); let (_, vol) = store .find_volume_mut(vid) .ok_or_else(|| Status::not_found(format!("not found volume id {}", vid)))?; @@ -7290,6 +7325,53 @@ mod tests { } } + #[tokio::test] + async fn test_read_needle_blob_negative_returns_invalid_argument() { + let (service, _tmp) = make_local_service_with_volume("", None); + let err = service + .read_needle_blob(Request::new(volume_server_pb::ReadNeedleBlobRequest { + volume_id: 1, + offset: 0, + size: -100, + })) + .await + .expect_err("negative size must be rejected"); + assert_eq!(err.code(), tonic::Code::InvalidArgument); + } + + #[tokio::test] + async fn test_write_needle_blob_negative_returns_invalid_argument() { + let (service, _tmp) = make_local_service_with_volume("", None); + let err = service + .write_needle_blob(Request::new(volume_server_pb::WriteNeedleBlobRequest { + volume_id: 1, + needle_id: 11, + size: -100, + needle_blob: vec![], + })) + .await + .expect_err("negative size must be rejected"); + assert_eq!(err.code(), tonic::Code::InvalidArgument); + } + + #[tokio::test] + async fn test_read_needle_blob_zero_size_passes_gate() { + let (service, _tmp) = make_local_service_with_volume("", None); + // Size(0) must pass the wire-size gate (may 404/Internal downstream, + // but must NOT be rejected as InvalidArgument). + match service + .read_needle_blob(Request::new(volume_server_pb::ReadNeedleBlobRequest { + volume_id: 1, + offset: 0, + size: 0, + })) + .await + { + Ok(_) => {} + Err(e) => assert_ne!(e.code(), tonic::Code::InvalidArgument, "got {e:?}"), + } + } + // Regression test for comparing the wrong compaction-revision field. // last_compact_revision() is bookkeeping recorded just before a compaction // starts (for makeup-diff catch-up) and is intentionally left behind diff --git a/seaweed-volume/src/storage/needle/needle.rs b/seaweed-volume/src/storage/needle/needle.rs index 278c944f5..2ad129716 100644 --- a/seaweed-volume/src/storage/needle/needle.rs +++ b/seaweed-volume/src/storage/needle/needle.rs @@ -615,6 +615,30 @@ pub fn get_actual_size(size: Size, version: Version) -> i64 { NEEDLE_HEADER_SIZE as i64 + needle_body_length(size, version) } +/// Validate a wire-supplied needle body size before any `as usize` cast. +/// Rejects negative/deleted sizes and bodies larger than the gRPC max message. +/// Size(0) is allowed: empty/anomalous entries and tombstones read as size 0 +/// (actual_size = header+checksum+pad > 0, safe alloc, no wrap). +/// Transport cap only: storage paths must NOT use this cap — see volume.rs +/// guards (a >1GiB stored needle from a high-limit cluster must remain +/// readable/compaction-safe). Keep `get_actual_size` unchanged (it +/// intentionally returns negative for deleted index entries). +pub fn validate_wire_size(size: Size) -> Result<(), String> { + if size.0 < 0 { + return Err(format!("invalid needle size {}", size.0)); + } + // Keep in sync with canonical `GRPC_MAX_MESSAGE_SIZE` in server/grpc_client.rs:10 + // (duplicated here to avoid a storage->server import and prevent drift). + const WIRE_MAX_NEEDLE_SIZE: i32 = 1 << 30; + if size.0 > WIRE_MAX_NEEDLE_SIZE { + return Err(format!( + "needle size {} exceeds max {}", + size.0, WIRE_MAX_NEEDLE_SIZE + )); + } + Ok(()) +} + /// Read 5 bytes as a u64 (big-endian, zero-padded high bytes). fn bytes_to_u64_5(bytes: &[u8]) -> u64 { assert!(bytes.len() >= 5); @@ -980,4 +1004,16 @@ mod tests { assert_eq!(fid.key, NeedleId(0x123)); assert_eq!(fid.cookie, Cookie(0)); } + + #[test] + fn test_validate_wire_size_boundaries() { + assert!(validate_wire_size(Size(-100)).is_err()); + assert!(validate_wire_size(Size(-1)).is_err()); + assert!(validate_wire_size(Size(0)).is_ok()); + assert!(validate_wire_size(Size(1024)).is_ok()); + assert!(validate_wire_size(Size(1)).is_ok()); + assert!(validate_wire_size(Size(1 << 30)).is_ok()); + assert!(validate_wire_size(Size((1 << 30) + 1)).is_err()); + assert!(validate_wire_size(Size(i32::MAX)).is_err()); + } } diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index 76c144f14..a2668604b 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -1574,6 +1574,15 @@ impl Volume { size: Size, ) -> Result<(), VolumeError> { let version = self.version(); + // Storage guard: negativity-only (Go parity — storage allocates what the + // index says). Size(0) and >1GiB map sizes must still read; only + // negative wraps/panics. Transport cap lives in RPC handlers only. + if size.0 < 0 { + return Err(VolumeError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("invalid needle size {}", size.0), + ))); + } let actual_size = get_actual_size(size, version); let mut buf = vec![0u8; actual_size as usize]; @@ -1591,6 +1600,13 @@ impl Volume { fn read_needle_blob_unlocked(&self, offset: i64, size: Size) -> Result, VolumeError> { let version = self.version(); + // Storage guard: negativity-only (Go parity). See read_needle_blob_and_parse. + if size.0 < 0 { + return Err(VolumeError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("invalid needle size {}", size.0), + ))); + } let actual_size = get_actual_size(size, version); let mut buf = vec![0u8; actual_size as usize]; self.read_exact_at_backend(&mut buf, offset as u64)?; @@ -3473,6 +3489,13 @@ impl Volume { if self.is_read_only() { return Err(VolumeError::ReadOnly); } + // Storage guard: negativity-only (Go parity). See read_needle_blob_and_parse. + if size.0 < 0 { + return Err(VolumeError::Io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("invalid needle size {}", size.0), + ))); + } // size indexes the needle and places the v3 append timestamp, so a caller using // the payload-only DataSize corrupts both, silently until the needle is read back. @@ -6537,6 +6560,38 @@ mod tests { .unwrap(); } + #[test] + fn test_read_blob_negative_does_not_panic() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let v = make_test_volume(dir); + let res = v.read_needle_blob(0, Size(-100)); + assert!(res.is_err(), "negative size must return Err, not panic"); + match res.unwrap_err() { + VolumeError::Io(e) => assert_eq!(e.kind(), std::io::ErrorKind::InvalidData), + e => panic!("expected Io InvalidData, got {e:?}"), + } + } + + #[test] + fn test_read_blob_zero_size_not_rejected_by_validation() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let v = make_test_volume(dir); + // Size(0) must not be rejected by the negativity guard: the read may + // Ok or fail on empty-volume IO, but never with our validation message. + match v.read_needle_blob(0, Size(0)) { + Ok(_) => {} + Err(VolumeError::Io(e)) => { + assert!( + !e.to_string().contains("invalid needle size"), + "Size(0) must pass validation, got {e}" + ); + } + Err(e) => panic!("unexpected error kind for Size(0): {e:?}"), + } + } + #[test] fn test_volume_destroy() { let tmp = TempDir::new().unwrap();