mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-07 23:07:48 +02:00
A draining filer stayed in the master's lock ring until its process exited, while gRPC GracefulStop was already refusing new connections. S3 gateways and peer filers kept routing object-write locks and owner-routed writes to it for the whole graceful-stop window. On shutdown the filer now sends leave_lock_ring on its open KeepConnected stream. The master removes it from the lock ring only; it stays a cluster member so peers keep following its metadata log through the drain. Once the ring update without it arrives, the filer has already transferred its locks to the new owners, and it keeps serving through the prior-owner window (plus a second for peers that apply the update later) before gRPC and HTTP begin draining. The whole leave is bounded at 10s so slow lock transfers or a stuck stream cannot hold up the drain. A lone filer, or a master that ignores the message, falls back to the previous behavior.
109 lines
3.4 KiB
Go
109 lines
3.4 KiB
Go
package wdclient
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
|
"google.golang.org/grpc"
|
|
"google.golang.org/grpc/credentials/insecure"
|
|
)
|
|
|
|
type fakeKeepConnectedServer struct {
|
|
master_pb.UnimplementedSeaweedServer
|
|
requests chan *master_pb.KeepConnectedRequest
|
|
// closeAfterLeave ends the stream once a leave arrives, forcing a reconnect.
|
|
closeAfterLeave bool
|
|
}
|
|
|
|
func (s *fakeKeepConnectedServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServer) error {
|
|
if err := stream.Send(&master_pb.KeepConnectedResponse{VolumeLocation: &master_pb.VolumeLocation{}}); err != nil {
|
|
return err
|
|
}
|
|
for {
|
|
req, err := stream.Recv()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
s.requests <- req
|
|
if req.LeaveLockRing && req.ClientType == "" {
|
|
if s.closeAfterLeave {
|
|
return nil
|
|
}
|
|
if err := stream.Send(&master_pb.KeepConnectedResponse{LockRingUpdate: &master_pb.LockRingUpdate{Servers: []string{"filer2:18888"}, Version: 2}}); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func nextRequest(t *testing.T, requests <-chan *master_pb.KeepConnectedRequest) *master_pb.KeepConnectedRequest {
|
|
t.Helper()
|
|
select {
|
|
case req := <-requests:
|
|
return req
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("master received no KeepConnected message")
|
|
return nil
|
|
}
|
|
}
|
|
|
|
func startLeaveTestClient(t *testing.T, srv *fakeKeepConnectedServer, onLockRingUpdate func(*master_pb.LockRingUpdate)) *MasterClient {
|
|
t.Helper()
|
|
addr := startFakeMasterServer(t, srv)
|
|
mc := NewMasterClient(
|
|
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
|
"", "filer", "filer1:18888", "", "",
|
|
*pb.NewServiceDiscoveryFromMap(map[string]pb.ServerAddress{"m": addr}),
|
|
)
|
|
mc.SetOnLockRingUpdateFn(onLockRingUpdate)
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
t.Cleanup(cancel)
|
|
go mc.KeepConnectedToMaster(ctx)
|
|
return mc
|
|
}
|
|
|
|
func TestLeaveLockRingIsSentOnTheOpenStream(t *testing.T) {
|
|
srv := &fakeKeepConnectedServer{requests: make(chan *master_pb.KeepConnectedRequest, 4)}
|
|
ringUpdates := make(chan *master_pb.LockRingUpdate, 1)
|
|
mc := startLeaveTestClient(t, srv, func(update *master_pb.LockRingUpdate) { ringUpdates <- update })
|
|
|
|
if first := nextRequest(t, srv.requests); first.LeaveLockRing {
|
|
t.Fatal("registration asked to leave the lock ring before LeaveLockRing")
|
|
}
|
|
if err := mc.LeaveLockRing(); err != nil {
|
|
t.Fatalf("LeaveLockRing: %v", err)
|
|
}
|
|
if leave := nextRequest(t, srv.requests); !leave.LeaveLockRing {
|
|
t.Fatalf("expected a leave message, got %+v", leave)
|
|
}
|
|
select {
|
|
case update := <-ringUpdates:
|
|
if update.Version != 2 {
|
|
t.Fatalf("unexpected ring update %+v", update)
|
|
}
|
|
case req := <-srv.requests:
|
|
t.Fatalf("leaving reconnected or re-registered: %+v", req)
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("the stream stopped delivering ring updates after leaving")
|
|
}
|
|
}
|
|
|
|
func TestReconnectAfterLeaveStaysOutOfLockRing(t *testing.T) {
|
|
srv := &fakeKeepConnectedServer{requests: make(chan *master_pb.KeepConnectedRequest, 4), closeAfterLeave: true}
|
|
mc := startLeaveTestClient(t, srv, nil)
|
|
|
|
nextRequest(t, srv.requests)
|
|
if err := mc.LeaveLockRing(); err != nil {
|
|
t.Fatalf("LeaveLockRing: %v", err)
|
|
}
|
|
nextRequest(t, srv.requests)
|
|
|
|
reconnect := nextRequest(t, srv.requests)
|
|
if reconnect.ClientType != "filer" || !reconnect.LeaveLockRing {
|
|
t.Fatalf("reconnect registration must stay out of the lock ring, got %+v", reconnect)
|
|
}
|
|
}
|