From 43fd5b8d827d14baebe495e00c7287d8ecca7770 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Mon, 28 Sep 2026 21:55:25 +0800 Subject: [PATCH] volume: reclaim staged EC shard generations left by the 2PC switch (#11501) * volume: remove staged EC generation files on teardown and shard delete The 2PC generation switch stages each run as .ecNN.v plus versioned .ecx/.ecj/.vif files. Nothing on the volume server removes them: isEcDataShardFile only recognises the exact .ecNN name, so the staged files are invisible to every bookkeeping pass, and even full_teardown's wipe-all path left them behind. Each re-encode therefore leaks a full shard set per shard-holding disk. RemoveEcGenerationFiles sweeps .ec*.v and .vif.v, optionally keeping generations at or above a threshold; teardown and the reconcile wipe remove every generation, and a per-shard delete removes that shard's staged generations too. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * volume: delete staged EC generations older than N via VolumeEcShardsDelete After a 2PC generation switch commits, the superseded generation's .*.v files sit on disk with no cleanup path: teardown removes everything, and a per-shard delete only touches the named shards, so the executor had no RPC that reclaims just the staged leftovers. delete_generations_older_than removes staged generation files strictly below the threshold on every disk. Versioned files are never mounted, so nothing is unloaded first; the committed generation and the canonical files are preserved. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * rust volume: mirror staged EC generation cleanup Parity with the Go volume server: remove_ec_generation_files sweeps .ec*.v and .vif.v staged by the 2PC switch, called by remove_ec_volume_files (which covers both teardown paths) and the new delete_generations_older_than request field; delete_ec_shards removes a shard's staged generations along with the canonical file. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * volume: match staged generation filenames literally filepath.Glob interprets metacharacters in the collection part of the base name, so a collection like a[bc] could match another volume's staged files (or miss its own). Scan the directory and compare names literally instead, mirroring the Rust read_dir implementation. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * rust volume: report generation-sweep errors and drop the store lock first - snapshot the location base names under the read lock and run the filesystem sweep after dropping it, so a slow disk cannot stall the store; - record per-entry read_dir errors in remove_ec_generation_files and propagate them from remove_ec_shard_generations instead of flatten() skipping them; - warn when a staged-shard generation fails to delete rather than reporting success with files left behind. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * volume: fail shard delete when the staged-generation listing fails A transient ReadDir failure fell back to removing canonical shard names only: staged .v files survived while the RPC still reported success, leaving the leak invisible to retrying callers. ENOENT still means the disk simply has no such directory; other listing errors now propagate. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * rust volume: propagate staged-generation removal failures delete_ec_shards logged remove_ec_shard_generations errors and the RPC returned success while staged .v files remained, diverging from the Go handler which surfaces the failure. The sweep keeps processing the remaining shards, retains the first error, and volume_ec_shards_delete maps it to Status::internal so callers can retry. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * rust volume: notify state change even when the shard sweep errors delete_ec_shards already deletes and unmounts the shards before returning a staged-generation failure, so returning early skipped volume_state_notify and the master kept routing to them until the next heartbeat. Notify before propagating the error. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- seaweed-volume/proto/volume_server.proto | 1 + seaweed-volume/src/server/grpc_server.rs | 175 +++++++++++++++++- seaweed-volume/src/storage/disk_location.rs | 7 + .../src/storage/erasure_coding/ec_shard.rs | 98 ++++++++++ seaweed-volume/src/storage/store.rs | 73 +++++++- .../src/storage/store_ec_reconcile.rs | 2 +- weed/pb/volume_server.proto | 1 + weed/pb/volume_server_pb/volume_server.pb.go | 34 ++-- .../volume_server_pb/volume_server_grpc.pb.go | 6 +- .../volume_grpc_ec_generation_fence_test.go | 86 +++++++++ weed/server/volume_grpc_erasure_coding.go | 62 ++++++- weed/storage/disk_location_ec.go | 8 + weed/storage/erasure_coding/ec_teardown.go | 58 ++++++ 13 files changed, 582 insertions(+), 29 deletions(-) diff --git a/seaweed-volume/proto/volume_server.proto b/seaweed-volume/proto/volume_server.proto index 1b983b652..f23f05863 100644 --- a/seaweed-volume/proto/volume_server.proto +++ b/seaweed-volume/proto/volume_server.proto @@ -475,6 +475,7 @@ message VolumeEcShardsDeleteRequest { repeated uint32 shard_ids = 3; bool full_teardown = 4; // pre-encode cleanup: wipe every EC artifact + generation for this volume, not just shard_ids int64 encode_ts_ns = 5; // full_teardown generation fence: delete only a disk whose .vif generation is strictly OLDER than this; preserve same-or-newer, generation 0, and an unreadable .vif. 0 => wipe-all (shell pre-encode / pre-upgrade) + uint32 delete_generations_older_than = 6; // post-commit cleanup: delete only staged .*.v artifacts with N strictly below this; 0 disables } message VolumeEcShardsDeleteResponse { bool full_teardown_done = 1; // set by a new server that performed full_teardown; absent from an old server lets the caller detect the silent no-op diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 1b82443bf..3c59d8e38 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -3410,14 +3410,68 @@ impl VolumeServer for VolumeGrpcService { )); } + if req.delete_generations_older_than > 0 { + // Post-commit cleanup of a 2PC generation switch: the committed + // generation has been promoted to the canonical names, so only the + // staged .*.v files strictly older than the threshold are + // superseded and safe to remove. Versioned files are never mounted, + // so nothing needs to be unloaded first. + // Snapshot the base names under the lock; the filesystem sweep runs + // after it is dropped so a slow disk cannot stall the store. + let mut bases = Vec::new(); + { + let store = self.state.store.read().unwrap(); + for loc in &store.locations { + bases.push(crate::storage::volume::volume_file_name( + &loc.directory, + &req.collection, + vid, + )); + if loc.idx_directory != loc.directory { + bases.push(crate::storage::volume::volume_file_name( + &loc.idx_directory, + &req.collection, + vid, + )); + } + } + } + for base in &bases { + crate::storage::erasure_coding::ec_shard::remove_ec_generation_files( + base, + req.delete_generations_older_than, + ) + .map_err(|e| { + Status::internal(format!( + "ec generation cleanup of volume {} on {}: {}", + req.volume_id, base, e + )) + })?; + } + self.state.volume_state_notify.notify_one(); + return Ok(Response::new( + volume_server_pb::VolumeEcShardsDeleteResponse { + full_teardown_done: false, + }, + )); + } + let mut store = self.state.store.write().unwrap(); let mut shard_ids: Vec = Vec::with_capacity(req.shard_ids.len()); for &sid in &req.shard_ids { shard_ids.push(shard_id_try_from(sid).map_err(Status::invalid_argument)?); } - store.delete_ec_shards(vid, &req.collection, &shard_ids); + let delete_result = store.delete_ec_shards(vid, &req.collection, &shard_ids); + // The shards are already deleted and unmounted even when a staged- + // generation sweep failed, so the state notification must still go out. drop(store); self.state.volume_state_notify.notify_one(); + delete_result.map_err(|e| { + Status::internal(format!( + "delete ec shards of volume {} in {}: {}", + req.volume_id, req.collection, e + )) + })?; Ok(Response::new( volume_server_pb::VolumeEcShardsDeleteResponse { full_teardown_done: false, @@ -9202,6 +9256,7 @@ mod tests { shard_ids: Vec::new(), full_teardown: true, encode_ts_ns: 0, + delete_generations_older_than: 0, }, )) .await @@ -9279,6 +9334,7 @@ mod tests { shard_ids: Vec::new(), full_teardown: true, encode_ts_ns: 150, + delete_generations_older_than: 0, }, )) .await @@ -9323,6 +9379,123 @@ mod tests { ); } + /// A full teardown wipes the 2PC-staged .*.v files together with + /// the canonical ones: they are invisible to EC bookkeeping, so anything + /// left behind leaks forever (Go removeStaleEcArtifacts). + #[tokio::test] + async fn test_volume_ec_shards_delete_teardown_removes_staged_generations() { + let collection = "ecdel-gens"; + let vid = VolumeId(7062); + let (service, tmp) = SplitDiskEcFixture { + collection, + ..SplitDiskEcFixture::new(vid.0) + } + .build(); + + for ext in [".ec00.v3", ".ecx.v3", ".vif.v3", ".ecsum.v3"] { + std::fs::write( + tmp.path() + .join("data0") + .join(format!("{}_{}{}", collection, vid.0, ext)), + b"staged", + ) + .unwrap(); + } + + service + .volume_ec_shards_delete(Request::new( + volume_server_pb::VolumeEcShardsDeleteRequest { + volume_id: vid.0, + collection: collection.to_string(), + shard_ids: Vec::new(), + full_teardown: true, + encode_ts_ns: 0, + delete_generations_older_than: 0, + }, + )) + .await + .unwrap(); + + for dir in ["data0", "data1"] { + let leftover = std::fs::read_dir(tmp.path().join(dir)) + .unwrap() + .flatten() + .filter(|e| { + e.file_name() + .to_string_lossy() + .starts_with(&format!("{}_{}.", collection, vid.0)) + }) + .count(); + assert_eq!(leftover, 0, "teardown must leave no EC files in {dir}"); + } + } + + /// Post-commit cleanup removes only staged generations strictly below the + /// threshold; the committed generation and live canonical files stay + /// (Go VolumeEcShardsDelete delete_generations_older_than). + #[tokio::test] + async fn test_volume_ec_shards_delete_generations_older_than() { + let collection = "ecdel-gc"; + let vid = VolumeId(7063); + let (service, tmp) = SplitDiskEcFixture { + collection, + ..SplitDiskEcFixture::new(vid.0) + } + .build(); + + let staged = |ext: &str| { + tmp.path() + .join("data0") + .join(format!("{}_{}{}", collection, vid.0, ext)) + }; + for ext in [".ec00.v3", ".ecx.v3", ".vif.v3", ".ec00.v7"] { + std::fs::write(staged(ext), b"staged").unwrap(); + } + + service + .volume_ec_shards_delete(Request::new( + volume_server_pb::VolumeEcShardsDeleteRequest { + volume_id: vid.0, + collection: collection.to_string(), + shard_ids: Vec::new(), + full_teardown: false, + encode_ts_ns: 0, + delete_generations_older_than: 5, + }, + )) + .await + .unwrap(); + + for ext in [".ec00.v3", ".ecx.v3", ".vif.v3"] { + assert!( + !staged(ext).exists(), + "{} must be removed", + staged(ext).display() + ); + } + assert!( + staged(".ec00.v7").exists(), + "the committed generation must be preserved" + ); + assert!( + tmp.path() + .join("data0") + .join(format!("{}_{}.ec00", collection, vid.0)) + .exists(), + "canonical shards must be preserved" + ); + assert!( + service + .state + .store + .read() + .unwrap() + .locations[0] + .has_ec_volume(vid), + "mounted shards must stay mounted" + ); + } + /// REGRESSION: a node-wide scrub must survive a volume that legitimately /// disappears while it runs. /// diff --git a/seaweed-volume/src/storage/disk_location.rs b/seaweed-volume/src/storage/disk_location.rs index 2a3109551..d588447e0 100644 --- a/seaweed-volume/src/storage/disk_location.rs +++ b/seaweed-volume/src/storage/disk_location.rs @@ -479,6 +479,13 @@ impl DiskLocation { if self.idx_directory != self.directory { remove_bitrot_sidecars(&idx_base)?; } + + // Staged 2PC generations (.ecNN.v, versioned .ecx/.ecj/.vif) + // belong to this volume's EC state too; leaving them orphans the files. + crate::storage::erasure_coding::ec_shard::remove_ec_generation_files(&base, 0)?; + if self.idx_directory != self.directory { + crate::storage::erasure_coding::ec_shard::remove_ec_generation_files(&idx_base, 0)?; + } Ok(()) } diff --git a/seaweed-volume/src/storage/erasure_coding/ec_shard.rs b/seaweed-volume/src/storage/erasure_coding/ec_shard.rs index 14195a7da..36038f4a6 100644 --- a/seaweed-volume/src/storage/erasure_coding/ec_shard.rs +++ b/seaweed-volume/src/storage/erasure_coding/ec_shard.rs @@ -195,6 +195,104 @@ impl ShardBits { } } +/// Parses the generation of a 2PC-staged `.v` file: `None` means the +/// name is not a generation file of `base`. +pub fn ec_file_generation(name: &str, base: &str) -> Option { + let suffix = name.strip_prefix(&format!("{}.v", base))?; + match suffix.parse::() { + Ok(g) if g > 0 => Some(g), + _ => None, + } +} + +/// Removes 2PC generation files staged under `base`: +/// `.ecNN.v`, `.ecx.v`, `.ecj.v`, `.ecsum.v` +/// and `.vif.v`. `generations_older_than == 0` removes every +/// generation; otherwise only generations strictly below it. Returns the +/// first real removal failure. Mirrors Go's `RemoveEcGenerationFiles`. +pub fn remove_ec_generation_files(base: &str, generations_older_than: u32) -> io::Result<()> { + let path = std::path::Path::new(base); + let (Some(parent), Some(fname)) = (path.parent(), path.file_name()) else { + return Ok(()); + }; + let ec_prefix = format!("{}.ec", fname.to_string_lossy()); + let vif_name = format!("{}.vif", fname.to_string_lossy()); + let mut first_err: Option = None; + let mut record = |res: io::Result<()>| { + if let Err(e) = res + && first_err.is_none() + { + first_err = Some(e); + } + }; + match fs::read_dir(parent) { + Ok(entries) => { + for entry in entries { + let entry = match entry { + Ok(entry) => entry, + Err(e) => { + // A skipped entry means an incomplete sweep; report it + // instead of pretending the cleanup finished. + record(Err(e)); + continue; + } + }; + let name = entry.file_name().to_string_lossy().into_owned(); + let Some((artifact, _)) = name.rsplit_once(".v") else { + continue; + }; + if artifact != vif_name && !artifact.starts_with(&ec_prefix) { + continue; + } + let Some(generation) = ec_file_generation(&name, artifact) else { + continue; + }; + if generations_older_than > 0 && generation >= generations_older_than { + continue; + } + record(match fs::remove_file(entry.path()) { + Err(e) if e.kind() != io::ErrorKind::NotFound => Err(e), + _ => Ok(()), + }); + } + } + Err(e) if e.kind() != io::ErrorKind::NotFound => record(Err(e)), + Err(_) => {} + } + match first_err { + Some(e) => Err(e), + None => Ok(()), + } +} + +/// Removes every staged generation `.v` of one shard file. +/// Returns true when at least one generation file was removed. +pub fn remove_ec_shard_generations(shard_file: &str) -> io::Result { + let path = std::path::Path::new(shard_file); + let (Some(parent), Some(fname)) = (path.parent(), path.file_name()) else { + return Ok(false); + }; + let fname = fname.to_string_lossy().into_owned(); + let mut removed = false; + match fs::read_dir(parent) { + Ok(entries) => { + for entry in entries { + let entry = entry?; + let name = entry.file_name().to_string_lossy().into_owned(); + if ec_file_generation(&name, &fname).is_some() { + match fs::remove_file(entry.path()) { + Err(e) if e.kind() != io::ErrorKind::NotFound => return Err(e), + _ => removed = true, + } + } + } + } + Err(e) if e.kind() != io::ErrorKind::NotFound => return Err(e), + Err(_) => {} + } + Ok(removed) +} + #[cfg(test)] mod tests { use super::*; diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index 533dd8ab4..0c656888a 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -1334,16 +1334,42 @@ impl Store { None } - /// Delete EC shard files from disk. - pub fn delete_ec_shards(&mut self, vid: VolumeId, collection: &str, shard_ids: &[ShardId]) { + /// Delete EC shard files from disk. Staged-generation removal failures are + /// retained and returned after every location has been processed, so a + /// failed sweep never masquerades as a successful delete. + pub fn delete_ec_shards( + &mut self, + vid: VolumeId, + collection: &str, + shard_ids: &[ShardId], + ) -> std::io::Result<()> { // Delete shard files from disk, tracking which locations actually held one. let mut deleted_at = vec![false; self.locations.len()]; + let mut first_err: Option = None; for (i, loc) in self.locations.iter().enumerate() { for &shard_id in shard_ids { let shard = EcVolumeShard::new(&loc.directory, collection, vid, shard_id); if std::fs::remove_file(shard.file_name()).is_ok() { deleted_at[i] = true; } + // The shard and every 2PC generation of it (.v) are + // removed: the shard must not live on this disk at all. + match crate::storage::erasure_coding::ec_shard::remove_ec_shard_generations( + &shard.file_name(), + ) { + Ok(true) => deleted_at[i] = true, + Ok(false) => {} + Err(e) => { + tracing::warn!( + "failed to remove staged generations of {}: {}", + shard.file_name(), + e + ); + if first_err.is_none() { + first_err = Some(e); + } + } + } } } @@ -1415,6 +1441,11 @@ impl Store { } } } + + match first_err { + Some(e) => Err(e), + None => Ok(()), + } } /// True if `loc` still has any on-disk EC shard for this volume. An @@ -2763,13 +2794,13 @@ mod tests { std::fs::write(format!("{}.ecsum", base1), b"x").unwrap(); // Disk 1 still has .ec01 afterwards: both sidecars survive. - store.delete_ec_shards(vid, collection, &[0]); + store.delete_ec_shards(vid, collection, &[0]).unwrap(); assert!(std::path::Path::new(&format!("{}.ecsum", base0)).exists()); assert!(std::path::Path::new(&format!("{}.ecsum", base1)).exists()); // Disk 1's last shard goes: its sidecar is orphaned and removed, but // disk 0 was never touched by either delete and keeps its sidecar. - store.delete_ec_shards(vid, collection, &[1]); + store.delete_ec_shards(vid, collection, &[1]).unwrap(); assert!(!std::path::Path::new(&format!("{}.ec01", base1)).exists()); assert!( !std::path::Path::new(&format!("{}.ecsum", base1)).exists(), @@ -2814,13 +2845,13 @@ mod tests { std::fs::write(format!("{}.ec01", base1), b"x").unwrap(); std::fs::write(format!("{}.ecsum", idx_base), b"x").unwrap(); - store.delete_ec_shards(vid, collection, &[0]); + store.delete_ec_shards(vid, collection, &[0]).unwrap(); assert!( std::path::Path::new(&format!("{}.ecsum", idx_base)).exists(), "shared idx sidecar must survive while a sibling disk still has shards" ); - store.delete_ec_shards(vid, collection, &[1]); + store.delete_ec_shards(vid, collection, &[1]).unwrap(); assert!( !std::path::Path::new(&format!("{}.ecsum", idx_base)).exists(), "shared idx sidecar should go with the last location's last shard" @@ -2843,7 +2874,7 @@ mod tests { std::fs::write(format!("{}.vif", base1), b"x").unwrap(); std::fs::write(format!("{}.idx", base1), b"x").unwrap(); - store.delete_ec_shards(vid, collection, &[0, 1]); + store.delete_ec_shards(vid, collection, &[0, 1]).unwrap(); assert!( !std::path::Path::new(&format!("{}.vif", base0)).exists(), @@ -2855,6 +2886,34 @@ mod tests { ); } + /// Deleting a shard removes its staged 2PC generations (.v) too: + /// a shard evicted off a disk leaves nothing (Go + /// deleteEcShardIdsForEachLocation). + #[test] + fn test_delete_ec_shards_removes_staged_generations() { + let (mut store, _tmp) = make_ec_target_test_store(1); + let collection = "c"; + let vid = VolumeId(7); + let base = volume_file_name(&store.locations[0].directory, collection, vid); + std::fs::write(format!("{}.ec00", base), b"x").unwrap(); + std::fs::write(format!("{}.ec05", base), b"x").unwrap(); + std::fs::write(format!("{}.ec05.v2", base), b"x").unwrap(); + std::fs::write(format!("{}.ecx", base), b"x").unwrap(); + + store.delete_ec_shards(vid, collection, &[5]).unwrap(); + + assert!(!std::path::Path::new(&format!("{}.ec05", base)).exists()); + assert!( + !std::path::Path::new(&format!("{}.ec05.v2", base)).exists(), + "staged generations of a deleted shard must go with it" + ); + assert!(std::path::Path::new(&format!("{}.ec00", base)).exists()); + assert!( + std::path::Path::new(&format!("{}.ecx", base)).exists(), + "index must survive while shards remain" + ); + } + /// An already-mounted EC volume on disk 1 must win over a stray /// `.ecx` on disk 2. Protects the post-startup steady state from /// being perturbed by leftover index files from a prior failed move. diff --git a/seaweed-volume/src/storage/store_ec_reconcile.rs b/seaweed-volume/src/storage/store_ec_reconcile.rs index 38ac80a1c..0c81e4244 100644 --- a/seaweed-volume/src/storage/store_ec_reconcile.rs +++ b/seaweed-volume/src/storage/store_ec_reconcile.rs @@ -1476,7 +1476,7 @@ mod tests { let vid = VolumeId(7004); let collection = "grafana-loki"; - store.delete_ec_shards(vid, collection, &[1]); + store.delete_ec_shards(vid, collection, &[1]).unwrap(); // Shard 1 file is gone on disk 1. let p1 = format!( diff --git a/weed/pb/volume_server.proto b/weed/pb/volume_server.proto index 2a66eb0a2..abbf6de00 100644 --- a/weed/pb/volume_server.proto +++ b/weed/pb/volume_server.proto @@ -477,6 +477,7 @@ message VolumeEcShardsDeleteRequest { repeated uint32 shard_ids = 3; bool full_teardown = 4; // pre-encode cleanup: wipe every EC artifact + generation for this volume, not just shard_ids int64 encode_ts_ns = 5; // full_teardown generation fence: delete only a disk whose .vif generation is strictly OLDER than this; preserve same-or-newer, generation 0, and an unreadable .vif. 0 => wipe-all (shell pre-encode / pre-upgrade) + uint32 delete_generations_older_than = 6; // post-commit cleanup: delete only staged .*.v artifacts with N strictly below this; 0 disables } message VolumeEcShardsDeleteResponse { bool full_teardown_done = 1; // set by a new server that performed full_teardown; absent from an old server lets the caller detect the silent no-op diff --git a/weed/pb/volume_server_pb/volume_server.pb.go b/weed/pb/volume_server_pb/volume_server.pb.go index 3ea68edfc..0c568858d 100644 --- a/weed/pb/volume_server_pb/volume_server.pb.go +++ b/weed/pb/volume_server_pb/volume_server.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go. DO NOT EDIT. // versions: // protoc-gen-go v1.36.6 -// protoc v6.33.4 +// protoc v3.21.12 // source: volume_server.proto package volume_server_pb @@ -1469,7 +1469,8 @@ type VolumeDeleteRequest struct { // used for moves where another server is taking over the same .vif. KeepRemoteData bool `protobuf:"varint,3,opt,name=keep_remote_data,json=keepRemoteData,proto3" json:"keep_remote_data,omitempty"` // when true, delete only if every needle is deleted: the volume held - // data once but nothing is live anymore. + // data once but nothing is live anymore. Passing either check, + // only_empty or this one, is enough to delete. OnlyGarbage bool `protobuf:"varint,4,opt,name=only_garbage,json=onlyGarbage,proto3" json:"only_garbage,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache @@ -3732,14 +3733,15 @@ func (*VolumeEcShardsCopyResponse) Descriptor() ([]byte, []int) { } type VolumeEcShardsDeleteRequest struct { - state protoimpl.MessageState `protogen:"open.v1"` - VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"` - Collection string `protobuf:"bytes,2,opt,name=collection,proto3" json:"collection,omitempty"` - ShardIds []uint32 `protobuf:"varint,3,rep,packed,name=shard_ids,json=shardIds,proto3" json:"shard_ids,omitempty"` - FullTeardown bool `protobuf:"varint,4,opt,name=full_teardown,json=fullTeardown,proto3" json:"full_teardown,omitempty"` // pre-encode cleanup: wipe every EC artifact + generation for this volume, not just shard_ids - EncodeTsNs int64 `protobuf:"varint,5,opt,name=encode_ts_ns,json=encodeTsNs,proto3" json:"encode_ts_ns,omitempty"` // full_teardown generation fence: delete only a disk whose .vif generation is strictly OLDER than this; preserve same-or-newer, generation 0, and an unreadable .vif. 0 => wipe-all (shell pre-encode / pre-upgrade) - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + state protoimpl.MessageState `protogen:"open.v1"` + VolumeId uint32 `protobuf:"varint,1,opt,name=volume_id,json=volumeId,proto3" json:"volume_id,omitempty"` + Collection string `protobuf:"bytes,2,opt,name=collection,proto3" json:"collection,omitempty"` + ShardIds []uint32 `protobuf:"varint,3,rep,packed,name=shard_ids,json=shardIds,proto3" json:"shard_ids,omitempty"` + FullTeardown bool `protobuf:"varint,4,opt,name=full_teardown,json=fullTeardown,proto3" json:"full_teardown,omitempty"` // pre-encode cleanup: wipe every EC artifact + generation for this volume, not just shard_ids + EncodeTsNs int64 `protobuf:"varint,5,opt,name=encode_ts_ns,json=encodeTsNs,proto3" json:"encode_ts_ns,omitempty"` // full_teardown generation fence: delete only a disk whose .vif generation is strictly OLDER than this; preserve same-or-newer, generation 0, and an unreadable .vif. 0 => wipe-all (shell pre-encode / pre-upgrade) + DeleteGenerationsOlderThan uint32 `protobuf:"varint,6,opt,name=delete_generations_older_than,json=deleteGenerationsOlderThan,proto3" json:"delete_generations_older_than,omitempty"` // post-commit cleanup: delete only staged .*.v artifacts with N strictly below this; 0 disables + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *VolumeEcShardsDeleteRequest) Reset() { @@ -3807,6 +3809,13 @@ func (x *VolumeEcShardsDeleteRequest) GetEncodeTsNs() int64 { return 0 } +func (x *VolumeEcShardsDeleteRequest) GetDeleteGenerationsOlderThan() uint32 { + if x != nil { + return x.DeleteGenerationsOlderThan + } + return 0 +} + type VolumeEcShardsDeleteResponse struct { state protoimpl.MessageState `protogen:"open.v1"` FullTeardownDone bool `protobuf:"varint,1,opt,name=full_teardown_done,json=fullTeardownDone,proto3" json:"full_teardown_done,omitempty"` // set by a new server that performed full_teardown; absent from an old server lets the caller detect the silent no-op @@ -7505,7 +7514,7 @@ const file_volume_server_proto_rawDesc = "" + "\x0fcopy_ecsum_file\x18\t \x01(\bR\rcopyEcsumFile\x12+\n" + "\x12io_byte_per_second\x18\n" + " \x01(\x03R\x0fioBytePerSecond\"\x1c\n" + - "\x1aVolumeEcShardsCopyResponse\"\xbe\x01\n" + + "\x1aVolumeEcShardsCopyResponse\"\x81\x02\n" + "\x1bVolumeEcShardsDeleteRequest\x12\x1b\n" + "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12\x1e\n" + "\n" + @@ -7514,7 +7523,8 @@ const file_volume_server_proto_rawDesc = "" + "\tshard_ids\x18\x03 \x03(\rR\bshardIds\x12#\n" + "\rfull_teardown\x18\x04 \x01(\bR\ffullTeardown\x12 \n" + "\fencode_ts_ns\x18\x05 \x01(\x03R\n" + - "encodeTsNs\"L\n" + + "encodeTsNs\x12A\n" + + "\x1ddelete_generations_older_than\x18\x06 \x01(\rR\x1adeleteGenerationsOlderThan\"L\n" + "\x1cVolumeEcShardsDeleteResponse\x12,\n" + "\x12full_teardown_done\x18\x01 \x01(\bR\x10fullTeardownDone\"\xd4\x01\n" + "\x1aVolumeEcShardsMountRequest\x12\x1b\n" + diff --git a/weed/pb/volume_server_pb/volume_server_grpc.pb.go b/weed/pb/volume_server_pb/volume_server_grpc.pb.go index 72d98971a..2104ed8ff 100644 --- a/weed/pb/volume_server_pb/volume_server_grpc.pb.go +++ b/weed/pb/volume_server_pb/volume_server_grpc.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go-grpc. DO NOT EDIT. // versions: // - protoc-gen-go-grpc v1.6.2 -// - protoc v6.33.4 +// - protoc v3.21.12 // source: volume_server.proto package volume_server_pb @@ -74,7 +74,7 @@ const ( // // For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. type VolumeServerClient interface { - // Experts only: takes multiple fid parameters. This function does not propagate deletes to replicas. + //Experts only: takes multiple fid parameters. This function does not propagate deletes to replicas. BatchDelete(ctx context.Context, in *BatchDeleteRequest, opts ...grpc.CallOption) (*BatchDeleteResponse, error) VacuumVolumeCheck(ctx context.Context, in *VacuumVolumeCheckRequest, opts ...grpc.CallOption) (*VacuumVolumeCheckResponse, error) VacuumVolumeCompact(ctx context.Context, in *VacuumVolumeCompactRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[VacuumVolumeCompactResponse], error) @@ -727,7 +727,7 @@ func (c *volumeServerClient) Ping(ctx context.Context, in *PingRequest, opts ... // All implementations must embed UnimplementedVolumeServerServer // for forward compatibility. type VolumeServerServer interface { - // Experts only: takes multiple fid parameters. This function does not propagate deletes to replicas. + //Experts only: takes multiple fid parameters. This function does not propagate deletes to replicas. BatchDelete(context.Context, *BatchDeleteRequest) (*BatchDeleteResponse, error) VacuumVolumeCheck(context.Context, *VacuumVolumeCheckRequest) (*VacuumVolumeCheckResponse, error) VacuumVolumeCompact(*VacuumVolumeCompactRequest, grpc.ServerStreamingServer[VacuumVolumeCompactResponse]) error diff --git a/weed/server/volume_grpc_ec_generation_fence_test.go b/weed/server/volume_grpc_ec_generation_fence_test.go index 2b6c8b1e9..6ec9edff7 100644 --- a/weed/server/volume_grpc_ec_generation_fence_test.go +++ b/weed/server/volume_grpc_ec_generation_fence_test.go @@ -149,6 +149,92 @@ func TestUnmountEcShardsFencedByGeneration(t *testing.T) { require.False(t, mountedEcShardIds(t, vs, vid)[1], "a strictly-older generation shard must be unmounted") } +// TestTeardownRemovesStagedGenerations pins that a full teardown wipes the +// 2PC-staged .*.v files together with the canonical ones: they are +// invisible to EC bookkeeping, so anything left behind leaks forever. +func TestTeardownRemovesStagedGenerations(t *testing.T) { + const collection = "ec-gen-leak" + vid := needle.VolumeId(57) + dir := t.TempDir() + vs := &VolumeServer{store: buildEcStoreWithGeneration(t, dir, collection, vid, 100, []erasure_coding.ShardId{0, 1})} + + base := erasure_coding.EcShardFileName(collection, dir, int(vid)) + for _, name := range []string{ + base + ".ec00.v3", base + ".ec01.v3", base + ".ecx.v3", base + ".vif.v3", base + ".ecsum.v3", + } { + require.NoError(t, os.WriteFile(name, []byte("staged"), 0o644)) + } + + _, err := vs.VolumeEcShardsDelete(context.Background(), &volume_server_pb.VolumeEcShardsDeleteRequest{ + VolumeId: uint32(vid), + Collection: collection, + FullTeardown: true, + EncodeTsNs: 200, + }) + require.NoError(t, err) + left, err := filepath.Glob(base + "*") + require.NoError(t, err) + require.Empty(t, left, "teardown must leave no EC files, staged generations included: %v", left) +} + +// TestDeleteGenerationsOlderThan covers the post-commit cleanup: staged +// generations below the threshold are removed while the committed generation +// and the live canonical files stay untouched. +func TestDeleteGenerationsOlderThan(t *testing.T) { + const collection = "ec-gen-gc" + vid := needle.VolumeId(58) + dir := t.TempDir() + vs := &VolumeServer{store: buildEcStoreWithGeneration(t, dir, collection, vid, 100, []erasure_coding.ShardId{0, 1})} + + base := erasure_coding.EcShardFileName(collection, dir, int(vid)) + stale := []string{base + ".ec00.v3", base + ".ecx.v3", base + ".vif.v3"} + fresh := []string{base + ".ec00.v7", base + ".vif.v7"} + for _, name := range append(stale, fresh...) { + require.NoError(t, os.WriteFile(name, []byte("staged"), 0o644)) + } + + _, err := vs.VolumeEcShardsDelete(context.Background(), &volume_server_pb.VolumeEcShardsDeleteRequest{ + VolumeId: uint32(vid), + Collection: collection, + DeleteGenerationsOlderThan: 5, + }) + require.NoError(t, err) + for _, name := range stale { + require.False(t, util.FileExists(name), "%s must be removed", name) + } + for _, name := range fresh { + require.True(t, util.FileExists(name), "%s must be preserved", name) + } + require.True(t, util.FileExists(base+".ec00"), "canonical shards must be preserved") + require.True(t, util.FileExists(base+".ecx")) + require.True(t, util.FileExists(base+".vif")) + require.True(t, mountedEcShardIds(t, vs, vid)[0], "mounted shards must stay mounted") +} + +// TestShardDeleteRemovesStagedGenerations pins that deleting a shard removes +// its staged generations too: a shard evicted off a disk leaves nothing. +func TestShardDeleteRemovesStagedGenerations(t *testing.T) { + const collection = "ec-shard-gen" + vid := needle.VolumeId(59) + dir := t.TempDir() + vs := &VolumeServer{store: buildEcStoreWithGeneration(t, dir, collection, vid, 100, []erasure_coding.ShardId{0, 5})} + + base := erasure_coding.EcShardFileName(collection, dir, int(vid)) + require.NoError(t, os.WriteFile(base+".ec05.v2", []byte("staged"), 0o644)) + require.NoError(t, vs.store.UnmountEcShards(vid, 5, 0)) + + _, err := vs.VolumeEcShardsDelete(context.Background(), &volume_server_pb.VolumeEcShardsDeleteRequest{ + VolumeId: uint32(vid), + Collection: collection, + ShardIds: []uint32{5}, + }) + require.NoError(t, err) + require.False(t, util.FileExists(base+".ec05")) + require.False(t, util.FileExists(base+".ec05.v2")) + require.True(t, util.FileExists(base+".ec00"), "sibling shards must be preserved") + require.True(t, util.FileExists(base+".ecx"), "index must survive while shards remain") +} + // TestReadEcGenerationTsNs covers the per-disk .vif generation read used by the // fenced teardown: a present .vif yields its generation (or 0 when it has no EC // config), and a missing .vif is reported unreadable (preserved, fail-safe). diff --git a/weed/server/volume_grpc_erasure_coding.go b/weed/server/volume_grpc_erasure_coding.go index 98a41e755..011b18ecd 100644 --- a/weed/server/volume_grpc_erasure_coding.go +++ b/weed/server/volume_grpc_erasure_coding.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "math" + "io/fs" "os" "path" "path/filepath" @@ -517,6 +518,27 @@ func (vs *VolumeServer) VolumeEcShardsDelete(ctx context.Context, req *volume_se return &volume_server_pb.VolumeEcShardsDeleteResponse{FullTeardownDone: true}, nil } + if req.DeleteGenerationsOlderThan > 0 { + // Post-commit cleanup of a 2PC generation switch: the committed + // generation has been promoted to the canonical names, so only the + // staged .*.v files strictly older than the threshold are + // superseded and safe to remove. Versioned files are never mounted, + // so nothing needs to be unloaded first. + for _, location := range vs.store.Locations { + dataBase := storage.VolumeFileName(location.Directory, req.Collection, int(req.VolumeId)) + idxBase := storage.VolumeFileName(location.IdxDirectory, req.Collection, int(req.VolumeId)) + if err := erasure_coding.RemoveEcGenerationFiles(dataBase, req.DeleteGenerationsOlderThan); err != nil { + return nil, fmt.Errorf("ec generation cleanup of volume %d on %s: %w", req.VolumeId, location.Directory, err) + } + if dataBase != idxBase { + if err := erasure_coding.RemoveEcGenerationFiles(idxBase, req.DeleteGenerationsOlderThan); err != nil { + return nil, fmt.Errorf("ec generation cleanup of volume %d on %s: %w", req.VolumeId, location.IdxDirectory, err) + } + } + } + return &volume_server_pb.VolumeEcShardsDeleteResponse{}, nil + } + glog.V(0).Infof("ec volume %s shard delete %v", bName, req.ShardIds) // Pass 1: delete the requested shard files (and any now-orphaned per-disk bitrot @@ -569,14 +591,37 @@ func deleteEcShardIdsForEachLocation(bName string, location *storage.DiskLocatio // Delete the requested shard files unconditionally. Gating on a local .ecx // (still used for index-file routing below) would leak an orphan shard left // by a failed copy that reconciliation later mounts under a foreign index. + shardFileNames := make([]string, 0, len(shardIds)) for _, shardId := range shardIds { - shardFileName := dataBaseFilename + erasure_coding.ToExt(int(shardId)) - if util.FileExists(shardFileName) { - found = true - if err := removeFileIfExists(shardFileName); err != nil { - return fmt.Errorf("remove ec shard %s: %w", shardFileName, err) + shardFileNames = append(shardFileNames, dataBaseFilename+erasure_coding.ToExt(int(shardId))) + } + // The shard and every 2PC generation of it (.v) are removed: + // the shard must not live on this disk at all. Names match literally — + // a glob would let glob metacharacters in the collection part of bName + // leak into another volume's files. + entries, readErr := os.ReadDir(location.Directory) + switch { + case readErr == nil: + for _, entry := range entries { + for _, shardFileName := range shardFileNames { + name := filepath.Join(location.Directory, entry.Name()) + if name != shardFileName && erasure_coding.EcFileGeneration(entry.Name(), filepath.Base(shardFileName)) < 0 { + continue + } + if util.FileExists(name) { + found = true + if err := removeFileIfExists(name); err != nil { + return fmt.Errorf("remove ec shard %s: %w", name, err) + } + } } } + case errors.Is(readErr, fs.ErrNotExist): + // No such directory means no shard files on this disk. + default: + // A listing failure must not fall back to canonical names only: + // staged .v files would survive while the RPC reports success. + return fmt.Errorf("list %s for ec shards of %s: %w", location.Directory, bName, readErr) } if !found { @@ -711,6 +756,13 @@ func removeStaleEcArtifacts(dataBaseFileName, indexBaseFileName string, total in record(removeBitrotSidecars(dataBaseFileName)) } + // Generations staged by the 2PC switch are .ecNN.v plus the + // versioned .ecx/.ecj/.vif/.ecsum: a teardown of this disk leaves none. + record(erasure_coding.RemoveEcGenerationFiles(dataBaseFileName, 0)) + if dataBaseFileName != indexBaseFileName { + record(erasure_coding.RemoveEcGenerationFiles(indexBaseFileName, 0)) + } + // Canonical .vif. A shard copy installs shards + .ecx before .vif, so an // interrupted copy can leave a stale .vif whose run identity / shard ratio / // dat_file_size a fresh generation would inherit. Remove it only on a shard-only diff --git a/weed/storage/disk_location_ec.go b/weed/storage/disk_location_ec.go index 4af53e05e..dd3b0f4c7 100644 --- a/weed/storage/disk_location_ec.go +++ b/weed/storage/disk_location_ec.go @@ -650,4 +650,12 @@ func (l *DiskLocation) removeEcVolumeFiles(collection string, vid needle.VolumeI for i := 0; i < erasure_coding.MaxShardCount; i++ { removeFile(baseFileName+erasure_coding.ToExt(i), "EC shard file") } + + // Staged 2PC generations (.ecNN.v, versioned .ecx/.ecj/.vif) + // belong to this volume's EC state too; leaving them orphans the files. + for _, dir := range []string{indexBaseFileName, baseFileName} { + if err := erasure_coding.RemoveEcGenerationFiles(dir, 0); err != nil { + glog.Warningf("Failed to remove EC generation files for %s: %v", dir, err) + } + } } diff --git a/weed/storage/erasure_coding/ec_teardown.go b/weed/storage/erasure_coding/ec_teardown.go index b0e98af26..72c74584c 100644 --- a/weed/storage/erasure_coding/ec_teardown.go +++ b/weed/storage/erasure_coding/ec_teardown.go @@ -4,6 +4,10 @@ import ( "context" "errors" "fmt" + "os" + "path/filepath" + "strconv" + "strings" "github.com/seaweedfs/seaweedfs/weed/operation" "github.com/seaweedfs/seaweedfs/weed/pb" @@ -74,3 +78,57 @@ func UnmountAndDeleteEcShards( return nil }) } + +// EcFileGeneration parses the generation of a 2PC-staged .v file: +// -1 means the name is not a generation file of base. +func EcFileGeneration(name, base string) int64 { + suffix, ok := strings.CutPrefix(name, base+".v") + if !ok { + return -1 + } + generation, err := strconv.ParseInt(suffix, 10, 64) + if err != nil || generation <= 0 { + return -1 + } + return generation +} + +// RemoveEcGenerationFiles removes 2PC generation files staged under base: +// .ecNN.v, .ecx.v, .ecj.v, .ecsum.v and +// .vif.v. generationsOlderThan == 0 removes every generation; +// otherwise only generations strictly below it. Returns the first real +// removal failure. +func RemoveEcGenerationFiles(baseFileName string, generationsOlderThan uint32) error { + var firstErr error + record := func(err error) { + if err != nil && firstErr == nil { + firstErr = err + } + } + dir, fileName := filepath.Dir(baseFileName), filepath.Base(baseFileName) + ecPrefix, vifName := fileName+".ec", fileName+".vif" + entries, err := os.ReadDir(dir) + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + for _, entry := range entries { + name := entry.Name() + // A generation file is .v; the last dot separates the + // staged-generation suffix from the artifact name. + artifact := name[:max(strings.LastIndexByte(name, '.'), 0)] + if artifact != vifName && !strings.HasPrefix(artifact, ecPrefix) { + continue + } + generation := EcFileGeneration(name, artifact) + if generation < 0 || (generationsOlderThan > 0 && generation >= int64(generationsOlderThan)) { + continue + } + if err := os.Remove(filepath.Join(dir, name)); err != nil && !os.IsNotExist(err) { + record(err) + } + } + return firstErr +}