From 574d0d7d1900629f4f904df52549b22488415f06 Mon Sep 17 00:00:00 2001 From: Matthew Keeler Date: Mon, 21 Sep 2026 09:23:53 -0400 Subject: [PATCH 1/2] feat: Keep retrying a rejected auto-configuration key A 401 or 403 on the auto-configuration stream stopped the stream and reported a fatal error, which made Relay exit. An operator can make the key valid again without Relay knowing, so recovery took a process restart, and a persistent store holding usable configuration was thrown away with the process. The stream now keeps retrying in every case. A failure unlikely to correct itself soon moves to a second eventsource retry profile, 5 minutes base and an hour ceiling, which bounds how much load a fleet of Relay Proxy instances puts on a service rejecting every request. The classification comes from internal/retry, so the auto-config stream and the big segment synchronizer sort failures the same way. A rejected key now reports INTERRUPTED rather than OFF, because the stream is still trying. OFF is left for Close. A stream that never connected stays INITIALIZING, matching the SDK data source status. Close now cancels the stream context as well as closing halt. halt is only observed when the next attempt fails, so on its own it cannot end a backoff wait already pending; without the cancel, Close would block for the remainder of a delay that now reaches an hour. Relay no longer exits when LaunchDarkly rejects the key. It serves from the persistent cache if one holds data, and answers 503 otherwise, which is what it already did for every other failure to reach LaunchDarkly. Anyone relying on the process exiting to detect a bad key should read the status resource instead. docs/proxy-mode.md says so. Three tests pinned the old stop-on-rejection behavior and are rewritten rather than renamed. stream_manager_known_limits_test.go pins an eventsource limitation instead: a server-sent "retry:" hint overrides the extended profile's base delay and is never cleared, so after a hint the extended backoff does not slow retries down. That is tolerated, not fixed, and the test exists so a future library fix surfaces deliberately. --- docs/proxy-mode.md | 16 ++ internal/autoconfig/stream_manager.go | 181 +++++++++++++----- .../autoconfig/stream_manager_errors_test.go | 83 ++++++-- .../stream_manager_known_limits_test.go | 101 ++++++++++ .../stream_manager_retry_profile_test.go | 64 +++++++ .../stream_manager_shutdown_test.go | 114 +++++++++++ .../autoconfig/stream_manager_status_test.go | 40 ++-- .../stream_manager_stop_policy_test.go | 87 +++++++++ 8 files changed, 611 insertions(+), 75 deletions(-) create mode 100644 internal/autoconfig/stream_manager_known_limits_test.go create mode 100644 internal/autoconfig/stream_manager_retry_profile_test.go create mode 100644 internal/autoconfig/stream_manager_shutdown_test.go create mode 100644 internal/autoconfig/stream_manager_stop_policy_test.go diff --git a/docs/proxy-mode.md b/docs/proxy-mode.md index 6a1524893..78c690f39 100644 --- a/docs/proxy-mode.md +++ b/docs/proxy-mode.md @@ -42,6 +42,22 @@ If the Relay Proxy receives requests from SDKs by this time, the behavior depend If you're an Enterprise customer using [automatic configuration](https://docs.launchdarkly.com/home/advanced/relay-proxy-enterprise/automatic-configuration), the first thing the Relay Proxy does on startup is request the configuration data from LaunchDarkly. During this time, the Relay Proxy does not yet know what the configured environments are, so it has no way to know if an SDK key or other credential in a request is valid. Therefore it returns a `503` error for all requests, indicating that it isn't ready yet. In this case, all LaunchDarkly SDKs will retry after a backoff delay. +### Relay Proxy receives a request when LaunchDarkly has rejected its auto-configuration key + +If LaunchDarkly rejects the auto-configuration stream with a `401` or `403` error, the auto-configuration key is not valid. The Relay Proxy does not give up. It keeps retrying on a slower schedule, which backs off to as long as one hour between attempts, so the connection recovers on its own if the key becomes valid again. Restart the Relay Proxy if you need it to reconnect immediately. + +This is the same thing the Relay Proxy already did for every other failure to reach LaunchDarkly, such as a `404`, a `5xx`, a network error, or a certificate problem. A rejected key is no longer treated differently. + +What happens to SDK requests in the meantime depends on whether the Relay Proxy has a configuration: + +* If you configure a [persistent store](./persistent-storage.md) and it holds configuration data from a previous run, the Relay Proxy loads that data and serves those environments while it keeps retrying. It logs `AutoConfig loaded from persistent cache`. + +* If the Relay Proxy has no configuration from any source, it cannot serve any request, because it does not know what its environments are. It answers every request with a `503` error until the key becomes valid. + +Read the Relay Proxy logs, or the [status resource](./endpoints.md), to tell a rejected key apart from an unreachable LaunchDarkly service. A rejected key logs `invalid auto-configuration key; will keep retrying in case it becomes valid`, and the status resource records the `401` or `403` under the auto-configuration stream's `lastError`. The reported state is `INTERRUPTED` if the stream had connected before the key was rejected, and `INITIALIZING` if it never connected at all. + +Versions of the Relay Proxy before 9.0.0 shut the process down immediately when LaunchDarkly rejected the auto-configuration key, even when a persistent store held usable configuration data. If you relied on the process exiting to detect a bad key, use the status resource instead. + ### Relay Proxy receives a request with invalid credentials If a server-side or mobile SDK connects to the Relay Proxy with an invalid SDK key or mobile key, the response is a `401` error. diff --git a/internal/autoconfig/stream_manager.go b/internal/autoconfig/stream_manager.go index bc62c559f..43722ff2a 100644 --- a/internal/autoconfig/stream_manager.go +++ b/internal/autoconfig/stream_manager.go @@ -21,6 +21,7 @@ import ( "github.com/launchdarkly/ld-relay/v9/internal/envfactory" "github.com/launchdarkly/ld-relay/v9/internal/httpconfig" "github.com/launchdarkly/ld-relay/v9/internal/logging" + "github.com/launchdarkly/ld-relay/v9/internal/retry" ) const ( @@ -31,6 +32,14 @@ const ( streamRetryResetInterval = 60 * time.Second streamJitterRatio = 0.5 defaultStreamRetryDelay = 1 * time.Second + + // Delays for a failure that is unlikely to correct itself soon, such as a rejected + // auto-configuration key. The stream keeps retrying on these instead of giving up, + // because an operator can make the key valid again without Relay knowing. The ceiling + // bounds how much load a fleet of Relay Proxy instances puts on a service that is + // rejecting every request. + streamExtendedRetryDelay = 5 * time.Minute + streamExtendedMaxRetryDelay = 1 * time.Hour ) var ( @@ -91,10 +100,19 @@ type StreamManager struct { lastKnownEnvs map[config.EnvironmentID]envfactory.EnvironmentRep httpConfig httpconfig.HTTPConfig initialRetryDelay time.Duration - logger *slog.Logger - halt chan struct{} - done chan struct{} // closed when the subscribe goroutine exits - closeOnce sync.Once + // extendedRetryDelay is the base delay used once a failure looks unlikely to correct + // itself soon. + extendedRetryDelay time.Duration + logger *slog.Logger + halt chan struct{} + done chan struct{} // closed when the subscribe goroutine exits + closeOnce sync.Once + + // streamCtx bounds the lifetime of the SSE connection, including any backoff wait between + // attempts. Cancelling it is the only way to interrupt that wait: eventsource's retry loop + // selects on the request's context and the delay timer, and nothing else. + streamCtx context.Context + streamCancel context.CancelFunc // statusLock guards status and failures. The eventsource error handler, the stream-consuming // goroutine, and the HTTP handlers that serve the status endpoint all touch them. @@ -130,16 +148,22 @@ func NewStreamManager( protocolVersionParam: []string{strconv.Itoa(protocolVersion)}, }.Encode() } + streamCtx, streamCancel := context.WithCancel(context.Background()) s := &StreamManager{ key: key, + streamCtx: streamCtx, + streamCancel: streamCancel, uri: streamURI, handler: handler, cache: cache, lastKnownEnvs: make(map[config.EnvironmentID]envfactory.EnvironmentRep), httpConfig: httpConfig, initialRetryDelay: initialRetryDelay, - logger: logger, - halt: make(chan struct{}), + // The extended delay has no configuration key. Its value bounds the load a fleet of + // Relay Proxy instances puts on a service that is rejecting its key. + extendedRetryDelay: streamExtendedRetryDelay, + logger: logger, + halt: make(chan struct{}), status: StreamStatus{ State: interfaces.DataSourceStateInitializing, StateSince: time.Now(), @@ -202,6 +226,11 @@ func (s *StreamManager) Start() <-chan error { func (s *StreamManager) Close() { s.closeOnce.Do(func() { close(s.halt) + // halt is only observed when the next attempt fails, so on its own it cannot end a + // backoff wait that is already pending. Cancelling the stream context does. Without + // this, Close would block on s.done for the remainder of a delay that now reaches an + // hour. + s.streamCancel() s.updateStatus(interfaces.DataSourceStateOff, interfaces.DataSourceErrorInfo{}) }) if s.done != nil { @@ -262,8 +291,8 @@ func (s *StreamManager) setStatus(state interfaces.DataSourceState, errorInfo in } if s.status.State == interfaces.DataSourceStateOff { - // OFF is terminal: either Close was called, or the key was rejected and the stream will not - // be retried. An event that was already in flight must not report the stream as working. + // OFF is terminal: Close was called, so the stream is finished. An event that was + // already in flight must not report the stream as working. return } @@ -304,43 +333,22 @@ func (s *StreamManager) subscribe(readyCh chan<- error) { var readyOnce sync.Once signalReady := func(err error) { readyOnce.Do(func() { readyCh <- err }) } - errorHandler := func(err error) es.StreamErrorHandlerResult { - // If Close() has been called, stop retrying so the SSE goroutine can exit. - select { - case <-s.halt: - return es.StreamErrorHandlerResult{CloseNow: true} - default: - } - - if se, ok := err.(es.SubscriptionError); ok { - errorInfo := interfaces.DataSourceErrorInfo{ - Kind: interfaces.DataSourceErrorKindErrorResponse, - StatusCode: se.Code, - Time: time.Now(), - } - if se.Code == 401 || se.Code == 403 { - s.logger.Error("invalid auto-configuration key; cannot get environments") - s.updateStatus(interfaces.DataSourceStateOff, errorInfo) - signalReady(errors.New("invalid auto-configuration key")) - return es.StreamErrorHandlerResult{CloseNow: true} - } - s.logger.Warn("HTTP error on auto-configuration stream", "statusCode", se.Code) - s.updateStatus(interfaces.DataSourceStateInterrupted, errorInfo) - return es.StreamErrorHandlerResult{CloseNow: false} - } - - s.logger.Warn("unexpected error on auto-configuration stream", "error", err) - s.updateStatus(interfaces.DataSourceStateInterrupted, interfaces.DataSourceErrorInfo{ - Kind: interfaces.DataSourceErrorKindNetworkError, - Time: time.Now(), - }) - return es.StreamErrorHandlerResult{CloseNow: false} + retryDelay := s.initialRetryDelay + if retryDelay <= 0 { + retryDelay = defaultStreamRetryDelay // COVERAGE: never happens in unit tests } - retry := s.initialRetryDelay - if retry <= 0 { - retry = defaultStreamRetryDelay // COVERAGE: never happens in unit tests - } + normalProfile := es.NewRetryProfile( + es.RetryProfileBaseDelay(retryDelay), + es.RetryProfileMaxDelay(streamMaxRetryDelay), + es.RetryProfileJitter(streamJitterRatio), + ) + extendedProfile := es.NewRetryProfile( + es.RetryProfileBaseDelay(s.extendedRetryDelay), + es.RetryProfileMaxDelay(streamExtendedMaxRetryDelay), + es.RetryProfileJitter(streamJitterRatio), + ) + errorHandler := s.newStreamErrorHandler(extendedProfile) rpacEndpoint, err := url.JoinPath(s.uri.String(), autoConfigStreamPath) if err != nil { @@ -353,7 +361,7 @@ func (s *StreamManager) subscribe(readyCh chan<- error) { return } - req, _ := http.NewRequest("GET", rpacEndpoint, nil) + req, _ := http.NewRequestWithContext(s.streamCtx, "GET", rpacEndpoint, nil) req.Header.Set("Authorization", string(s.key)) s.logger.Info("connecting to auto-configuration stream", "url", rpacEndpoint) @@ -368,9 +376,8 @@ func (s *StreamManager) subscribe(readyCh chan<- error) { stream, err := es.SubscribeWithRequestAndOptions(req, es.StreamOptionHTTPClient(client), es.StreamOptionReadTimeout(streamReadTimeout), - es.StreamOptionInitialRetry(retry), - es.StreamOptionUseBackoff(streamMaxRetryDelay), - es.StreamOptionUseJitter(streamJitterRatio), + es.StreamOptionDefaultRetryProfile(normalProfile), + es.StreamOptionRegisterRetryProfile(extendedProfile), es.StreamOptionRetryResetInterval(streamRetryResetInterval), es.StreamOptionErrorHandler(errorHandler), es.StreamOptionCanRetryFirstConnection(-1), @@ -615,6 +622,86 @@ func (s *StreamManager) dispatchEnvAction(id config.EnvironmentID, rep envfactor } } +// newStreamErrorHandler builds the SSE error handler for one subscribe cycle. +// +// Nothing stops the stream. A rejected key can become valid again without Relay knowing, and +// every other failure could clear at any time, so the stream keeps retrying in all cases. A +// failure that is unlikely to correct itself soon moves to the longer delays, which bounds the +// load a fleet puts on a service that is rejecting every request. +func (s *StreamManager) newStreamErrorHandler(extendedProfile *es.RetryProfile) func(error) es.StreamErrorHandlerResult { + // loggedExtended keeps the notice about the longer delays to once per subscribe cycle. + // The library returns to the normal delays itself once the connection has been healthy + // for streamRetryResetInterval, and does not report that, so re-logging would mislead. + loggedExtended := false + + return func(err error) es.StreamErrorHandlerResult { + // If Close() has been called, stop retrying so the SSE goroutine can exit. + select { + case <-s.halt: + return es.StreamErrorHandlerResult{CloseNow: true} + default: + } + + // Interrupted rather than Off even for a rejected key: the stream keeps retrying, so + // it is not finished. Off is left for Close. + s.updateStatus(interfaces.DataSourceStateInterrupted, streamErrorInfo(err)) + + result := es.StreamErrorHandlerResult{CloseNow: false} + if s.classifyAndLogStreamError(err) == retry.Unexpected { + if !loggedExtended { + s.logger.Info("classified failure as unexpected; engaging extended backoff") + loggedExtended = true + } + result.ActivateProfile = extendedProfile + } + return result + } +} + +// classifyAndLogStreamError sorts a stream failure into a retry class and logs it. The class +// decides how long to wait before the next attempt; it never decides whether to keep trying. +// +// A failure that is unlikely to correct itself soon is worth an error, because it nearly +// always means a real configuration problem, even though the stream recovers on its own once +// that problem is fixed. +func (s *StreamManager) classifyAndLogStreamError(err error) retry.FailureClass { + var se es.SubscriptionError + if !errors.As(err, &se) { + // No transport-level failure is unexpected, so these keep the short delays. + s.logger.Warn("unexpected error on auto-configuration stream", "error", err) + return retry.Normal + } + + class := retry.ClassifyHTTPStatus(se.Code) + switch { + case se.Code == http.StatusUnauthorized || se.Code == http.StatusForbidden: + s.logger.Error("invalid auto-configuration key; will keep retrying in case it becomes valid") + case class == retry.Unexpected: + s.logger.Error("HTTP error on auto-configuration stream", "statusCode", se.Code) + default: + s.logger.Warn("HTTP error on auto-configuration stream", "statusCode", se.Code) + } + return class +} + +// streamErrorInfo describes a stream failure in the shape the SDK data source status uses, so +// the status resource reports both the same way. An HTTP failure carries its status code; a +// transport failure has none to carry. +func streamErrorInfo(err error) interfaces.DataSourceErrorInfo { + var se es.SubscriptionError + if errors.As(err, &se) { + return interfaces.DataSourceErrorInfo{ + Kind: interfaces.DataSourceErrorKindErrorResponse, + StatusCode: se.Code, + Time: time.Now(), + } + } + return interfaces.DataSourceErrorInfo{ + Kind: interfaces.DataSourceErrorKindNetworkError, + Time: time.Now(), + } +} + func (s *StreamManager) applyCachedContent(content *PutContent) { s.handlePut(PutContent{ Environments: content.Environments, diff --git a/internal/autoconfig/stream_manager_errors_test.go b/internal/autoconfig/stream_manager_errors_test.go index 04b3e5ba9..ec566c77c 100644 --- a/internal/autoconfig/stream_manager_errors_test.go +++ b/internal/autoconfig/stream_manager_errors_test.go @@ -13,6 +13,9 @@ import ( helpers "github.com/launchdarkly/go-test-helpers/v3" "github.com/launchdarkly/go-test-helpers/v3/httphelpers" + + "github.com/launchdarkly/ld-relay/v9/config" + "github.com/launchdarkly/ld-relay/v9/internal/envfactory" ) func eventShouldCauseStreamRestart(t *testing.T, event httphelpers.SSEEvent) { @@ -125,29 +128,83 @@ func TestReconnectAfterNetworkError(t *testing.T) { errorShouldCauseReconnect(t, httphelpers.BrokenConnectionHandler(), "unexpected error", 0) } -func TestNoReconnectAfterUnrecoverableHTTPError(t *testing.T) { +func TestRecoversAfterUnrecoverableHTTPError(t *testing.T) { + // A rejected key no longer stops the stream. An operator can make the key valid again + // without Relay knowing, so the stream keeps trying and recovers on its own, which + // previously took a process restart. for _, status := range []int{401, 403} { t.Run(fmt.Sprintf("status %d", status), func(t *testing.T) { initialEvent := makeEnvPutEvent(testEnv1) streamHandler, stream := httphelpers.SSEHandler(&initialEvent) defer stream.Close() - errorProducingHandler := httphelpers.HandlerWithStatus(status) handler := httphelpers.SequentialHandler( - errorProducingHandler, // first request will get this - streamHandler, // request after reconnect will get this + httphelpers.HandlerWithStatus(status), // first request is rejected + streamHandler, // the retry succeeds ) streamManagerTestWithStreamHandler(t, handler, stream, func(p streamManagerTestParams) { + // Shorten the extended delay so the retry happens within the test. + p.streamManager.extendedRetryDelay = time.Millisecond p.startStream() - <-p.requestsCh // first request - select { - case <-p.requestsCh: // got expected stream restart - require.Fail(t, "got unexpected stream restart") - case <-p.messageHandler.received: - require.Fail(t, "got unexpected event") - case <-time.After(time.Millisecond * 200): - assert.True(t, p.mockLog.HasMessage(slog.LevelError, "invalid auto-configuration key")) - } + + <-p.requestsCh // the rejected request + <-p.requestsCh // the retry + + p.requireMessage() // the environment from the recovered stream + p.requireReceivedAllMessage() + + assert.True(t, p.mockLog.HasMessage(slog.LevelError, "will keep retrying")) + assert.True(t, p.mockLog.HasMessage(slog.LevelInfo, "engaging extended backoff")) }) }) } } + +func TestServesTheCachedConfigurationWhileTheKeyIsRejected(t *testing.T) { + // A rejected key leaves Relay running, so a cached configuration is worth having: Relay + // serves those environments while it keeps trying the key. Before this change the process + // exited and threw the cache away. + handler := httphelpers.HandlerWithStatus(401) + _, stream := httphelpers.SSEHandler(nil) + defer stream.Close() + + cache := cacheWithContent{content: &PutContent{ + Environments: map[config.EnvironmentID]envfactory.EnvironmentRep{testEnv1.EnvID: testEnv1}, + }} + + streamManagerTestWithCache(t, handler, stream, cache, func(p streamManagerTestParams) { + p.streamManager.extendedRetryDelay = time.Millisecond + + readyCh := p.streamManager.Start() + + // The cached environment reaches the handler, which is what makes Relay serviceable. + p.requireMessage() + p.requireReceivedAllMessage() + + if !helpers.AssertNoMoreValues(t, readyCh, time.Second, "Relay reported a failure") { + t.FailNow() + } + assert.True(t, p.mockLog.HasMessage(slog.LevelInfo, "loaded from persistent cache")) + }) +} + +func TestKeepsRunningWithNoConfigurationAtAll(t *testing.T) { + // The case that used to exit the process. With a rejected key and nothing cached, Relay has + // nothing to serve, and it still keeps running and retrying rather than reporting a failure. + // It answers 503 until the key becomes valid, which is what it already did for every other + // failure to reach LaunchDarkly. + handler := httphelpers.HandlerWithStatus(401) + _, stream := httphelpers.SSEHandler(nil) + defer stream.Close() + + streamManagerTestWithStreamHandler(t, handler, stream, func(p streamManagerTestParams) { + p.streamManager.extendedRetryDelay = time.Millisecond + + readyCh := p.streamManager.Start() + + if !helpers.AssertNoMoreValues(t, readyCh, 500*time.Millisecond, + "Relay reported a failure on a rejected key") { + t.FailNow() + } + assert.True(t, p.mockLog.HasMessage(slog.LevelError, "will keep retrying")) + }) +} diff --git a/internal/autoconfig/stream_manager_known_limits_test.go b/internal/autoconfig/stream_manager_known_limits_test.go new file mode 100644 index 000000000..ed33a51d6 --- /dev/null +++ b/internal/autoconfig/stream_manager_known_limits_test.go @@ -0,0 +1,101 @@ +package autoconfig + +// Tests in this file pin behavior Relay currently has and that we have decided not to change +// yet. They are not proofs of pending bugs: each one asserts today's behavior so that a future +// fix shows up as a failing test somebody has to update deliberately, rather than a silent +// change. The comment on each test says what the limitation is and why it is tolerated. + +import ( + "fmt" + "log/slog" + "net/http" + "net/http/httptest" + "net/url" + "strconv" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/launchdarkly/ld-relay/v9/config" + "github.com/launchdarkly/ld-relay/v9/internal/httpconfig" + "github.com/launchdarkly/ld-relay/v9/internal/logging/logtest" +) + +// reconnectDelays returns every delay eventsource logged, from both the first-connection loop +// ("retrying in N secs") and the post-connection loop ("Reconnecting in N secs"). The library +// writes these through the eventsource logger bridge, which records them as slog messages. +func reconnectDelays(mockLog *logtest.MockHandler) []time.Duration { + var out []time.Duration + for _, e := range mockLog.EntriesForLevel(slog.LevelInfo) { + for _, marker := range []string{"retrying in ", "Reconnecting in "} { + _, after, found := strings.Cut(e.Message, marker) + if !found { + continue + } + secs, err := strconv.ParseFloat(strings.TrimSuffix(strings.TrimSpace(after), " secs"), 64) + if err == nil { + out = append(out, time.Duration(secs*float64(time.Second))) + } + } + } + return out +} + +// TestKnownLimitServerRetryHintDefeatsTheExtendedProfile: eventsource's baseDelayOverride (set by an +// SSE "retry:" field) replaces the active profile's base delay, including the extended +// profile's, and is never cleared (retry_delay.go NextRetryDelay: "if +// activeState.baseDelayOverride != nil { effectiveBase = *activeState.baseDelayOverride }"). +// So once the stream has seen a retry hint, engaging the extended profile does not slow +// retries down to the intended 5 min: Relay keeps hammering the rejecting service on the +// hinted base delay instead. +func TestKnownLimitServerRetryHintDefeatsTheExtendedProfile(t *testing.T) { + const hintMillis = 10 + var requestCount int64 + + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + n := atomic.AddInt64(&requestCount, 1) + if n > 1 { + w.WriteHeader(401) // the key is rejected from now on + return + } + // One good connection that carries a server-directed retry hint, then it drops. + w.Header().Set("Content-Type", "text/event-stream") + w.WriteHeader(200) + fmt.Fprintf(w, "retry: %d\nevent: put\ndata: {\"path\":\"/\",\"data\":{\"environments\":{}}}\n\n", hintMillis) + w.(http.Flusher).Flush() + })) + defer server.Close() + + logger, mockLog := logtest.NewMockLogger() + + httpConfig, err := httpconfig.NewHTTPConfig(config.ProxyConfig{}, config.HTTPConfig{}, nil, "", logger) + require.NoError(t, err) + serverURL, err := url.Parse(server.URL) + require.NoError(t, err) + + sm := NewStreamManager( + testConfigKey, serverURL, newTestMessageHandler(), httpConfig, + time.Millisecond, rpacProtocolVersion, logger, noopTestCache{}, + ) + defer sm.Close() + sm.extendedRetryDelay = 10 * time.Minute // the extended base delay we are supposed to get + + sm.Start() + + // Wait until the key has been rejected several times. + require.Eventually(t, func() bool { + return atomic.LoadInt64(&requestCount) >= 5 + }, 3*time.Second, 10*time.Millisecond, "the rejecting service was contacted fewer than 5 times") + + assert.True(t, mockLog.HasMessage(slog.LevelInfo, "engaging extended backoff")) + delays := reconnectDelays(mockLog) + t.Logf("delays eventsource used after the extended profile was engaged: %v", delays) + for _, d := range delays { + assert.Less(t, d, 5*time.Second, + "every delay stayed on the 10ms hint, not the 10 min extended base") + } +} diff --git a/internal/autoconfig/stream_manager_retry_profile_test.go b/internal/autoconfig/stream_manager_retry_profile_test.go new file mode 100644 index 000000000..303423709 --- /dev/null +++ b/internal/autoconfig/stream_manager_retry_profile_test.go @@ -0,0 +1,64 @@ +package autoconfig + +import ( + "net/http" + "net/http/httptest" + "net/url" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/launchdarkly/ld-relay/v9/config" + "github.com/launchdarkly/ld-relay/v9/internal/httpconfig" + "github.com/launchdarkly/ld-relay/v9/internal/logging/logtest" +) + +// TestExtendedDelaysDoubleFromTheExtendedBase checks that engaging the extended profile really +// slows retries down, which is the point of the change and the part no other test covers: +// removing the profile activation leaves every other test in the package green. +// +// The activation has to survive repeated passes through eventsource's +// CanRetryFirstConnection(-1) loop, and the per-profile attempt counter has to double from the +// extended base rather than the normal one. Both hold. +func TestExtendedDelaysDoubleFromTheExtendedBase(t *testing.T) { + const base = 20 * time.Millisecond + + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(401) + })) + defer server.Close() + + logger, mockLog := logtest.NewMockLogger() + + httpConfig, err := httpconfig.NewHTTPConfig(config.ProxyConfig{}, config.HTTPConfig{}, nil, "", logger) + require.NoError(t, err) + serverURL, err := url.Parse(server.URL) + require.NoError(t, err) + + sm := NewStreamManager( + testConfigKey, serverURL, newTestMessageHandler(), httpConfig, + time.Millisecond, rpacProtocolVersion, logger, noopTestCache{}, + ) + defer sm.Close() + sm.extendedRetryDelay = base + + sm.Start() + + var delays []time.Duration + require.Eventually(t, func() bool { + delays = reconnectDelays(mockLog) + return len(delays) >= 4 + }, 5*time.Second, 10*time.Millisecond, "expected at least 4 logged delays") + + t.Logf("first four extended delays: %v", delays[:4]) + // Jitter removes up to half, so attempt n's delay is in [base*2^(n-1)/2, base*2^(n-1)]. + for i := 0; i < 4; i++ { + want := time.Duration(1<= 1 }, 2*time.Second, 10*time.Millisecond, + "expected the first attempt") + + sm.Close() + atClose := requests() + + // Several extended delays' worth of wall clock, so a surviving wait would have fired. + time.Sleep(700 * time.Millisecond) + assert.Equal(t, atClose, requests(), "Relay contacted the service after Close() returned") +} diff --git a/internal/autoconfig/stream_manager_status_test.go b/internal/autoconfig/stream_manager_status_test.go index c83622789..1d801b551 100644 --- a/internal/autoconfig/stream_manager_status_test.go +++ b/internal/autoconfig/stream_manager_status_test.go @@ -232,25 +232,34 @@ func TestStreamStatusRecordsInvalidDataError(t *testing.T) { }) } -func TestStreamStatusIsOffAfterUnrecoverableHTTPError(t *testing.T) { +func TestStreamStatusIsNotOffAfterARejectedKey(t *testing.T) { + // The stream keeps retrying a rejected key, so it is not finished and must not report OFF. + // OFF is reserved for Close. for _, status := range []int{401, 403} { t.Run(fmt.Sprintf("status %d", status), func(t *testing.T) { - streamHandler, stream := httphelpers.SSEHandler(nil) + // Every attempt is rejected, so the stream never becomes ready and stays in the + // retrying state this test is about. Recovery is covered separately by + // TestRecoversAfterUnrecoverableHTTPError. + handler := httphelpers.HandlerWithStatus(status) + _, stream := httphelpers.SSEHandler(nil) defer stream.Close() - handler := httphelpers.SequentialHandler( - httphelpers.HandlerWithStatus(status), - streamHandler, - ) streamManagerTestWithStreamHandler(t, handler, stream, func(p streamManagerTestParams) { + p.streamManager.extendedRetryDelay = time.Millisecond readyCh := p.streamManager.Start() - err := helpers.RequireValue(t, readyCh, time.Second, "timed out waiting for stream failure") - require.Error(t, err) - got := requireStatusEventually(t, p, "expected OFF", func(s StreamStatus) bool { - return s.State == interfaces.DataSourceStateOff - }) + got := requireStatusEventually(t, p, "expected the rejection to be recorded", + func(s StreamStatus) bool { + return s.LastError.StatusCode == status + }) + assert.NotEqual(t, interfaces.DataSourceStateOff, got.State, + "the stream is still retrying, so it is not finished") + // It never connected, so an interruption keeps the initializing state. + assert.Equal(t, interfaces.DataSourceStateInitializing, got.State) assert.Equal(t, interfaces.DataSourceErrorKindErrorResponse, got.LastError.Kind) - assert.Equal(t, status, got.LastError.StatusCode) + + if !helpers.AssertNoMoreValues(t, readyCh, 300*time.Millisecond, "Relay reported a failure") { + t.FailNow() + } }) }) } @@ -545,7 +554,7 @@ func TestStreamStatusRecoversOnAnUnrecognizedEvent(t *testing.T) { // A key revoked mid-stream is the case where this status is the only signal. The error never reaches // the channel Relay exits on, because signalReady has already fired, so Relay keeps serving with a // configuration stream that will not be retried. -func TestStreamStatusIsOffAfterMidStreamKeyRevocation(t *testing.T) { +func TestStreamStatusIsInterruptedAfterMidStreamKeyRevocation(t *testing.T) { initialEvent := makeEnvPutEvent(testEnv1) streamHandler, stream := httphelpers.SSEHandler(&initialEvent) defer stream.Close() @@ -563,9 +572,10 @@ func TestStreamStatusIsOffAfterMidStreamKeyRevocation(t *testing.T) { stream.EndAll() - got := requireStatusEventually(t, p, "expected OFF after the key was rejected", + // The stream had connected, so a rejection now interrupts it rather than ending it. + got := requireStatusEventually(t, p, "expected INTERRUPTED after the key was rejected", func(s StreamStatus) bool { - return s.State == interfaces.DataSourceStateOff + return s.State == interfaces.DataSourceStateInterrupted }) assert.Equal(t, interfaces.DataSourceErrorKindErrorResponse, got.LastError.Kind) assert.Equal(t, 401, got.LastError.StatusCode) diff --git a/internal/autoconfig/stream_manager_stop_policy_test.go b/internal/autoconfig/stream_manager_stop_policy_test.go new file mode 100644 index 000000000..ef2036649 --- /dev/null +++ b/internal/autoconfig/stream_manager_stop_policy_test.go @@ -0,0 +1,87 @@ +package autoconfig + +import ( + "fmt" + "log/slog" + "net/http/httptest" + "net/url" + "testing" + "time" + + helpers "github.com/launchdarkly/go-test-helpers/v3" + "github.com/launchdarkly/go-test-helpers/v3/httphelpers" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/launchdarkly/ld-relay/v9/config" + "github.com/launchdarkly/ld-relay/v9/internal/httpconfig" + "github.com/launchdarkly/ld-relay/v9/internal/logging/logtest" + "github.com/launchdarkly/ld-relay/v9/internal/retry" +) + +// Every status in the slow-retry class keeps Relay running and keeps it trying. The class +// decides how long to wait, and nothing else: a misconfigured stream URI or a load balancer +// answering 404 mid-deploy must not take a fleet down, and neither must a rejected key. +func TestUnexpectedStatusKeepsRelayRunning(t *testing.T) { + for _, status := range []int{404, 405, 410, 418, 451} { + t.Run(fmt.Sprintf("status %d", status), func(t *testing.T) { + // These are all in the slow-retry class, which is the point: the class picks the + // delay and never decides whether to keep trying. + require.Equal(t, retry.Unexpected, retry.ClassifyHTTPStatus(status)) + + handler, requestsCh := httphelpers.RecordingHandler(httphelpers.HandlerWithStatus(status)) + _, stream := httphelpers.SSEHandler(nil) + defer stream.Close() + + streamManagerTestWithStreamHandler(t, handler, stream, func(p streamManagerTestParams) { + p.streamManager.extendedRetryDelay = time.Millisecond + + readyCh := p.streamManager.Start() + + // It keeps trying rather than reporting a failure. + helpers.RequireValue(t, requestsCh, time.Second, "expected the first attempt") + helpers.RequireValue(t, requestsCh, time.Second, "expected a retry") + + if !helpers.AssertNoMoreValues(t, readyCh, 300*time.Millisecond, + "Relay reported a fatal error for a status that is not a rejected credential") { + t.FailNow() + } + assert.False(t, p.mockLog.HasMessage(slog.LevelError, "invalid auto-configuration key")) + }) + }) + } +} + +// A certificate failure is in the slow-retry class too, and likewise must not stop Relay. A +// machine with a skewed clock, or a TLS-terminating proxy mid-cert-rollover, produces this and +// resolves without Relay's involvement. +func TestCertificateFailureDoesNotStopRelay(t *testing.T) { + server := httptest.NewTLSServer(httphelpers.HandlerWithStatus(200)) + defer server.Close() + + logger, mockLog := logtest.NewMockLogger() + + httpConfig, err := httpconfig.NewHTTPConfig(config.ProxyConfig{}, config.HTTPConfig{}, nil, "", logger) + require.NoError(t, err) + serverURL, err := url.Parse(server.URL) + require.NoError(t, err) + + sm := NewStreamManager( + testConfigKey, serverURL, newTestMessageHandler(), httpConfig, + time.Millisecond, rpacProtocolVersion, logger, noopTestCache{}, + ) + defer sm.Close() + sm.extendedRetryDelay = time.Millisecond + + readyCh := sm.Start() + if !helpers.AssertNoMoreValues(t, readyCh, time.Second, + "Relay reported a fatal error for a certificate failure") { + t.FailNow() + } + + // A certificate failure is a transport failure, and no transport failure is unexpected, so + // it must stay on the short delays. Waiting minutes would not help: the failure resolves + // the moment an operator fixes the certificate. + assert.False(t, mockLog.HasMessage(slog.LevelInfo, "engaging extended backoff")) + assert.True(t, mockLog.HasMessage(slog.LevelWarn, "unexpected error on auto-configuration stream")) +} From 8075c61340a454b5a0eeb837c56db985b1088f11 Mon Sep 17 00:00:00 2001 From: Matthew Keeler Date: Mon, 21 Sep 2026 15:05:25 -0400 Subject: [PATCH 2/2] fix: Do not report a shutdown as a fatal auto-configuration failure Close cancels the stream context, and eventsource returns that cancellation from SubscribeWithRequestAndOptions without consulting the error handler. Two arms of subscribe's select could then mishandle it, and both were reachable because Close closes halt and cancels the context, leaving the arms ready at once so which wins is decided per-run. The streamCh arm treated the cancellation as a stream failure and passed it to signalReady. ld-relay.go exits the process on a non-nil value from that channel, so a graceful shutdown of a stream that had never connected -- the rejected-key path this change exists to keep alive -- could exit 1. The halt arm returned without signalling readyCh at all, so a caller still waiting for the first connection waited forever. That is the arm that wins in practice: the new test reproduced it 5 times out of 5 before this fix, while the exit needs the cancellation to land in the same scheduling window. Both now call signalShutdown, which closes readyCh under the same sync.Once as signalReady, so exactly one of the two ever touches the channel. Closing rather than sending nil is the point: a non-nil error means fatal, and a shutdown is not that. Reported by Cursor Bugbot and by the review agent on #883, and confirmed by kinyoklion. The test is the one the review suggested. --- internal/autoconfig/stream_manager.go | 21 ++++++++++++++++ .../stream_manager_shutdown_test.go | 25 +++++++++++++++++++ 2 files changed, 46 insertions(+) diff --git a/internal/autoconfig/stream_manager.go b/internal/autoconfig/stream_manager.go index 43722ff2a..65f4c1211 100644 --- a/internal/autoconfig/stream_manager.go +++ b/internal/autoconfig/stream_manager.go @@ -332,6 +332,13 @@ func (s *StreamManager) subscribe(readyCh chan<- error) { var readyOnce sync.Once signalReady := func(err error) { readyOnce.Do(func() { readyCh <- err }) } + // signalShutdown releases a caller waiting on readyCh without reporting a failure, for the + // case where Close is what ended the attempt. It shares readyOnce with signalReady, so + // exactly one of the two ever touches the channel and neither can send on a closed one. + // + // Closing rather than sending nil matters: a non-nil error on this channel makes the caller + // treat the failure as fatal and exit the process, and a shutdown is not that. + signalShutdown := func() { readyOnce.Do(func() { close(readyCh) }) } retryDelay := s.initialRetryDelay if retryDelay <= 0 { @@ -404,6 +411,17 @@ func (s *StreamManager) subscribe(readyCh chan<- error) { case result := <-streamCh: if result.err != nil { + // Close cancels the stream context, and eventsource returns that cancellation from + // here without consulting the error handler. It is the shutdown's own doing rather + // than a stream failure, so it must not be reported as one: the caller exits the + // process on a non-nil error, and this arm and the halt arm below are both ready + // once Close has run, so which one wins is decided per-run. + select { + case <-s.halt: + signalShutdown() + return + default: + } s.logger.Error("unexpected error on auto-configuration stream", "error", result.err) // The error handler has already recorded why the connection failed, so this reports // only that the stream is permanently off, and keeps that specific error. @@ -427,6 +445,9 @@ func (s *StreamManager) subscribe(readyCh chan<- error) { result.stream.Close() } }() + // Nothing signalled readyCh on this path before, so a caller that was still waiting + // for the first connection waited forever. + signalShutdown() return } } diff --git a/internal/autoconfig/stream_manager_shutdown_test.go b/internal/autoconfig/stream_manager_shutdown_test.go index 35492dfc1..d3fd8b691 100644 --- a/internal/autoconfig/stream_manager_shutdown_test.go +++ b/internal/autoconfig/stream_manager_shutdown_test.go @@ -95,6 +95,31 @@ func TestCloseInterruptsAPendingBackoffWait(t *testing.T) { stacksContaining("abandonStreamGoroutine")...)) } +// Close() while the stream is stuck in its first-connection retry loop must not report the +// shutdown's own context cancellation as a fatal stream failure. Eventsource returns that +// cancellation without consulting the error handler, so if it reaches the ready channel here it +// makes ld-relay.go call os.Exit(1) in the middle of a graceful shutdown. Close() must instead +// close the ready channel, so whoever waits on it unblocks without an error. +func TestCloseWhileRetryingDoesNotSignalAFatalStreamError(t *testing.T) { + sm, mockLog, _, closeServer := newRejectingStreamManager(t, 5*time.Minute) + defer closeServer() + + readyCh := sm.Start() + require.Eventually(t, func() bool { + return mockLog.HasMessage(slog.LevelError, "invalid auto-configuration key") + }, 2*time.Second, 10*time.Millisecond, "expected the first rejection") + + sm.Close() + + select { + case err := <-readyCh: + assert.NoError(t, err, + "the shutdown's context cancellation must not be reported as a fatal stream error") + case <-time.After(2 * time.Second): + t.Fatal("the ready channel was not closed after Close() returned") + } +} + // The same guarantee stated in terms of observable behavior: nothing reaches LaunchDarkly after // Close() returns. A short extended delay is used so a surviving wait would actually fire. func TestNoRequestIsMadeAfterClose(t *testing.T) {