From dfd3732c6bbd7446ffb538b509dde69abd1e5b37 Mon Sep 17 00:00:00 2001 From: Forrest Marshall Date: Mon, 6 Dec 2021 14:37:46 -0800 Subject: [PATCH] improve concurrent watcher registration perf --- lib/cache/cache.go | 4 +- lib/services/fanout.go | 117 +++++++++++++++++++++++++++++++++--- lib/services/fanout_test.go | 71 ++++++++++++++++++++++ 3 files changed, 180 insertions(+), 12 deletions(-) diff --git a/lib/cache/cache.go b/lib/cache/cache.go index cf9603d9748..5d44f660d98 100644 --- a/lib/cache/cache.go +++ b/lib/cache/cache.go @@ -350,7 +350,7 @@ type Cache struct { webSessionCache types.WebSessionInterface webTokenCache types.WebTokenInterface windowsDesktopsCache services.WindowsDesktops - eventsFanout *services.Fanout + eventsFanout *services.FanoutSet // closed indicates that the cache has been closed closed *atomic.Bool @@ -624,7 +624,7 @@ func New(config Config) (*Cache, error) { webSessionCache: local.NewIdentityService(wrapper).WebSessions(), webTokenCache: local.NewIdentityService(wrapper).WebTokens(), windowsDesktopsCache: local.NewWindowsDesktopService(wrapper), - eventsFanout: services.NewFanout(), + eventsFanout: services.NewFanoutSet(), Entry: log.WithFields(log.Fields{ trace.Component: config.Component, }), diff --git a/lib/services/fanout.go b/lib/services/fanout.go index b4699d40da7..313c8586869 100644 --- a/lib/services/fanout.go +++ b/lib/services/fanout.go @@ -23,6 +23,8 @@ import ( "github.com/gravitational/teleport/api/types" "github.com/gravitational/trace" + + "go.uber.org/atomic" ) const defaultQueueSize = 64 @@ -201,7 +203,7 @@ func (f *Fanout) Emit(events ...types.Event) { func (f *Fanout) Reset() { f.mu.Lock() defer f.mu.Unlock() - f.closeWatchers() + f.closeWatchersAsync() f.init = false } @@ -210,19 +212,31 @@ func (f *Fanout) Reset() { func (f *Fanout) Close() { f.mu.Lock() defer f.mu.Unlock() - f.closeWatchers() + f.closeWatchersAsync() f.closed = true } -func (f *Fanout) closeWatchers() { - for _, entries := range f.watchers { - for _, entry := range entries { - entry.watcher.cancel() - } - } - // watcher map was potentially quite large, so - // relenguish that memory. +// closeWatchersAsync moves ownership of the watcher mapping to a background goroutine +// for asynchronous cancellation and sets up a new empty mapping. +func (f *Fanout) closeWatchersAsync() { + watchersToClose := f.watchers f.watchers = make(map[string][]fanoutEntry) + // goroutines run with a "happens after" releationship to the + // expressions that create them. since we move ownership of the + // old watcher mapping prior to spawning this goroutine, we are + // "safe" to modify it without worrying about locking. because + // we don't continue to hold the lock in the foreground goroutine, + // this fanout instance may permit new events/registrations/inits/resets + // while the old watchers are still being closed. this is fine, since + // the aformentioned move guarantees that these old watchers aren't + // going to observe any of the new state transitions. + go func() { + for _, entries := range watchersToClose { + for _, entry := range entries { + entry.watcher.cancel() + } + } + }() } func (f *Fanout) addWatcher(w *fanoutWatcher) { @@ -354,3 +368,86 @@ func (w *fanoutWatcher) Error() error { return nil } } + +// fanoutSetSize is the number of members in a fanout set. selected based on some experimentation with +// the FanoutSetRegistration benchmark. This value keeps 100K concurrent registrations well under 1s. +const fanoutSetSize = 128 + +// FanoutSet is a collection of separate Fanout instances. It exposes an identical API, and "load balances" +// watcher registration across the enclosed instances. In very large clusters it is possible for tens of +// thousands of nodes to simultaneously request watchers. This can cause serious contention issues. FanoutSet is +// a simple but effective solution to that problem. +type FanoutSet struct { + // rw mutex is used to ensure that Close and Reset operations are exclusive, + // since these operations close watchers. Enforcing this property isn't strictly + // necessary, but it prevents a scenario where watchers might observe a reset/close, + // attempt re-registration, and observe the *same* reset/close again. This isn't + // necessarily a problem, but it might confuse attempts to debug other event-system + // issues, so we choose to avoid it. + rw sync.RWMutex + counter *atomic.Uint64 + members []*Fanout +} + +// NewFanoutSet creates a new FanoutSet instance in an uninitialized +// state. Until initialized, watchers will be queued but no +// events will be sent. +func NewFanoutSet() *FanoutSet { + members := make([]*Fanout, 0, fanoutSetSize) + for i := 0; i < fanoutSetSize; i++ { + members = append(members, NewFanout()) + } + return &FanoutSet{ + counter: atomic.NewUint64(0), + members: members, + } +} + +// NewWatcher attaches a new watcher to a fanout instance. +func (s *FanoutSet) NewWatcher(ctx context.Context, watch types.Watch) (types.Watcher, error) { + s.rw.RLock() // see field-level docks for locking model + defer s.rw.RUnlock() + fi := int(s.counter.Inc() % uint64(len(s.members))) + return s.members[fi].NewWatcher(ctx, watch) +} + +// SetInit sets the Fanout instances into an initialized state, sending OpInit +// events to any watchers which were added prior to initialization. +func (s *FanoutSet) SetInit() { + s.rw.RLock() // see field-level docks for locking model + defer s.rw.RUnlock() + for _, f := range s.members { + f.SetInit() + } +} + +// Emit broadcasts events to all matching watchers that have been attached +// to this fanout set. +func (s *FanoutSet) Emit(events ...types.Event) { + s.rw.RLock() // see field-level docks for locking model + defer s.rw.RUnlock() + for _, f := range s.members { + f.Emit(events...) + } +} + +// Reset closes all attached watchers and places the fanout instances +// into an uninitialized state. Reset may be called on an uninitialized +// fanout set to remove "queued" watchers. +func (s *FanoutSet) Reset() { + s.rw.Lock() // see field-level docks for locking model + defer s.rw.Unlock() + for _, f := range s.members { + f.Reset() + } +} + +// Close permanently closes the fanout. Existing watchers will be +// closed and no new watchers will be added. +func (s *FanoutSet) Close() { + s.rw.Lock() // see field-level docks for locking model + defer s.rw.Unlock() + for _, f := range s.members { + f.Close() + } +} diff --git a/lib/services/fanout_test.go b/lib/services/fanout_test.go index f9d9fe4068c..4c774ad5d7d 100644 --- a/lib/services/fanout_test.go +++ b/lib/services/fanout_test.go @@ -18,6 +18,7 @@ package services import ( "context" + "sync" "testing" "time" @@ -70,3 +71,73 @@ func TestFanoutInit(t *testing.T) { default: } } + +/* +goos: linux +goarch: amd64 +pkg: github.com/gravitational/teleport/lib/services +cpu: Intel(R) Core(TM) i9-10885H CPU @ 2.40GHz +BenchmarkFanoutRegistration-16 1 118856478045 ns/op +*/ +// NOTE: this benchmark exists primarily to "contrast" with the set registration +// benchmark below, and demonstrate why the set-based strategy is necessary. +func BenchmarkFanoutRegistration(b *testing.B) { + const iterations = 100_000 + ctx := context.Background() + + for n := 0; n < b.N; n++ { + f := NewFanout() + f.SetInit() + + var wg sync.WaitGroup + + for i := 0; i < iterations; i++ { + wg.Add(1) + go func() { + defer wg.Done() + w, err := f.NewWatcher(ctx, types.Watch{ + Name: "test", + Kinds: []types.WatchKind{{Name: "spam"}, {Name: "eggs"}}, + }) + require.NoError(b, err) + w.Close() + }() + } + + wg.Wait() + } +} + +/* +goos: linux +goarch: amd64 +pkg: github.com/gravitational/teleport/lib/services +cpu: Intel(R) Core(TM) i9-10885H CPU @ 2.40GHz +BenchmarkFanoutSetRegistration-16 3 394211563 ns/op +*/ +func BenchmarkFanoutSetRegistration(b *testing.B) { + const iterations = 100_000 + ctx := context.Background() + + for n := 0; n < b.N; n++ { + f := NewFanoutSet() + f.SetInit() + + var wg sync.WaitGroup + + for i := 0; i < iterations; i++ { + wg.Add(1) + go func() { + defer wg.Done() + w, err := f.NewWatcher(ctx, types.Watch{ + Name: "test", + Kinds: []types.WatchKind{{Name: "spam"}, {Name: "eggs"}}, + }) + require.NoError(b, err) + w.Close() + }() + } + + wg.Wait() + } +}