Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions docs/proxy-mode.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
202 changes: 155 additions & 47 deletions internal/autoconfig/stream_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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 (
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
}

Expand Down Expand Up @@ -303,44 +332,30 @@ 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) }) }

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 {
Expand All @@ -353,7 +368,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)

Expand All @@ -368,9 +383,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),
Comment thread
cursor[bot] marked this conversation as resolved.
es.StreamOptionRetryResetInterval(streamRetryResetInterval),
es.StreamOptionErrorHandler(errorHandler),
es.StreamOptionCanRetryFirstConnection(-1),
Expand All @@ -397,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.
Expand All @@ -420,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
}
}
Expand Down Expand Up @@ -615,6 +643,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,
Expand Down
Loading
Loading