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
This commit is contained in:
Forrest
2023-03-03 17:39:00 +00:00
committed by GitHub
parent 948c1db915
commit aaee695927
8 changed files with 211 additions and 56 deletions
+1 -1
View File
@@ -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
)
+1 -1
View File
@@ -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(),
+1 -1
View File
@@ -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 (
+1 -1
View File
@@ -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(),
+73 -36
View File
@@ -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
+127 -5
View File
@@ -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{}
-1
View File
@@ -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 {
+7 -10
View File
@@ -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,
})