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