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.
This commit is contained in:
Peter Dodd authored and GitHub committed 2026-10-06 09:53:40 +08:00
1 parent aa5b337716
commit 52c6df3bec
12 files changed
+402 -11

No files matched your search

+3
View File
@@ -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 {
+11
View File
@@ -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()
+6 -3
View File
@@ -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
+41 -2
View File
@@ -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")
}
}
+3
View File
@@ -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 {
+13 -2
View File
@@ -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" +
+59
View File
@@ -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...")
+90
View File
@@ -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)
}
}
+10 -2
View File
@@ -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
+23
View File
@@ -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),
+35 -2
View File
@@ -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
+108
View File
@@ -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)
}
}