fix(ec): correct EC FULL scrub for deleted needles, shard-location cache, and parity coverage (#10152)

* fix(ec): correct EC FULL scrub for deleted needles + shard-location cache

Addresses review findings on the EC FULL distributed scrub:
- Remote EC reads now thread Go's (bytes, is_deleted) contract. A runtime EC
  delete keeps the .ecx size positive (the delete lives in .ecj/memory), so the
  raw-index walk verifies the needle, and its header interval is usually remote;
  the peer answers is_deleted with no payload. The scrub zero-fills that interval
  (so the needle reaches read_bytes -> SizeMismatch{0} -> the delete-state
  suppression), the serving direct read short-circuits to not-found, and
  reconstruction EXCLUDES the shard instead of feeding zeros into Reed-Solomon.
- The walk skips size.is_deleted() (not just is_tombstone), so a -originalSize
  .ecx entry (pre-encode delete) can't yield empty intervals or panic parse_header.
- Restore Go's < data_shards completeness guard (per-volume, custom-ratio aware)
  and per-shard merge in the location cache instead of clobber-with-partial.
- Abort the scrub with an error on mid-scan unmount instead of a false-CLEAN.
- Hoist the refreshed location map once instead of cloning it per needle.

Claude-Session: https://claude.ai/code/session_015EE9Sc9EvNp8BCVva4RKdo

* feat(scrub): keep RS parity check in EC FULL until CHECKSUM lands

The per-needle FULL walk only reads live data-shard intervals, so it can't catch
bitrot in a parity shard or an unwalked cold region. Run verify_ec_shards
alongside the walk, gated on all-shards-local (single-node EC), via spawn_blocking.
A deliberate temporary divergence from Go FULL; moves to mode 4 (CHECKSUM) once
the .ecsum subsystem lands.

Claude-Session: https://claude.ai/code/session_015EE9Sc9EvNp8BCVva4RKdo
This commit is contained in:
Chris Lu
2026-06-30 10:18:31 -07:00
committed by GitHub
parent a85318111c
commit 8f2a2abae4
4 changed files with 266 additions and 72 deletions
+65 -12
View File
@@ -3962,22 +3962,75 @@ impl VolumeServer for VolumeGrpcService {
}
}
2 => {
// FULL: verify every needle across local AND remote shards.
{
// Match Go: a missing requested volume is a hard error.
// FULL: Go-parity per-needle local+remote walk, PLUS a TEMPORARY
// local Reed-Solomon parity check. The needle walk only reads
// DATA-shard intervals of LIVE needles, so on its own it can't
// catch silent bitrot in a PARITY shard or an unwalked cold
// region. Go closes that gap with a separate CHECKSUM mode over
// .ecsum, which Rust does not have yet; running both here is a
// deliberate divergence from Go FULL to preserve coverage. Drop
// verify_ec_shards from this arm once mode 4 (CHECKSUM) lands.
//
// The RS recompute needs every shard co-located, so only run it
// when this node holds all data+parity shards (single-node EC);
// on a distributed layout it would report every non-local shard
// as missing. Snapshot under a brief lock; release before await.
let (dir, collection, data_shards, parity_shards, all_local) = {
let store = self.state.store.read().unwrap();
if store.find_ec_volume(vid).is_none() {
return Err(Status::not_found(format!(
"EC volume id {} not found",
vid.0
)));
}
}
let ecv = store.find_ec_volume(vid).ok_or_else(|| {
Status::not_found(format!("EC volume id {} not found", vid.0))
})?;
let total = (ecv.data_shards + ecv.parity_shards) as usize;
let local = ecv.shards.iter().filter(|s| s.is_some()).count();
(
ecv.dir.clone(),
ecv.collection.clone(),
ecv.data_shards as usize,
ecv.parity_shards as usize,
local == total,
)
};
total_volumes += 1;
let (files, shard_infos, errs) =
// (1) Per-needle local+remote walk (Go ScrubEcVolume parity).
let (files, mut shard_infos, mut errs) =
crate::server::store_ec::scrub_ec_volume_distributed(&self.state, vid, false)
.await;
total_files += files as u64;
total_files += files as u64; // count comes from the needle walk only
// (2) Local parity check, gated on all-shards-local. Blocking RS
// verify -> spawn_blocking; inputs are owned, no lock held.
if all_local && !dir.is_empty() {
let collection_pc = collection.clone();
let (parity_broken, parity_details) = tokio::task::spawn_blocking(move || {
crate::storage::erasure_coding::ec_encoder::verify_ec_shards(
&dir,
&collection_pc,
vid,
data_shards,
parity_shards,
)
})
.await
.map_err(|e| Status::internal(format!("verify_ec_shards join: {}", e)))?
.unwrap_or_else(|e| (Vec::new(), vec![format!("verify_ec_shards: {}", e)]));
let mut seen: std::collections::HashSet<u32> =
shard_infos.iter().map(|s| s.shard_id).collect();
for sid in parity_broken {
if seen.insert(sid) {
shard_infos.push(volume_server_pb::EcShardInfo {
shard_id: sid,
collection: collection.clone(),
volume_id: vid.0,
..Default::default()
});
}
}
shard_infos.sort_by_key(|s| s.shard_id);
errs.extend(parity_details);
}
if !errs.is_empty() || !shard_infos.is_empty() {
broken_volume_ids.push(vid.0);
broken_shard_infos.extend(shard_infos);
+168 -49
View File
@@ -115,8 +115,13 @@ pub async fn read_ec_shard_needle_distributed(
{
match cached_lookup_ec_shard_locations(state, vid).await {
Ok(fresh) => {
shard_locations = fresh.clone();
write_back_shard_locations(state, vid, fresh);
// A complete reply merges into the cache; an incomplete one
// (< data_shards) is left unwritten — keep the prior cache.
if let Some(merged) =
write_back_shard_locations(state, vid, fresh, snapshot.data_shards as usize)
{
shard_locations = merged;
}
}
Err(e) => {
// Lookup failed — proceed with cached values. If cache
@@ -144,7 +149,7 @@ pub async fn read_ec_shard_needle_distributed(
shard_offset,
size,
} => {
let buf = fetch_one_interval(
let (buf, is_deleted) = fetch_one_interval(
state,
vid,
needle_id,
@@ -157,6 +162,12 @@ pub async fn read_ec_shard_needle_distributed(
snapshot.encode_ts_ns,
)
.await?;
// 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);
}
}
@@ -241,11 +252,23 @@ pub async fn scrub_ec_volume_distributed(
let mut errs = seed_errs;
// Refresh the shard-location cache once up front (mirrors Go's
// cachedLookupEcShardLocations); a failed lookup is a hard error rather than
// a per-needle connect storm against a down master.
// cachedLookupEcShardLocations). A partial reply (< data_shards locations, a
// master mid-recovery) or a failed lookup is a hard, retryable error — never
// overwrite a good cache with a partial map or storm a down master per needle.
if needs_refresh(&cached_locations, cache_refreshed_at, data_shards, total_shards) {
match cached_lookup_ec_shard_locations(state, vid).await {
Ok(fresh) => write_back_shard_locations(state, vid, fresh),
Ok(fresh) => {
if write_back_shard_locations(state, vid, fresh, data_shards).is_none() {
return (
0,
Vec::new(),
vec![format!(
"failed to locate shard via master grpc: fewer than {} data-shard locations returned",
data_shards
)],
);
}
}
Err(e) => {
return (
0,
@@ -256,6 +279,24 @@ pub async fn scrub_ec_volume_distributed(
}
}
// Hoist the post-refresh shard-location map once; it is stable for the whole
// walk, so per-needle snapshots no longer clone it.
let locations: HashMap<ShardId, Vec<String>> = {
let store = state.store.read().unwrap();
let ecv = match store.find_ec_volume(vid) {
Some(v) => v,
None => {
return (
0,
Vec::new(),
vec![format!("EC volume id {} not found", vid.0)],
)
}
};
let map = ecv.shard_locations.read().unwrap().clone();
map
};
// Walk the .ecx (private fd, no lock) for the row count + live (id, offset, size).
let mut count: i64 = 0;
let mut needles: Vec<(NeedleId, Offset, Size)> = Vec::new();
@@ -263,7 +304,12 @@ pub async fn scrub_ec_volume_distributed(
Ok(mut f) => {
if let Err(e) = crate::storage::idx::walk_index_file(&mut f, 0, |id, offset, size| {
count += 1;
if !size.is_tombstone() {
// Skip ALL deleted entries: -1 tombstones (runtime delete folded
// into .ecx) and -originalSize entries (a needle deleted on the
// regular volume before EC encode). get_actual_size uses the raw
// signed size, so a negative would yield empty intervals
// (false-positive) or an under-16-byte buffer (parse panic).
if !size.is_deleted() {
needles.push((id, offset, size));
}
Ok(())
@@ -281,8 +327,13 @@ pub async fn scrub_ec_volume_distributed(
// Per-needle snapshot under the lock from the RAW .ecx (offset, size) so
// logically-deleted needles are still verified; lock dropped before await.
let snapshot = match scrub_snapshot_under_lock(state, vid, offset, size) {
Ok(Some(s)) => s,
Ok(None) => continue, // volume gone since the walk
Ok(s) => s,
// Volume unmounted mid-scan: abort with an error rather than skipping
// every remaining needle, which would report a false-CLEAN result.
Err(e) if e.kind() == io::ErrorKind::NotFound => {
errs.push(format!("EC volume {} unmounted during scrub: {}", vid.0, e));
break;
}
Err(e) => {
errs.push(format!("needle {} on EC volume {}: {}", id.0, vid.0, e));
continue;
@@ -301,14 +352,11 @@ pub async fn scrub_ec_volume_distributed(
shard_offset,
size: ssize,
} => {
let sources = snapshot
.cached_locations
.get(shard_id)
.cloned()
.unwrap_or_default();
let sources: &[String] =
locations.get(shard_id).map(Vec::as_slice).unwrap_or(&[]);
match read_remote_ec_shard_interval(
state,
&sources,
sources,
vid,
id,
*shard_id,
@@ -318,7 +366,11 @@ pub async fn scrub_ec_volume_distributed(
)
.await
{
Ok(buf) => data.extend_from_slice(&buf),
// A deleted shard yields no bytes; zero-fill the interval so
// the assembled needle reaches read_bytes -> SizeMismatch{0}
// -> the delete-state suppression (mirrors Go's pre-zeroed buffer).
Ok((_, true)) => data.resize(data.len() + *ssize, 0),
Ok((buf, false)) => data.extend_from_slice(&buf),
Err(_) => {
errs.push(format!(
"failed to read EC shard {} for needle {} on volume {} (interval {}/{})",
@@ -408,33 +460,52 @@ fn scrub_snapshot_under_lock(
vid: VolumeId,
offset: Offset,
size: Size,
) -> io::Result<Option<Snapshot>> {
) -> io::Result<ScrubSnapshot> {
let store = state.store.read().unwrap();
let ecv = match store.find_ec_volume(vid) {
Some(v) => v,
None => return Ok(None),
// Volume unmounted mid-scan: a distinct NotFound so the caller aborts
// with an error rather than silently skipping (which would false-CLEAN).
None => {
return Err(io::Error::new(
io::ErrorKind::NotFound,
format!("EC volume {} not found (unmounted mid-scan)", vid.0),
))
}
};
let intervals = ecv.locate_ec_shard_needle_interval(offset.to_actual_offset(), size);
build_snapshot(ecv, offset, size, &intervals).map(Some)
}
/// Read any locally-held shard intervals and assemble the rest as `NeedRemote`,
/// then snapshot the scalars + shard-location cache so the caller can drop the
/// store lock before awaiting remote reads. Shared by the read and scrub paths.
fn build_snapshot(
ecv: &crate::storage::erasure_coding::ec_volume::EcVolume,
offset: Offset,
size: Size,
intervals: &[crate::storage::erasure_coding::ec_locate::Interval],
) -> io::Result<Snapshot> {
if intervals.is_empty() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"no intervals for needle",
));
}
let actual = get_actual_size(size, ecv.version);
Ok(ScrubSnapshot {
version: ecv.version,
actual_size: get_actual_size(size, ecv.version) as usize,
size_for_parse: size,
intervals: read_local_intervals(ecv, &intervals),
encode_ts_ns: ecv.encode_ts_ns,
})
}
/// Scalars + locally-read intervals for the FULL scrub. Unlike `Snapshot` it
/// omits the shard-location cache: the scrub hoists the refreshed map once up
/// front instead of cloning it per needle.
struct ScrubSnapshot {
version: Version,
actual_size: usize,
size_for_parse: Size,
intervals: Vec<IntervalResult>,
encode_ts_ns: i64,
}
/// Read any locally-held shard intervals, marking the rest `NeedRemote`. Shared
/// by the read-path `build_snapshot` and the scrub snapshot.
fn read_local_intervals(
ecv: &crate::storage::erasure_coding::ec_volume::EcVolume,
intervals: &[crate::storage::erasure_coding::ec_locate::Interval],
) -> Vec<IntervalResult> {
let mut interval_results = Vec::with_capacity(intervals.len());
for interval in intervals {
let (shard_id, shard_offset) = interval.to_shard_id_and_offset(ecv.data_shards);
@@ -459,7 +530,25 @@ fn build_snapshot(
}),
}
}
interval_results
}
/// Read local intervals, then snapshot the scalars + shard-location cache so the
/// read path can drop the store lock before awaiting remote reads.
fn build_snapshot(
ecv: &crate::storage::erasure_coding::ec_volume::EcVolume,
offset: Offset,
size: Size,
intervals: &[crate::storage::erasure_coding::ec_locate::Interval],
) -> io::Result<Snapshot> {
if intervals.is_empty() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"no intervals for needle",
));
}
let actual = get_actual_size(size, ecv.version);
let interval_results = read_local_intervals(ecv, intervals);
let cached_locations = ecv.shard_locations.read().unwrap().clone();
let cache_refreshed_at = *ecv.shard_locations_refresh_time.lock().unwrap();
@@ -555,19 +644,29 @@ async fn cached_lookup_ec_shard_locations(
Ok(out)
}
/// Merge a fresh `LookupEcVolume` reply into the per-EcVolume shard-location
/// cache. Returns the merged map on success, or `None` when the reply is
/// incomplete (and was therefore NOT written) or the volume is gone.
///
/// Completeness guard (mirrors Go's `cachedLookupEcShardLocations`): a reply
/// carrying fewer than `data_shards` shard locations is a master-mid-recovery
/// partial, not ground truth — Go aborts the lookup with a retryable error and
/// leaves the cache + refresh time untouched. Returning `None` without writing
/// does the same, so a previously-complete cache is never clobbered with a
/// partial map. On a complete reply the write is a per-shard MERGE (see
/// `merge_shard_locations`), not a full replace.
fn write_back_shard_locations(
state: &Arc<VolumeServerState>,
vid: VolumeId,
locations: HashMap<ShardId, Vec<String>>,
) {
let store = state.store.read().unwrap();
if let Some(ecv) = store.find_ec_volume(vid) {
// Atomic swap + freshness stamp so a concurrent reader sees
// either the prior cache or the fresh one — never an
// intermediate half-replaced map with the freshness flag
// already flipped.
ecv.replace_shard_locations(locations);
data_shards: usize,
) -> Option<HashMap<ShardId, Vec<String>>> {
if locations.len() < data_shards {
return None;
}
let store = state.store.read().unwrap();
let ecv = store.find_ec_volume(vid)?;
Some(ecv.merge_shard_locations(locations))
}
/// Build a SeaweedFS-style `host:httpPort.grpcPort` address from a
@@ -599,7 +698,7 @@ async fn fetch_one_interval(
data_shards: usize,
parity_shards: usize,
expected_encode_ts_ns: i64,
) -> io::Result<Vec<u8>> {
) -> io::Result<(Vec<u8>, bool)> {
// Direct peer read against the cached locations for this shard.
if let Some(sources) = shard_locations.get(&shard_id) {
if !sources.is_empty() {
@@ -615,7 +714,9 @@ async fn fetch_one_interval(
)
.await
{
Ok(buf) => return Ok(buf),
// A deleted needle short-circuits: don't reconstruct (every shard
// would report deleted), let the caller return "deleted".
Ok((buf, is_deleted)) => return Ok((buf, is_deleted)),
Err(e) => {
tracing::debug!(
"direct read ec shard {}.{} from {:?} failed: {} — will reconstruct",
@@ -631,7 +732,7 @@ async fn fetch_one_interval(
// Reconstruct: fan-out reads to every other shard at the same
// (shard_offset, size). Mirrors `recoverOneRemoteEcShardInterval`.
recover_one_remote_ec_shard_interval(
let buf = recover_one_remote_ec_shard_interval(
state,
vid,
needle_id,
@@ -643,7 +744,8 @@ async fn fetch_one_interval(
parity_shards,
expected_encode_ts_ns,
)
.await
.await?;
Ok((buf, false))
}
async fn read_remote_ec_shard_interval(
@@ -655,7 +757,7 @@ async fn read_remote_ec_shard_interval(
shard_offset: i64,
size: usize,
expected_encode_ts_ns: i64,
) -> io::Result<Vec<u8>> {
) -> io::Result<(Vec<u8>, bool)> {
let mut last_err: Option<io::Error> = None;
for src in sources {
match do_read_remote_ec_shard_interval(
@@ -670,7 +772,7 @@ async fn read_remote_ec_shard_interval(
)
.await
{
Ok(buf) => return Ok(buf),
Ok(res) => return Ok(res),
Err(e) => last_err = Some(e),
}
}
@@ -691,7 +793,7 @@ async fn do_read_remote_ec_shard_interval(
shard_offset: i64,
size: usize,
expected_encode_ts_ns: i64,
) -> io::Result<Vec<u8>> {
) -> io::Result<(Vec<u8>, bool)> {
let grpc_addr =
parse_grpc_address(source).map_err(|e| io::Error::new(io::ErrorKind::InvalidInput, e))?;
let endpoint = build_grpc_endpoint(&grpc_addr, state.outgoing_grpc_tls.as_ref())
@@ -740,6 +842,7 @@ async fn do_read_remote_ec_shard_interval(
let mut stream = resp.into_inner();
let mut out = Vec::with_capacity(size);
let mut is_deleted = false;
while let Some(msg) = stream
.message()
.await
@@ -757,10 +860,23 @@ async fn do_read_remote_ec_shard_interval(
),
));
}
if msg.is_deleted {
is_deleted = true;
}
if !msg.data.is_empty() {
out.extend_from_slice(&msg.data);
}
}
// A runtime EC delete keeps the .ecx size positive; the holder masks the delete
// at read time and answers is_deleted with no payload. Signal the deletion to
// the caller (Go's `(bytes, is_deleted)` contract) instead of synthesizing bytes
// here: the scrub zero-fills the interval (so the assembled needle hits read_bytes
// -> SizeMismatch{found:0} -> suppression), the serving direct read short-circuits
// to "deleted", and reconstruction EXCLUDES the shard rather than feeding zeros
// into Reed-Solomon. Exempt from the short-read guard below.
if is_deleted {
return Ok((Vec::new(), true));
}
if out.len() < size {
return Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
@@ -775,7 +891,7 @@ async fn do_read_remote_ec_shard_interval(
));
}
out.truncate(size);
Ok(out)
Ok((out, false))
}
async fn recover_one_remote_ec_shard_interval(
@@ -859,8 +975,11 @@ async fn recover_one_remote_ec_shard_interval(
for (sid, res) in results {
match res {
Ok(buf) => {
if (sid as usize) < total_shards {
// Exclude a deleted shard from reconstruction (Go gates on a full
// read): feeding the empty/zero buffer into Reed-Solomon would
// corrupt the recovered shard.
Ok((buf, is_deleted)) => {
if !is_deleted && (sid as usize) < total_shards {
bufs[sid as usize] = Some(buf);
}
}
@@ -207,13 +207,12 @@ pub fn rebuild_ec_files(
Ok(())
}
/// Verify EC shards by computing parity against the existing data and identifying corrupted shards.
///
/// No longer wired to a scrub mode: FULL is now Go's per-needle local+remote walk
/// (`store_ec::scrub_ec_volume_distributed`). Retained as a standalone Reed-Solomon
/// parity check — the only parity/cold-region verification available until the
/// `.ecsum` bitrot CHECKSUM path lands.
#[allow(dead_code)]
/// Reed-Solomon parity check over the locally-held shards: recompute parity from
/// the data shards and flag any shard whose bytes disagree. Wired into FULL
/// (scrub mode 2) as a TEMPORARY local parity/cold-region check — the per-needle
/// FULL walk only reads live data-shard intervals, so on its own it can't catch
/// bitrot in a parity shard or an unwalked region. Move to mode 4 (CHECKSUM) and
/// drop it from mode 2 once the `.ecsum` subsystem lands.
pub fn verify_ec_shards(
dir: &str,
collection: &str,
@@ -464,10 +464,10 @@ impl EcVolume {
/// against a half-populated cache and return NotFound for the
/// not-yet-inserted shards.
///
/// Callers populating the whole cache atomically should use
/// [`Self::replace_shard_locations`] instead — it swaps the
/// entire map under the write lock and advances the refresh
/// timestamp in one step.
/// Callers writing back a whole `LookupEcVolume` reply should use
/// [`Self::merge_shard_locations`] instead — it upserts the reply's
/// shards under the write lock and advances the refresh timestamp in
/// one step, retaining cached shards the reply omits.
pub fn set_shard_locations(&self, shard_id: ShardId, locations: Vec<String>) {
self.shard_locations
.write()
@@ -486,6 +486,29 @@ impl EcVolume {
*self.shard_locations_refresh_time.lock().unwrap() = Some(std::time::Instant::now());
}
/// Merge a fresh `LookupEcVolume` reply into the shard-locations cache and
/// stamp the refresh time, returning a clone of the resulting map.
///
/// Per-shard upsert (mirrors Go's `cachedLookupEcShardLocations`): each shard
/// id present in the reply overwrites its cached entry, while shard ids absent
/// from the reply keep their previously-cached locations — so a reply that
/// passes the data-shard completeness guard but omits a shard already in cache
/// does not drop that shard's known location (unlike a full replace).
pub fn merge_shard_locations(
&self,
locations: HashMap<ShardId, Vec<String>>,
) -> HashMap<ShardId, Vec<String>> {
let merged = {
let mut guard = self.shard_locations.write().unwrap();
for (shard_id, addrs) in locations {
guard.insert(shard_id, addrs);
}
guard.clone()
};
*self.shard_locations_refresh_time.lock().unwrap() = Some(std::time::Instant::now());
merged
}
/// Get a cloned list of server addresses for a given shard ID.
pub fn get_shard_locations(&self, shard_id: ShardId) -> Vec<String> {
self.shard_locations