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) }