From aaee69592725558cff2144aceef43523e22ddb4c Mon Sep 17 00:00:00 2001 From: Forrest <30576607+fspmarshall@users.noreply.github.com> Date: Fri, 3 Mar 2023 09:39:00 -0800 Subject: [PATCH] node hb and watcher scalability improvements (#21490) * node hb and watcher scalability improvements * ensure that node has ample time to hb * fix degraded state * hbv2 test cov --- integration/proxy/proxy_helpers.go | 2 +- lib/cache/cache.go | 2 +- lib/defaults/defaults.go | 2 +- lib/services/watcher.go | 2 +- lib/srv/heartbeatv2.go | 109 ++++++++++++++++-------- lib/srv/heartbeatv2_test.go | 132 +++++++++++++++++++++++++++-- lib/srv/regular/sshserver.go | 1 - lib/utils/retry.go | 17 ++-- 8 files changed, 211 insertions(+), 56 deletions(-) diff --git a/integration/proxy/proxy_helpers.go b/integration/proxy/proxy_helpers.go index 3893572f9b7..2278c9ec190 100644 --- a/integration/proxy/proxy_helpers.go +++ b/integration/proxy/proxy_helpers.go @@ -204,7 +204,7 @@ func (p *Suite) addNodeToLeafCluster(t *testing.T, tunnelNodeHostname string) { func (p *Suite) mustConnectToClusterAndRunSSHCommand(t *testing.T, config helpers.ClientConfig) { const ( - deadline = time.Second * 5 + deadline = time.Second * 20 nextIterWaitTime = time.Millisecond * 100 ) diff --git a/lib/cache/cache.go b/lib/cache/cache.go index 1d1a4583677..bca2aea7101 100644 --- a/lib/cache/cache.go +++ b/lib/cache/cache.go @@ -820,7 +820,7 @@ func New(config Config) (*Cache, error) { // Start the cache. Should only be called once. func (c *Cache) Start() error { retry, err := retryutils.NewLinear(retryutils.LinearConfig{ - First: utils.HalfJitter(c.MaxRetryPeriod / 10), + First: utils.FullJitter(c.MaxRetryPeriod / 10), Step: c.MaxRetryPeriod / 5, Max: c.MaxRetryPeriod, Jitter: retryutils.NewHalfJitter(), diff --git a/lib/defaults/defaults.go b/lib/defaults/defaults.go index 7c833d36815..4eeca393d7e 100644 --- a/lib/defaults/defaults.go +++ b/lib/defaults/defaults.go @@ -307,7 +307,7 @@ const ( // MaxWatcherBackoff is the maximum retry time a watcher should use in // the event of connection issues - MaxWatcherBackoff = time.Minute + MaxWatcherBackoff = 90 * time.Second ) const ( diff --git a/lib/services/watcher.go b/lib/services/watcher.go index ffd8200dd6b..dd80fad0a7e 100644 --- a/lib/services/watcher.go +++ b/lib/services/watcher.go @@ -98,7 +98,7 @@ func (cfg *ResourceWatcherConfig) CheckAndSetDefaults() error { // incl. cfg.CheckAndSetDefaults. func newResourceWatcher(ctx context.Context, collector resourceCollector, cfg ResourceWatcherConfig) (*resourceWatcher, error) { retry, err := retryutils.NewLinear(retryutils.LinearConfig{ - First: utils.HalfJitter(cfg.MaxRetryPeriod / 10), + First: utils.FullJitter(cfg.MaxRetryPeriod / 10), Step: cfg.MaxRetryPeriod / 5, Max: cfg.MaxRetryPeriod, Jitter: retryutils.NewHalfJitter(), diff --git a/lib/srv/heartbeatv2.go b/lib/srv/heartbeatv2.go index 0bccf39aed2..f5cbdc0d1ae 100644 --- a/lib/srv/heartbeatv2.go +++ b/lib/srv/heartbeatv2.go @@ -41,14 +41,14 @@ type SSHServerHeartbeatConfig struct { InventoryHandle inventory.DownstreamHandle // GetServer gets the latest server spec. GetServer func() *types.ServerV2 + + // -- below values are all optional + // Announcer is a fallback used to perform basic upsert-style heartbeats // if the control stream is unavailable. // // DELETE IN: 11.0 (only exists for back-compat with v9 auth servers) Announcer auth.Announcer - - // -- below values are all optional - // OnHeartbeat is a per-attempt callback (optional). OnHeartbeat func(error) // AnnounceInterval is the interval at which heartbeats are attempted (optional). @@ -64,9 +64,6 @@ func (c *SSHServerHeartbeatConfig) Check() error { if c.GetServer == nil { return trace.BadParameter("missing required parameter GetServer for ssh heartbeat") } - if c.Announcer == nil { - return trace.BadParameter("missing required parameter Announcer for ssh heartbeat") - } return nil } @@ -105,9 +102,14 @@ const ( hbv2Start hbv2TestEvent = "hb-start" hbv2Close hbv2TestEvent = "hb-close" - hbv2AnnounceInterval = "hb-announce-interval" + hbv2AnnounceInterval hbv2TestEvent = "hb-announce-interval" - hbv2FallbackBackoff = "hb-fallback-backoff" + hbv2FallbackBackoff hbv2TestEvent = "hb-fallback-backoff" + + hbv2NoFallback hbv2TestEvent = "no-fallback" + + hbv2OnHeartbeatOk = "on-heartbeat-ok" + hbv2OnHeartbeatErr = "on-heartbeat-err" ) // newHeartbeatV2 configures a new HeartbeatV2 instance to wrap a given implementation. @@ -143,10 +145,11 @@ type HeartbeatV2 struct { announceFailed error fallbackFailed error + icsUnavailable error - announce *interval.Interval - poll *interval.Interval - dc *interval.Interval + announce *interval.Interval + poll *interval.Interval + degradedCheck *interval.Interval // fallbackBackoffTime approximately replicate the backoff used by heartbeat V1 when an announce // fails. It can be removed once we remove the fallback announce operation, since control-stream @@ -171,8 +174,9 @@ type heartbeatV2Config struct { // -- below values only used in tests - fallbackBackoff time.Duration - testEvents chan hbv2TestEvent + fallbackBackoff time.Duration + testEvents chan hbv2TestEvent + degradedCheckInterval time.Duration } func (c *heartbeatV2Config) SetDefaults() { @@ -189,6 +193,12 @@ func (c *heartbeatV2Config) SetDefaults() { // only set externally during tests c.fallbackBackoff = time.Minute } + + if c.degradedCheckInterval == 0 { + // a lot of integration tests rely on overriding ServerKeepAliveTTL to modify how + // quickly teleport detects that it is in a degraded state. + c.degradedCheckInterval = apidefaults.ServerKeepAliveTTL() + } } // noSenderErr is used to periodically trigger "degraded state" events when the control @@ -200,6 +210,7 @@ func (h *HeartbeatV2) run() { // so we just allocate something reasonably descriptive once. h.announceFailed = trace.Errorf("control stream heartbeat failed (variant=%T)", h.inner) h.fallbackFailed = trace.Errorf("upsert fallback heartbeat failed (variant=%T)", h.inner) + h.icsUnavailable = trace.Errorf("ics unavailable for heartbeat (variant=%T)", h.inner) // set up interval for forced announcement (i.e. heartbeat even if state is unchanged). h.announce = interval.New(interval.Config{ @@ -223,40 +234,45 @@ func (h *HeartbeatV2) run() { // down. Since we no longer perform keepalives, we instead simply emit an error on this // interval when we don't have a healthy control stream. // TODO(fspmarshall): find a more elegant solution to this problem. - h.dc = interval.New(interval.Config{ - Duration: apidefaults.ServerKeepAliveTTL(), + h.degradedCheck = interval.New(interval.Config{ + Duration: h.degradedCheckInterval, }) - defer h.dc.Stop() + defer h.degradedCheck.Stop() h.testEvent(hbv2Start) defer h.testEvent(hbv2Close) for { // outer loop performs announcement via the fallback method (used for backwards compatibility - // with older auth servers). + // with older auth servers). Not all drivers support fallback. if h.shouldAnnounce { - if time.Now().After(h.fallbackBackoffTime) { - if ok := h.inner.FallbackAnnounce(h.closeContext); ok { - h.testEvent(hbv2FallbackOk) - // reset announce interval and state on successful announce - h.announce.Reset() - h.shouldAnnounce = false - h.onHeartbeat(nil) + if h.inner.SupportsFallback() { + if time.Now().After(h.fallbackBackoffTime) { + if ok := h.inner.FallbackAnnounce(h.closeContext); ok { + h.testEvent(hbv2FallbackOk) + // reset announce interval and state on successful announce + h.announce.Reset() + h.degradedCheck.Reset() + h.shouldAnnounce = false + h.onHeartbeat(nil) - // unblock tests waiting on an announce operation - for _, waiter := range h.announceWaiters { - close(waiter) + // unblock tests waiting on an announce operation + for _, waiter := range h.announceWaiters { + close(waiter) + } + h.announceWaiters = nil + } else { + h.testEvent(hbv2FallbackErr) + // announce failed, enter a backoff state. + h.fallbackBackoffTime = time.Now().Add(utils.SeventhJitter(h.fallbackBackoff)) + h.onHeartbeat(h.fallbackFailed) } - h.announceWaiters = nil } else { - h.testEvent(hbv2FallbackErr) - // announce failed, enter a backoff state. - h.fallbackBackoffTime = time.Now().Add(utils.SeventhJitter(h.fallbackBackoff)) - h.onHeartbeat(h.fallbackFailed) + h.testEvent(hbv2FallbackBackoff) } } else { - h.testEvent(hbv2FallbackBackoff) + h.testEvent(hbv2NoFallback) } } @@ -267,7 +283,7 @@ func (h *HeartbeatV2) run() { case sender := <-h.handle.Sender(): // sender is available, hand off to the primary run loop h.runWithSender(sender) - h.dc.Reset() + h.degradedCheck.Reset() case <-h.announce.Next(): h.testEvent(hbv2AnnounceInterval) h.shouldAnnounce = true @@ -278,8 +294,11 @@ func (h *HeartbeatV2) run() { } else { h.testEvent(hbv2PollSame) } - case <-h.dc.Next(): - if !h.inner.Poll() && !h.shouldAnnounce { + case <-h.degradedCheck.Next(): + if !h.inner.SupportsFallback() || (!h.inner.Poll() && !h.shouldAnnounce) { + // if we don't have fallback and/or aren't planning to hit the fallback + // soon, then we need to emit a heartbeat error in order to inform the + // rest of teleport that we are in a degraded state. h.onHeartbeat(noSenderErr) } case ch := <-h.testAnnounce: @@ -303,6 +322,7 @@ func (h *HeartbeatV2) runWithSender(sender inventory.DownstreamSender) { h.testEvent(hbv2AnnounceOk) // reset announce interval and state on successful announce h.announce.Reset() + h.degradedCheck.Reset() h.shouldAnnounce = false h.onHeartbeat(nil) @@ -332,6 +352,12 @@ func (h *HeartbeatV2) runWithSender(sender inventory.DownstreamSender) { } else { h.testEvent(hbv2PollSame) } + case <-h.degradedCheck.Next(): + if !h.inner.Poll() && !h.shouldAnnounce { + // its been a while since we announced and we are not in a retry/announce + // state now, so clear up any degraded state. + h.onHeartbeat(nil) + } case waiter := <-h.testAnnounce: h.shouldAnnounce = true h.announceWaiters = append(h.announceWaiters, waiter) @@ -378,6 +404,11 @@ func (h *HeartbeatV2) ForceSend(timeout time.Duration) error { } func (h *HeartbeatV2) onHeartbeat(err error) { + if err != nil { + h.testEvent(hbv2OnHeartbeatErr) + } else { + h.testEvent(hbv2OnHeartbeatOk) + } if h.onHeartbeatInner == nil { return } @@ -397,6 +428,8 @@ type heartbeatV2Driver interface { FallbackAnnounce(ctx context.Context) (ok bool) // Announce attempts to heartbeat via the inventory control stream. Announce(ctx context.Context, sender inventory.DownstreamSender) (ok bool) + // SupportsFallback checks if the driver supports fallback. + SupportsFallback() bool } // sshServerHeartbeatV2 is the heartbeatV2 implementation for ssh servers. @@ -413,6 +446,10 @@ func (h *sshServerHeartbeatV2) Poll() (changed bool) { return services.CompareServers(h.getServer(), h.prev) == services.Different } +func (h *sshServerHeartbeatV2) SupportsFallback() bool { + return h.announcer != nil +} + func (h *sshServerHeartbeatV2) FallbackAnnounce(ctx context.Context) (ok bool) { if h.announcer == nil { return false diff --git a/lib/srv/heartbeatv2_test.go b/lib/srv/heartbeatv2_test.go index 7882e0197c8..38e57a50fb6 100644 --- a/lib/srv/heartbeatv2_test.go +++ b/lib/srv/heartbeatv2_test.go @@ -48,6 +48,8 @@ type fakeHeartbeatDriver struct { pollChanged int fallbackErr int announceErr int + + disableFallback bool } func (h *fakeHeartbeatDriver) Poll() (changed bool) { @@ -64,6 +66,9 @@ func (h *fakeHeartbeatDriver) Poll() (changed bool) { func (h *fakeHeartbeatDriver) FallbackAnnounce(ctx context.Context) (ok bool) { h.mu.Lock() defer h.mu.Unlock() + if !h.SupportsFallback() { + panic("FallbackAnnounce called when SupportsFallback is false") + } h.fallbackCount++ if h.fallbackErr > 0 { h.fallbackErr-- @@ -72,6 +77,10 @@ func (h *fakeHeartbeatDriver) FallbackAnnounce(ctx context.Context) (ok bool) { return true } +func (h *fakeHeartbeatDriver) SupportsFallback() bool { + return !h.disableFallback +} + func (h *fakeHeartbeatDriver) Announce(ctx context.Context, sender inventory.DownstreamSender) (ok bool) { h.mu.Lock() defer h.mu.Unlock() @@ -176,8 +185,8 @@ func TestHeartbeatV2Basics(t *testing.T) { // use the control-stream announce. First poll always reads // as different, so expect that too. awaitEvents(t, hb.testEvents, - expect(hbv2PollDiff, hbv2FallbackOk, hbv2Start), - deny(hbv2FallbackErr, hbv2FallbackBackoff, hbv2AnnounceOk, hbv2AnnounceErr), + expect(hbv2PollDiff, hbv2FallbackOk, hbv2Start, hbv2OnHeartbeatOk), + deny(hbv2FallbackErr, hbv2FallbackBackoff, hbv2AnnounceOk, hbv2AnnounceErr, hbv2OnHeartbeatErr), ) // verify that we're now polling "same" and that time-based announces @@ -196,7 +205,7 @@ func TestHeartbeatV2Basics(t *testing.T) { // wait for fallback errors to happen, and confirm that we see fallback backoff // come into effect. we still expect no proper announce events. awaitEvents(t, hb.testEvents, - expect(hbv2FallbackErr, hbv2FallbackErr, hbv2FallbackBackoff, hbv2FallbackOk), + expect(hbv2FallbackErr, hbv2FallbackErr, hbv2FallbackBackoff, hbv2FallbackOk, hbv2OnHeartbeatErr, hbv2OnHeartbeatOk), deny(hbv2AnnounceOk, hbv2AnnounceErr), ) @@ -224,8 +233,8 @@ func TestHeartbeatV2Basics(t *testing.T) { // in case we refactor anything later). Take this opportunity to re-check that our announces // are internval and not poll based. awaitEvents(t, hb.testEvents, - expect(hbv2AnnounceOk, hbv2AnnounceOk, hbv2PollSame, hbv2AnnounceInterval), - deny(hbv2AnnounceErr, hbv2FallbackOk, hbv2FallbackErr), + expect(hbv2AnnounceOk, hbv2AnnounceOk, hbv2PollSame, hbv2AnnounceInterval, hbv2OnHeartbeatOk), + deny(hbv2AnnounceErr, hbv2FallbackOk, hbv2FallbackErr, hbv2OnHeartbeatErr), ) // set up a "changed" poll since we haven't traversed that path @@ -279,6 +288,119 @@ func TestHeartbeatV2Basics(t *testing.T) { ) } +func TestHeartbeatV2NoFallbackUnchecked(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // set up fake hb driver that lets us easily inject failures for + // the diff steps and assists w/ faking inventory control handles. + driver := newFakeHeartbeatDriver(t) + + driver.mu.Lock() + driver.disableFallback = true + driver.mu.Unlock() + + // set up hb with very long degraded check interval to confirm expected + // OnHeartbeat behavior outside of periodic checks. + hb := newHeartbeatV2(driver.handle, driver, heartbeatV2Config{ + announceInterval: time.Millisecond * 200, + pollInterval: time.Millisecond * 50, + fallbackBackoff: time.Millisecond * 400, + degradedCheckInterval: time.Hour, + testEvents: make(chan hbv2TestEvent, 1028), + }) + go hb.Run() + defer hb.Close() + + // verify that we tick but don't ever emit a degraded state + awaitEvents(t, hb.testEvents, + expect(hbv2NoFallback, hbv2NoFallback), + deny(hbv2OnHeartbeatOk, hbv2OnHeartbeatErr), + ) + + // make a stream available to the heartbeat instance + // (note: we don't need to pull from our half of the stream since + // fakeHeartbeatDriverInner doesn't actually send any messages across it). + stream := driver.newStream(ctx, t) + + // verify heartbeats + awaitEvents(t, hb.testEvents, + expect(hbv2OnHeartbeatOk), + deny(hbv2OnHeartbeatErr), + ) + + // verify that while heartbeating, we don't hit the "NoFallback" case + awaitEvents(t, hb.testEvents, + expect(hbv2OnHeartbeatOk, hbv2OnHeartbeatOk), + deny(hbv2OnHeartbeatErr, hbv2NoFallback), + ) + + stream.Close() + + awaitEvents(t, hb.testEvents, + expect(hbv2NoFallback), + deny(hbv2OnHeartbeatErr), + ) +} + +func TestHeartbeatV2NoFallbackChecked(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // set up fake hb driver that lets us easily inject failures for + // the diff steps and assists w/ faking inventory control handles. + driver := newFakeHeartbeatDriver(t) + + driver.mu.Lock() + driver.disableFallback = true + driver.mu.Unlock() + + // set up hb with fast degraded check interval + hb := newHeartbeatV2(driver.handle, driver, heartbeatV2Config{ + announceInterval: time.Millisecond * 200, + pollInterval: time.Millisecond * 50, + fallbackBackoff: time.Millisecond * 400, + degradedCheckInterval: time.Millisecond * 100, + testEvents: make(chan hbv2TestEvent, 1028), + }) + go hb.Run() + defer hb.Close() + + // verify that we tick and emit hb errs + awaitEvents(t, hb.testEvents, + expect(hbv2NoFallback, hbv2NoFallback, hbv2OnHeartbeatErr), + deny(hbv2OnHeartbeatOk), + ) + + // make a stream available to the heartbeat instance + // (note: we don't need to pull from our half of the stream since + // fakeHeartbeatDriverInner doesn't actually send any messages across it). + stream := driver.newStream(ctx, t) + + // verify heartbeats start + awaitEvents(t, hb.testEvents, + expect(hbv2OnHeartbeatOk), + ) + + // verify that the hb errs have stopped + awaitEvents(t, hb.testEvents, + expect(hbv2OnHeartbeatOk, hbv2OnHeartbeatOk), + deny(hbv2OnHeartbeatErr), + ) + + stream.Close() + + // verify that closed stream means errs resume + awaitEvents(t, hb.testEvents, + expect(hbv2NoFallback, hbv2OnHeartbeatErr), + deny(hbv2OnHeartbeatOk), + ) +} + type eventOpts struct { expect map[hbv2TestEvent]int deny map[hbv2TestEvent]struct{} diff --git a/lib/srv/regular/sshserver.go b/lib/srv/regular/sshserver.go index cc064cf1352..c72103952e9 100644 --- a/lib/srv/regular/sshserver.go +++ b/lib/srv/regular/sshserver.go @@ -845,7 +845,6 @@ func New( heartbeat, err = srv.NewSSHServerHeartbeat(srv.SSHServerHeartbeatConfig{ InventoryHandle: s.inventoryHandle, GetServer: s.getServerInfo, - Announcer: s.authService, OnHeartbeat: s.onHeartbeat, }) } else { diff --git a/lib/utils/retry.go b/lib/utils/retry.go index 9c1088c2860..01a21660cf4 100644 --- a/lib/utils/retry.go +++ b/lib/utils/retry.go @@ -40,18 +40,15 @@ var SeventhJitter = retryutils.NewSeventhJitter() // any usecases that might scale with cluster size or request count. var FullJitter = retryutils.NewFullJitter() -// NewDefaultLinear creates a linear retry using a half jitter, 10s step, and maxing out -// at 1 minute. These values were selected by reviewing commonly used parameters elsewhere -// in the code base, which (at the time of writing) seem to converge on approximately this -// configuration for "critical but potentially load-inducing" operations like cache watcher -// registration and auth connector setup. It also includes an auto-reset value of 5m. Auto-reset -// is less commonly used, and if used should probably be shorter, but 5m is a reasonable -// safety net to reduce the impact of accidental misuse. +// NewDefaultLinear creates a linear retry with reasonable default parameters for +// attempting to restart "critical but potentially load-inducing" operations, such +// as watcher or control stream resume. Exact parameters are subject to change, +// but this retry will always be configured for automatic reset. func NewDefaultLinear() *retryutils.Linear { retry, err := retryutils.NewLinear(retryutils.LinearConfig{ - First: HalfJitter(time.Second * 5), - Step: time.Second * 10, - Max: time.Minute, + First: FullJitter(time.Second * 10), + Step: time.Second * 15, + Max: time.Second * 90, Jitter: retryutils.NewHalfJitter(), AutoReset: 5, })