From b8f074b7d3cb026f5d8cc32da156a9230d88eb20 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Thu, 1 Oct 2026 15:18:34 +0300 Subject: [PATCH] volume server: VolumeNeedleStatus reads remote EC shards and reports deleted needles like Go (#11535) * volume server: VolumeNeedleStatus reads remote EC shards and reports deleted needles like Go For an EC volume the handler read only locally mounted shards, so a node that did not hold the shard with the needle's bytes answered Internal "ec shard N not available locally". Go's ReadEcShardNeedle fetches the interval from a peer or reconstructs it. It also mapped every regular volume read error, including a tombstone, to NotFound "needle not found", which fs.verify treats as a missing needle; Go returns ErrorDeleted as a plain error ("already deleted"), which fs.verify skips. The EC branch now drops the store guard and uses the distributed EC read the HTTP GET path uses. Errors map like Go: needle absent -> NotFound "needle not found ", tombstoned (regular or EC .ecx/.ecj) -> Unknown "already deleted", anything else -> Unknown with the error text. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: tell EC deletions and vanished volumes apart in VolumeNeedleStatus The distributed EC reader returned Ok(None) for an absent needle, a needle a peer reported deleted, and a volume unmounted after the handler's own existence check. VolumeNeedleStatus answered all three NotFound "needle not found", which fs.verify -pruneEntries counts as lost data. A reported deletion was also lost when an earlier interval failed. The reader now says why it has no needle (EcMiss: NotFound, Deleted, VolumeNotFound), classifying the local tombstone itself and letting a reported deletion outrank other interval errors, as Go's ReadEcShardNeedle does. VolumeNeedleStatus maps Deleted to Unknown "already deleted" and VolumeNotFound to "volume not found", and drops its separate EC pre-check. read_ec_shard_needle_distributed keeps its Ok(None) for every miss, so the other callers are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- seaweed-volume/src/server/grpc_server.rs | 323 +++++++++++++++++++---- seaweed-volume/src/server/store_ec.rs | 96 +++++-- 2 files changed, 342 insertions(+), 77 deletions(-) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 1e98ec462..202e4b7f4 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -47,6 +47,11 @@ fn scrub_mode_label(mode: i32) -> &'static str { } } +/// Go formats the id in decimal; fs.verify matches the "needle not found " prefix. +fn needle_not_found(needle_id: NeedleId) -> Status { + Status::not_found(format!("needle not found {}", needle_id.0)) +} + fn unix_now_seconds() -> f64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -5023,64 +5028,72 @@ impl VolumeServer for VolumeGrpcService { let vid = VolumeId(req.volume_id); let needle_id = NeedleId(req.needle_id); - let store = self.state.store.read().unwrap(); + { + let store = self.state.store.read().unwrap(); - // Try normal volume first - if store.find_volume(vid).is_some() { - let mut n = Needle { - id: needle_id, - ..Needle::default() - }; - match store.read_volume_needle(vid, &mut n) { - Ok(_) => { - let ttl_str = n.ttl.as_ref().map_or(String::new(), |t| t.to_string()); - return Ok(Response::new( - volume_server_pb::VolumeNeedleStatusResponse { - needle_id: n.id.0, - cookie: n.cookie.0, - size: n.size.0 as u32, - last_modified: n.last_modified, - crc: n.checksum.0, - ttl: ttl_str, - }, - )); - } - Err(_) => return Err(Status::not_found(format!("needle not found {}", needle_id))), - } - } - - // Fall back to EC shards — read full needle from local shards - if let Some(ec_vol) = store.find_ec_volume(vid) { - match ec_vol.read_ec_shard_needle(needle_id) { - Ok(Some(n)) => { - let ttl_str = match &n.ttl { - Some(t) if n.has_ttl() => t.to_string(), - _ => String::new(), - }; - return Ok(Response::new( - volume_server_pb::VolumeNeedleStatusResponse { - needle_id: n.id.0, - cookie: n.cookie.0, - size: n.size.0 as u32, - last_modified: n.last_modified, - crc: n.checksum.0, - ttl: ttl_str, - }, - )); - } - Ok(None) => { - return Err(Status::not_found(format!("needle not found {}", needle_id))); - } - Err(e) => { - return Err(Status::internal(format!( - "read ec shard needle {} from volume {}: {}", - needle_id, vid, e - ))); + // Try normal volume first + if store.find_volume(vid).is_some() { + let mut n = Needle { + id: needle_id, + ..Needle::default() + }; + match store.read_volume_needle(vid, &mut n) { + Ok(_) => { + let ttl_str = n.ttl.as_ref().map_or(String::new(), |t| t.to_string()); + return Ok(Response::new( + volume_server_pb::VolumeNeedleStatusResponse { + needle_id: n.id.0, + cookie: n.cookie.0, + size: n.size.0 as u32, + last_modified: n.last_modified, + crc: n.checksum.0, + ttl: ttl_str, + }, + )); + } + Err(crate::storage::volume::VolumeError::NotFound) => { + return Err(needle_not_found(needle_id)); + } + // fs.verify skips "already deleted" by message, like Go's plain ErrorDeleted. + Err(e) => return Err(Status::unknown(e.to_string())), } } } - Err(Status::not_found(format!("volume not found {}", vid))) + // Intervals on shards held by other nodes are fetched or reconstructed, as in Go. + use crate::server::store_ec::EcMiss; + match crate::server::store_ec::read_ec_shard_needle_or_miss(&self.state, vid, needle_id) + .await + { + Ok(Ok(n)) => { + let ttl_str = match &n.ttl { + Some(t) if n.has_ttl() => t.to_string(), + _ => String::new(), + }; + Ok(Response::new( + volume_server_pb::VolumeNeedleStatusResponse { + needle_id: n.id.0, + cookie: n.cookie.0, + size: n.size.0 as u32, + last_modified: n.last_modified, + crc: n.checksum.0, + ttl: ttl_str, + }, + )) + } + Ok(Err(EcMiss::NotFound)) => Err(needle_not_found(needle_id)), + // fs.verify skips "already deleted" by message, like Go's plain ErrorDeleted. + Ok(Err(EcMiss::Deleted)) => Err(Status::unknown( + crate::storage::volume::VolumeError::Deleted.to_string(), + )), + Ok(Err(EcMiss::VolumeNotFound)) => { + Err(Status::not_found(format!("volume not found {}", vid))) + } + Err(e) => Err(Status::unknown(format!( + "read ec shard needle {} from volume {}: {}", + needle_id, vid, e + ))), + } } async fn ping( @@ -9750,6 +9763,210 @@ mod tests { assert_eq!(resp.total_files, 1); } + async fn needle_status( + service: &VolumeGrpcService, + needle_id: u64, + ) -> Result { + service + .volume_needle_status(Request::new(volume_server_pb::VolumeNeedleStatusRequest { + volume_id: 1, + needle_id, + })) + .await + .map(Response::into_inner) + } + + async fn generate_and_mount_ec_1(service: &VolumeGrpcService, shard_ids: Vec) { + service + .volume_ec_shards_generate(Request::new( + volume_server_pb::VolumeEcShardsGenerateRequest { + volume_id: 1, + collection: String::new(), + }, + )) + .await + .unwrap(); + mount_ec_1(service, shard_ids).await; + } + + async fn mount_ec_1(service: &VolumeGrpcService, shard_ids: Vec) { + service + .volume_ec_shards_mount(Request::new(volume_server_pb::VolumeEcShardsMountRequest { + volume_id: 1, + collection: String::new(), + shard_ids, + source_disk_type: String::new(), + recover_missing_index: false, + })) + .await + .unwrap(); + } + + /// Leaves `service` holding only EC shard 1 of volume 1, with shard 0 (the + /// one with needle 11's bytes) served by a peer; `peer_deletes` journals + /// the needle's delete on the peer alone. + async fn put_ec_1_needle_shard_on_a_peer( + service: &VolumeGrpcService, + tmp: &TempDir, + peer_deletes: bool, + ) -> (TempDir, tokio::sync::oneshot::Sender<()>) { + // One local shard, not the needle's: only a peer read can answer. + generate_and_mount_ec_1(service, vec![1]).await; + + let (peer, peer_tmp) = make_local_service_with_volume("", None); + peer.state + .store + .write() + .unwrap() + .unmount_volume(VolumeId(1)) + .unwrap(); + for entry in std::fs::read_dir(tmp.path()).unwrap() { + let name = entry.unwrap().file_name().into_string().unwrap(); + if name.starts_with("1.ec") || name == "1.vif" { + std::fs::copy(tmp.path().join(&name), peer_tmp.path().join(&name)).unwrap(); + } + } + mount_ec_1(&peer, vec![0]).await; + if peer_deletes { + peer.state + .store + .write() + .unwrap() + .find_ec_volume_mut(VolumeId(1)) + .unwrap() + .journal_delete(NeedleId(11)) + .unwrap(); + } + let (port, shutdown) = serve_source(peer).await; + + { + let mut store = service.state.store.write().unwrap(); + store.unmount_volume(VolumeId(1)).unwrap(); + store + .find_ec_volume(VolumeId(1)) + .unwrap() + .merge_shard_locations( + (0u8..14) + .map(|sid| { + let addr = if sid == 0 { + format!("127.0.0.1:{port}.{port}") + } else { + "127.0.0.1:255.1".to_string() + }; + (sid, vec![addr]) + }) + .collect(), + ); + } + (peer_tmp, shutdown) + } + + /// fs.verify asks every EC shard holder, so a node that does not hold the + /// shard with the needle's bytes must fetch them from a peer, as Go does. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_volume_needle_status_reads_an_ec_needle_from_a_peer_shard() { + let (service, tmp) = make_local_service_with_volume("", None); + let want = needle_status(&service, 11).await.unwrap(); + let _peer = put_ec_1_needle_shard_on_a_peer(&service, &tmp, false).await; + + let got = needle_status(&service, 11) + .await + .expect("the needle's shard is on a reachable peer"); + assert_eq!(got, want); + } + + /// The local .ecx still shows the needle live; the peer's answer that it is + /// deleted must reach fs.verify as Go's ErrorDeleted, not as a missing needle. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_volume_needle_status_reports_a_peer_reported_ec_deletion() { + let (service, tmp) = make_local_service_with_volume("", None); + let _peer = put_ec_1_needle_shard_on_a_peer(&service, &tmp, true).await; + + let err = needle_status(&service, 11).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::Unknown, "{err}"); + assert_eq!(err.message(), "already deleted"); + } + + /// A volume unmounted before the EC read resolves it is "volume not found", + /// which fs.verify does not count as a lost needle. + #[tokio::test] + async fn test_volume_needle_status_reports_a_vanished_volume_as_not_found() { + let (service, _tmp) = make_local_service_with_volume("", None); + service + .state + .store + .write() + .unwrap() + .unmount_volume(VolumeId(1)) + .unwrap(); + + let miss = crate::server::store_ec::read_ec_shard_needle_or_miss( + &service.state, + VolumeId(1), + NeedleId(11), + ) + .await + .unwrap() + .unwrap_err(); + assert_eq!(miss, crate::server::store_ec::EcMiss::VolumeNotFound); + + let err = needle_status(&service, 11).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::NotFound, "{err}"); + assert_eq!(err.message(), "volume not found 1"); + } + + /// Go returns ErrorDeleted as a plain error, which fs.verify skips by + /// message; a NotFound "needle not found" would class it as missing. + #[tokio::test] + async fn test_volume_needle_status_reports_deleted_needles_like_go() { + let (service, _tmp) = make_local_service_with_volume("", None); + + let err = needle_status(&service, 12345).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::NotFound, "{err}"); + assert_eq!(err.message(), "needle not found 12345"); + + generate_and_mount_ec_1(&service, (0..14).collect()).await; + { + let mut store = service.state.store.write().unwrap(); + let mut n = Needle { + id: NeedleId(11), + cookie: Cookie(0x3344), + ..Needle::default() + }; + store.delete_volume_needle(VolumeId(1), &mut n).unwrap(); + } + let err = needle_status(&service, 11).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::Unknown, "{err}"); + assert_eq!(err.message(), "already deleted"); + + // The EC copy still has the needle live until its own delete lands. + service + .state + .store + .write() + .unwrap() + .unmount_volume(VolumeId(1)) + .unwrap(); + assert_eq!(needle_status(&service, 11).await.unwrap().cookie, 0x3344); + + service + .state + .store + .write() + .unwrap() + .find_ec_volume_mut(VolumeId(1)) + .unwrap() + .journal_delete(NeedleId(11)) + .unwrap(); + let err = needle_status(&service, 11).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::Unknown, "{err}"); + assert_eq!(err.message(), "already deleted"); + + let err = needle_status(&service, 12345).await.unwrap_err(); + assert_eq!(err.code(), tonic::Code::NotFound, "{err}"); + assert_eq!(err.message(), "needle not found 12345"); + } + /// Batch atomicity: mount pre-validates the ENTIRE shard_ids before /// acquiring the write lock or mounting anything. A batch like [0, 32] /// must fail with InvalidArgument and mount NOTHING — not the valid diff --git a/seaweed-volume/src/server/store_ec.rs b/seaweed-volume/src/server/store_ec.rs index a90effccc..6d8b3ca12 100644 --- a/seaweed-volume/src/server/store_ec.rs +++ b/seaweed-volume/src/server/store_ec.rs @@ -95,20 +95,41 @@ struct Snapshot { encode_ts_ns: i64, } -/// Top-level entry point. Returns `Ok(None)` for "not found" (matches -/// Go's `ReadEcShardNeedle`); errors propagate as `io::Error`. +/// Why a distributed EC read has no needle to return. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum EcMiss { + NotFound, + /// Tombstoned in the local `.ecx`/`.ecj`, or reported deleted by a peer. + Deleted, + VolumeNotFound, +} + +/// Top-level entry point. Returns `Ok(None)` for any miss — absent, +/// deleted, or volume gone; errors propagate as `io::Error`. pub async fn read_ec_shard_needle_distributed( state: &Arc, vid: VolumeId, needle_id: NeedleId, ) -> io::Result> { + Ok(read_ec_shard_needle_or_miss(state, vid, needle_id) + .await? + .ok()) +} + +/// Like `read_ec_shard_needle_distributed`, but says why there is no needle, +/// as Go's `ReadEcShardNeedle` tells `ErrorDeleted` from not-found. +pub async fn read_ec_shard_needle_or_miss( + state: &Arc, + vid: VolumeId, + needle_id: NeedleId, +) -> io::Result> { // Phase A — under the Store read lock, locate the needle, compute // intervals, and read any locally-mounted shard intervals. We must // not `.await` while holding this guard (std::sync::RwLockReadGuard // is !Send). let mut snapshot = match snapshot_under_lock(state, vid, needle_id)? { - Some(s) => s, - None => return Ok(None), + Ok(s) => s, + Err(miss) => return Ok(Err(miss)), }; // Phase B — refresh the shard_locations cache from the master if @@ -210,17 +231,11 @@ pub async fn read_ec_shard_needle_distributed( .collect() .await; - let mut assembled: Vec> = Vec::with_capacity(fetched.len()); - for res in fetched { - let (buf, is_deleted) = res?; - // A peer reports the needle deleted (a cross-server window where the - // local index still shows it live): treat as not-found rather than - // serving zeros, mirroring Go's ErrorDeleted. - if is_deleted { - return Ok(None); - } - assembled.push(buf); - } + // A peer reports the needle deleted (a cross-server window where the + // local index still shows it live): answer deleted rather than serving zeros. + let Some(assembled) = gather_intervals(fetched)? else { + return Ok(Err(EcMiss::Deleted)); + }; // Phase D — assemble and parse the Needle. Mirrors the tail of // `EcVolume::read_ec_shard_needle`. @@ -252,7 +267,20 @@ pub async fn read_ec_shard_needle_distributed( snapshot.version, ) .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, format!("{}", e)))?; - Ok(Some(n)) + Ok(Ok(n)) +} + +/// `None` when any holder reported the needle deleted. That outranks another +/// interval's error: deletes are never invented, so the needle is gone either way. +fn gather_intervals(fetched: Vec, bool)>>) -> io::Result>>> { + if fetched.iter().any(|r| matches!(r, Ok((_, true)))) { + return Ok(None); + } + fetched + .into_iter() + .map(|r| r.map(|(buf, _)| buf)) + .collect::>() + .map(Some) } /// What one EC delete RPC carries — `VolumeEcBlobDeleteRequest` minus tonic. @@ -864,11 +892,10 @@ fn snapshot_under_lock( state: &Arc, vid: VolumeId, needle_id: NeedleId, -) -> io::Result> { +) -> io::Result> { let store = state.store.read().unwrap(); - let ecv = match store.find_ec_volume(vid) { - Some(v) => v, - None => return Ok(None), + let Some(ecv) = store.find_ec_volume(vid) else { + return Ok(Err(EcMiss::VolumeNotFound)); }; // Reuse EcVolume::locate_needle for offset/size resolution AND @@ -876,11 +903,17 @@ fn snapshot_under_lock( // local-only read path uses, so we stay byte-identical on the // shard-size + interval boundaries. locate_needle applies the runtime // delete mask, which is correct for serving reads. - let (offset, size, intervals) = match ecv.locate_needle(needle_id)? { - Some(v) => v, - None => return Ok(None), + let Some((offset, size, intervals)) = ecv.locate_needle(needle_id)? else { + // locate_needle folds a tombstone into not-found. + let deleted = + matches!(ecv.find_needle_from_ecx(needle_id)?, Some((_, s)) if s.is_deleted()); + return Ok(Err(if deleted { + EcMiss::Deleted + } else { + EcMiss::NotFound + })); }; - build_snapshot(ecv, offset, size, &intervals).map(Some) + build_snapshot(ecv, offset, size, &intervals).map(Ok) } /// Like `snapshot_under_lock`, but locates intervals from the RAW .ecx @@ -1790,6 +1823,21 @@ async fn drain_copy_stream( mod tests { use super::*; + #[test] + fn gather_intervals_puts_a_reported_deletion_ahead_of_errors() { + let failed = || Err(io::Error::other("shard unreachable")); + assert!( + gather_intervals(vec![failed(), Ok((Vec::new(), true))]) + .unwrap() + .is_none() + ); + assert!(gather_intervals(vec![failed(), Ok((vec![1], false))]).is_err()); + assert_eq!( + gather_intervals(vec![Ok((vec![1], false)), Ok((vec![2], false))]).unwrap(), + Some(vec![vec![1], vec![2]]) + ); + } + fn locations(count: usize) -> HashMap> { (0..count) .map(|sid| (sid as ShardId, vec!["127.0.0.1:8080".to_string()]))