mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 22:41:56 +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.
142 lines
5.2 KiB
Go
142 lines
5.2 KiB
Go
package weed_server
|
|
|
|
import (
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/cluster"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
func TestInitialLockRingUpdateReturnsLastBroadcastForFilers(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")
|
|
|
|
resp := ms.initialLockRingUpdate(cluster.FilerType, "group-a")
|
|
require.NotNil(t, resp)
|
|
require.NotNil(t, resp.LockRingUpdate)
|
|
assert.Equal(t, "group-a", resp.LockRingUpdate.FilerGroup)
|
|
assert.ElementsMatch(t, []string{"filer1:8888", "filer2:8888"}, resp.LockRingUpdate.Servers)
|
|
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),
|
|
}
|
|
|
|
ms.LockRingManager.AddServer("group-a", "filer1:8888")
|
|
ms.LockRingManager.FlushPending("group-a")
|
|
|
|
assert.Nil(t, ms.initialLockRingUpdate(cluster.BrokerType, "group-a"))
|
|
}
|
|
|
|
// TestReconnectedClientSurvivesOldHandlerCleanup covers a client reconnecting
|
|
// KeepConnected under the same name before the old handler has exited: the old
|
|
// handler's deferred cleanup must not close the channel the reconnected stream
|
|
// registered, or that stream reads nil from its closed channel forever and
|
|
// floods the client with empty responses.
|
|
func TestReconnectedClientSurvivesOldHandlerCleanup(t *testing.T) {
|
|
ms := &MasterServer{
|
|
clientChans: make(map[string]chan *master_pb.KeepConnectedResponse),
|
|
}
|
|
|
|
clientName, oldChan := ms.addClient("", cluster.MasterType, "peer:19333")
|
|
_, newChan := ms.addClient("", cluster.MasterType, "peer:19333")
|
|
|
|
// the old handler's deferred cleanup runs after the reconnect registered
|
|
ms.deleteClient(clientName, oldChan)
|
|
|
|
// the old channel is closed so its drain goroutine can exit
|
|
select {
|
|
case _, ok := <-oldChan:
|
|
assert.False(t, ok, "old channel should be closed")
|
|
default:
|
|
t.Fatal("old channel left open")
|
|
}
|
|
|
|
// the reconnected stream's channel is untouched and still receives broadcasts
|
|
ms.broadcastToClients(&master_pb.KeepConnectedResponse{
|
|
VolumeLocation: &master_pb.VolumeLocation{Url: "volume-a:8080"},
|
|
})
|
|
select {
|
|
case message, ok := <-newChan:
|
|
require.True(t, ok, "reconnected client's channel was closed by the old handler's cleanup")
|
|
require.NotNil(t, message.GetVolumeLocation())
|
|
default:
|
|
t.Fatal("no broadcast reached the reconnected client")
|
|
}
|
|
|
|
// the reconnected handler's own cleanup still removes the registration
|
|
ms.deleteClient(clientName, newChan)
|
|
ms.clientChansLock.RLock()
|
|
_, found := ms.clientChans[clientName]
|
|
ms.clientChansLock.RUnlock()
|
|
assert.False(t, found)
|
|
}
|
|
|
|
// TestBroadcastVolumeLocationsToClients verifies grown volume locations are sent to registered clients.
|
|
func TestBroadcastVolumeLocationsToClients(t *testing.T) {
|
|
clientChan := make(chan *master_pb.KeepConnectedResponse, 2)
|
|
ms := &MasterServer{
|
|
clientChans: map[string]chan *master_pb.KeepConnectedResponse{
|
|
"default.filer@127.0.0.1:8888": clientChan,
|
|
},
|
|
}
|
|
|
|
ms.broadcastVolumeLocationsToClients([]*master_pb.VolumeLocation{
|
|
{Url: "volume-a:8080", NewVids: []uint32{7}},
|
|
{Url: "volume-b:8080", NewVids: []uint32{8}},
|
|
})
|
|
|
|
var first *master_pb.KeepConnectedResponse
|
|
select {
|
|
case first = <-clientChan:
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("timed out waiting for first broadcast")
|
|
}
|
|
require.NotNil(t, first.GetVolumeLocation())
|
|
assert.Equal(t, []uint32{7}, first.GetVolumeLocation().GetNewVids())
|
|
assert.Equal(t, "volume-a:8080", first.GetVolumeLocation().GetUrl())
|
|
|
|
var second *master_pb.KeepConnectedResponse
|
|
select {
|
|
case second = <-clientChan:
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("timed out waiting for second broadcast")
|
|
}
|
|
require.NotNil(t, second.GetVolumeLocation())
|
|
assert.Equal(t, []uint32{8}, second.GetVolumeLocation().GetNewVids())
|
|
assert.Equal(t, "volume-b:8080", second.GetVolumeLocation().GetUrl())
|
|
}
|