fix(volume): reject negative Size, recover poisoned store lock (#11345)

This commit is contained in:
Eliah Rusin
2026-09-16 08:12:11 -07:00
committed by GitHub
parent 4fc9ada2ec
commit 701e397337
3 changed files with 175 additions and 2 deletions
+84 -2
View File
@@ -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
@@ -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());
}
}
+55
View File
@@ -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<Vec<u8>, 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();