From 52fb9f93ff4bdff2ee3f9f787424971fb62fb1cf Mon Sep 17 00:00:00 2001 From: Peter Dodd Date: Sat, 3 Oct 2026 02:10:43 +0100 Subject: [PATCH] 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 --- weed/s3api/s3api_server.go | 7 +- weed/wdclient/filer_client.go | 97 +++++++++++++++- weed/wdclient/filer_client_test.go | 176 +++++++++++++++++++++++++++++ 3 files changed, 273 insertions(+), 7 deletions(-) diff --git a/weed/s3api/s3api_server.go b/weed/s3api/s3api_server.go index aea2a558a..48886f199 100644 --- a/weed/s3api/s3api_server.go +++ b/weed/s3api/s3api_server.go @@ -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) diff --git a/weed/wdclient/filer_client.go b/weed/wdclient/filer_client.go index d4f150ddc..957be5f0d 100644 --- a/weed/wdclient/filer_client.go +++ b/weed/wdclient/filer_client.go @@ -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{}{} diff --git a/weed/wdclient/filer_client_test.go b/weed/wdclient/filer_client_test.go index 1ab8b9a87..b7c589ecc 100644 --- a/weed/wdclient/filer_client_test.go +++ b/weed/wdclient/filer_client_test.go @@ -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) + } +}