From a261f90e188df5851fb659f08eac530627d6f62a Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Sun, 27 Sep 2026 07:05:22 +0800 Subject: [PATCH] vacuum: let the sweep release volumes that stay empty and quiet (#11477) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * vacuum: let the sweep release volumes that stay empty and quiet Vacuuming reclaims bytes but not slots: a fully emptied volume stays registered to its collection forever, and since growth is gated only on slot count a store at 99% free disk can still refuse writes to other collections (#11429). volume.deleteEmpty exists but is manual-only. With -vacuumDeleteEmptyAfterSeconds (or master.vacuumDeleteEmptyAfterSeconds under weed server/mini; default 0, off) the automatic sweep now deletes replica copies that have stayed empty and quiet for that long, the same rule volume.deleteEmpty applies on demand: remote-backed copies are skipped, and every delete carries the volume server's onlyEmpty / onlyGarbage guards so a copy written since the last report is refused rather than removed. Copies that still hold data or were written recently stay; only a volume whose every copy is deleted leaves the sweep's work map, sparing a compaction of bytes that are all deleted. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * vacuum: harden empty-volume sweep against partial and racing deletes Review follow-up on #11477: - delete a volume only when every replica copy is a verifiable empty-and-quiet candidate; deleting the empty copy of a volume whose sibling holds live files would silently cut its replica count (greptile P1). - drain the volume out of the writable list before deleting, the same drain the compact pass uses, so PickForWrite stops assigning it and pending writes settle (devin). - bound the VolumeDelete RPC so one stalled server cannot hold the vacuum lock indefinitely (greptile P1, reusing allocateVolumeTimeout). The vid2location panic scenario raised in review does not exist: VolumeLocationList methods are nil-receiver safe and a missing vid just fails enoughCopies, so a partially deleted volume skips compaction instead of crashing the sweep. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * vacuum: unregister deleted empty replicas and prune the sweep list A successful VolumeDelete only updates the volume server; the master still tracked the replica and kept it in the sweep's location list for the compaction pass (coderabbit on #11477). Unregister the replica right after its delete succeeds and drop it from the sweep copy, so a partially deleted volume only compacts copies that still exist. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * vacuum: pin deleting volumes out of the writable list across heartbeats Review follow-up on #11477 (greptile): DrainAndRemoveFromWritable only removed the volume once; a heartbeat landing between the drain and the replica deletes re-evaluated writability and re-added it, so a client write could reach a replica whose siblings were already gone and leave the volume under-replicated when the last copy refused its onlyEmpty delete. MarkDeleting records the vid in deletingVolumes — checked inside setVolumeWritable so heartbeat, capacity-recovery, and admin re-add paths all hold it out — and UnmarkDeleting releases it once the sweep finishes the copy pass. A partially deleted volume's surviving replicas then return to writable through the normal heartbeat path. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * vacuum: restore writability when a sweep delete survives 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> --- weed/command/master.go | 41 ++-- weed/command/master_follower.go | 1 + weed/command/mini.go | 1 + weed/command/server.go | 1 + weed/server/master_grpc_server_volume.go | 2 +- weed/server/master_server.go | 16 +- weed/server/master_server_handlers_admin.go | 2 +- weed/topology/topology_event_handling.go | 4 +- weed/topology/topology_vacuum.go | 85 ++++++- weed/topology/topology_vacuum_test.go | 255 ++++++++++++++++++++ weed/topology/volume_layout.go | 29 +++ weed/topology/volume_layout_drain_test.go | 37 +++ 12 files changed, 441 insertions(+), 33 deletions(-) create mode 100644 weed/topology/topology_vacuum_test.go diff --git a/weed/command/master.go b/weed/command/master.go index 5b29b914f..81ffa6db9 100644 --- a/weed/command/master.go +++ b/weed/command/master.go @@ -46,18 +46,19 @@ const ( ) type MasterOptions struct { - port *int - portGrpc *int - ip *string - ipBind *string - metaFolder *string - peers *string - mastersDeprecated *string // deprecated, for backward compatibility in master.follower - volumeSizeLimitMB *uint - fileSizeLimitMB *int - volumePreallocate *bool - maxParallelVacuumPerServer *int - vacuumIntervalSeconds *int + port *int + portGrpc *int + ip *string + ipBind *string + metaFolder *string + peers *string + mastersDeprecated *string // deprecated, for backward compatibility in master.follower + volumeSizeLimitMB *uint + fileSizeLimitMB *int + volumePreallocate *bool + maxParallelVacuumPerServer *int + vacuumIntervalSeconds *int + vacuumDeleteEmptyAfterSeconds *int // pulseSeconds *int defaultReplication *string garbageThreshold *float64 @@ -95,6 +96,7 @@ func init() { m.volumePreallocate = cmdMaster.Flag.Bool("volumePreallocate", false, "Preallocate disk space for volumes.") m.maxParallelVacuumPerServer = cmdMaster.Flag.Int("maxParallelVacuumPerServer", 1, "maximum number of volumes to vacuum in parallel per volume server") m.vacuumIntervalSeconds = cmdMaster.Flag.Int("vacuumIntervalSeconds", 840, "seconds between automatic vacuum sweeps") + m.vacuumDeleteEmptyAfterSeconds = cmdMaster.Flag.Int("vacuumDeleteEmptyAfterSeconds", 0, "automatic sweep deletes volume copies that stay empty this many seconds; 0 disables") // m.pulseSeconds = cmdMaster.Flag.Int("pulseSeconds", 5, "number of seconds between heartbeats") m.defaultReplication = cmdMaster.Flag.String("defaultReplication", "", "Default replication type if not specified.") m.garbageThreshold = cmdMaster.Flag.Float64("garbageThreshold", 0.3, "threshold to vacuum and reclaim spaces") @@ -469,13 +471,14 @@ func peerIndex(self pb.ServerAddress, peers []pb.ServerAddress) int { func (m *MasterOptions) toMasterOption(whiteList []string) *weed_server.MasterOption { masterAddress := pb.NewServerAddress(*m.ip, *m.port, *m.portGrpc) return &weed_server.MasterOption{ - Master: masterAddress, - MetaFolder: *m.metaFolder, - VolumeSizeLimitMB: uint32(*m.volumeSizeLimitMB), - FileSizeLimitMB: *m.fileSizeLimitMB, - VolumePreallocate: *m.volumePreallocate, - MaxParallelVacuumPerServer: *m.maxParallelVacuumPerServer, - VacuumIntervalSeconds: *m.vacuumIntervalSeconds, + Master: masterAddress, + MetaFolder: *m.metaFolder, + VolumeSizeLimitMB: uint32(*m.volumeSizeLimitMB), + FileSizeLimitMB: *m.fileSizeLimitMB, + VolumePreallocate: *m.volumePreallocate, + MaxParallelVacuumPerServer: *m.maxParallelVacuumPerServer, + VacuumIntervalSeconds: *m.vacuumIntervalSeconds, + VacuumDeleteEmptyAfterSeconds: *m.vacuumDeleteEmptyAfterSeconds, // PulseSeconds: *m.pulseSeconds, DefaultReplicaPlacement: *m.defaultReplication, GarbageThreshold: *m.garbageThreshold, diff --git a/weed/command/master_follower.go b/weed/command/master_follower.go index 1a8ffefba..4e0983eb9 100644 --- a/weed/command/master_follower.go +++ b/weed/command/master_follower.go @@ -48,6 +48,7 @@ func init() { mf.raftResumeState = aws.Bool(false) mf.maxParallelVacuumPerServer = aws.Int(1) mf.vacuumIntervalSeconds = aws.Int(840) + mf.vacuumDeleteEmptyAfterSeconds = aws.Int(0) mf.telemetryUrl = aws.String("https://telemetry.seaweedfs.com/api/collect") mf.telemetryEnabled = aws.Bool(false) } diff --git a/weed/command/mini.go b/weed/command/mini.go index 2cf14a491..e5c2460dd 100644 --- a/weed/command/mini.go +++ b/weed/command/mini.go @@ -428,6 +428,7 @@ func initMiniMasterFlags() { miniMasterOptions.volumePreallocate = cmdMini.Flag.Bool("master.volumePreallocate", false, "Preallocate disk space for volumes.") miniMasterOptions.maxParallelVacuumPerServer = cmdMini.Flag.Int("master.maxParallelVacuumPerServer", 1, "maximum number of volumes to vacuum in parallel on one volume server") miniMasterOptions.vacuumIntervalSeconds = cmdMini.Flag.Int("master.vacuumIntervalSeconds", 840, "seconds between automatic vacuum sweeps") + miniMasterOptions.vacuumDeleteEmptyAfterSeconds = cmdMini.Flag.Int("master.vacuumDeleteEmptyAfterSeconds", 0, "automatic sweep deletes volume copies that stay empty this many seconds; 0 disables") miniMasterOptions.defaultReplication = cmdMini.Flag.String("master.defaultReplication", "", "Default replication type if not specified.") miniMasterOptions.garbageThreshold = cmdMini.Flag.Float64("master.garbageThreshold", 0.3, "threshold to vacuum and reclaim spaces") miniMasterOptions.metricsAddress = cmdMini.Flag.String("master.metrics.address", "", "Prometheus gateway address") diff --git a/weed/command/server.go b/weed/command/server.go index b987fe9a3..6fbacae0d 100644 --- a/weed/command/server.go +++ b/weed/command/server.go @@ -99,6 +99,7 @@ func init() { masterOptions.volumePreallocate = cmdServer.Flag.Bool("master.volumePreallocate", false, "Preallocate disk space for volumes.") masterOptions.maxParallelVacuumPerServer = cmdServer.Flag.Int("master.maxParallelVacuumPerServer", 1, "maximum number of volumes to vacuum in parallel on one volume server") masterOptions.vacuumIntervalSeconds = cmdServer.Flag.Int("master.vacuumIntervalSeconds", 840, "seconds between automatic vacuum sweeps") + masterOptions.vacuumDeleteEmptyAfterSeconds = cmdServer.Flag.Int("master.vacuumDeleteEmptyAfterSeconds", 0, "automatic sweep deletes volume copies that stay empty this many seconds; 0 disables") masterOptions.defaultReplication = cmdServer.Flag.String("master.defaultReplication", "", "Default replication type if not specified.") masterOptions.garbageThreshold = cmdServer.Flag.Float64("master.garbageThreshold", 0.3, "threshold to vacuum and reclaim spaces") masterOptions.metricsAddress = cmdServer.Flag.String("master.metrics.address", "", "Prometheus gateway address") diff --git a/weed/server/master_grpc_server_volume.go b/weed/server/master_grpc_server_volume.go index 1a3ea6c09..deb4a4f0e 100644 --- a/weed/server/master_grpc_server_volume.go +++ b/weed/server/master_grpc_server_volume.go @@ -382,7 +382,7 @@ func (ms *MasterServer) VacuumVolume(ctx context.Context, req *master_pb.VacuumV resp := &master_pb.VacuumVolumeResponse{} - ms.Topo.Vacuum(ms.grpcDialOption, float64(req.GarbageThreshold), ms.option.MaxParallelVacuumPerServer, req.VolumeId, req.Collection, ms.preallocateSize, false) + ms.Topo.Vacuum(ms.grpcDialOption, float64(req.GarbageThreshold), ms.option.MaxParallelVacuumPerServer, req.VolumeId, req.Collection, ms.preallocateSize, false, 0) return resp, nil } diff --git a/weed/server/master_server.go b/weed/server/master_server.go index 73106c99b..dbf746997 100644 --- a/weed/server/master_server.go +++ b/weed/server/master_server.go @@ -45,13 +45,14 @@ const ( ) type MasterOption struct { - Master pb.ServerAddress - MetaFolder string - VolumeSizeLimitMB uint32 - FileSizeLimitMB int - VolumePreallocate bool - MaxParallelVacuumPerServer int - VacuumIntervalSeconds int + Master pb.ServerAddress + MetaFolder string + VolumeSizeLimitMB uint32 + FileSizeLimitMB int + VolumePreallocate bool + MaxParallelVacuumPerServer int + VacuumIntervalSeconds int + VacuumDeleteEmptyAfterSeconds int // PulseSeconds int DefaultReplicaPlacement string GarbageThreshold float64 @@ -204,6 +205,7 @@ func NewMasterServer(r *mux.Router, option *MasterOption, peers map[string]pb.Se topology.VolumeGrowStrategy.Threshold, ms.preallocateSize, time.Duration(ms.option.VacuumIntervalSeconds)*time.Second, + time.Duration(ms.option.VacuumDeleteEmptyAfterSeconds)*time.Second, ) ms.ProcessGrowRequest() diff --git a/weed/server/master_server_handlers_admin.go b/weed/server/master_server_handlers_admin.go index 54802d93d..293043d67 100644 --- a/weed/server/master_server_handlers_admin.go +++ b/weed/server/master_server_handlers_admin.go @@ -61,7 +61,7 @@ func (ms *MasterServer) volumeVacuumHandler(w http.ResponseWriter, r *http.Reque } } // glog.Infoln("garbageThreshold =", gcThreshold) - ms.Topo.Vacuum(ms.grpcDialOption, gcThreshold, ms.option.MaxParallelVacuumPerServer, 0, "", ms.preallocateSize, false) + ms.Topo.Vacuum(ms.grpcDialOption, gcThreshold, ms.option.MaxParallelVacuumPerServer, 0, "", ms.preallocateSize, false, 0) ms.dirStatusHandler(w, r) } diff --git a/weed/topology/topology_event_handling.go b/weed/topology/topology_event_handling.go index c5c031a67..a2825345b 100644 --- a/weed/topology/topology_event_handling.go +++ b/weed/topology/topology_event_handling.go @@ -13,7 +13,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/storage" ) -func (t *Topology) StartRefreshWritableVolumes(grpcDialOption grpc.DialOption, garbageThreshold float64, concurrentVacuumLimitPerVolumeServer int, growThreshold float64, preallocate int64, vacuumInterval time.Duration) { +func (t *Topology) StartRefreshWritableVolumes(grpcDialOption grpc.DialOption, garbageThreshold float64, concurrentVacuumLimitPerVolumeServer int, growThreshold float64, preallocate int64, vacuumInterval time.Duration, deleteEmptyAfter time.Duration) { if vacuumInterval <= 0 { vacuumInterval = 14 * time.Minute } @@ -39,7 +39,7 @@ func (t *Topology) StartRefreshWritableVolumes(grpcDialOption grpc.DialOption, g t.EnableVacuumByPlugin() } if !t.IsVacuumDisabled() { - t.Vacuum(grpcDialOption, garbageThreshold, concurrentVacuumLimitPerVolumeServer, 0, "", preallocate, true) + t.Vacuum(grpcDialOption, garbageThreshold, concurrentVacuumLimitPerVolumeServer, 0, "", preallocate, true, deleteEmptyAfter) } } else { stats.MasterReplicaPlacementMismatch.Reset() diff --git a/weed/topology/topology_vacuum.go b/weed/topology/topology_vacuum.go index 171f848ce..0a4537a94 100644 --- a/weed/topology/topology_vacuum.go +++ b/weed/topology/topology_vacuum.go @@ -13,7 +13,9 @@ import ( "google.golang.org/grpc" + "github.com/seaweedfs/seaweedfs/weed/storage" "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/operation" @@ -214,7 +216,7 @@ func (t *Topology) batchVacuumVolumeCleanup(grpcDialOption grpc.DialOption, vl * } } -func (t *Topology) Vacuum(grpcDialOption grpc.DialOption, garbageThreshold float64, maxParallelVacuumPerServer int, volumeId uint32, collection string, preallocate int64, automatic bool) { +func (t *Topology) Vacuum(grpcDialOption grpc.DialOption, garbageThreshold float64, maxParallelVacuumPerServer int, volumeId uint32, collection string, preallocate int64, automatic bool, deleteEmptyAfter time.Duration) { // if there is vacuum going on, return immediately swapped := atomic.CompareAndSwapInt64(&t.vacuumLockCounter, 0, 1) @@ -248,7 +250,7 @@ func (t *Topology) Vacuum(grpcDialOption grpc.DialOption, garbageThreshold float t.vacuumOneVolumeId(grpcDialOption, volumeLayout, c, garbageThreshold, locationList, vid, preallocate, false) } } else { - t.vacuumOneVolumeLayout(grpcDialOption, volumeLayout, c, garbageThreshold, maxParallelVacuumPerServer, preallocate, automatic) + t.vacuumOneVolumeLayout(grpcDialOption, volumeLayout, c, garbageThreshold, maxParallelVacuumPerServer, preallocate, automatic, deleteEmptyAfter) } } if automatic && t.IsVacuumDisabled() { @@ -262,7 +264,7 @@ func (t *Topology) Vacuum(grpcDialOption grpc.DialOption, garbageThreshold float } } -func (t *Topology) vacuumOneVolumeLayout(grpcDialOption grpc.DialOption, volumeLayout *VolumeLayout, c *Collection, garbageThreshold float64, maxParallelVacuumPerServer int, preallocate int64, automatic bool) { +func (t *Topology) vacuumOneVolumeLayout(grpcDialOption grpc.DialOption, volumeLayout *VolumeLayout, c *Collection, garbageThreshold float64, maxParallelVacuumPerServer int, preallocate int64, automatic bool, deleteEmptyAfter time.Duration) { volumeLayout.accessLock.RLock() todoVolumeMap := make(map[needle.VolumeId]*VolumeLocationList) @@ -271,6 +273,12 @@ func (t *Topology) vacuumOneVolumeLayout(grpcDialOption grpc.DialOption, volumeL } volumeLayout.accessLock.RUnlock() + // Empty volumes hold their slots forever: deleting them here also spares + // compacting bytes that are all deleted already. + if deleteEmptyAfter > 0 { + t.deleteEmptyVolumes(grpcDialOption, volumeLayout, todoVolumeMap, deleteEmptyAfter) + } + // limiter for each volume server limiter := make(map[NodeId]int) var limiterLock sync.Mutex @@ -377,3 +385,74 @@ func (t *Topology) vacuumOneVolumeId(grpcDialOption grpc.DialOption, volumeLayou } } } + +// deleteEmptyVolumes removes a volume whose every replica copy has stayed +// empty and quiet for quietPeriod, the same rule volume.deleteEmpty applies +// on demand. A copy that still holds data or was written recently keeps the +// whole volume: deleting only the empty copies would silently cut the +// surviving copy's replica count. Fully deleted vids leave the sweep's work +// map; anything else falls through to the normal compaction path. +func (t *Topology) deleteEmptyVolumes(grpcDialOption grpc.DialOption, vl *VolumeLayout, todoVolumeMap map[needle.VolumeId]*VolumeLocationList, quietPeriod time.Duration) { + quietSeconds := int64(quietPeriod / time.Second) + nowUnixSeconds := time.Now().Unix() + for vid, locationList := range todoVolumeMap { + eligible := true + for _, dn := range locationList.list { + v, err := dn.GetVolumesById(vid) + if err != nil || !isEmptyVolumeDeleteCandidate(v, quietSeconds, nowUnixSeconds) { + eligible = false + break + } + } + if !eligible { + continue + } + // keep the volume out of assignments for the whole delete: a heartbeat + // landing mid-delete must not re-add it while replicas are dropping + vl.MarkDeleting(vid) + remaining := 0 + kept := locationList.list[:0] + for _, dn := range locationList.list { + v, err := dn.GetVolumesById(vid) + if err == nil { + onlyGarbage := v.FileCount > 0 && v.FileCount <= v.DeleteCount + glog.V(0).Infof("deleting empty volume %d on %s", vid, dn.ServerAddress()) + if err = t.deleteEmptyVolume(grpcDialOption, dn, vid, onlyGarbage); err == nil { + t.UnRegisterVolumeLayout(v, dn) + continue + } + } + glog.Warningf("delete empty volume %d on %s: %v", vid, dn.ServerAddress(), err) + kept = append(kept, dn) + remaining++ + } + vl.UnmarkDeleting(vid) + locationList.list = kept + if remaining == 0 { + delete(todoVolumeMap, vid) + } + } +} + +func (t *Topology) deleteEmptyVolume(grpcDialOption grpc.DialOption, dn *DataNode, vid needle.VolumeId, onlyGarbage bool) error { + return operation.WithVolumeServerClient(false, dn.ServerAddress(), grpcDialOption, func(client volume_server_pb.VolumeServerClient) error { + // onlyEmpty stays set so a pre-upgrade server checks emptiness and + // refuses instead of deleting a volume that changed since the report. + ctx, cancel := context.WithTimeout(context.Background(), allocateVolumeTimeout) + defer cancel() + _, err := client.VolumeDelete(ctx, &volume_server_pb.VolumeDeleteRequest{ + VolumeId: uint32(vid), + OnlyEmpty: true, + OnlyGarbage: onlyGarbage, + }) + return err + }) +} + +func isEmptyVolumeDeleteCandidate(v storage.VolumeInfo, quietSeconds, nowUnixSeconds int64) bool { + return 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/topology/topology_vacuum_test.go b/weed/topology/topology_vacuum_test.go new file mode 100644 index 000000000..fd193b495 --- /dev/null +++ b/weed/topology/topology_vacuum_test.go @@ -0,0 +1,255 @@ +package topology + +import ( + "context" + "errors" + "fmt" + "net" + "sync" + "testing" + "time" + + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" + "github.com/seaweedfs/seaweedfs/weed/sequence" + "github.com/seaweedfs/seaweedfs/weed/storage" + "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/super_block" + "github.com/seaweedfs/seaweedfs/weed/storage/types" +) + +func TestIsEmptyVolumeDeleteCandidate(t *testing.T) { + now := time.Now().Unix() + old := now - 7200 + base := storage.VolumeInfo{ + Size: 0, + ModifiedAtSecond: old, + } + for _, tc := range []struct { + name string + v storage.VolumeInfo + quiet int64 + want bool + }{ + {"empty and quiet", base, 3600, true}, + {"all garbage and quiet", storage.VolumeInfo{Size: 1 << 20, FileCount: 100, DeleteCount: 100, ModifiedAtSecond: old}, 3600, true}, + {"remote-backed", storage.VolumeInfo{Size: 0, ModifiedAtSecond: old, RemoteStorageName: "s3"}, 3600, false}, + {"read-only without delete permission", storage.VolumeInfo{Size: 0, ModifiedAtSecond: old, ReadOnly: true}, 3600, false}, + {"read-only can delete", storage.VolumeInfo{Size: 0, ModifiedAtSecond: old, ReadOnly: true, ReadOnlyCanDelete: true}, 3600, true}, + {"live files", storage.VolumeInfo{Size: 1 << 20, FileCount: 100, DeleteCount: 50, ModifiedAtSecond: old}, 3600, false}, + {"no live files but oversized is fine", storage.VolumeInfo{Size: 1 << 20, FileCount: 100, DeleteCount: 100, ModifiedAtSecond: old}, 3600, true}, + {"unreported mtime", storage.VolumeInfo{Size: 0}, 3600, false}, + {"still quiet-recent", storage.VolumeInfo{Size: 0, ModifiedAtSecond: now - 60}, 3600, false}, + } { + t.Run(tc.name, func(t *testing.T) { + if got := isEmptyVolumeDeleteCandidate(tc.v, tc.quiet, now); got != tc.want { + t.Errorf("isEmptyVolumeDeleteCandidate() = %v, want %v", got, tc.want) + } + }) + } +} + +type fakeVolumeDeleteServer struct { + volume_server_pb.UnimplementedVolumeServerServer + mu sync.Mutex + deletes []*volume_server_pb.VolumeDeleteRequest +} + +func (f *fakeVolumeDeleteServer) VolumeDelete(ctx context.Context, req *volume_server_pb.VolumeDeleteRequest) (*volume_server_pb.VolumeDeleteResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.deletes = append(f.deletes, req) + return &volume_server_pb.VolumeDeleteResponse{}, nil +} + +func startFakeVolumeServer(t *testing.T, vs *fakeVolumeDeleteServer) (grpcPort int, dialOption grpc.DialOption) { + t.Helper() + lis, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + srv := grpc.NewServer() + volume_server_pb.RegisterVolumeServerServer(srv, vs) + serveErr := make(chan error, 1) + go func() { serveErr <- srv.Serve(lis) }() + t.Cleanup(func() { + srv.Stop() + if err := <-serveErr; err != nil && !errors.Is(err, grpc.ErrServerStopped) { + t.Errorf("fake volume server: %v", err) + } + }) + return lis.Addr().(*net.TCPAddr).Port, grpc.WithTransportCredentials(insecure.NewCredentials()) +} + +func TestDeleteEmptyVolumesDeletesOnlyQuietEmptyCopies(t *testing.T) { + fake := &fakeVolumeDeleteServer{} + grpcPort, dialOption := startFakeVolumeServer(t, fake) + + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + rack := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1") + dn := rack.GetOrCreateDataNode("127.0.0.1", 8080, grpcPort, "127.0.0.1", fmt.Sprintf("dn-%d", grpcPort), map[string]uint32{"": 10}) + + now := time.Now().Unix() + old := now - 7200 + volume := func(id int, size uint64, files, deleted uint32, mtime int64) storage.VolumeInfo { + return storage.VolumeInfo{ + Id: needle.VolumeId(id), + Size: size, + Collection: "c", + FileCount: files, + DeleteCount: deleted, + ModifiedAtSecond: mtime, + Version: needle.GetCurrentVersion(), + ReplicaPlacement: &super_block.ReplicaPlacement{}, + Ttl: needle.EMPTY_TTL, + } + } + quietEmpty := volume(1, 0, 0, 0, old) + quietGarbage := volume(2, 1<<20, 100, 100, old) + recentEmpty := volume(3, 0, 0, 0, now) + live := volume(4, 1<<20, 100, 50, old) + + todo := make(map[needle.VolumeId]*VolumeLocationList) + dn.UpdateVolumes([]storage.VolumeInfo{quietEmpty, quietGarbage, recentEmpty, live}) + for _, v := range []storage.VolumeInfo{quietEmpty, quietGarbage, recentEmpty, live} { + topo.RegisterVolumeLayout(v, dn) + ll := NewVolumeLocationList() + ll.list = append(ll.list, dn) + todo[v.Id] = ll + } + + vl := topo.GetVolumeLayout("c", &super_block.ReplicaPlacement{}, needle.EMPTY_TTL, types.ToDiskType("")) + topo.deleteEmptyVolumes(dialOption, vl, todo, time.Hour) + + if _, ok := todo[quietEmpty.Id]; ok { + t.Error("quiet empty volume stayed in the sweep map") + } + if _, ok := todo[quietGarbage.Id]; ok { + t.Error("quiet all-garbage volume stayed in the sweep map") + } + for _, v := range []storage.VolumeInfo{recentEmpty, live} { + if _, ok := todo[v.Id]; !ok { + t.Errorf("volume %d was deleted though not an empty-quiet candidate", v.Id) + } + } + + fake.mu.Lock() + defer fake.mu.Unlock() + if len(fake.deletes) != 2 { + t.Fatalf("VolumeDelete calls = %d, want 2", len(fake.deletes)) + } + got := map[uint32]*volume_server_pb.VolumeDeleteRequest{} + for _, d := range fake.deletes { + if !d.OnlyEmpty { + t.Errorf("VolumeDelete %d missing the onlyEmpty guard", d.VolumeId) + } + got[d.VolumeId] = d + } + if got[uint32(quietEmpty.Id)].OnlyGarbage { + t.Error("truly empty volume deleted with onlyGarbage") + } + if !got[uint32(quietGarbage.Id)].OnlyGarbage { + t.Error("all-garbage volume should carry onlyGarbage") + } +} + +// A volume whose sibling replica holds live files is not "empty": deleting +// the empty copy would silently cut the live copy's replica count, so the +// whole vid is left alone. +func TestDeleteEmptyVolumesSkipsVidWithLiveReplica(t *testing.T) { + fake := &fakeVolumeDeleteServer{} + grpcPort, dialOption := startFakeVolumeServer(t, fake) + + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + rack := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1") + emptyDn := rack.GetOrCreateDataNode("127.0.0.1", 8080, grpcPort, "127.0.0.1", "dn-empty", map[string]uint32{"": 10}) + liveDn := rack.GetOrCreateDataNode("127.0.0.2", 8080, 1, "127.0.0.2", "dn-live", map[string]uint32{"": 10}) + + now := time.Now().Unix() + v := storage.VolumeInfo{ + Id: needle.VolumeId(1), + Size: 0, + Collection: "c", + ModifiedAtSecond: now - 7200, + Version: needle.GetCurrentVersion(), + ReplicaPlacement: &super_block.ReplicaPlacement{}, + Ttl: needle.EMPTY_TTL, + } + live := storage.VolumeInfo{ + Id: v.Id, + Size: 1 << 20, + Collection: "c", + FileCount: 100, + DeleteCount: 50, + ModifiedAtSecond: now - 7200, + Version: needle.GetCurrentVersion(), + ReplicaPlacement: &super_block.ReplicaPlacement{}, + Ttl: needle.EMPTY_TTL, + } + emptyDn.UpdateVolumes([]storage.VolumeInfo{v}) + liveDn.UpdateVolumes([]storage.VolumeInfo{live}) + topo.RegisterVolumeLayout(v, emptyDn) + topo.RegisterVolumeLayout(live, liveDn) + + ll := NewVolumeLocationList() + ll.list = append(ll.list, emptyDn, liveDn) + todo := map[needle.VolumeId]*VolumeLocationList{v.Id: ll} + + vl := topo.GetVolumeLayout("c", &super_block.ReplicaPlacement{}, needle.EMPTY_TTL, types.ToDiskType("")) + topo.deleteEmptyVolumes(dialOption, vl, todo, time.Hour) + + if _, ok := todo[v.Id]; !ok { + t.Fatal("vid left the sweep map while a live replica keeps it") + } + fake.mu.Lock() + defer fake.mu.Unlock() + if len(fake.deletes) != 0 { + t.Fatalf("VolumeDelete calls = %d, want 0 (a live sibling keeps the whole vid)", len(fake.deletes)) + } +} + +// With every copy eligible a copy whose delete RPC fails keeps the vid in the +// sweep map; the deleted copy is gone but the vid falls through to the normal +// compaction path for its surviving replicas. +func TestDeleteEmptyVolumesKeepsVidWhenCopyDeleteFails(t *testing.T) { + fake := &fakeVolumeDeleteServer{} + grpcPort, dialOption := startFakeVolumeServer(t, fake) + + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false) + rack := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1") + emptyDn := rack.GetOrCreateDataNode("127.0.0.1", 8080, grpcPort, "127.0.0.1", "dn-empty", map[string]uint32{"": 10}) + deadDn := rack.GetOrCreateDataNode("127.0.0.2", 8080, 1, "127.0.0.2", "dn-dead", map[string]uint32{"": 10}) + + now := time.Now().Unix() + v := storage.VolumeInfo{ + Id: needle.VolumeId(1), + Size: 0, + Collection: "c", + ModifiedAtSecond: now - 7200, + Version: needle.GetCurrentVersion(), + ReplicaPlacement: &super_block.ReplicaPlacement{}, + Ttl: needle.EMPTY_TTL, + } + emptyDn.UpdateVolumes([]storage.VolumeInfo{v}) + deadDn.UpdateVolumes([]storage.VolumeInfo{v}) + topo.RegisterVolumeLayout(v, emptyDn) + topo.RegisterVolumeLayout(v, deadDn) + + ll := NewVolumeLocationList() + ll.list = append(ll.list, emptyDn, deadDn) + todo := map[needle.VolumeId]*VolumeLocationList{v.Id: ll} + + vl := topo.GetVolumeLayout("c", &super_block.ReplicaPlacement{}, needle.EMPTY_TTL, types.ToDiskType("")) + topo.deleteEmptyVolumes(dialOption, vl, todo, time.Hour) + + if _, ok := todo[v.Id]; !ok { + t.Fatal("vid left the sweep map while a replica delete failed") + } + fake.mu.Lock() + defer fake.mu.Unlock() + if len(fake.deletes) != 1 { + t.Fatalf("VolumeDelete calls = %d, want 1 (only the reachable copy)", len(fake.deletes)) + } +} diff --git a/weed/topology/volume_layout.go b/weed/topology/volume_layout.go index 8a36f79cc..b8da55d5f 100644 --- a/weed/topology/volume_layout.go +++ b/weed/topology/volume_layout.go @@ -65,6 +65,7 @@ type VolumeLayout struct { writableMembers map[needle.VolumeId]struct{} crowded map[needle.VolumeId]struct{} vacuumedVolumes map[needle.VolumeId]time.Time + deletingVolumes map[needle.VolumeId]struct{} volumeSizeLimit uint64 replicationAsMin bool accessLock sync.RWMutex @@ -92,6 +93,7 @@ func NewVolumeLayout(rp *super_block.ReplicaPlacement, ttl *needle.TTL, diskType writableMembers: make(map[needle.VolumeId]struct{}), crowded: make(map[needle.VolumeId]struct{}), vacuumedVolumes: make(map[needle.VolumeId]time.Time), + deletingVolumes: make(map[needle.VolumeId]struct{}), volumeSizeLimit: volumeSizeLimit, replicationAsMin: replicationAsMin, sizeTracking: make(map[needle.VolumeId]*volumeSizeTracking), @@ -553,6 +555,30 @@ func (vl *VolumeLayout) DrainAndRemoveFromWritable(vid needle.VolumeId) { vl.waitForPendingDrain(context.Background(), vid) } +// MarkDeleting pins a volume out of the writable list for the duration of a +// sweep delete. DrainAndRemoveFromWritable alone is not enough here: unlike a +// compaction the volume disappears, so a heartbeat landing mid-delete must not +// re-add it and let a write reach a replica whose siblings are already gone. +func (vl *VolumeLayout) MarkDeleting(vid needle.VolumeId) { + vl.accessLock.Lock() + vl.deletingVolumes[vid] = struct{}{} + vl.removeFromWritable(vid) + vl.accessLock.Unlock() + vl.waitForPendingDrain(context.Background(), vid) +} + +func (vl *VolumeLayout) UnmarkDeleting(vid needle.VolumeId) { + vl.accessLock.Lock() + delete(vl.deletingVolumes, vid) + // Re-run the standard writable gate so a surviving volume regains + // assignment immediately; digest heartbeats only re-report changed + // volumes, so waiting on them could strand it unwritable indefinitely. + if !vl.vid2location[vid].AnyOversized() && vl.enoughCopies(vid) && vl.isAllWritable(vid) { + vl.setVolumeWritable(vid) + } + vl.accessLock.Unlock() +} + func (vl *VolumeLayout) isEmpty() bool { vl.accessLock.RLock() defer vl.accessLock.RUnlock() @@ -920,6 +946,9 @@ func (vl *VolumeLayout) removeFromWritable(vid needle.VolumeId) bool { return false } func (vl *VolumeLayout) setVolumeWritable(vid needle.VolumeId) bool { + if _, ok := vl.deletingVolumes[vid]; ok { + return false + } if _, ok := vl.writableMembers[vid]; ok { return false } diff --git a/weed/topology/volume_layout_drain_test.go b/weed/topology/volume_layout_drain_test.go index 00490b234..0199266b2 100644 --- a/weed/topology/volume_layout_drain_test.go +++ b/weed/topology/volume_layout_drain_test.go @@ -4,6 +4,7 @@ import ( "testing" "time" + "github.com/seaweedfs/seaweedfs/weed/storage" "github.com/seaweedfs/seaweedfs/weed/storage/needle" "github.com/seaweedfs/seaweedfs/weed/storage/super_block" "github.com/seaweedfs/seaweedfs/weed/storage/types" @@ -259,3 +260,39 @@ func TestSetVolumeReadOnly_PreservesPending(t *testing.T) { t.Errorf("expected 5000 pending (not drained), got %d", p) } } + +// A heartbeat landing mid-delete re-evaluates writability; a marked volume +// must stay out of the writable list until the sweep finishes with it. +func TestMarkDeletingKeepsVolumeUnwritableAcrossHeartbeats(t *testing.T) { + layout := ` +{ + "dc1":{ + "rack1":{ + "server1":{ + "volumes":[ + {"id":1, "size":1000, "replication":"000"} + ], + "limit":10 + } + } + } +} +` + _, vl := setupPickTest(t, layout, 10000) + + if writable, _ := vl.GetWritableVolumeCount(); writable != 1 { + t.Fatalf("volume not writable before mark, got %d", writable) + } + + vl.MarkDeleting(1) + + vl.EnsureCorrectWritables(&storage.VolumeInfo{Id: 1}) + if writable, _ := vl.GetWritableVolumeCount(); writable != 0 { + t.Fatal("marked volume became writable again") + } + + vl.UnmarkDeleting(1) + if writable, _ := vl.GetWritableVolumeCount(); writable != 1 { + t.Fatal("unmarked volume did not regain writability without a heartbeat") + } +}