From 52c6df3beca29c5f5716bf23b786663e73fb1b60 Mon Sep 17 00:00:00 2001 From: Peter Dodd Date: Tue, 6 Oct 2026 02:53:40 +0100 Subject: [PATCH] fix(filer): leave the lock ring before stopping gRPC on shutdown (#11604) 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. --- seaweed-volume/proto/master.proto | 3 + weed/cluster/lock_manager/lock_ring.go | 11 +++ weed/command/filer.go | 9 +- weed/command/filer_shutdown_test.go | 43 ++++++++- weed/pb/master.proto | 3 + weed/pb/master_pb/master.pb.go | 15 +++- weed/server/filer_server.go | 59 +++++++++++++ weed/server/filer_server_leave_test.go | 90 +++++++++++++++++++ weed/server/master_grpc_server.go | 12 ++- weed/server/master_grpc_server_test.go | 23 +++++ weed/wdclient/masterclient.go | 37 +++++++- weed/wdclient/masterclient_leave_test.go | 108 +++++++++++++++++++++++ 12 files changed, 402 insertions(+), 11 deletions(-) create mode 100644 weed/server/filer_server_leave_test.go create mode 100644 weed/wdclient/masterclient_leave_test.go diff --git a/seaweed-volume/proto/master.proto b/seaweed-volume/proto/master.proto index 21987ccc1..153ed085b 100644 --- a/seaweed-volume/proto/master.proto +++ b/seaweed-volume/proto/master.proto @@ -213,6 +213,9 @@ message KeepConnectedRequest { string filer_group = 5; string data_center = 6; string rack = 7; + // A draining filer leaves the lock ring but stays a cluster member, so peers + // keep following its metadata until the stream closes. + bool leave_lock_ring = 8; } message VolumeLocation { diff --git a/weed/cluster/lock_manager/lock_ring.go b/weed/cluster/lock_manager/lock_ring.go index 8b87e6243..625b8d234 100644 --- a/weed/cluster/lock_manager/lock_ring.go +++ b/weed/cluster/lock_manager/lock_ring.go @@ -152,6 +152,17 @@ func (r *LockRing) GetSnapshot() (servers []pb.ServerAddress) { return r.snapshots[0].servers } +// PriorOwnerWindowEnd is when the latest ring change stops routing moved keys +// to their prior owner. +func (r *LockRing) PriorOwnerWindowEnd() time.Time { + r.RLock() + defer r.RUnlock() + if len(r.snapshots) == 0 { + return time.Time{} + } + return r.snapshots[0].ts.Add(r.snapshotInterval) +} + // WaitForCleanup waits for all pending cleanup operations to complete func (r *LockRing) WaitForCleanup() { r.cleanupWg.Wait() diff --git a/weed/command/filer.go b/weed/command/filer.go index a23f31544..e58a672bd 100644 --- a/weed/command/filer.go +++ b/weed/command/filer.go @@ -588,7 +588,7 @@ func (fo *FilerOptions) startFiler() { } httpS := newHttpServer(defaultHandler, tlsConfig) httpServers = append(httpServers, httpS) - shutdown := newFilerShutdown(stopGrpcServer, fs.Shutdown, httpServers...) + shutdown := newFilerShutdown(fs.LeaveLockRing, stopGrpcServer, fs.Shutdown, httpServers...) grace.OnInterrupt(shutdown) @@ -617,7 +617,7 @@ func (fo *FilerOptions) startFiler() { } httpS := newHttpServer(defaultHandler, nil) httpServers = append(httpServers, httpS) - shutdown := newFilerShutdown(stopGrpcServer, fs.Shutdown, httpServers...) + shutdown := newFilerShutdown(fs.LeaveLockRing, stopGrpcServer, fs.Shutdown, httpServers...) grace.OnInterrupt(shutdown) @@ -637,8 +637,11 @@ func (fo *FilerOptions) startFiler() { } // newFilerShutdown joins shutdown callers while gRPC and HTTP drain concurrently. -func newFilerShutdown(stopGrpc, shutdownFiler func(), httpServers ...*http.Server) func() { +// The filer leaves the lock ring first: peers and S3 gateways route keys to it +// until the ring changes, and would hit refused connections once gRPC stops. +func newFilerShutdown(leaveLockRing, stopGrpc, shutdownFiler func(), httpServers ...*http.Server) func() { return sync.OnceFunc(func() { + leaveLockRing() shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() var drained sync.WaitGroup diff --git a/weed/command/filer_shutdown_test.go b/weed/command/filer_shutdown_test.go index f115b8f61..60f32c7ad 100644 --- a/weed/command/filer_shutdown_test.go +++ b/weed/command/filer_shutdown_test.go @@ -56,7 +56,7 @@ func TestFilerShutdownJoinsServeDuringParallelDrain(t *testing.T) { grpcStopped := make(chan struct{}) filerClosed := make(chan struct{}) var closes atomic.Int32 - shutdown := newFilerShutdown(func() { + shutdown := newFilerShutdown(func() {}, func() { close(grpcStarted) <-grpcRelease close(grpcStopped) @@ -138,7 +138,7 @@ func TestFilerShutdownWaitsForHTTPAfterGrpcStops(t *testing.T) { server.Config.RegisterOnShutdown(func() { close(httpClosing) }) grpcStopped := make(chan struct{}) filerClosed := make(chan struct{}) - shutdown := newFilerShutdown(func() { close(grpcStopped) }, func() { close(filerClosed) }, server.Config) + shutdown := newFilerShutdown(func() {}, func() { close(grpcStopped) }, func() { close(filerClosed) }, server.Config) joined := make(chan struct{}) go func() { shutdown(); close(joined) }() <-httpClosing @@ -157,3 +157,42 @@ func TestFilerShutdownWaitsForHTTPAfterGrpcStops(t *testing.T) { t.Error("filer was not closed after HTTP request completed") } } + +func TestFilerShutdownLeavesLockRingBeforeDraining(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {})) + t.Cleanup(server.Close) + httpClosing := make(chan struct{}) + server.Config.RegisterOnShutdown(func() { close(httpClosing) }) + + leaving := make(chan struct{}) + releaseLeave := make(chan struct{}) + finishLeave := sync.OnceFunc(func() { close(releaseLeave) }) + t.Cleanup(finishLeave) + grpcStopping := make(chan struct{}) + shutdown := newFilerShutdown(func() { + close(leaving) + <-releaseLeave + }, func() { close(grpcStopping) }, func() {}, server.Config) + joined := make(chan struct{}) + go func() { shutdown(); close(joined) }() + + select { + case <-leaving: + case <-time.After(5 * time.Second): + t.Fatal("shutdown did not leave the lock ring") + } + select { + case <-grpcStopping: + t.Fatal("gRPC stopped accepting while the filer was still leaving the lock ring") + case <-httpClosing: + t.Fatal("HTTP stopped accepting while the filer was still leaving the lock ring") + case <-time.After(100 * time.Millisecond): + } + finishLeave() + <-joined + select { + case <-grpcStopping: + default: + t.Error("gRPC was not stopped after leaving the lock ring") + } +} diff --git a/weed/pb/master.proto b/weed/pb/master.proto index 21987ccc1..153ed085b 100644 --- a/weed/pb/master.proto +++ b/weed/pb/master.proto @@ -213,6 +213,9 @@ message KeepConnectedRequest { string filer_group = 5; string data_center = 6; string rack = 7; + // A draining filer leaves the lock ring but stays a cluster member, so peers + // keep following its metadata until the stream closes. + bool leave_lock_ring = 8; } message VolumeLocation { diff --git a/weed/pb/master_pb/master.pb.go b/weed/pb/master_pb/master.pb.go index 62928abc3..69833d906 100644 --- a/weed/pb/master_pb/master.pb.go +++ b/weed/pb/master_pb/master.pb.go @@ -998,6 +998,9 @@ type KeepConnectedRequest struct { FilerGroup string `protobuf:"bytes,5,opt,name=filer_group,json=filerGroup,proto3" json:"filer_group,omitempty"` DataCenter string `protobuf:"bytes,6,opt,name=data_center,json=dataCenter,proto3" json:"data_center,omitempty"` Rack string `protobuf:"bytes,7,opt,name=rack,proto3" json:"rack,omitempty"` + // A draining filer leaves the lock ring but stays a cluster member, so peers + // keep following its metadata until the stream closes. + LeaveLockRing bool `protobuf:"varint,8,opt,name=leave_lock_ring,json=leaveLockRing,proto3" json:"leave_lock_ring,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1074,6 +1077,13 @@ func (x *KeepConnectedRequest) GetRack() string { return "" } +func (x *KeepConnectedRequest) GetLeaveLockRing() bool { + if x != nil { + return x.LeaveLockRing + } + return false +} + type VolumeLocation struct { state protoimpl.MessageState `protogen:"open.v1"` Url string `protobuf:"bytes,1,opt,name=url,proto3" json:"url,omitempty"` @@ -5137,7 +5147,7 @@ const file_master_proto_rawDesc = "" + "\x04data\x18\x01 \x01(\rR\x04data\x12\x16\n" + "\x06parity\x18\x02 \x01(\rR\x06parity\x12\x1d\n" + "\n" + - "volume_ids\x18\x03 \x03(\rR\tvolumeIds\"\xce\x01\n" + + "volume_ids\x18\x03 \x03(\rR\tvolumeIds\"\xf6\x01\n" + "\x14KeepConnectedRequest\x12\x1f\n" + "\vclient_type\x18\x01 \x01(\tR\n" + "clientType\x12%\n" + @@ -5147,7 +5157,8 @@ const file_master_proto_rawDesc = "" + "filerGroup\x12\x1f\n" + "\vdata_center\x18\x06 \x01(\tR\n" + "dataCenter\x12\x12\n" + - "\x04rack\x18\a \x01(\tR\x04rack\"\x9e\x03\n" + + "\x04rack\x18\a \x01(\tR\x04rack\x12&\n" + + "\x0fleave_lock_ring\x18\b \x01(\bR\rleaveLockRing\"\x9e\x03\n" + "\x0eVolumeLocation\x12\x10\n" + "\x03url\x18\x01 \x01(\tR\x03url\x12\x1d\n" + "\n" + diff --git a/weed/server/filer_server.go b/weed/server/filer_server.go index 2b2eec443..49c99a87a 100644 --- a/weed/server/filer_server.go +++ b/weed/server/filer_server.go @@ -5,11 +5,13 @@ import ( "fmt" "net/http" "os" + "slices" "strings" "sync" "sync/atomic" "time" + "github.com/seaweedfs/seaweedfs/weed/cluster" "github.com/seaweedfs/seaweedfs/weed/credential" "github.com/seaweedfs/seaweedfs/weed/stats" "golang.org/x/sync/singleflight" @@ -383,6 +385,63 @@ func (fs *FilerServer) Shutdown() { fs.filer.Shutdown() } +// priorOwnerWindowSkew covers peers and gateways that start their prior-owner +// window after this filer, since each times it from its own ring update. +const priorOwnerWindowSkew = time.Second + +// leaveLockRingBudget bounds the leave: installing the ring update runs lock +// transfers under the ring lock, which can outlast the removal timeout. +const leaveLockRingBudget = 10 * time.Second + +// LeaveLockRing hands this filer's lock-ring keys to its peers while its own +// servers still accept the lock transfers and prior-owner writes that follow. +func (fs *FilerServer) LeaveLockRing() { + fs.leaveLockRingWithin(3*cluster.LockRingStabilizationInterval, leaveLockRingBudget) +} + +func (fs *FilerServer) leaveLockRingWithin(removalTimeout, budget time.Duration) { + left := make(chan struct{}) + go func() { + defer close(left) + fs.leaveLockRing(removalTimeout) + }() + select { + case <-left: + case <-time.After(budget): + glog.Warningf("LockRing: %s did not finish leaving within %v, shutting down anyway", fs.option.Host, budget) + } +} + +func (fs *FilerServer) leaveLockRing(removalTimeout time.Duration) { + if fs.filer.Dlm == nil { + return + } + ring := fs.filer.Dlm.LockRing + self := fs.option.Host + members := ring.GetSnapshot() + // A lone filer has no peer to take its keys: leaving would strand lock + // requests arriving during the drain without an owner. + if len(members) < 2 || !slices.Contains(members, self) { + return + } + // A failed send still leaves: the broken stream's close drops this filer + // from the ring and the reconnect registers without it. + if err := fs.filer.MasterClient.LeaveLockRing(); err != nil { + glog.Warningf("LockRing: %s leave request failed, waiting for reconnect: %v", self, err) + } + deadline := time.Now().Add(removalTimeout) + for slices.Contains(ring.GetSnapshot(), self) { + if time.Now().After(deadline) { + glog.Warningf("LockRing: %s still in the ring after %v, shutting down anyway", self, removalTimeout) + return + } + time.Sleep(50 * time.Millisecond) + } + until := ring.PriorOwnerWindowEnd().Add(priorOwnerWindowSkew) + glog.V(0).Infof("LockRing: %s left, serving prior-owner requests until %v", self, until) + time.Sleep(time.Until(until)) +} + func (fs *FilerServer) Reload() { glog.V(0).Infoln("Reload filer server...") diff --git a/weed/server/filer_server_leave_test.go b/weed/server/filer_server_leave_test.go new file mode 100644 index 000000000..b2bee6dcc --- /dev/null +++ b/weed/server/filer_server_leave_test.go @@ -0,0 +1,90 @@ +package weed_server + +import ( + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager" + "github.com/seaweedfs/seaweedfs/weed/filer" + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/wdclient" + "google.golang.org/grpc" +) + +func newLeaveTestFilerServer(priorOwnerWindow time.Duration, ring ...pb.ServerAddress) *FilerServer { + const self = pb.ServerAddress("filer1:18888") + dlm := lock_manager.NewDistributedLockManager(self) + dlm.LockRing = lock_manager.NewLockRing(priorOwnerWindow) + dlm.LockRing.SetSnapshot(ring, 1) + return &FilerServer{ + option: &FilerOption{Host: self}, + filer: &filer.Filer{ + Dlm: dlm, + MasterClient: wdclient.NewMasterClient(grpc.EmptyDialOption{}, "", "filer", self, "", "", pb.ServerDiscovery{}), + }, + } +} + +func TestLeaveLockRingWaitsForRemovalThenPriorOwnerWindow(t *testing.T) { + const window = 300 * time.Millisecond + fs := newLeaveTestFilerServer(window, "filer1:18888", "filer2:18888") + + removed := make(chan time.Time, 1) + go func() { + time.Sleep(100 * time.Millisecond) + before := time.Now() + fs.filer.Dlm.LockRing.SetSnapshot([]pb.ServerAddress{"filer2:18888"}, 2) + removed <- before + }() + + fs.leaveLockRing(5 * time.Second) + returned := time.Now() + + select { + case removedAt := <-removed: + if want := window + priorOwnerWindowSkew; returned.Sub(removedAt) < want-20*time.Millisecond { + t.Errorf("returned %v after the ring dropped this filer, want the %v prior-owner window plus peer skew", returned.Sub(removedAt), want) + } + default: + t.Fatal("returned before the master dropped this filer from the lock ring") + } +} + +func TestLeaveLockRingGivesUpWhenTheMasterKeepsTheFiler(t *testing.T) { + fs := newLeaveTestFilerServer(5*time.Second, "filer1:18888", "filer2:18888") + + start := time.Now() + fs.leaveLockRing(200 * time.Millisecond) + if elapsed := time.Since(start); elapsed < 200*time.Millisecond || elapsed > 2*time.Second { + t.Errorf("leave took %v, want about the 200ms removal timeout", elapsed) + } +} + +func TestLeaveLockRingSkipsALoneFiler(t *testing.T) { + fs := newLeaveTestFilerServer(5*time.Second, "filer1:18888") + + start := time.Now() + fs.leaveLockRing(5 * time.Second) + if elapsed := time.Since(start); elapsed > time.Second { + t.Errorf("a lone filer waited %v to leave a ring with no peer to take its keys", elapsed) + } +} + +func TestLeaveLockRingIsBoundedBySlowLockTransfers(t *testing.T) { + fs := newLeaveTestFilerServer(5*time.Second, "filer1:18888", "filer2:18888") + transferring := make(chan struct{}) + release := make(chan struct{}) + t.Cleanup(func() { close(release) }) + fs.filer.Dlm.LockRing.SetTakeSnapshotCallback(func([]pb.ServerAddress) { + close(transferring) + <-release + }) + go fs.filer.Dlm.LockRing.SetSnapshot([]pb.ServerAddress{"filer2:18888"}, 2) + <-transferring + + start := time.Now() + fs.leaveLockRingWithin(5*time.Second, 300*time.Millisecond) + if elapsed := time.Since(start); elapsed > 2*time.Second { + t.Errorf("leave blocked %v behind lock transfers, want the 300ms budget", elapsed) + } +} diff --git a/weed/server/master_grpc_server.go b/weed/server/master_grpc_server.go index 976435545..058ecacd5 100644 --- a/weed/server/master_grpc_server.go +++ b/weed/server/master_grpc_server.go @@ -422,7 +422,7 @@ func (ms *MasterServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServ glog.V(1).Infof("Cluster: %s node %s added to group '%s'", req.ClientType, peerAddress, req.FilerGroup) ms.broadcastToClients(update) } - if req.ClientType == cluster.FilerType { + if req.ClientType == cluster.FilerType && !req.LeaveLockRing { ms.LockRingManager.AddServer(cluster.FilerGroupName(req.FilerGroup), peerAddress) } if req.ClientType == cluster.MasterType { @@ -484,7 +484,7 @@ func (ms *MasterServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServ go func() { for { - _, err := stream.Recv() + message, err := stream.Recv() if err != nil { glog.V(2).Infof("- client %v: %v", clientName, err) go func() { @@ -496,6 +496,7 @@ func (ms *MasterServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServ close(stopChan) return } + ms.onKeepConnectedMessage(req, peerAddress, message) } }() @@ -532,6 +533,13 @@ func (ms *MasterServer) KeepConnected(stream master_pb.Seaweed_KeepConnectedServ } +func (ms *MasterServer) onKeepConnectedMessage(registered *master_pb.KeepConnectedRequest, peerAddress pb.ServerAddress, message *master_pb.KeepConnectedRequest) { + if registered.ClientType == cluster.FilerType && message.LeaveLockRing { + glog.V(0).Infof("LockRing: filer %s leaving group '%s'", peerAddress, registered.FilerGroup) + ms.LockRingManager.RemoveServer(cluster.FilerGroupName(registered.FilerGroup), peerAddress) + } +} + func (ms *MasterServer) initialLockRingUpdate(clientType string, filerGroup string) *master_pb.KeepConnectedResponse { if ms.LockRingManager == nil { return nil diff --git a/weed/server/master_grpc_server_test.go b/weed/server/master_grpc_server_test.go index 7822078a4..b692c14b0 100644 --- a/weed/server/master_grpc_server_test.go +++ b/weed/server/master_grpc_server_test.go @@ -27,6 +27,29 @@ func TestInitialLockRingUpdateReturnsLastBroadcastForFilers(t *testing.T) { assert.Greater(t, resp.LockRingUpdate.Version, int64(0)) } +func TestLeaveLockRingDropsOnlyTheLeavingFiler(t *testing.T) { + ms := &MasterServer{ + LockRingManager: cluster.NewLockRingManager(nil), + } + ms.LockRingManager.AddServer("group-a", "filer1:8888") + ms.LockRingManager.AddServer("group-a", "filer2:8888") + ms.LockRingManager.FlushPending("group-a") + + registered := &master_pb.KeepConnectedRequest{ClientType: cluster.FilerType, FilerGroup: "group-a"} + ms.onKeepConnectedMessage(registered, "filer1:8888", &master_pb.KeepConnectedRequest{}) + ms.LockRingManager.FlushPending("group-a") + assert.ElementsMatch(t, []string{"filer1:8888", "filer2:8888"}, ms.LockRingManager.GetServers("group-a")) + + ms.onKeepConnectedMessage(registered, "filer1:8888", &master_pb.KeepConnectedRequest{LeaveLockRing: true}) + ms.LockRingManager.FlushPending("group-a") + assert.ElementsMatch(t, []string{"filer2:8888"}, ms.LockRingManager.GetServers("group-a")) + + s3 := &master_pb.KeepConnectedRequest{ClientType: cluster.S3Type, FilerGroup: "group-a"} + ms.onKeepConnectedMessage(s3, "filer2:8888", &master_pb.KeepConnectedRequest{LeaveLockRing: true}) + ms.LockRingManager.FlushPending("group-a") + assert.ElementsMatch(t, []string{"filer2:8888"}, ms.LockRingManager.GetServers("group-a")) +} + func TestInitialLockRingUpdateSkipsNonFilers(t *testing.T) { ms := &MasterServer{ LockRingManager: cluster.NewLockRingManager(nil), diff --git a/weed/wdclient/masterclient.go b/weed/wdclient/masterclient.go index 6d638b6d3..e5ab8075a 100644 --- a/weed/wdclient/masterclient.go +++ b/weed/wdclient/masterclient.go @@ -161,6 +161,12 @@ type MasterClient struct { OnLockRingUpdateLock sync.RWMutex OnMasterChange func(previous, current pb.ServerAddress) OnMasterChangeLock sync.RWMutex + + // streamLock publishes the stream only after its registration send, so a + // leave sent afterwards never races that send. + streamLock sync.Mutex + stream master_pb.Seaweed_KeepConnectedClient + leavingLockRing bool } func NewMasterClient(grpcDialOption grpc.DialOption, filerGroup string, clientType string, clientHost pb.ServerAddress, clientDataCenter string, rack string, masters pb.ServerDiscovery) *MasterClient { @@ -248,18 +254,32 @@ func (mc *MasterClient) tryConnectToMaster(ctx context.Context, master pb.Server } glog.V(1).Infof("%s.%s masterClient gRPC stream established to %s in %v", mc.FilerGroup, mc.clientType, master, time.Since(connectStartTime)) - if err = stream.Send(&master_pb.KeepConnectedRequest{ + mc.streamLock.Lock() + err = stream.Send(&master_pb.KeepConnectedRequest{ FilerGroup: mc.FilerGroup, DataCenter: mc.GetDataCenter(), Rack: mc.rack, ClientType: mc.clientType, ClientAddress: string(mc.clientHost), Version: version.Version(), - }); err != nil { + LeaveLockRing: mc.leavingLockRing, + }) + if err == nil { + mc.stream = stream + } + mc.streamLock.Unlock() + if err != nil { glog.V(0).Infof("%s.%s masterClient failed to send to %s: %v", mc.FilerGroup, mc.clientType, master, err) stats.MasterClientConnectCounter.WithLabelValues(stats.FailedToSend).Inc() return err } + defer func() { + mc.streamLock.Lock() + if mc.stream == stream { + mc.stream = nil + } + mc.streamLock.Unlock() + }() glog.V(1).Infof("%s.%s masterClient Connected to %v", mc.FilerGroup, mc.clientType, master) resp, err := stream.Recv() @@ -370,6 +390,19 @@ func (mc *MasterClient) tryConnectToMaster(ctx context.Context, master pb.Server return nextHintedLeader } +// LeaveLockRing asks the master to drop this client from the lock ring while it +// stays connected, and keeps it out of the ring on any later reconnect. +func (mc *MasterClient) LeaveLockRing() error { + mc.streamLock.Lock() + mc.leavingLockRing = true + stream := mc.stream + mc.streamLock.Unlock() + if stream == nil { + return nil + } + return stream.Send(&master_pb.KeepConnectedRequest{LeaveLockRing: true}) +} + // addedVids indexes added ids that are also being removed. A volume moved // between a server's disks is reported both ways in one message, and the server // still has it -- acting on the removal would drop a good location, whichever diff --git a/weed/wdclient/masterclient_leave_test.go b/weed/wdclient/masterclient_leave_test.go new file mode 100644 index 000000000..18d5e47c9 --- /dev/null +++ b/weed/wdclient/masterclient_leave_test.go @@ -0,0 +1,108 @@ +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) + } +}