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