From 8f2a2abae43035150dc99ad63eda1d3c6b4243d1 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Tue, 30 Jun 2026 10:18:31 -0700 Subject: [PATCH] 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 --- seaweed-volume/src/server/grpc_server.rs | 77 ++++++- seaweed-volume/src/server/store_ec.rs | 217 ++++++++++++++---- .../src/storage/erasure_coding/ec_encoder.rs | 13 +- .../src/storage/erasure_coding/ec_volume.rs | 31 ++- 4 files changed, 266 insertions(+), 72 deletions(-) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index f9c76f22d..e1535f7b1 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -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 = + 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); diff --git a/seaweed-volume/src/server/store_ec.rs b/seaweed-volume/src/server/store_ec.rs index a0fca844c..4c6d6438f 100644 --- a/seaweed-volume/src/server/store_ec.rs +++ b/seaweed-volume/src/server/store_ec.rs @@ -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> = { + 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> { +) -> io::Result { 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 { 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, + 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 { 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 { + 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, vid: VolumeId, locations: HashMap>, -) { - 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>> { + 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> { +) -> io::Result<(Vec, 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> { +) -> io::Result<(Vec, bool)> { let mut last_err: Option = 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> { +) -> io::Result<(Vec, 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); } } diff --git a/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs b/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs index aa8108b15..ad01f2012 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_encoder.rs @@ -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, diff --git a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs index 0c2d17c4a..eff6d3047 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_volume.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_volume.rs @@ -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) { 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>, + ) -> HashMap> { + 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 { self.shard_locations