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