From 0ca484c3547598edd348cf38b42bd6017891ce6b Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Sat, 3 Oct 2026 12:32:42 +0800 Subject: [PATCH] vacuum: bound master vacuum RPCs with phase deadlines (#11579) * vacuum: bound the commit RPC with a phase deadline VacuumVolumeCommit ran on context.Background(), so a volume server that keeps the call pending would hold the topology-wide vacuum guard forever and every later sweep would be skipped. Give the call a deadline scaled like the existing phase waits (one minute per GB of the volume size limit) so a stalled commit ends as an error instead of blocking the sweep; the timeout is a var so tests can shrink it. * vacuum: bound the replica status probe with a phase deadline The VolumeStatus call on replicas that were not compacted also ran on context.Background(), so a stalled replica could pin the sweep the same way a stalled commit can. Give it the same per-phase deadline. * vacuum: bound the cleanup RPC with a phase deadline VacuumVolumeCleanup also ran on context.Background(); a stalled server would keep the sweep worker and the shared vacuum guard pending forever. Give it the same per-phase deadline. * vacuum: let the check and compact phase waits cancel their RPCs The coordinator wait timers fired while the check and compact calls still ran on context.Background(), so the sweep gave up but the RPC goroutine stayed until the server answered, and a compact stream kept writing on the server. Share one deadline context between the wait and the calls so an expired wait actually cancels them. * vacuum: test that a stalled volume server releases the vacuum guard A fake volume server keeps one vacuum-phase RPC pending until the client context is cancelled. Before the phase deadlines, Vacuum never returned and vacuumLockCounter stayed held; now each phase cancels on its deadline and the guard is free for the next request. * volume: stop compaction at the next needle when the client cancels The progress callback only noticed a gone client when a 128 MiB report failed to send, so an aborted VacuumVolumeCompact kept copying for up to a whole interval while the master had already moved on to cleanup. Check the stream context on every needle, the same early return the Rust volume server does with tx.is_closed(). * vacuum: assert the stalled phase RPC is cancelled, not just bypassed The check and compact coordinator waits already returned on timeout before the deadlines existed, so a regression that put the calls back on context.Background() would pass unnoticed. Wait for the fake server to report that the phase RPC context ended. * vacuum: give the stalled-RPC test room to reach the handler The 50ms phase budget starts before goroutine scheduling and the gRPC dial, so a busy test host could expire it before the fake server saw the call. Raise the override to 250ms; the test still finishes in about a second. * vacuum: describe the phase deadline as scaled, not per-GB The formula keeps the exact expression the check and compact waits already used (floor plus one at 1 GiB granularity); it is a backstop, not a per-GB SLO. --- weed/server/volume_grpc_vacuum.go | 3 + weed/topology/topology_vacuum.go | 39 ++++--- weed/topology/topology_vacuum_test.go | 145 +++++++++++++++++++++++--- 3 files changed, 162 insertions(+), 25 deletions(-) diff --git a/weed/server/volume_grpc_vacuum.go b/weed/server/volume_grpc_vacuum.go index e56c48405..968966164 100644 --- a/weed/server/volume_grpc_vacuum.go +++ b/weed/server/volume_grpc_vacuum.go @@ -55,6 +55,9 @@ func (vs *VolumeServer) VacuumVolumeCompact(req *volume_server_pb.VacuumVolumeCo fs, fsErr := procfs.NewDefaultFS() var sendErr error err := vs.store.CompactVolume(needle.VolumeId(req.VolumeId), req.Preallocate, vs.compactionBytePerSecond, func(processed int64) bool { + if stream.Context().Err() != nil { + return false + } if processed > nextReportTarget { resp.ProcessedBytes = processed if fsErr == nil && numCPU > 0 { diff --git a/weed/topology/topology_vacuum.go b/weed/topology/topology_vacuum.go index 4480637d2..7babdceb0 100644 --- a/weed/topology/topology_vacuum.go +++ b/weed/topology/topology_vacuum.go @@ -22,14 +22,25 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" ) +// vacuumPhaseTimeout bounds one synchronous vacuum RPC, scaled with the +// volume size limit, so a stalled volume server cannot hold the vacuum +// guard forever. +var vacuumPhaseTimeout = time.Minute + +func (t *Topology) vacuumRPCTimeout() time.Duration { + return vacuumPhaseTimeout * time.Duration(t.volumeSizeLimit/1024/1024/1000+1) +} + func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid needle.VolumeId, locationlist *VolumeLocationList, garbageThreshold float64, skipReadOnly bool) (*VolumeLocationList, bool) { ch := make(chan int, locationlist.Length()) errCount := int32(0) + ctx, cancel := context.WithTimeout(context.Background(), t.vacuumRPCTimeout()) + defer cancel() for index, dn := range locationlist.list { go func(index int, dn *DataNode, url pb.ServerAddress, vid needle.VolumeId) { err := operation.WithVolumeServerClient(false, url, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error { - resp, err := volumeServerClient.VacuumVolumeCheck(context.Background(), &volume_server_pb.VacuumVolumeCheckRequest{ + resp, err := volumeServerClient.VacuumVolumeCheck(ctx, &volume_server_pb.VacuumVolumeCheckRequest{ VolumeId: uint32(vid), }) if err != nil { @@ -66,16 +77,13 @@ func (t *Topology) batchVacuumVolumeCheck(grpcDialOption grpc.DialOption, vid ne } vacuumLocationList := NewVolumeLocationList() - waitTimeout := time.NewTimer(time.Minute * time.Duration(t.volumeSizeLimit/1024/1024/1000+1)) - defer waitTimeout.Stop() - for range locationlist.list { select { case index := <-ch: if index != -1 { vacuumLocationList.list = append(vacuumLocationList.list, locationlist.list[index]) } - case <-waitTimeout.C: + case <-ctx.Done(): return vacuumLocationList, false } } @@ -87,11 +95,13 @@ func (t *Topology) batchVacuumVolumeCompact(grpcDialOption grpc.DialOption, vl * vl.DrainAndRemoveFromWritable(vid) ch := make(chan bool, locationlist.Length()) + ctx, cancel := context.WithTimeout(context.Background(), 3*t.vacuumRPCTimeout()) + defer cancel() for index, dn := range locationlist.list { go func(index int, url pb.ServerAddress, vid needle.VolumeId) { glog.V(0).Infoln(index, "Start vacuuming", vid, "on", url) err := operation.WithVolumeServerClient(true, url, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error { - stream, err := volumeServerClient.VacuumVolumeCompact(context.Background(), &volume_server_pb.VacuumVolumeCompactRequest{ + stream, err := volumeServerClient.VacuumVolumeCompact(ctx, &volume_server_pb.VacuumVolumeCompactRequest{ VolumeId: uint32(vid), Preallocate: preallocate, }) @@ -124,14 +134,11 @@ func (t *Topology) batchVacuumVolumeCompact(grpcDialOption grpc.DialOption, vl * } isVacuumSuccess := true - waitTimeout := time.NewTimer(3 * time.Minute * time.Duration(t.volumeSizeLimit/1024/1024/1000+1)) - defer waitTimeout.Stop() - for range locationlist.list { select { case canCommit := <-ch: isVacuumSuccess = isVacuumSuccess && canCommit - case <-waitTimeout.C: + case <-ctx.Done(): return false } } @@ -145,7 +152,9 @@ func (t *Topology) batchVacuumVolumeCommit(grpcDialOption grpc.DialOption, vl *V for _, dn := range vacuumLocationList.list { glog.V(0).Infoln("Start Committing vacuum", vid, "on", dn.Url()) err := operation.WithVolumeServerClient(false, dn.ServerAddress(), grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error { - resp, err := volumeServerClient.VacuumVolumeCommit(context.Background(), &volume_server_pb.VacuumVolumeCommitRequest{ + ctx, cancel := context.WithTimeout(context.Background(), t.vacuumRPCTimeout()) + defer cancel() + resp, err := volumeServerClient.VacuumVolumeCommit(ctx, &volume_server_pb.VacuumVolumeCommitRequest{ VolumeId: uint32(vid), }) if resp != nil { @@ -178,7 +187,9 @@ func (t *Topology) batchVacuumVolumeCommit(grpcDialOption grpc.DialOption, vl *V } if !isFound { err := operation.WithVolumeServerClient(false, dn.ServerAddress(), grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error { - resp, err := volumeServerClient.VolumeStatus(context.Background(), &volume_server_pb.VolumeStatusRequest{ + ctx, cancel := context.WithTimeout(context.Background(), t.vacuumRPCTimeout()) + defer cancel() + resp, err := volumeServerClient.VolumeStatus(ctx, &volume_server_pb.VolumeStatusRequest{ VolumeId: uint32(vid), }) if resp != nil { @@ -218,7 +229,9 @@ func (t *Topology) batchVacuumVolumeCleanup(grpcDialOption grpc.DialOption, vl * for _, dn := range locationlist.list { glog.V(0).Infoln("Start cleaning up", vid, "on", dn.Url()) err := operation.WithVolumeServerClient(false, dn.ServerAddress(), grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error { - _, err := volumeServerClient.VacuumVolumeCleanup(context.Background(), &volume_server_pb.VacuumVolumeCleanupRequest{ + ctx, cancel := context.WithTimeout(context.Background(), t.vacuumRPCTimeout()) + defer cancel() + _, err := volumeServerClient.VacuumVolumeCleanup(ctx, &volume_server_pb.VacuumVolumeCleanupRequest{ VolumeId: uint32(vid), }) return err diff --git a/weed/topology/topology_vacuum_test.go b/weed/topology/topology_vacuum_test.go index 2b504b516..6d96add77 100644 --- a/weed/topology/topology_vacuum_test.go +++ b/weed/topology/topology_vacuum_test.go @@ -6,6 +6,7 @@ import ( "fmt" "net" "sync" + "sync/atomic" "testing" "time" @@ -256,12 +257,37 @@ func TestDeleteEmptyVolumesKeepsVidWhenCopyDeleteFails(t *testing.T) { type fakeVacuumServer struct { volume_server_pb.UnimplementedVolumeServerServer - mu sync.Mutex - checks map[uint32]*volume_server_pb.VacuumVolumeCheckResponse - committed []uint32 + mu sync.Mutex + checks map[uint32]*volume_server_pb.VacuumVolumeCheckResponse + committed []uint32 + compactErr error + hang string + entered chan string + cancelled chan struct{} } -func (f *fakeVacuumServer) VacuumVolumeCheck(_ context.Context, req *volume_server_pb.VacuumVolumeCheckRequest) (*volume_server_pb.VacuumVolumeCheckResponse, error) { +func (f *fakeVacuumServer) stall(ctx context.Context, phase string) error { + select { + case f.entered <- phase: + default: + } + <-ctx.Done() + f.mu.Lock() + if f.cancelled != nil { + select { + case <-f.cancelled: + default: + close(f.cancelled) + } + } + f.mu.Unlock() + return ctx.Err() +} + +func (f *fakeVacuumServer) VacuumVolumeCheck(ctx context.Context, req *volume_server_pb.VacuumVolumeCheckRequest) (*volume_server_pb.VacuumVolumeCheckResponse, error) { + if f.hang == "check" { + return nil, f.stall(ctx, "check") + } resp, ok := f.checks[req.VolumeId] if !ok { return nil, fmt.Errorf("volume %d not found", req.VolumeId) @@ -269,21 +295,37 @@ func (f *fakeVacuumServer) VacuumVolumeCheck(_ context.Context, req *volume_serv return resp, nil } -func (f *fakeVacuumServer) VacuumVolumeCompact(_ *volume_server_pb.VacuumVolumeCompactRequest, _ volume_server_pb.VolumeServer_VacuumVolumeCompactServer) error { - return nil +func (f *fakeVacuumServer) VacuumVolumeCompact(_ *volume_server_pb.VacuumVolumeCompactRequest, stream volume_server_pb.VolumeServer_VacuumVolumeCompactServer) error { + if f.hang == "compact" { + return f.stall(stream.Context(), "compact") + } + return f.compactErr } -func (f *fakeVacuumServer) VacuumVolumeCommit(_ context.Context, req *volume_server_pb.VacuumVolumeCommitRequest) (*volume_server_pb.VacuumVolumeCommitResponse, error) { +func (f *fakeVacuumServer) VacuumVolumeCommit(ctx context.Context, req *volume_server_pb.VacuumVolumeCommitRequest) (*volume_server_pb.VacuumVolumeCommitResponse, error) { + if f.hang == "commit" { + return nil, f.stall(ctx, "commit") + } f.mu.Lock() defer f.mu.Unlock() f.committed = append(f.committed, req.VolumeId) return &volume_server_pb.VacuumVolumeCommitResponse{IsReadOnly: true}, nil } -func (f *fakeVacuumServer) VacuumVolumeCleanup(_ context.Context, _ *volume_server_pb.VacuumVolumeCleanupRequest) (*volume_server_pb.VacuumVolumeCleanupResponse, error) { +func (f *fakeVacuumServer) VacuumVolumeCleanup(ctx context.Context, _ *volume_server_pb.VacuumVolumeCleanupRequest) (*volume_server_pb.VacuumVolumeCleanupResponse, error) { + if f.hang == "cleanup" { + return nil, f.stall(ctx, "cleanup") + } return &volume_server_pb.VacuumVolumeCleanupResponse{}, nil } +func (f *fakeVacuumServer) VolumeStatus(ctx context.Context, _ *volume_server_pb.VolumeStatusRequest) (*volume_server_pb.VolumeStatusResponse, error) { + if f.hang == "status" { + return nil, f.stall(ctx, "status") + } + return &volume_server_pb.VolumeStatusResponse{}, nil +} + // A sweep keeps a read-only replica eligible only when the volume server // reports disk_space_low; other read-only causes stay skipped unless the // request names the volume explicitly. @@ -336,10 +378,10 @@ func TestVacuumReadOnlyDiskLowVolume(t *testing.T) { topo.vacuumOneVolumeId(dialOption, vl, c, 0.3, ll, vid, 0, skipReadOnly) } - vacuum(1, true) // read-only but disk_low: compacted - vacuum(2, true) // read-only otherwise: skipped by sweep - vacuum(3, true) // writable: compacted - vacuum(4, true) // disk_low but below threshold: skipped + vacuum(1, true) // read-only but disk_low: compacted + vacuum(2, true) // read-only otherwise: skipped by sweep + vacuum(3, true) // writable: compacted + vacuum(4, true) // disk_low but below threshold: skipped vacuum(5, false) // read-only, explicit request: compacted fake.mu.Lock() @@ -360,3 +402,82 @@ func TestVacuumReadOnlyDiskLowVolume(t *testing.T) { } } } + +// A volume server that keeps a vacuum RPC pending must not hold the +// topology-wide vacuum guard forever: the phase deadline cancels the call, +// Vacuum returns, and the next request is not skipped. +func TestVacuumStalledVolumeServerReleasesGuard(t *testing.T) { + for _, phase := range []string{"check", "compact", "commit", "status", "cleanup"} { + t.Run(phase, func(t *testing.T) { + old := vacuumPhaseTimeout + vacuumPhaseTimeout = 250 * time.Millisecond + defer func() { vacuumPhaseTimeout = old }() + + hangFake := &fakeVacuumServer{ + checks: map[uint32]*volume_server_pb.VacuumVolumeCheckResponse{1: {GarbageRatio: 0.9}}, + hang: phase, + entered: make(chan string, 1), + cancelled: make(chan struct{}), + } + if phase == "cleanup" { + hangFake.compactErr = errors.New("compact failed") + } + if phase == "status" { + hangFake.checks[1] = &volume_server_pb.VacuumVolumeCheckResponse{GarbageRatio: 0.1} + } + + topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 1, 5, true) + rack := topo.GetOrCreateDataCenter("dc1").GetOrCreateRack("rack1") + + var dns []*DataNode + if phase == "status" { + okFake := &fakeVacuumServer{ + checks: map[uint32]*volume_server_pb.VacuumVolumeCheckResponse{1: {GarbageRatio: 0.9}}, + } + okPort, _ := startFakeVolumeServer(t, okFake) + dns = append(dns, rack.GetOrCreateDataNode("127.0.0.1", 8080, okPort, "127.0.0.1", "dn-ok", map[string]uint32{"": 10})) + } + hangPort, dialOption := startFakeVolumeServer(t, hangFake) + dns = append(dns, rack.GetOrCreateDataNode("127.0.0.1", 8081, hangPort, "127.0.0.1", "dn-hang", map[string]uint32{"": 10})) + + v := storage.VolumeInfo{ + Id: needle.VolumeId(1), + Size: 1 << 20, + Collection: "c", + ModifiedAtSecond: time.Now().Unix(), + Version: needle.GetCurrentVersion(), + ReplicaPlacement: &super_block.ReplicaPlacement{}, + Ttl: needle.EMPTY_TTL, + } + for _, dn := range dns { + dn.UpdateVolumes([]storage.VolumeInfo{v}) + topo.RegisterVolumeLayout(v, dn) + } + + done := make(chan struct{}) + go func() { + topo.Vacuum(dialOption, 0.3, 1, 1, "c", 0, false, 0) + close(done) + }() + + select { + case <-hangFake.entered: + case <-time.After(15 * time.Second): + t.Fatalf("vacuum never reached the %s RPC", phase) + } + select { + case <-done: + case <-time.After(15 * time.Second): + t.Fatalf("vacuum did not return while the %s RPC was pending", phase) + } + select { + case <-hangFake.cancelled: + case <-time.After(15 * time.Second): + t.Fatalf("the %s RPC was not cancelled", phase) + } + if c := atomic.LoadInt64(&topo.vacuumLockCounter); c != 0 { + t.Fatalf("vacuumLockCounter = %d, want 0", c) + } + }) + } +}