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.
This commit is contained in:
Chris Lu authored and GitHub committed 2026-10-03 12:32:42 +08:00
1 parent 3e679e925e
commit 0ca484c354
3 files changed
+162 -25

No files matched your search

+3
View File
@@ -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 {
+26 -13
View File
@@ -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
+133 -12
View File
@@ -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)
}
})
}
}