From 8ad2f29e3ef46d9cd47e73a1de6578bc681aee93 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Sat, 26 Sep 2026 15:53:39 +0800 Subject: [PATCH] shell: let volume.deleteEmpty drop volumes with no live needles (#11437) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * shell: let volume.deleteEmpty drop volumes with no live needles The candidate check only accepted a .dat at superblock size, so a volume whose every needle was deleted still had to be vacuumed first — minutes of compaction to rewrite bytes that were all garbage anyway. FileCount counts every indexed entry and DeleteCount every entry made garbage by overwrite or delete, so FileCount <= DeleteCount means nothing live remains and the volume can be unlinked directly. The quietFor guard is unchanged. * volume server: add only_garbage VolumeDelete guard VolumeDelete(only_empty) refuses every volume that ever held data, so a volume whose needles are all deleted could only be removed after a vacuum rewrote it. The new only_garbage flag deletes only when the byte counters show nothing live: DeletedSize covering all of ContentSize, the same all-garbage state vacuum measures. Byte counters are used because the file/delete counts drift on index reload. * rust volume: mirror only_garbage VolumeDelete guard Same check as the Go server: a volume deletes under only_garbage when its deleted bytes cover all content bytes. The grpc handler rejects before the store drops the volume from its map, since destroy errors after removal would still unmount it. * volume delete: let either enabled check pass, keep onlyEmpty on the wire An upgraded shell sending only_garbage to a pre-upgrade server would be read as an unconditional delete (field ignored, only_empty false). The request now keeps only_empty set so old servers check emptiness and refuse, while new servers delete when either check passes. * volume.deleteEmpty: skip remote-backed and protected read-only volumes A remote-tiered replica shares its cloud object with the other replicas, so keepRemoteData=false on one delete removes data they still reference. Protected read-only volumes are quarantined or under maintenance, which is exactly when a replica should not be dropped. * volume delete: validate guarded copies across disks before deleting * volume delete: hold copy locks across guarded validate-and-delete CheckVolumeDeletable released each copy's locks before Destroy ran, so a write landing on a later copy between the two passes refused its destroy after earlier copies were already removed. Pin every copy's dataFileAccessLock (and its location's volumesLock) across validation and removal so a refused delete leaves all copies intact. * volume delete: send deleted-volume notices after releasing locks A blocking send on a full DeletedVolumesChan under volumesLock can stall the heartbeat loop that drains it while it waits on the same locks. Collect the notices under the lock span and send after release. * pb: restore generated-file cosmetics to match the repo's protoc version 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 | 4 + seaweed-volume/src/server/grpc_server.rs | 108 ++++++++++++++++-- seaweed-volume/src/server/heartbeat.rs | 2 +- seaweed-volume/src/storage/disk_location.rs | 7 +- seaweed-volume/src/storage/store.rs | 3 +- seaweed-volume/src/storage/volume.rs | 25 ++-- weed/command/backup.go | 2 +- weed/ec/ec_encode.go | 2 +- weed/operation/volume_move/volume_move.go | 7 +- weed/pb/volume_server.proto | 4 + weed/pb/volume_server_pb/volume_server.pb.go | 19 ++- weed/server/volume_grpc_admin.go | 2 +- weed/server/volume_grpc_copy.go | 2 +- weed/shell/command_volume_delete.go | 2 +- weed/shell/command_volume_delete_empty.go | 14 ++- .../shell/command_volume_delete_empty_test.go | 57 ++++----- weed/shell/command_volume_fix_replication.go | 2 +- weed/shell/command_volume_merge.go | 4 +- weed/shell/command_volume_move.go | 4 +- weed/shell/command_volume_tier_move.go | 2 +- weed/storage/disk_location.go | 10 +- weed/storage/remote_tier_integration_test.go | 4 +- weed/storage/store.go | 84 +++++++++++++- weed/storage/store_delete_volume_test.go | 36 +++++- weed/storage/store_duplicate_vid_test.go | 81 ++++++++++++- weed/storage/volume.go | 11 ++ weed/storage/volume_async_worker_test.go | 2 +- weed/storage/volume_destroy_ec_vif_test.go | 4 +- weed/storage/volume_vacuum_crash_safe_test.go | 2 +- weed/storage/volume_write.go | 41 +++++-- weed/storage/volume_write_test.go | 8 +- 31 files changed, 451 insertions(+), 104 deletions(-) diff --git a/seaweed-volume/proto/volume_server.proto b/seaweed-volume/proto/volume_server.proto index 39b8f1af9..1b983b652 100644 --- a/seaweed-volume/proto/volume_server.proto +++ b/seaweed-volume/proto/volume_server.proto @@ -259,6 +259,10 @@ message VolumeDeleteRequest { // when true, do not remove the cloud-tier object backing the volume. // used for moves where another server is taking over the same .vif. bool keep_remote_data = 3; + // when true, delete only if every needle is deleted: the volume held + // data once but nothing is live anymore. Passing either check, + // only_empty or this one, is enough to delete. + bool only_garbage = 4; } message VolumeDeleteResponse { } diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 8c57fbf7c..a984d5036 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -1597,16 +1597,20 @@ impl VolumeServer for VolumeGrpcService { let req = request.into_inner(); let vid = VolumeId(req.volume_id); let mut store = self.state.store.write().unwrap(); - if req.only_empty { + if req.only_empty || req.only_garbage { let (_, vol) = store.find_volume(vid).ok_or_else(|| { Status::from(crate::storage::volume::VolumeError::VolumeNotFound(vid)) })?; - if vol.file_count() > 0 { + let empty_ok = req.only_empty && vol.file_count() == 0; + let garbage_ok = req.only_garbage + && vol.content_size() > 0 + && vol.deleted_size() >= vol.content_size(); + if !empty_ok && !garbage_ok { return Err(Status::from(crate::storage::volume::VolumeError::NotEmpty)); } } store - .delete_volume(vid, req.only_empty, req.keep_remote_data) + .delete_volume(vid, req.only_empty, req.only_garbage, req.keep_remote_data) .map_err(|e| crate::server::status_with_context(&format!("delete volume {vid}"), e))?; self.state.volume_state_notify.notify_one(); Ok(Response::new(volume_server_pb::VolumeDeleteResponse {})) @@ -1903,7 +1907,7 @@ impl VolumeServer for VolumeGrpcService { let mut store = self.state.store.write().unwrap(); // keep remote data: the inbound copy carries a .vif that may point // at the same cloud-tier object the existing volume references. - store.delete_volume(vid, false, true).map_err(|e| { + store.delete_volume(vid, false, false, true).map_err(|e| { Status::internal(format!("failed to delete existing volume {}: {}", vid, e)) })?; drop(store); @@ -2224,7 +2228,7 @@ impl VolumeServer for VolumeGrpcService { // false would delete the source's remote data. if mounted { let mut store = state.store.write().unwrap(); - let _ = store.delete_volume(vid, false, true); + let _ = store.delete_volume(vid, false, false, true); state.volume_state_notify.notify_one(); } let _ = std::fs::remove_file(format!("{}.dat", data_base_name)); @@ -7115,7 +7119,7 @@ mod tests { let (dest_service, dest_tmp) = make_local_service_with_volume("", None); { let mut store = dest_service.state.store.write().unwrap(); - store.delete_volume(VolumeId(1), false, false).unwrap(); + store.delete_volume(VolumeId(1), false, false, false).unwrap(); // available_space is filled in by the periodic disk check, which // does not run in a unit test; without it VolumeCopy finds no // location with room and never gets as far as copying. @@ -7262,7 +7266,7 @@ mod tests { let (dest_service, dest_tmp) = make_local_service_with_volume("", None); { let mut store = dest_service.state.store.write().unwrap(); - store.delete_volume(VolumeId(1), false, false).unwrap(); + store.delete_volume(VolumeId(1), false, false, false).unwrap(); for loc in &store.locations { loc.check_disk_space(); } @@ -9252,6 +9256,7 @@ mod tests { .volume_delete(Request::new(volume_server_pb::VolumeDeleteRequest { volume_id: 4242, only_empty: false, + only_garbage: false, keep_remote_data: false, })) .await @@ -9267,6 +9272,7 @@ mod tests { .volume_delete(Request::new(volume_server_pb::VolumeDeleteRequest { volume_id: 1, only_empty: true, + only_garbage: false, keep_remote_data: false, })) .await @@ -9285,6 +9291,94 @@ mod tests { ); } + #[tokio::test] + async fn volume_delete_only_garbage_deletes_a_fully_deleted_volume() { + let (service, _tmp) = make_local_service_with_volume("", None); + { + let mut store = service.state.store.write().unwrap(); + let (_, vol) = store.find_volume_mut(VolumeId(1)).unwrap(); + vol.delete_needle(&mut Needle { + id: NeedleId(11), + cookie: Cookie(0x3344), + ..Needle::default() + }) + .unwrap(); + vol.sync_to_disk().unwrap(); + } + + service + .volume_delete(Request::new(volume_server_pb::VolumeDeleteRequest { + volume_id: 1, + only_empty: false, + only_garbage: true, + keep_remote_data: false, + })) + .await + .expect("a fully deleted volume is garbage and must delete"); + assert!( + service + .state + .store + .read() + .unwrap() + .find_volume(VolumeId(1)) + .is_none(), + "the garbage volume must be gone" + ); + } + + #[tokio::test] + async fn volume_delete_only_garbage_refuses_a_volume_with_live_needles() { + let (service, _tmp) = make_local_service_with_volume("", None); + let err = service + .volume_delete(Request::new(volume_server_pb::VolumeDeleteRequest { + volume_id: 1, + only_empty: false, + only_garbage: true, + keep_remote_data: false, + })) + .await + .expect_err("only_garbage must refuse a volume holding a live needle"); + assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}"); + assert!( + service + .state + .store + .read() + .unwrap() + .find_volume(VolumeId(1)) + .is_some(), + "refused delete must leave the volume mounted" + ); + } + + #[tokio::test] + async fn volume_delete_empty_or_garbage_uses_either_check() { + // Both flags set, volume fully deleted: the garbage check passes even + // though the empty check would not. + let (service, _tmp) = make_local_service_with_volume("", None); + { + let mut store = service.state.store.write().unwrap(); + let (_, vol) = store.find_volume_mut(VolumeId(1)).unwrap(); + vol.delete_needle(&mut Needle { + id: NeedleId(11), + cookie: Cookie(0x3344), + ..Needle::default() + }) + .unwrap(); + vol.sync_to_disk().unwrap(); + } + service + .volume_delete(Request::new(volume_server_pb::VolumeDeleteRequest { + volume_id: 1, + only_empty: true, + only_garbage: true, + keep_remote_data: false, + })) + .await + .expect("a fully deleted volume deletes under either check"); + } + // A collection or ext carrying a separator or ".." is folded into a path // CopyFile then opens; it must be rejected rather than climbed out of the // volume directory. Mirrors Go's checkVolumeFileExtension / diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index 082475791..64e9dc5d9 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -1048,7 +1048,7 @@ fn build_heartbeat_with_ec_status( } for vid in delete_vids { - let _ = loc.delete_volume(vid, false, false); + let _ = loc.delete_volume(vid, false, false, false); } for vid in quarantine_vids { diff --git a/seaweed-volume/src/storage/disk_location.rs b/seaweed-volume/src/storage/disk_location.rs index 787beb84c..286d414ba 100644 --- a/seaweed-volume/src/storage/disk_location.rs +++ b/seaweed-volume/src/storage/disk_location.rs @@ -598,13 +598,14 @@ impl DiskLocation { &mut self, vid: VolumeId, only_empty: bool, + only_garbage: bool, keep_remote_data: bool, ) -> Result<(), VolumeError> { if let Some(mut v) = self.volumes.remove(&vid) { crate::metrics::VOLUME_GAUGE .with_label_values(&[&v.collection, "volume"]) .dec(); - v.destroy(only_empty, keep_remote_data)?; + v.destroy(only_empty, only_garbage, keep_remote_data)?; Ok(()) } else { Err(VolumeError::NotFound) @@ -625,7 +626,7 @@ impl DiskLocation { crate::metrics::VOLUME_GAUGE .with_label_values(&[&v.collection, "volume"]) .dec(); - if let Err(e) = v.destroy(false, false) { + if let Err(e) = v.destroy(false, false, false) { warn!(volume_id = vid.0, error = %e, "delete collection: failed to destroy volume"); } } @@ -1747,7 +1748,7 @@ mod tests { .unwrap(); assert_eq!(loc.volumes_len(), 2); - loc.delete_volume(VolumeId(1), false, false).unwrap(); + loc.delete_volume(VolumeId(1), false, false, false).unwrap(); assert_eq!(loc.volumes_len(), 1); assert!(loc.find_volume(VolumeId(1)).is_none()); } diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index e0a610953..5a7553048 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -395,11 +395,12 @@ impl Store { &mut self, vid: VolumeId, only_empty: bool, + only_garbage: bool, keep_remote_data: bool, ) -> Result<(), VolumeError> { for loc in &mut self.locations { if loc.find_volume(vid).is_some() { - return loc.delete_volume(vid, only_empty, keep_remote_data); + return loc.delete_volume(vid, only_empty, only_garbage, keep_remote_data); } } Err(VolumeError::NotFound) diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index df9941e1e..288d57611 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -4424,8 +4424,19 @@ impl Volume { /// Destroy removes everything related to this volume. When keep_remote_data /// is true the cloud-tier object backing the volume is left intact — used /// by moves where another server is taking over the same .vif. - pub fn destroy(&mut self, only_empty: bool, keep_remote_data: bool) -> Result<(), VolumeError> { - if only_empty && self.file_count() > 0 { + pub fn destroy( + &mut self, + only_empty: bool, + only_garbage: bool, + keep_remote_data: bool, + ) -> Result<(), VolumeError> { + // Either enabled check may pass: a volume with no live data qualifies + // whether it reads empty or as all garbage. Byte counters, not counts: + // index rows and live tallies drift apart on reload. + let empty_ok = only_empty && self.file_count() == 0; + let garbage_ok = + only_garbage && self.content_size() > 0 && self.deleted_size() >= self.content_size(); + if (only_empty || only_garbage) && !empty_ok && !garbage_ok { return Err(VolumeError::NotEmpty); } if self.is_compacting { @@ -5238,7 +5249,7 @@ mod tests { let plan = v.dat_scan_plan(sb_size).unwrap(); let dat_path = v.file_name(".dat"); - v.destroy(false, false).unwrap(); + v.destroy(false, false, false).unwrap(); assert!( !Path::new(&dat_path).exists(), "precondition: destroy removed .dat" @@ -6995,7 +7006,7 @@ mod tests { dat_path = v.file_name(".dat"); idx_path = v.file_name(".idx"); assert!(Path::new(&dat_path).exists()); - v.destroy(false, false).unwrap(); + v.destroy(false, false, false).unwrap(); } assert!(!Path::new(&dat_path).exists()); @@ -8789,7 +8800,7 @@ mod tests { assert!(std::path::Path::new(&idx_path).exists()); // Destroy the volume - v.destroy(false, false).unwrap(); + v.destroy(false, false, false).unwrap(); // .dat and .idx should be gone assert!( @@ -8842,7 +8853,7 @@ mod tests { assert!(std::path::Path::new(&dat_path).exists()); assert!(std::path::Path::new(&idx_path).exists()); - v.destroy(false, false).unwrap(); + v.destroy(false, false, false).unwrap(); assert!( !std::path::Path::new(&dat_path).exists(), @@ -8881,7 +8892,7 @@ mod tests { let ecx_path = format!("{}/1.ecx", dir); std::fs::write(&ecx_path, b"ec-index").unwrap(); - v.destroy(false, false).unwrap(); + v.destroy(false, false, false).unwrap(); let dat_path = format!("{}/1.dat", dir); let idx_path = format!("{}/1.idx", dir); diff --git a/weed/command/backup.go b/weed/command/backup.go index 4778ae961..cfeabcde2 100644 --- a/weed/command/backup.go +++ b/weed/command/backup.go @@ -156,7 +156,7 @@ func backupFromLocation(volumeServer pb.ServerAddress, grpcDialOption grpc.DialO // If local volume is larger than remote, recreate it if datSize > stats.TailOffset { - if err := v.Destroy(false, false); err != nil { + if err := v.Destroy(false, false, false); err != nil { v.Close() return fmt.Errorf("destroying volume: %w", err), false } diff --git a/weed/ec/ec_encode.go b/weed/ec/ec_encode.go index c6f77fdaa..75f4a73ad 100644 --- a/weed/ec/ec_encode.go +++ b/weed/ec/ec_encode.go @@ -43,7 +43,7 @@ func markVolumeReplicaWritable(ctx context.Context, grpcDialOption grpc.DialOpti // deleteVolume removes the volume from sourceVolumeServer via the canonical // volume_move helper. func deleteVolume(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer pb.ServerAddress, onlyEmpty bool, keepRemoteData bool) (err error) { - return volume_move.NewMover(grpcDialOption).DeleteVolume(ctx, volumeId, sourceVolumeServer, onlyEmpty, keepRemoteData) + return volume_move.NewMover(grpcDialOption).DeleteVolume(ctx, volumeId, sourceVolumeServer, onlyEmpty, false, keepRemoteData) } func ChunkVolumeIds(volumeIds []needle.VolumeId, batchSize int) [][]needle.VolumeId { diff --git a/weed/operation/volume_move/volume_move.go b/weed/operation/volume_move/volume_move.go index c7c4754f0..e195a0f84 100644 --- a/weed/operation/volume_move/volume_move.go +++ b/weed/operation/volume_move/volume_move.go @@ -116,7 +116,7 @@ func (m *Mover) LiveMoveVolume(ctx context.Context, volumeId needle.VolumeId, so if cleanupTarget { // The target copy may be missing tailed entries; remove it so // the source stays the only replica. - if dErr := m.DeleteVolume(cleanupCtx, volumeId, target, false, true); dErr != nil { + if dErr := m.DeleteVolume(cleanupCtx, volumeId, target, false, false, true); dErr != nil { // Restoring the source while the stale target stays mounted // risks divergent replicas; keep the source readonly. A // re-run refuses while the copy exists, so name the fix. @@ -232,7 +232,7 @@ func (m *Mover) LiveMoveVolume(ctx context.Context, volumeId needle.VolumeId, so opts.Progress(90, fmt.Sprintf("deleting volume %d from %s", volumeId, source)) fmt.Fprintf(opts.Writer, "deleting volume %d from %s\n", volumeId, source) sourceDeleteStarted = true - if err = m.DeleteVolume(ctx, volumeId, source, false, true); err != nil { + if err = m.DeleteVolume(ctx, volumeId, source, false, false, true); err != nil { return fmt.Errorf("delete volume %d from %s: %v", volumeId, source, err) } @@ -415,11 +415,12 @@ func (m *Mover) ReadVolumeFileStatus(ctx context.Context, volumeId needle.Volume // DeleteVolume removes the volume from server. When keepRemoteData is true, the // cloud-tier object backing the volume is left intact — used on the source side // of a move where another server is taking over the same .vif. -func (m *Mover) DeleteVolume(ctx context.Context, volumeId needle.VolumeId, server pb.ServerAddress, onlyEmpty bool, keepRemoteData bool) error { +func (m *Mover) DeleteVolume(ctx context.Context, volumeId needle.VolumeId, server pb.ServerAddress, onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) error { return m.withClient(false, server, func(client volume_server_pb.VolumeServerClient) error { _, deleteErr := client.VolumeDelete(ctx, &volume_server_pb.VolumeDeleteRequest{ VolumeId: uint32(volumeId), OnlyEmpty: onlyEmpty, + OnlyGarbage: onlyGarbage, KeepRemoteData: keepRemoteData, }) return deleteErr diff --git a/weed/pb/volume_server.proto b/weed/pb/volume_server.proto index 4fddb2f4e..2a66eb0a2 100644 --- a/weed/pb/volume_server.proto +++ b/weed/pb/volume_server.proto @@ -259,6 +259,10 @@ message VolumeDeleteRequest { // when true, do not remove the cloud-tier object backing the volume. // used for moves where another server is taking over the same .vif. bool keep_remote_data = 3; + // when true, delete only if every needle is deleted: the volume held + // data once but nothing is live anymore. Passing either check, + // only_empty or this one, is enough to delete. + bool only_garbage = 4; } message VolumeDeleteResponse { } diff --git a/weed/pb/volume_server_pb/volume_server.pb.go b/weed/pb/volume_server_pb/volume_server.pb.go index 6e7c220eb..3ea68edfc 100644 --- a/weed/pb/volume_server_pb/volume_server.pb.go +++ b/weed/pb/volume_server_pb/volume_server.pb.go @@ -1468,8 +1468,11 @@ type VolumeDeleteRequest struct { // when true, do not remove the cloud-tier object backing the volume. // 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"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // when true, delete only if every needle is deleted: the volume held + // data once but nothing is live anymore. + OnlyGarbage bool `protobuf:"varint,4,opt,name=only_garbage,json=onlyGarbage,proto3" json:"only_garbage,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *VolumeDeleteRequest) Reset() { @@ -1523,6 +1526,13 @@ func (x *VolumeDeleteRequest) GetKeepRemoteData() bool { return false } +func (x *VolumeDeleteRequest) GetOnlyGarbage() bool { + if x != nil { + return x.OnlyGarbage + } + return false +} + type VolumeDeleteResponse struct { state protoimpl.MessageState `protogen:"open.v1"` unknownFields protoimpl.UnknownFields @@ -7324,12 +7334,13 @@ const file_volume_server_proto_rawDesc = "" + "\x15VolumeUnmountResponse\"<\n" + "\x1dVolumeConsolidateIndexRequest\x12\x1b\n" + "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\" \n" + - "\x1eVolumeConsolidateIndexResponse\"{\n" + + "\x1eVolumeConsolidateIndexResponse\"\x9e\x01\n" + "\x13VolumeDeleteRequest\x12\x1b\n" + "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12\x1d\n" + "\n" + "only_empty\x18\x02 \x01(\bR\tonlyEmpty\x12(\n" + - "\x10keep_remote_data\x18\x03 \x01(\bR\x0ekeepRemoteData\"\x16\n" + + "\x10keep_remote_data\x18\x03 \x01(\bR\x0ekeepRemoteData\x12!\n" + + "\fonly_garbage\x18\x04 \x01(\bR\vonlyGarbage\"\x16\n" + "\x14VolumeDeleteResponse\"q\n" + "\x19VolumeMarkReadonlyRequest\x12\x1b\n" + "\tvolume_id\x18\x01 \x01(\rR\bvolumeId\x12\x18\n" + diff --git a/weed/server/volume_grpc_admin.go b/weed/server/volume_grpc_admin.go index f2b1eb286..d4d96c412 100644 --- a/weed/server/volume_grpc_admin.go +++ b/weed/server/volume_grpc_admin.go @@ -198,7 +198,7 @@ func (vs *VolumeServer) VolumeDelete(ctx context.Context, req *volume_server_pb. return resp, err } - err := vs.store.DeleteVolume(needle.VolumeId(req.VolumeId), req.OnlyEmpty, req.KeepRemoteData) + err := vs.store.DeleteVolume(needle.VolumeId(req.VolumeId), req.OnlyEmpty, req.OnlyGarbage, req.KeepRemoteData) if err != nil { glog.Errorf("volume delete %v: %v", req, err) diff --git a/weed/server/volume_grpc_copy.go b/weed/server/volume_grpc_copy.go index fd8e0deed..fdf24ab95 100644 --- a/weed/server/volume_grpc_copy.go +++ b/weed/server/volume_grpc_copy.go @@ -101,7 +101,7 @@ func (vs *VolumeServer) VolumeCopy(req *volume_server_pb.VolumeCopyRequest, stre glog.V(0).Infof("volume %d already exists. deleting before copying from %s...", req.VolumeId, req.SourceDataNode) // keep remote data: the inbound copy carries a .vif that may point at // the same cloud-tier object the existing volume references. - if delErr := vs.store.DeleteVolume(needle.VolumeId(req.VolumeId), false, true); delErr != nil { + if delErr := vs.store.DeleteVolume(needle.VolumeId(req.VolumeId), false, false, true); delErr != nil { return fmt.Errorf("failed to delete existing volume %d: %v", req.VolumeId, delErr) } glog.V(0).Infof("deleted existing volume %d before copying.", req.VolumeId) diff --git a/weed/shell/command_volume_delete.go b/weed/shell/command_volume_delete.go index 7b45554d0..e6a118073 100644 --- a/weed/shell/command_volume_delete.go +++ b/weed/shell/command_volume_delete.go @@ -61,6 +61,6 @@ func (c *commandVolumeDelete) Do(args []string, commandEnv *CommandEnv, writer i defer cancel() } - return deleteVolume(ctx, commandEnv.option.GrpcDialOption, volumeId, sourceVolumeServer, false, false) + return deleteVolume(ctx, commandEnv.option.GrpcDialOption, volumeId, sourceVolumeServer, false, false, false) } diff --git a/weed/shell/command_volume_delete_empty.go b/weed/shell/command_volume_delete_empty.go index 507b70de4..b0874467e 100644 --- a/weed/shell/command_volume_delete_empty.go +++ b/weed/shell/command_volume_delete_empty.go @@ -32,6 +32,9 @@ func (c *commandVolumeDeleteEmpty) Help() string { volume.deleteEmpty -collectionPattern=important* -quietFor=24h -apply This command deletes all empty volumes from one volume server. + A volume with no live needles left is empty too, even when its + .dat file is still large: compacting it first would only rewrite + bytes that are all deleted already. ` } @@ -79,8 +82,11 @@ func (c *commandVolumeDeleteEmpty) Do(args []string, commandEnv *CommandEnv, wri if isEmptyVolumeDeleteCandidate(v, quietSeconds, nowUnixSeconds, collectionMatcher) { if *applyBalancing { log.Printf("deleting empty volume %d from %s", v.Id, dn.Id) + onlyGarbage := v.FileCount > 0 && v.FileCount <= v.DeleteCount + // onlyEmpty stays set so a pre-upgrade server checks + // emptiness and refuses instead of deleting unseen. if deleteErr := deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, needle.VolumeId(v.Id), - pb.NewServerAddressFromDataNode(dn), true, false); deleteErr != nil { + pb.NewServerAddressFromDataNode(dn), true, onlyGarbage, false); deleteErr != nil { err = deleteErr } continue @@ -97,7 +103,11 @@ func (c *commandVolumeDeleteEmpty) Do(args []string, commandEnv *CommandEnv, wri func isEmptyVolumeDeleteCandidate(v *master_pb.VolumeInformationMessage, quietSeconds, nowUnixSeconds int64, collectionMatcher *wildcard.CollectionMatcher) bool { return collectionMatcher.Matches(v.Collection) && - v.Size <= super_block.SuperBlockSize && + // A remote-backed replica shares its cloud object with the others, so + // deleting one cannot drop it without hurting the survivors. + v.RemoteStorageName == "" && + (!v.ReadOnly || v.ReadOnlyCanDelete) && + (v.Size <= super_block.SuperBlockSize || v.FileCount > 0 && v.FileCount <= v.DeleteCount) && v.ModifiedAtSecond > 0 && v.ModifiedAtSecond+quietSeconds < nowUnixSeconds } diff --git a/weed/shell/command_volume_delete_empty_test.go b/weed/shell/command_volume_delete_empty_test.go index 344724271..c831858d0 100644 --- a/weed/shell/command_volume_delete_empty_test.go +++ b/weed/shell/command_volume_delete_empty_test.go @@ -4,45 +4,38 @@ import ( "testing" "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" "github.com/seaweedfs/seaweedfs/weed/util/wildcard" + "github.com/stretchr/testify/assert" ) -func TestIsEmptyVolumeDeleteCandidateCollectionPattern(t *testing.T) { - now := int64(1000) - quietSeconds := int64(100) - quietEmptyVolume := &master_pb.VolumeInformationMessage{ - Size: 0, - ModifiedAtSecond: now - quietSeconds - 1, - Collection: "important-logs", - } +func TestIsEmptyVolumeDeleteCandidate(t *testing.T) { + matcher, err := wildcard.CompileCollectionMatcher("") + assert.NoError(t, err) + + const quietSeconds = int64(60) + const now = int64(1_000_000) + old := now - quietSeconds - 1 tests := []struct { - name string - pattern string - collection string - want bool + name string + v *master_pb.VolumeInformationMessage + want bool }{ - {name: "empty pattern matches named collection", pattern: "", collection: "important-logs", want: true}, - {name: "wildcard matches collection", pattern: "important*", collection: "important-logs", want: true}, - {name: "wildcard rejects collection", pattern: "important*", collection: "other-logs", want: false}, - {name: "default pattern matches empty collection", pattern: CollectionDefault, collection: "", want: true}, - {name: "default pattern rejects named collection", pattern: CollectionDefault, collection: "important-logs", want: false}, - {name: "list matches a listed collection", pattern: "other-logs,important-logs", collection: "important-logs", want: true}, - {name: "list rejects an unlisted collection", pattern: "other-logs,important-logs", collection: "audit-logs", want: false}, + {"small .dat is empty", &master_pb.VolumeInformationMessage{Size: super_block.SuperBlockSize, ModifiedAtSecond: old}, true}, + {"all needles deleted is empty", &master_pb.VolumeInformationMessage{Size: 1 << 30, FileCount: 100, DeleteCount: 100, ModifiedAtSecond: old}, true}, + {"live needles keep the volume", &master_pb.VolumeInformationMessage{Size: 1 << 30, FileCount: 100, DeleteCount: 99, ModifiedAtSecond: old}, false}, + {"overwrites alone are not empty", &master_pb.VolumeInformationMessage{Size: 1 << 30, FileCount: 200, DeleteCount: 100, ModifiedAtSecond: old}, false}, + {"recent all-deleted volume is kept", &master_pb.VolumeInformationMessage{Size: 1 << 30, FileCount: 100, DeleteCount: 100, ModifiedAtSecond: now}, false}, + {"never written volume is kept", &master_pb.VolumeInformationMessage{Size: super_block.SuperBlockSize, ModifiedAtSecond: 0}, false}, + {"remote-backed replica is kept", &master_pb.VolumeInformationMessage{Size: super_block.SuperBlockSize, RemoteStorageName: "s3", RemoteStorageKey: "v", ModifiedAtSecond: old}, false}, + {"remote-backed garbage is kept", &master_pb.VolumeInformationMessage{Size: 1 << 30, FileCount: 100, DeleteCount: 100, RemoteStorageName: "s3", ModifiedAtSecond: old}, false}, + {"protected read-only volume is kept", &master_pb.VolumeInformationMessage{Size: super_block.SuperBlockSize, ReadOnly: true, ModifiedAtSecond: old}, false}, + {"deletable read-only volume can go", &master_pb.VolumeInformationMessage{Size: super_block.SuperBlockSize, ReadOnly: true, ReadOnlyCanDelete: true, ModifiedAtSecond: old}, true}, } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - v := *quietEmptyVolume - v.Collection = tt.collection - matcher, err := wildcard.CompileCollectionMatcher(tt.pattern) - if err != nil { - t.Fatalf("CompileCollectionMatcher(%q): %v", tt.pattern, err) - } - if got := isEmptyVolumeDeleteCandidate(&v, quietSeconds, now, matcher); got != tt.want { - t.Fatalf("isEmptyVolumeDeleteCandidate(collection=%q, pattern=%q) = %v, want %v", - tt.collection, tt.pattern, got, tt.want) - } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + assert.Equal(t, tc.want, isEmptyVolumeDeleteCandidate(tc.v, quietSeconds, now, matcher)) }) } } diff --git a/weed/shell/command_volume_fix_replication.go b/weed/shell/command_volume_fix_replication.go index 8eaebc718..7cb38fd3b 100644 --- a/weed/shell/command_volume_fix_replication.go +++ b/weed/shell/command_volume_fix_replication.go @@ -390,7 +390,7 @@ func (c *commandVolumeFixReplication) deleteOneVolume(commandEnv *CommandEnv, wr // Surplus replica being trimmed; keep the remote object since other // replicas of the same .vif still reference it. if err := deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, needle.VolumeId(replica.info.Id), - pb.NewServerAddressFromDataNode(replica.location.dataNode), false, true); err != nil { + pb.NewServerAddressFromDataNode(replica.location.dataNode), false, false, true); err != nil { fmt.Fprintf(writer, "deleting volume %d from %s : %v", replica.info.Id, replica.location.dataNode.Id, err) } else { deleted++ diff --git a/weed/shell/command_volume_merge.go b/weed/shell/command_volume_merge.go index 296cdbe8c..003af4d04 100644 --- a/weed/shell/command_volume_merge.go +++ b/weed/shell/command_volume_merge.go @@ -120,7 +120,7 @@ func (c *commandVolumeMerge) Do(args []string, commandEnv *CommandEnv, writer io if !cleanupTarget { return } - if delErr := deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, volumeId, targetServer, false, false); delErr != nil { + if delErr := deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, volumeId, targetServer, false, false, false); delErr != nil { glog.Warningf("failed to clean up temporary merge volume %d on %s: %v", volumeId, targetServer, delErr) } }() @@ -192,7 +192,7 @@ func (c *commandVolumeMerge) Do(args []string, commandEnv *CommandEnv, writer io } } - if err = deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, volumeId, targetServer, false, false); err != nil { + if err = deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, volumeId, targetServer, false, false, false); err != nil { return err } diff --git a/weed/shell/command_volume_move.go b/weed/shell/command_volume_move.go index 41d0c5b40..d8c74beb5 100644 --- a/weed/shell/command_volume_move.go +++ b/weed/shell/command_volume_move.go @@ -114,8 +114,8 @@ func tailVolume(ctx context.Context, grpcDialOption grpc.DialOption, volumeId ne // deleteVolume removes the volume from sourceVolumeServer. When keepRemoteData // is true, the cloud-tier object backing the volume is left intact — used on // the source side of a move where another server is taking over the same .vif. -func deleteVolume(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer pb.ServerAddress, onlyEmpty bool, keepRemoteData bool) (err error) { - return volume_move.NewMover(grpcDialOption).DeleteVolume(ctx, volumeId, sourceVolumeServer, onlyEmpty, keepRemoteData) +func deleteVolume(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer pb.ServerAddress, onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) (err error) { + return volume_move.NewMover(grpcDialOption).DeleteVolume(ctx, volumeId, sourceVolumeServer, onlyEmpty, onlyGarbage, keepRemoteData) } func markVolumeWritable(ctx context.Context, grpcDialOption grpc.DialOption, volumeId needle.VolumeId, sourceVolumeServer pb.ServerAddress, writable, persist bool) (err error) { diff --git a/weed/shell/command_volume_tier_move.go b/weed/shell/command_volume_tier_move.go index 7b8135cad..4725c4470 100644 --- a/weed/shell/command_volume_tier_move.go +++ b/weed/shell/command_volume_tier_move.go @@ -389,7 +389,7 @@ func (c *commandVolumeTierMove) doMoveOneVolume(commandEnv *CommandEnv, writer i } // keepRemoteData=true: remote-tiered replicas share one cloud object, so // deleting a replica must not delete the object the survivors still point at. - if err = deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, vid, loc.ServerAddress(), false, true); err != nil { + if err = deleteVolume(context.Background(), commandEnv.option.GrpcDialOption, vid, loc.ServerAddress(), false, false, true); err != nil { fmt.Fprintf(writer, "failed to delete volume %d on %s: %v\n", vid, loc.Url, err) } } diff --git a/weed/storage/disk_location.go b/weed/storage/disk_location.go index 102689500..51afa1f4c 100644 --- a/weed/storage/disk_location.go +++ b/weed/storage/disk_location.go @@ -461,7 +461,7 @@ func (l *DiskLocation) DeleteCollectionFromDiskLocation(collection string) (dele wg.Add(2) go func() { for k, v := range delVolsMap { - if err := v.Destroy(false, false); err != nil { + if err := v.Destroy(false, false, false); err != nil { errChain <- err } else { l.volumesLock.Lock() @@ -497,12 +497,12 @@ func (l *DiskLocation) DeleteCollectionFromDiskLocation(collection string) (dele return } -func (l *DiskLocation) deleteVolumeById(vid needle.VolumeId, onlyEmpty bool, keepRemoteData bool) (found bool, e error) { +func (l *DiskLocation) deleteVolumeById(vid needle.VolumeId, onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) (found bool, e error) { v, ok := l.volumes[vid] if !ok { return } - e = v.Destroy(onlyEmpty, keepRemoteData) + e = v.Destroy(onlyEmpty, onlyGarbage, keepRemoteData) if e != nil { return } @@ -528,7 +528,7 @@ func (l *DiskLocation) LoadVolume(diskId uint32, vid needle.VolumeId, needleMapK var ErrVolumeNotFound = fmt.Errorf("volume not found") -func (l *DiskLocation) DeleteVolume(vid needle.VolumeId, onlyEmpty bool, keepRemoteData bool) error { +func (l *DiskLocation) DeleteVolume(vid needle.VolumeId, onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) error { l.volumesLock.Lock() defer l.volumesLock.Unlock() @@ -536,7 +536,7 @@ func (l *DiskLocation) DeleteVolume(vid needle.VolumeId, onlyEmpty bool, keepRem if !ok { return ErrVolumeNotFound } - _, err := l.deleteVolumeById(vid, onlyEmpty, keepRemoteData) + _, err := l.deleteVolumeById(vid, onlyEmpty, onlyGarbage, keepRemoteData) return err } diff --git a/weed/storage/remote_tier_integration_test.go b/weed/storage/remote_tier_integration_test.go index fc59258e5..d6356efb8 100644 --- a/weed/storage/remote_tier_integration_test.go +++ b/weed/storage/remote_tier_integration_test.go @@ -343,7 +343,7 @@ func TestRemoteTier_Move_KeepsRemoteObject(t *testing.T) { v := reloadVolume(t, dir, vid) require.True(t, v.HasRemoteFile()) - require.NoError(t, v.Destroy(false, true)) + require.NoError(t, v.Destroy(false, false, true)) require.True(t, b.objectExists(key), "Destroy(keepRemoteData=true) must not delete remote object") require.Empty(t, b.deleteHistory(), "no DeleteFile call expected on a move-style destroy") @@ -363,7 +363,7 @@ func TestRemoteTier_RealDelete_RemovesRemoteObject(t *testing.T) { v := reloadVolume(t, dir, vid) require.True(t, v.HasRemoteFile()) - require.NoError(t, v.Destroy(false, false)) + require.NoError(t, v.Destroy(false, false, false)) require.False(t, b.objectExists(key), "Destroy(keepRemoteData=false) must delete remote object") require.Equal(t, []string{key}, b.deleteHistory()) diff --git a/weed/storage/store.go b/weed/storage/store.go index 329208483..d4447cbc3 100644 --- a/weed/storage/store.go +++ b/weed/storage/store.go @@ -590,7 +590,7 @@ func (s *Store) CollectHeartbeat() *master_pb.Heartbeat { // delete expired volumes. location.volumesLock.Lock() for _, vid := range deleteVids { - found, err := location.deleteVolumeById(vid, false, false) + found, err := location.deleteVolumeById(vid, false, false, false) if err == nil { if found { glog.V(0).Infof("volume %d is deleted", vid) @@ -1126,10 +1126,13 @@ func RenameOrCopyFile(src, dst string) error { return nil } -func (s *Store) DeleteVolume(i needle.VolumeId, onlyEmpty bool, keepRemoteData bool) error { +func (s *Store) DeleteVolume(i needle.VolumeId, onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) error { // Delete every copy of the volume id across disks, not just the first match, so // a stale twin (e.g. a re-attached disk; NewStore has no cross-disk duplicate // guard) cannot survive a delete and re-register as the volume's content. + if onlyEmpty || onlyGarbage { + return s.deleteVolumeGuarded(i, onlyEmpty, onlyGarbage, keepRemoteData) + } deletedAny := false var errs []error for _, location := range s.Locations { @@ -1146,7 +1149,7 @@ func (s *Store) DeleteVolume(i needle.VolumeId, onlyEmpty bool, keepRemoteData b DiskType: string(location.DiskType), DiskId: v.diskId, } - err := location.DeleteVolume(i, onlyEmpty, keepRemoteData) + err := location.DeleteVolume(i, onlyEmpty, onlyGarbage, keepRemoteData) if err == nil { glog.V(0).Infof("DeleteVolume %d disk_id:%d", i, v.diskId) s.DeletedVolumesChan <- &message @@ -1173,6 +1176,81 @@ func (s *Store) DeleteVolume(i needle.VolumeId, onlyEmpty bool, keepRemoteData b return nil } +// deleteVolumeGuarded removes a volume only when every duplicate copy passes +// the emptiness guards. Each copy's locks are held across validation and +// removal, so a write or mount cannot slip between the check on one copy and +// the destroy of another and leave a partial delete. +func (s *Store) deleteVolumeGuarded(i needle.VolumeId, onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) error { + var lockedLocations []*DiskLocation + var lockedVolumes []*Volume + unlockAll := func() { + for _, v := range lockedVolumes { + v.dataFileAccessLock.Unlock() + } + for _, location := range lockedLocations { + location.volumesLock.Unlock() + } + lockedVolumes, lockedLocations = nil, nil + } + defer unlockAll() + for _, location := range s.Locations { + location.volumesLock.Lock() + if v, ok := location.volumes[i]; ok { + v.dataFileAccessLock.Lock() + lockedLocations = append(lockedLocations, location) + lockedVolumes = append(lockedVolumes, v) + } else { + location.volumesLock.Unlock() + } + } + + for _, v := range lockedVolumes { + if err := v.checkDeletableLocked(onlyEmpty, onlyGarbage); err != nil { + return fmt.Errorf("DeleteVolume %d: %w", i, err) + } + } + + deletedAny := false + var errs []error + var deletedMessages []*master_pb.VolumeShortInformationMessage + for _, location := range lockedLocations { + v := location.volumes[i] + message := &master_pb.VolumeShortInformationMessage{ + Id: uint32(v.Id), + Collection: v.Collection, + ReplicaPlacement: uint32(v.ReplicaPlacement.Byte()), + Version: uint32(v.Version()), + Ttl: v.Ttl.ToUint32(), + DiskType: string(location.DiskType), + DiskId: v.diskId, + } + if err := v.destroyLocked(onlyEmpty, onlyGarbage, keepRemoteData); err != nil { + // A real failure on one disk must not be masked by another copy's + // success: a stale copy left on the failing disk would re-register. + glog.Errorf("DeleteVolume %d: %v", i, err) + errs = append(errs, err) + continue + } + delete(location.volumes, i) + glog.V(0).Infof("DeleteVolume %d disk_id:%d", i, v.diskId) + deletedMessages = append(deletedMessages, message) + deletedAny = true + } + // Send after the locks are released: a full channel would otherwise block + // here while the draining heartbeat loop waits on these same locks. + unlockAll() + for _, m := range deletedMessages { + s.DeletedVolumesChan <- m + } + if len(errs) > 0 { + return fmt.Errorf("DeleteVolume %d failed on some disks: %w", i, errors.Join(errs...)) + } + if !deletedAny { + return fmt.Errorf("delete volume %d not found on disk: %w", i, ErrVolumeNotFound) + } + return nil +} + func (s *Store) ConfigureVolume(i needle.VolumeId, replication string) error { for _, location := range s.Locations { diff --git a/weed/storage/store_delete_volume_test.go b/weed/storage/store_delete_volume_test.go index 634c4b6e7..52d7aa31d 100644 --- a/weed/storage/store_delete_volume_test.go +++ b/weed/storage/store_delete_volume_test.go @@ -16,7 +16,7 @@ func TestDeleteVolumeErrorsAreInspectable(t *testing.T) { t.Fatal(err) } - err := store.DeleteVolume(5, true, false) + err := store.DeleteVolume(5, true, false, false) if !errors.Is(err, ErrVolumeNotEmpty) { t.Fatalf("only-empty delete of a non-empty volume = %v, want ErrVolumeNotEmpty", err) } @@ -24,12 +24,42 @@ func TestDeleteVolumeErrorsAreInspectable(t *testing.T) { t.Fatal("refused delete removed the volume") } - err = store.DeleteVolume(99, false, false) + err = store.DeleteVolume(99, false, false, false) if !errors.Is(err, ErrVolumeNotFound) { t.Fatalf("delete of an absent volume = %v, want ErrVolumeNotFound", err) } - if err := store.DeleteVolume(5, false, false); err != nil { + if err := store.DeleteVolume(5, false, false, false); err != nil { t.Fatalf("forced delete: %v", err) } } + +func TestDeleteVolumeOnlyGarbage(t *testing.T) { + store := newTestStore(t, 1) + mountCollectionVolume(t, store.Locations[0], 5, "") + n := &needle.Needle{Id: types.Uint64ToNeedleId(1), Data: []byte("keep")} + if _, err := store.WriteVolumeNeedle(5, n, false, false); err != nil { + t.Fatal(err) + } + + // A live needle refuses, and the volume stays mounted. + err := store.DeleteVolume(5, false, true, false) + if !errors.Is(err, ErrVolumeNotEmpty) { + t.Fatalf("only-garbage delete of a live volume = %v, want ErrVolumeNotEmpty", err) + } + if _, found := store.Locations[0].FindVolume(5); !found { + t.Fatal("refused delete removed the volume") + } + + if _, err := store.DeleteVolumeNeedle(5, n); err != nil { + t.Fatal(err) + } + // The shell sends both flags so a pre-upgrade server still refuses; either + // check passing must suffice here. + if err := store.DeleteVolume(5, true, true, false); err != nil { + t.Fatalf("delete of a fully deleted volume: %v", err) + } + if _, found := store.Locations[0].FindVolume(5); found { + t.Fatal("fully deleted volume survived") + } +} diff --git a/weed/storage/store_duplicate_vid_test.go b/weed/storage/store_duplicate_vid_test.go index 160c4a9ba..49169439c 100644 --- a/weed/storage/store_duplicate_vid_test.go +++ b/weed/storage/store_duplicate_vid_test.go @@ -2,9 +2,11 @@ package storage import ( "testing" + "time" "github.com/seaweedfs/seaweedfs/weed/storage/needle" "github.com/seaweedfs/seaweedfs/weed/storage/super_block" + "github.com/seaweedfs/seaweedfs/weed/storage/types" "github.com/stretchr/testify/require" ) @@ -42,10 +44,87 @@ func TestDeleteVolumeRemovesAllDuplicateCopies(t *testing.T) { loc.SetVolume(vid, v) } - require.NoError(t, store.DeleteVolume(vid, false, false)) + require.NoError(t, store.DeleteVolume(vid, false, false, false)) _, found0 := store.Locations[0].FindVolume(vid) _, found1 := store.Locations[1].FindVolume(vid) require.False(t, found0, "copy on disk 0 must be deleted") require.False(t, found1, "the stale twin on disk 1 must also be deleted") } + +// A guarded delete must check every copy before destroying any: an earlier +// garbage copy must survive when a later duplicate still holds live data. +func TestDeleteVolumeGuardChecksAllCopiesBeforeDeleting(t *testing.T) { + store := newTestStore(t, 2) + const vid = needle.VolumeId(4244) + + for _, loc := range store.Locations { + v, err := NewVolume(loc.Directory, loc.IdxDirectory, "", vid, NeedleMapInMemory, + &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0) + require.NoError(t, err) + loc.SetVolume(vid, v) + } + + // Disk 0's copy is fully deleted; disk 1's twin still holds a live needle. + copy0, _ := store.Locations[0].FindVolume(vid) + copy1, _ := store.Locations[1].FindVolume(vid) + n := &needle.Needle{Id: types.Uint64ToNeedleId(1), Data: []byte("x")} + _, _, _, err := copy0.writeNeedle2(n, false, false, false) + require.NoError(t, err) + _, _, _, err = copy1.writeNeedle2(n, false, false, false) + require.NoError(t, err) + _, err = copy0.deleteNeedle2(n) + require.NoError(t, err) + + err = store.DeleteVolume(vid, false, true, false) + require.ErrorIs(t, err, ErrVolumeNotEmpty) + + _, found0 := store.Locations[0].FindVolume(vid) + _, found1 := store.Locations[1].FindVolume(vid) + require.True(t, found0, "the garbage copy must survive a refused delete") + require.True(t, found1, "the live copy must survive a refused delete") +} + +// A guarded delete must serialize with writes in flight on any copy: while a +// copy's data lock is held (a write in progress), the delete cannot start +// destroying other copies, or a write landing between validation and removal +// would refuse the later copy and leave a partial delete. +func TestDeleteVolumeGuardWaitsForInFlightCopyWrite(t *testing.T) { + store := newTestStore(t, 2) + const vid = needle.VolumeId(4245) + + for _, loc := range store.Locations { + v, err := NewVolume(loc.Directory, loc.IdxDirectory, "", vid, NeedleMapInMemory, + &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0) + require.NoError(t, err) + loc.SetVolume(vid, v) + } + + copy1, _ := store.Locations[1].FindVolume(vid) + copy1.dataFileAccessLock.Lock() + + done := make(chan error, 1) + go func() { + done <- store.DeleteVolume(vid, false, true, false) + }() + + select { + case err := <-done: + t.Fatalf("guarded delete proceeded while a copy was locked for write: %v", err) + case <-time.After(200 * time.Millisecond): + } + + // The write wins; the delete must see the new needle and refuse, leaving + // every copy intact. + n := &needle.Needle{Id: types.Uint64ToNeedleId(1), Data: []byte("x")} + _, _, _, err := copy1.doWriteRequest(n, false) + require.NoError(t, err) + copy1.dataFileAccessLock.Unlock() + + require.ErrorIs(t, <-done, ErrVolumeNotEmpty) + + _, found0 := store.Locations[0].FindVolume(vid) + _, found1 := store.Locations[1].FindVolume(vid) + require.True(t, found0, "the earlier copy must survive a refused delete") + require.True(t, found1, "the written copy must survive a refused delete") +} diff --git a/weed/storage/volume.go b/weed/storage/volume.go index 129d69b89..0528b5029 100644 --- a/weed/storage/volume.go +++ b/weed/storage/volume.go @@ -225,6 +225,17 @@ func (v *Volume) doIsEmpty() (bool, error) { return true, nil } +// doIsGarbage reports whether every byte ever written is already deleted — +// the same all-garbage state vacuum measures, checked against the byte +// counters because the count counters drift on reload. +func (v *Volume) doIsGarbage() bool { + if v.nm == nil { + return false + } + contentSize := v.nm.ContentSize() + return contentSize > 0 && v.nm.DeletedSize() >= contentSize +} + func (v *Volume) DeletedSize() uint64 { v.dataFileAccessLock.RLock() defer v.dataFileAccessLock.RUnlock() diff --git a/weed/storage/volume_async_worker_test.go b/weed/storage/volume_async_worker_test.go index 9d8f26d4c..7b29f3e11 100644 --- a/weed/storage/volume_async_worker_test.go +++ b/weed/storage/volume_async_worker_test.go @@ -52,7 +52,7 @@ func TestDurableWriteAfterDestroyWritesInline(t *testing.T) { _, _, _, err = v.writeNeedle2(newRandomNeedle(1), true, true, false) require.NoError(t, err) - require.NoError(t, v.Destroy(false, false)) + require.NoError(t, v.Destroy(false, false, false)) require.Nil(t, v.asyncRequestsChan) require.False(t, v.asyncRequestAppend(needle.NewAsyncRequest(newRandomNeedle(2), true))) diff --git a/weed/storage/volume_destroy_ec_vif_test.go b/weed/storage/volume_destroy_ec_vif_test.go index 8a8fffbe4..c51aeea88 100644 --- a/weed/storage/volume_destroy_ec_vif_test.go +++ b/weed/storage/volume_destroy_ec_vif_test.go @@ -32,7 +32,7 @@ func TestDestroyKeepsVifWhenEcCoexists(t *testing.T) { ecxPath := erasure_coding.EcShardFileName("", dir, 1) + ".ecx" require.NoError(t, os.WriteFile(ecxPath, []byte("ec-index"), 0o644)) - require.NoError(t, v.Destroy(false, false)) + require.NoError(t, v.Destroy(false, false, false)) assertFileExist(t, false, base+".dat") assertFileExist(t, false, base+".idx") @@ -54,7 +54,7 @@ func TestDestroyRemovesVifWhenNoEc(t *testing.T) { vifPath := base + ".vif" require.NoError(t, os.WriteFile(vifPath, []byte("regular-volume-info"), 0o644)) - require.NoError(t, v.Destroy(false, false)) + require.NoError(t, v.Destroy(false, false, false)) assertFileExist(t, false, base+".dat") assertFileExist(t, false, base+".idx") diff --git a/weed/storage/volume_vacuum_crash_safe_test.go b/weed/storage/volume_vacuum_crash_safe_test.go index 223bc6126..74279017a 100644 --- a/weed/storage/volume_vacuum_crash_safe_test.go +++ b/weed/storage/volume_vacuum_crash_safe_test.go @@ -320,7 +320,7 @@ func TestDestroyRemovesCommitMarker(t *testing.T) { } } - if err := v.Destroy(false, false); err != nil { + if err := v.Destroy(false, false, false); err != nil { t.Fatalf("destroy: %v", err) } diff --git a/weed/storage/volume_write.go b/weed/storage/volume_write.go index 13b26a1b4..95f84f6e8 100644 --- a/weed/storage/volume_write.go +++ b/weed/storage/volume_write.go @@ -139,20 +139,18 @@ var ErrVolumeNotEmpty = fmt.Errorf("volume not empty") // Destroy removes everything related to this volume. When keepRemoteData is // true the cloud-tier object backing the volume is left intact — used by // moves where another server is taking over the same .vif. -func (v *Volume) Destroy(onlyEmpty bool, keepRemoteData bool) (err error) { +func (v *Volume) Destroy(onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) (err error) { v.dataFileAccessLock.Lock() defer v.dataFileAccessLock.Unlock() + return v.destroyLocked(onlyEmpty, onlyGarbage, keepRemoteData) +} - if onlyEmpty { - isEmpty, e := v.doIsEmpty() - if e != nil { - err = fmt.Errorf("failed to read isEmpty %v", e) - return - } - if !isEmpty { - err = ErrVolumeNotEmpty - return - } +// destroyLocked is Destroy for callers already holding dataFileAccessLock, +// e.g. a guarded multi-copy delete that pins every copy under one lock span +// so validation and removal cannot be split by a write. +func (v *Volume) destroyLocked(onlyEmpty bool, onlyGarbage bool, keepRemoteData bool) (err error) { + if err = v.checkDeletableLocked(onlyEmpty, onlyGarbage); err != nil { + return } if !v.isCompactionInProgress.CompareAndSwap(false, true) { err = fmt.Errorf("volume %d is compacting", v.Id) @@ -178,6 +176,27 @@ func (v *Volume) Destroy(onlyEmpty bool, keepRemoteData bool) (err error) { return } +// checkDeletableLocked enforces the guards Destroy applies before removing +// any file: either enabled check may pass, since a volume with no live data +// qualifies whether it reads empty or as all garbage. +func (v *Volume) checkDeletableLocked(onlyEmpty bool, onlyGarbage bool) (err error) { + if !onlyEmpty && !onlyGarbage { + return nil + } + emptyOk := false + if onlyEmpty { + isEmpty, e := v.doIsEmpty() + if e != nil { + return fmt.Errorf("failed to read isEmpty %v", e) + } + emptyOk = isEmpty + } + if !emptyOk && !(onlyGarbage && v.doIsGarbage()) { + return ErrVolumeNotEmpty + } + return nil +} + // sharesVifWithEcVolume reports whether an EC volume for this volume id lives // on the same disk, in which case its .vif is the same file as the regular // volume's and must outlive the regular volume's deletion. diff --git a/weed/storage/volume_write_test.go b/weed/storage/volume_write_test.go index 24beac8c2..f6e8b5270 100644 --- a/weed/storage/volume_write_test.go +++ b/weed/storage/volume_write_test.go @@ -92,7 +92,7 @@ func TestDestroyEmptyVolumeWithOnlyEmpty(t *testing.T) { // should can Destroy empty volume with onlyEmpty assertFileExist(t, true, path) - err = v.Destroy(true, false) + err = v.Destroy(true, false, false) if err != nil { t.Fatalf("destroy volume: %v", err) } @@ -110,7 +110,7 @@ func TestDestroyEmptyVolumeWithoutOnlyEmpty(t *testing.T) { // should can Destroy empty volume without onlyEmpty assertFileExist(t, true, path) - err = v.Destroy(false, false) + err = v.Destroy(false, false, false) if err != nil { t.Fatalf("destroy volume: %v", err) } @@ -135,7 +135,7 @@ func TestDestroyNonemptyVolumeWithOnlyEmpty(t *testing.T) { assert.Equal(t, uint64(1), v.FileCount()) assertFileExist(t, true, path) - err = v.Destroy(true, false) + err = v.Destroy(true, false, false) assert.EqualError(t, err, "volume not empty") assertFileExist(t, true, path) @@ -165,7 +165,7 @@ func TestDestroyNonemptyVolumeWithoutOnlyEmpty(t *testing.T) { assert.Equal(t, uint64(1), v.FileCount()) assertFileExist(t, true, path) - err = v.Destroy(false, false) + err = v.Destroy(false, false, false) if err != nil { t.Fatalf("destroy volume: %v", err) }