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 <base>.ecNN.v<N> 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 <base>.ec*.v<N> and <base>.vif.v<N>,
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
<base>.*.v<N> 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
<base>.ec*.v<N> and <base>.vif.v<N> 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<N> 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<N> 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>
This commit is contained in:
Chris Lu
2026-09-28 21:55:25 +08:00
committed by GitHub
co-authored by Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
parent 150a69fe11
commit 43fd5b8d82
13 changed files with 582 additions and 29 deletions
+1
View File
@@ -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 <base>.*.v<N> 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
+174 -1
View File
@@ -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 <base>.*.v<N> 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<ShardId> = 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 <base>.*.v<N> 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.
///
@@ -479,6 +479,13 @@ impl DiskLocation {
if self.idx_directory != self.directory {
remove_bitrot_sidecars(&idx_base)?;
}
// Staged 2PC generations (<base>.ecNN.v<N>, 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(())
}
@@ -195,6 +195,104 @@ impl ShardBits {
}
}
/// Parses the generation of a 2PC-staged `<base>.v<N>` file: `None` means the
/// name is not a generation file of `base`.
pub fn ec_file_generation(name: &str, base: &str) -> Option<u32> {
let suffix = name.strip_prefix(&format!("{}.v", base))?;
match suffix.parse::<u32>() {
Ok(g) if g > 0 => Some(g),
_ => None,
}
}
/// Removes 2PC generation files staged under `base`:
/// `<base>.ecNN.v<N>`, `<base>.ecx.v<N>`, `<base>.ecj.v<N>`, `<base>.ecsum.v<N>`
/// and `<base>.vif.v<N>`. `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<io::Error> = 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 `<shard_file>.v<N>` 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<bool> {
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::*;
+66 -7
View File
@@ -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<std::io::Error> = 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 (<name>.v<N>) 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 (<name>.v<N>) 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.
@@ -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!(
+1
View File
@@ -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 <base>.*.v<N> 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
+22 -12
View File
@@ -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 <base>.*.v<N> 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" +
@@ -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
@@ -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 <base>.*.v<N> 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).
+57 -5
View File
@@ -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 <base>.*.v<N> 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 (<name>.v<N>) 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<N> 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 <base>.ecNN.v<N> 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 <base>.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
+8
View File
@@ -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 (<base>.ecNN.v<N>, 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)
}
}
}
@@ -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 <base>.v<N> 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:
// <base>.ecNN.v<N>, <base>.ecx.v<N>, <base>.ecj.v<N>, <base>.ecsum.v<N> and
// <base>.vif.v<N>. 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 <artifact>.v<N>; 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
}