s3: track filer joins and leaves pushed by the master (#11563)

* fix(s3): track filer joins and leaves pushed by the master

The S3 FilerClient replaced its -filer seed with a master snapshot of
filer IPs at boot and refreshed it only every 5 minutes. A rolling
restart replaces every filer well inside that window, leaving S3
servers with only dead addresses and failing every write until the
next poll.

Apply the master's ClusterNodeUpdate pushes to the filer list as they
arrive, keeping the poll as a backstop. The last filer is never
removed, and a poll snapshot requested before a push was applied is
discarded rather than overwriting newer membership.

* Defer last-filer leaves; bump the generation only on real changes

* fix(s3): cancel deferred filer leaves on rejoin and on discovery

A deferred last-filer leave outlived the filer it was recorded for: a
rejoin at the same address looked like a duplicate add, and a discovery
snapshot left the entry behind. The next join then removed a live
filer until the following poll.

A join now cancels any deferred leave for its address, and an applied
snapshot clears them, since it is the master's current membership.

* Bump the push generation when a rejoin cancels a deferred leave

---------

Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
This commit is contained in:
Peter DoddandChris Lu authored and GitHub committed 2026-10-03 09:10:43 +08:00
1 parent 0ff7794c54
commit 52fb9f93ff
3 files changed
+273 -7

No files matched your search

+4 -3
View File
@@ -229,14 +229,15 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl
objectWriteLockClient.ResetRing()
})
}
// Start the master client connection loop - required for GetMaster() to work
go masterClient.KeepConnectedToMaster(context.Background())
filerClient = wdclient.NewFilerClient(option.Filers, option.GrpcDialOption, option.DataCenter, &wdclient.FilerClientOption{
MasterClient: masterClient,
FilerGroup: option.FilerGroup,
DiscoveryInterval: 5 * time.Minute,
})
masterClient.SetOnPeerUpdateFn(filerClient.OnPeerUpdate)
// Start the master client connection loop - required for GetMaster() to work
go masterClient.KeepConnectedToMaster(context.Background())
glog.V(1).Infof("S3 API initialized FilerClient with %d filer(s) and discovery enabled (group: %s, masters: %v)",
len(option.Filers), option.FilerGroup, option.Masters)
+93 -4
View File
@@ -17,6 +17,7 @@ import (
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
"github.com/seaweedfs/seaweedfs/weed/util"
)
@@ -42,9 +43,11 @@ type filerHealth struct {
type FilerClient struct {
*vidMapClient
filerAddresses []pb.ServerAddress
filerAddressesMu sync.RWMutex // Protects filerAddresses and filerHealth
filerIndex int32 // atomic: current filer index for round-robin
filerHealth []*filerHealth // health status per filer (same order as filerAddresses)
filerAddressesMu sync.RWMutex // Protects filerAddresses and filerHealth
filerIndex int32 // atomic: current filer index for round-robin
filerHealth []*filerHealth // health status per filer (same order as filerAddresses)
peerUpdates uint64 // pushed filer updates applied; protected by filerAddressesMu
deferredLeaves map[pb.ServerAddress]struct{} // leaves suppressed while the filer was the last known; protected by filerAddressesMu
grpcDialOption grpc.DialOption
urlPreference UrlPreference
grpcTimeout time.Duration
@@ -335,6 +338,8 @@ func (fc *FilerClient) refreshFilerList() {
return
}
generation := fc.peerUpdateGeneration()
// Query master for filers in our group
updates := cluster.ListExistingPeerUpdates(currentMaster, fc.grpcDialOption, fc.filerGroup, cluster.FilerType)
@@ -358,7 +363,26 @@ func (fc *FilerClient) refreshFilerList() {
return
}
fc.applyDiscoveredFilers(discoveredFilers)
fc.applyDiscoverySnapshot(discoveredFilers, generation)
}
func (fc *FilerClient) peerUpdateGeneration() uint64 {
fc.filerAddressesMu.RLock()
defer fc.filerAddressesMu.RUnlock()
return fc.peerUpdates
}
// applyDiscoverySnapshot drops a snapshot requested before a pushed update was
// applied: the push is newer, and the next poll reconciles anything it missed.
func (fc *FilerClient) applyDiscoverySnapshot(discoveredFilers map[pb.ServerAddress]struct{}, generation uint64) {
fc.filerAddressesMu.Lock()
defer fc.filerAddressesMu.Unlock()
if fc.peerUpdates != generation {
glog.V(1).Infof("FilerClient: discarding discovery snapshot for group '%s' superseded by pushed updates", fc.filerGroup)
return
}
fc.deferredLeaves = nil
fc.applyDiscoveredFilersLocked(discoveredFilers)
}
// applyDiscoveredFilers treats the master snapshot as authoritative: survivors
@@ -368,7 +392,72 @@ func (fc *FilerClient) refreshFilerList() {
func (fc *FilerClient) applyDiscoveredFilers(discoveredFilers map[pb.ServerAddress]struct{}) {
fc.filerAddressesMu.Lock()
defer fc.filerAddressesMu.Unlock()
fc.applyDiscoveredFilersLocked(discoveredFilers)
}
// OnPeerUpdate applies a filer join/leave pushed by the master, so the list
// tracks a rolling restart instead of waiting for the next discovery poll. A
// leave that would empty the list is deferred until another filer exists to
// take over: an empty list fails every request, while a stale one still
// recovers through the poll.
func (fc *FilerClient) OnPeerUpdate(update *master_pb.ClusterNodeUpdate, _ time.Time) {
if update.NodeType != cluster.FilerType || update.Address == "" {
return
}
addr := pb.ServerAddress(update.Address)
fc.filerAddressesMu.Lock()
defer fc.filerAddressesMu.Unlock()
filers := make(map[pb.ServerAddress]struct{}, len(fc.filerAddresses)+1)
for _, f := range fc.filerAddresses {
filers[f] = struct{}{}
}
changed := false
if update.IsAdd {
// A rejoin cancels its deferred leave; bumping the generation keeps an
// in-flight snapshot that lacks the rejoined filer from pruning it.
if _, ok := fc.deferredLeaves[addr]; ok {
delete(fc.deferredLeaves, addr)
changed = true
}
if _, ok := filers[addr]; !ok {
filers[addr] = struct{}{}
changed = true
}
} else {
if _, ok := filers[addr]; ok {
delete(filers, addr)
changed = true
if len(filers) == 0 {
filers[addr] = struct{}{}
if fc.deferredLeaves == nil {
fc.deferredLeaves = make(map[pb.ServerAddress]struct{})
}
fc.deferredLeaves[addr] = struct{}{}
}
}
}
// A deferred leave can be honored once another filer exists to take over.
for gone := range fc.deferredLeaves {
if _, ok := filers[gone]; ok && len(filers) <= 1 {
continue
}
delete(filers, gone)
delete(fc.deferredLeaves, gone)
changed = true
}
if !changed {
return
}
fc.peerUpdates++
fc.applyDiscoveredFilersLocked(filers)
}
func (fc *FilerClient) applyDiscoveredFilersLocked(discoveredFilers map[pb.ServerAddress]struct{}) {
existingFilers := make(map[pb.ServerAddress]struct{}, len(fc.filerAddresses))
for _, f := range fc.filerAddresses {
existingFilers[f] = struct{}{}
+176
View File
@@ -3,8 +3,11 @@ package wdclient
import (
"sync/atomic"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/cluster"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
)
func newTestFilerClient(addrs ...pb.ServerAddress) *FilerClient {
@@ -100,3 +103,176 @@ func TestApplyDiscoveredFilersNoChangeIsNoop(t *testing.T) {
t.Errorf("no-op refresh should not reallocate health entries")
}
}
func filerUpdate(addr pb.ServerAddress, isAdd bool) *master_pb.ClusterNodeUpdate {
return &master_pb.ClusterNodeUpdate{NodeType: cluster.FilerType, Address: string(addr), IsAdd: isAdd}
}
// A rolling restart replaces every filer within one discovery interval; the
// pushed updates must keep the list live without waiting for the next poll.
func TestOnPeerUpdateTracksRollingReplacement(t *testing.T) {
old1 := pb.ServerAddress("10.0.0.1:8888")
old2 := pb.ServerAddress("10.0.0.2:8888")
new1 := pb.ServerAddress("10.0.1.1:8888")
new2 := pb.ServerAddress("10.0.1.2:8888")
fc := newTestFilerClient(old1, old2)
fc.OnPeerUpdate(filerUpdate(old1, false), time.Now())
fc.OnPeerUpdate(filerUpdate(new1, true), time.Now())
fc.OnPeerUpdate(filerUpdate(old2, false), time.Now())
fc.OnPeerUpdate(filerUpdate(new2, true), time.Now())
got := filerAddressList(fc)
if len(got) != 2 || got[0] != new1 || got[1] != new2 {
t.Fatalf("expected [%s %s], got %v", new1, new2, got)
}
}
func TestOnPeerUpdateKeepsLastFiler(t *testing.T) {
only := pb.ServerAddress("10.0.0.1:8888")
fc := newTestFilerClient(only)
fc.OnPeerUpdate(filerUpdate(only, false), time.Now())
if got := filerAddressList(fc); len(got) != 1 || got[0] != only {
t.Fatalf("removing the last filer must keep it, got %v", got)
}
}
func TestOnPeerUpdateIgnoresOtherNodeTypes(t *testing.T) {
a := pb.ServerAddress("10.0.0.1:8888")
fc := newTestFilerClient(a)
fc.OnPeerUpdate(&master_pb.ClusterNodeUpdate{NodeType: cluster.S3Type, Address: "10.0.0.9:18333", IsAdd: true}, time.Now())
if got := filerAddressList(fc); len(got) != 1 || got[0] != a {
t.Fatalf("non-filer update changed the list: %v", got)
}
}
func TestOnPeerUpdateDuplicateAddKeepsHealth(t *testing.T) {
a := pb.ServerAddress("10.0.0.1:8888")
fc := newTestFilerClient(a)
atomic.StoreInt32(&fc.filerHealth[0].failureCount, 2)
fc.OnPeerUpdate(filerUpdate(a, true), time.Now())
if len(fc.filerHealth) != 1 || atomic.LoadInt32(&fc.filerHealth[0].failureCount) != 2 {
t.Fatalf("re-announced filer lost its health state")
}
}
func TestDiscoverySnapshotTakenBeforePushIsDiscarded(t *testing.T) {
old := pb.ServerAddress("10.0.0.1:8888")
joined := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(old)
generation := fc.peerUpdateGeneration()
fc.OnPeerUpdate(filerUpdate(joined, true), time.Now())
fc.OnPeerUpdate(filerUpdate(old, false), time.Now())
fc.applyDiscoverySnapshot(map[pb.ServerAddress]struct{}{old: {}}, generation)
if got := filerAddressList(fc); len(got) != 1 || got[0] != joined {
t.Fatalf("stale snapshot overwrote pushed membership: %v", got)
}
}
func TestDiscoverySnapshotWithoutInterveningPushIsApplied(t *testing.T) {
old := pb.ServerAddress("10.0.0.1:8888")
replacement := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(old)
fc.applyDiscoverySnapshot(map[pb.ServerAddress]struct{}{replacement: {}}, fc.peerUpdateGeneration())
if got := filerAddressList(fc); len(got) != 1 || got[0] != replacement {
t.Fatalf("expected snapshot to replace list, got %v", got)
}
}
// A leave suppressed to keep the last filer is deferred, then honored as soon
// as a replacement joins, so the departed address stops being a candidate.
func TestOnPeerUpdateDeferredLeaveFlushesOnJoin(t *testing.T) {
old := pb.ServerAddress("10.0.0.1:8888")
joined := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(old)
fc.OnPeerUpdate(filerUpdate(old, false), time.Now())
if got := filerAddressList(fc); len(got) != 1 || got[0] != old {
t.Fatalf("leave of the last filer must be deferred, got %v", got)
}
if len(fc.deferredLeaves) != 1 {
t.Fatalf("expected the leave to be deferred, got %v", fc.deferredLeaves)
}
fc.OnPeerUpdate(filerUpdate(joined, true), time.Now())
if got := filerAddressList(fc); len(got) != 1 || got[0] != joined {
t.Fatalf("deferred leave should flush on join, got %v", got)
}
if len(fc.deferredLeaves) != 0 {
t.Fatalf("deferred leaves should be empty after flush, got %v", fc.deferredLeaves)
}
}
// A pushed add for an already-known filer changes nothing and must not bump
// the generation that guards an in-flight discovery snapshot.
func TestOnPeerUpdateNoopDoesNotDiscardSnapshot(t *testing.T) {
old := pb.ServerAddress("10.0.0.1:8888")
replacement := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(old)
generation := fc.peerUpdateGeneration()
fc.OnPeerUpdate(filerUpdate(old, true), time.Now())
fc.applyDiscoverySnapshot(map[pb.ServerAddress]struct{}{replacement: {}}, generation)
if got := filerAddressList(fc); len(got) != 1 || got[0] != replacement {
t.Fatalf("no-op push should not discard the snapshot, got %v", got)
}
}
func TestOnPeerUpdateRejoinCancelsDeferredLeave(t *testing.T) {
restarted := pb.ServerAddress("10.0.0.1:8888")
joined := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(restarted)
fc.OnPeerUpdate(filerUpdate(restarted, false), time.Now())
fc.OnPeerUpdate(filerUpdate(restarted, true), time.Now())
fc.OnPeerUpdate(filerUpdate(joined, true), time.Now())
if got := filerAddressList(fc); len(got) != 2 || got[0] != restarted || got[1] != joined {
t.Fatalf("rejoined filer was dropped by its stale deferred leave: %v", got)
}
}
func TestDiscoverySnapshotClearsDeferredLeaves(t *testing.T) {
departed := pb.ServerAddress("10.0.0.1:8888")
replacement := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(departed)
fc.OnPeerUpdate(filerUpdate(departed, false), time.Now())
fc.applyDiscoverySnapshot(map[pb.ServerAddress]struct{}{replacement: {}}, fc.peerUpdateGeneration())
fc.OnPeerUpdate(filerUpdate(departed, true), time.Now())
if got := filerAddressList(fc); len(got) != 2 || got[0] != replacement || got[1] != departed {
t.Fatalf("rejoined filer was dropped by a deferred leave the poll should have cleared: %v", got)
}
}
// A poll that started after a deferred leave but before the filer rejoined can
// return a snapshot lacking the rejoined filer; the rejoin must bump the
// generation so that snapshot is discarded.
func TestOnPeerUpdateRejoinBumpsGeneration(t *testing.T) {
old := pb.ServerAddress("10.0.0.1:8888")
other := pb.ServerAddress("10.0.1.1:8888")
fc := newTestFilerClient(old)
fc.OnPeerUpdate(filerUpdate(old, false), time.Now())
generation := fc.peerUpdateGeneration()
fc.OnPeerUpdate(filerUpdate(old, true), time.Now())
fc.applyDiscoverySnapshot(map[pb.ServerAddress]struct{}{other: {}}, generation)
if got := filerAddressList(fc); len(got) != 1 || got[0] != old {
t.Fatalf("rejoined filer was pruned by an in-flight snapshot: %v", got)
}
}