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 <decimal id>", 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) <noreply@anthropic.com>

* 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) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Eliah RusinandClaude Opus 5.5 authored and GitHub committed 2026-10-01 20:18:34 +08:00
1 parent cc1ec48151
commit b8f074b7d3
2 files changed
+342 -77

No files matched your search

+270 -53
View File
@@ -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<volume_server_pb::VolumeNeedleStatusResponse, Status> {
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<u32>) {
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<u32>) {
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
+72 -24
View File
@@ -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<VolumeServerState>,
vid: VolumeId,
needle_id: NeedleId,
) -> io::Result<Option<Needle>> {
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<VolumeServerState>,
vid: VolumeId,
needle_id: NeedleId,
) -> io::Result<Result<Needle, EcMiss>> {
// 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<u8>> = 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<io::Result<(Vec<u8>, bool)>>) -> io::Result<Option<Vec<Vec<u8>>>> {
if fetched.iter().any(|r| matches!(r, Ok((_, true)))) {
return Ok(None);
}
fetched
.into_iter()
.map(|r| r.map(|(buf, _)| buf))
.collect::<io::Result<_>>()
.map(Some)
}
/// What one EC delete RPC carries — `VolumeEcBlobDeleteRequest` minus tonic.
@@ -864,11 +892,10 @@ fn snapshot_under_lock(
state: &Arc<VolumeServerState>,
vid: VolumeId,
needle_id: NeedleId,
) -> io::Result<Option<Snapshot>> {
) -> io::Result<Result<Snapshot, EcMiss>> {
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<ShardId, Vec<String>> {
(0..count)
.map(|sid| (sid as ShardId, vec!["127.0.0.1:8080".to_string()]))