diff --git a/contract-tests/src/main/java/com/launchdarkly/sdktest/TestService.java b/contract-tests/src/main/java/com/launchdarkly/sdktest/TestService.java index b6618d9e..50acf5e5 100644 --- a/contract-tests/src/main/java/com/launchdarkly/sdktest/TestService.java +++ b/contract-tests/src/main/java/com/launchdarkly/sdktest/TestService.java @@ -43,7 +43,9 @@ public class TestService extends NanoHTTPD { "evaluation-hooks", "track-hooks", "client-per-context-summaries", - "client-event-source-http-errors" + "client-event-source-http-errors", + "retry-conformance-fdv1-streaming", + "retry-conformance-fdv1-polling" }; private static final String MIME_JSON = "application/json"; static final Gson gson = new GsonBuilder() diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ConnectivityManager.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ConnectivityManager.java index cbeaf9fb..aae89b11 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ConnectivityManager.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/ConnectivityManager.java @@ -136,8 +136,6 @@ public void setStatus(@NonNull DataSourceState state, Throwable failure) { @Override public void shutDown() { - // The DataSource will call this method if it receives an error such as HTTP 401 that - // indicates the mobile key is invalid. ConnectivityManager.this.shutDown(); setStatus(ConnectionInformation.ConnectionMode.SHUTDOWN, null); } diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSource.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSource.java index 9b8d4056..3f9efac0 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSource.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSource.java @@ -125,6 +125,29 @@ public interface DataSourceFactory { @NonNull LDLogger logger, long fallbackTimeoutSeconds, long recoveryTimeoutSeconds + ) { + this(evaluationContext, initializers, synchronizers, fdv1FallbackSynchronizer, + dataSourceUpdateSink, sharedExecutor, logger, fallbackTimeoutSeconds, + recoveryTimeoutSeconds, RetryRegime.EXTENDED); + } + + /** + * This constructor allows tests to shorten the backoff after a synchronizer's unexpected + * error. See the other constructor for the remaining parameters. + * + * @param synchronizerBackoff the delay bounds for that backoff + */ + FDv2DataSource( + @NonNull LDContext evaluationContext, + @NonNull List> initializers, + @NonNull List> synchronizers, + @Nullable DataSourceFactory fdv1FallbackSynchronizer, + @NonNull DataSourceUpdateSinkV2 dataSourceUpdateSink, + @NonNull ScheduledExecutorService sharedExecutor, + @NonNull LDLogger logger, + long fallbackTimeoutSeconds, + long recoveryTimeoutSeconds, + @NonNull RetryRegime synchronizerBackoff ) { this.evaluationContext = evaluationContext; this.dataSourceUpdateSink = dataSourceUpdateSink; @@ -140,16 +163,17 @@ public interface DataSourceFactory { List allSynchronizers = new ArrayList<>(); for (DataSourceFactory factory : synchronizers) { - allSynchronizers.add(new SynchronizerFactoryWithState(factory)); + allSynchronizers.add(new SynchronizerFactoryWithState(factory, false, synchronizerBackoff)); } if (fdv1FallbackSynchronizer != null) { - SynchronizerFactoryWithState fdv1 = new SynchronizerFactoryWithState(fdv1FallbackSynchronizer, true); + SynchronizerFactoryWithState fdv1 = + new SynchronizerFactoryWithState(fdv1FallbackSynchronizer, true, synchronizerBackoff); fdv1.block(); allSynchronizers.add(fdv1); } // note that the source manager only uses the initializers after the cache initializers and not the cache initializers - this.sourceManager = new SourceManager(allSynchronizers, generalInitializers); + this.sourceManager = new SourceManager(allSynchronizers, generalInitializers, sharedExecutor); this.fallbackTimeoutSeconds = fallbackTimeoutSeconds; this.recoveryTimeoutSeconds = recoveryTimeoutSeconds; this.sharedExecutor = sharedExecutor; @@ -475,11 +499,20 @@ private List getConditions(int synchronizerC List list = new ArrayList<>(); list.add(new FDv2DataSourceConditions.FallbackCondition(sharedExecutor, fallbackTimeoutSeconds)); if (!isPrime) { - list.add(new FDv2DataSourceConditions.RecoveryCondition(sharedExecutor, recoveryTimeoutSeconds)); + // Recovery only goes ahead once a higher-priority synchronizer is available, not + // merely backing off. + list.add(new FDv2DataSourceConditions.RecoveryCondition(sharedExecutor, recoveryTimeoutSeconds, + sourceManager::hasAvailableSynchronizerBeforeCurrent)); } return list; } + /** + * Returns the next available synchronizer, waiting for a backoff to end if that is the only + * way to get one. Returns null once the data source is stopped or there is nothing left to + * try. + */ + @Nullable private static String detailForThrowable(@Nullable Throwable error) { if (error == null) { return "unknown error"; @@ -515,12 +548,12 @@ private void runSynchronizers( @NonNull DataSourceUpdateSinkV2 sink ) { try { - Synchronizer synchronizer = sourceManager.getNextAvailableSynchronizerAndSetActive(); + Synchronizer synchronizer = sourceManager.nextAvailableSynchronizer().get(); while (synchronizer != null) { String synchronizerName = synchronizer.name(); logger.info("Synchronizer '{}' is starting.", synchronizerName); resetSynchronizerStatusDedupe(); - int synchronizerCount = sourceManager.getAvailableSynchronizerCount(); + int synchronizerCount = sourceManager.getUsableSynchronizerCount(); boolean isPrime = sourceManager.isPrimeSynchronizer(); try { boolean running = true; @@ -570,6 +603,7 @@ private void runSynchronizers( if (changeSet != null) { sink.apply(context, changeSet); sink.setStatus(DataSourceState.VALID, null); + sourceManager.recordCurrentSynchronizerHealthy(System.currentTimeMillis()); tryCompleteStart(true, null); } break; @@ -595,15 +629,22 @@ private void runSynchronizers( running = false; break; case TERMINAL_ERROR: + // Move on to the next synchronizer now, and put this one + // into a backoff so that it is tried again later. maybeLogSynchronizerStatusChange( synchronizer.name(), status.getState() ); - sourceManager.blockCurrentSynchronizer(); + long backoffMillis = sourceManager.backOffCurrentSynchronizer( + System.currentTimeMillis()); logger.warn( - "Synchronizer '{}' permanently failed and will not be used again until application restart.", - synchronizer.name() + "Synchronizer '{}' reported an unexpected error and will not be tried again for {} seconds.", + synchronizer.name(), + backoffMillis / 1000 ); + if (sourceManager.getAvailableSynchronizerCount() == 0) { + logger.info("All synchronizers are waiting out a backoff after unexpected errors; the first to become available will be tried next."); + } running = false; sink.setStatus(DataSourceState.INTERRUPTED, status.getError()); break; @@ -651,7 +692,7 @@ private void runSynchronizers( Thread.currentThread().interrupt(); return; } - synchronizer = sourceManager.getNextAvailableSynchronizerAndSetActive(); + synchronizer = sourceManager.nextAvailableSynchronizer().get(); } if (!stopCalled.get()) { logger.warn("No more synchronizers available."); diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSourceConditions.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSourceConditions.java index fd9541d2..43983646 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSourceConditions.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2DataSourceConditions.java @@ -1,6 +1,7 @@ package com.launchdarkly.sdk.android; import androidx.annotation.NonNull; +import androidx.annotation.Nullable; import com.launchdarkly.sdk.android.subsystems.FDv2SourceResult; import com.launchdarkly.sdk.fdv2.SourceResultType; @@ -99,16 +100,48 @@ public ConditionType getType() { } /** - * Recovery: timer starts when built. Future completes with RECOVERY when timer fires. + * Decides, when a recovery timer fires, whether there is anything to recover to. + */ + interface RecoveryGate { + boolean canRecover(); + } + + /** + * Recovery: timer starts when built. Future completes with RECOVERY when timer fires, unless + * the gate says there is nothing to recover to yet, in which case the timer is re-armed for + * another timeout. */ static final class RecoveryCondition extends TimedCondition { + @Nullable + private final RecoveryGate gate; + private volatile boolean closed = false; RecoveryCondition(@NonNull ScheduledExecutorService executor, long timeoutSeconds) { + this(executor, timeoutSeconds, null); + } + + RecoveryCondition(@NonNull ScheduledExecutorService executor, long timeoutSeconds, @Nullable RecoveryGate gate) { super(executor, timeoutSeconds); - this.timerFuture = executor.schedule( - () -> resultFuture.set(ConditionType.RECOVERY), - timeoutSeconds, - TimeUnit.SECONDS); + this.gate = gate; + arm(); + } + + private void arm() { + this.timerFuture = sharedExecutor.schedule(this::onTimer, timeoutSeconds, TimeUnit.SECONDS); + } + + private void onTimer() { + if (gate == null || gate.canRecover()) { + resultFuture.set(ConditionType.RECOVERY); + } else if (!closed) { + arm(); + } + } + + @Override + public void close() { + closed = true; + super.close(); } @Override diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizer.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizer.java index ef8e0d5f..12048e7e 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizer.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizer.java @@ -5,12 +5,9 @@ import androidx.annotation.VisibleForTesting; import com.launchdarkly.eventsource.ConnectStrategy; -import com.launchdarkly.eventsource.ErrorStrategy; import com.launchdarkly.eventsource.EventSource; -import com.launchdarkly.eventsource.FaultEvent; import com.launchdarkly.eventsource.HttpConnectStrategy; import com.launchdarkly.eventsource.MessageEvent; -import com.launchdarkly.eventsource.RetryDelayStrategy; import com.launchdarkly.eventsource.StreamClosedByCallerException; import com.launchdarkly.eventsource.StreamEvent; import com.launchdarkly.eventsource.ResponseHeaders; @@ -34,8 +31,9 @@ import java.io.IOException; import java.net.URI; import java.util.Map; -import java.util.concurrent.Executor; import java.util.concurrent.Future; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -51,12 +49,15 @@ * If an optional {@link FDv2Requestor} is supplied, {@code ping} SSE events are handled by * issuing a poll request. If no requestor is supplied, {@code ping} events are ignored. *

+ * Reconnection is managed here rather than by the EventSource library. A transport failure, a + * recoverable HTTP status, or a bad payload is reported as INTERRUPTED and followed by a + * reconnect after a backoff. An HTTP status that is not expected to resolve soon, such as 401, is + * reported as TERMINAL_ERROR and ends this synchronizer. */ final class FDv2StreamingSynchronizer implements Synchronizer { private static final String METHOD_REPORT = "REPORT"; private static final String PING = "ping"; private static final long READ_TIMEOUT_MS = 300_000; // 5 minutes - private static final long MAX_RECONNECT_TIME_MS = 300_000; // 5 minutes private final HttpProperties httpProperties; private final URI streamBaseUri; @@ -67,21 +68,28 @@ final class FDv2StreamingSynchronizer implements Synchronizer { @Nullable private final FDv2Requestor requestor; private final boolean evaluationReasons; - private final int initialReconnectDelayMillis; @Nullable private final DiagnosticStore diagnosticStore; private final LDLogger logger; - private final Executor executor; + private final ScheduledExecutorService executor; private final LDAsyncQueue resultQueue = new LDAsyncQueue<>(); private final LDAwaitFuture shutdownFuture = new LDAwaitFuture<>(); private final AtomicBoolean started = new AtomicBoolean(false); + + // The following are only touched by the current connection attempt. Attempts never overlap. private final FDv2ProtocolHandler protocolHandler = new FDv2ProtocolHandler(); + private final StreamingRetryState retryState; + // Set while handling a message when the current connection must be dropped and a new one + // made, along with the wait before doing so. + private boolean restartRequested = false; + private long restartDelayMillis = 0; - // closeLock guards: closed and the eventSource assignment in startStream. + // closeLock guards closed, eventSource, and pendingAttempt. private final Object closeLock = new Object(); private boolean closed = false; - private volatile EventSource eventSource; + private EventSource eventSource; + private ScheduledFuture pendingAttempt; private volatile long streamStarted = 0; /** @@ -90,12 +98,13 @@ final class FDv2StreamingSynchronizer implements Synchronizer { * @param streamBaseUri base URI for the stream endpoint * @param streamRequestPath path appended to the base URI for the stream request * @param requestor optional requestor for handling ping events via poll; may be null - * @param initialReconnectDelayMillis delay before reconnecting after an error, in milliseconds + * @param initialReconnectDelayMillis initial delay before reconnecting after a failure, in + * milliseconds * @param evaluationReasons true to request evaluation reasons in the stream * @param useReport true to use HTTP REPORT for the request body * @param httpProperties HTTP configuration for the stream request - * @param executor executor used to run the streaming loop on a background - * thread; should use background-priority threads + * @param executor executor for the connection attempts and the waits + * between them. Should use background-priority threads. * @param logger logger * @param diagnosticStore optional store for stream diagnostics; may be null */ @@ -109,7 +118,7 @@ final class FDv2StreamingSynchronizer implements Synchronizer { boolean evaluationReasons, boolean useReport, @NonNull HttpProperties httpProperties, - @NonNull Executor executor, + @NonNull ScheduledExecutorService executor, @NonNull LDLogger logger, @Nullable DiagnosticStore diagnosticStore ) { @@ -118,13 +127,13 @@ final class FDv2StreamingSynchronizer implements Synchronizer { this.streamBaseUri = streamBaseUri; this.streamRequestPath = streamRequestPath; this.requestor = requestor; - this.initialReconnectDelayMillis = initialReconnectDelayMillis; this.evaluationReasons = evaluationReasons; this.useReport = useReport; this.httpProperties = httpProperties; this.executor = executor; this.logger = logger; this.diagnosticStore = diagnosticStore; + this.retryState = new StreamingRetryState(initialReconnectDelayMillis); } @Override @@ -134,7 +143,7 @@ public Future next() { shouldStart = !closed && !started.getAndSet(true); } if (shouldStart) { - startStream(); + executor.execute(this::runConnectionAttempt); } return LDFutures.anyOf(shutdownFuture, resultQueue.take()); } @@ -149,6 +158,10 @@ public void close() { closed = true; esToClose = eventSource; eventSource = null; + if (pendingAttempt != null) { + pendingAttempt.cancel(false); + pendingAttempt = null; + } } if (esToClose != null) { esToClose.close(); @@ -162,7 +175,51 @@ public void close() { shutdownFuture.set(FDv2SourceResult.status(FDv2SourceResult.Status.shutdown(), false)); } - private void startStream() { + private boolean isClosed() { + synchronized (closeLock) { + return closed; + } + } + + /** + * Runs one connection attempt and schedules the next one after the backoff delay. Only + * {@link #close()} stops the sequence. + */ + private void runConnectionAttempt() { + EventSource es = buildEventSource(); + synchronized (closeLock) { + if (closed) { + es.close(); + return; + } + eventSource = es; + } + streamStarted = System.currentTimeMillis(); + + long delay; + try { + delay = readUntilDisconnected(es); + } finally { + es.close(); + synchronized (closeLock) { + if (eventSource == es) { + eventSource = null; + } + } + } + + if (delay < 0) { + return; + } + synchronized (closeLock) { + if (closed) { + return; + } + pendingAttempt = executor.schedule(this::runConnectionAttempt, delay, TimeUnit.MILLISECONDS); + } + } + + private EventSource buildEventSource() { HttpConnectStrategy connectStrategy = ConnectStrategy.http(getStreamUri()) .clientBuilderActions(clientBuilder -> { httpProperties.applyToHttpClientBuilder(clientBuilder); @@ -194,52 +251,57 @@ private void startStream() { RequestBody.create(JsonSerialization.serialize(evaluationContext), JSON)); } - EventSource es = new EventSource.Builder(connectStrategy) - .retryDelay(initialReconnectDelayMillis, TimeUnit.MILLISECONDS) - .retryDelayStrategy(RetryDelayStrategy.defaultStrategy() - .maxDelay(MAX_RECONNECT_TIME_MS, TimeUnit.MILLISECONDS)) - .errorStrategy(ErrorStrategy.alwaysContinue()) - .build(); - - synchronized (closeLock) { - if (closed) { - es.close(); - return; - } - eventSource = es; - } + // The default error strategy throws every failure from readAnyEvent() instead of + // reconnecting, which leaves the backoff to this class. + return new EventSource.Builder(connectStrategy).build(); + } - executor.execute(() -> { - streamStarted = System.currentTimeMillis(); - try { - for (StreamEvent event : es.anyEvents()) { - if (closed) { - break; - } - if (event instanceof MessageEvent) { - handleMessage((MessageEvent) event); - } else if (event instanceof FaultEvent) { - handleError((FaultEvent) event); - } - // CommentEvent (SSE comment/heartbeat line) — no action needed + /** + * Reads events from one connection until it ends. + * + * @return the wait in milliseconds before the next connection attempt, or -1 if the + * synchronizer has been closed and must not reconnect + */ + private long readUntilDisconnected(EventSource es) { + try { + while (true) { + StreamEvent event = es.readAnyEvent(); + if (isClosed()) { + return -1; } - } catch (Exception e) { - synchronized (closeLock) { - if (closed) { - return; + if (event instanceof MessageEvent) { + handleMessage((MessageEvent) event); + if (restartRequested) { + restartRequested = false; + return restartDelayMillis; } } - LDUtil.logExceptionAtErrorLevel(logger, e, "Stream thread ended with unexpected exception"); - recordStreamInit(true); - resultQueue.put(FDv2SourceResult.status( - FDv2SourceResult.Status.interrupted( - new LDFailure("Stream thread ended unexpectedly", e, - LDFailure.FailureType.UNKNOWN_ERROR)), - false)); - } finally { - es.close(); + // StartedEvent and CommentEvent (SSE comment/heartbeat line) need no action. } - }); + } catch (StreamException e) { + if (isClosed() || e instanceof StreamClosedByCallerException) { + return -1; + } + return handleError(e); + } catch (RuntimeException e) { + if (isClosed()) { + return -1; + } + LDUtil.logExceptionAtErrorLevel(logger, e, "Unexpected exception while reading stream"); + recordStreamInit(true); + protocolHandler.reset(); + resultQueue.put(FDv2SourceResult.status( + FDv2SourceResult.Status.interrupted( + new LDFailure("Unexpected exception while reading stream", e, + LDFailure.FailureType.UNKNOWN_ERROR)), + false)); + return recordFailureAndGetDelay(false); + } + } + + private long recordFailureAndGetDelay(boolean unexpected) { + retryState.recordFailure(unexpected, System.currentTimeMillis()); + return retryState.nextDelayMillis(); } private URI getStreamUri() { @@ -274,6 +336,8 @@ void handleMessage(MessageEvent event) { String eventData = event.getData(); logger.debug("onMessage: {}: {}", eventName, eventData); + retryState.recordSuccess(System.currentTimeMillis()); + if (PING.equalsIgnoreCase(eventName)) { handlePing(); return; @@ -387,37 +451,33 @@ private void handlePing() { resultQueue.put(result); } - private void handleError(FaultEvent event) { - StreamException t = event.getCause(); - if (t instanceof StreamClosedByCallerException) { - return; - } - + /** + * Classifies a connection failure and reports it. + * + * @return the wait in milliseconds before the next attempt, or -1 if this synchronizer has + * ended and must not reconnect + */ + private long handleError(StreamException t) { recordStreamInit(true); protocolHandler.reset(); - boolean fdv1Fallback = isFdv1Fallback(event.getHeaders()); + boolean fdv1Fallback = false; + LDFailure failure; if (t instanceof StreamHttpErrorException) { StreamHttpErrorException httpError = (StreamHttpErrorException) t; - fdv1Fallback = fdv1Fallback || isFdv1Fallback(httpError.getHeaders()); + fdv1Fallback = isFdv1Fallback(httpError.getHeaders()); int code = httpError.getCode(); boolean recoverable = LDUtil.isHttpErrorRecoverable(code); - LDFailure failure = new LDInvalidResponseCodeFailure( + failure = new LDInvalidResponseCodeFailure( "Unexpected response code from stream", t, code, recoverable); if (!recoverable) { logger.error("Encountered non-retriable error: {}. Aborting connection to stream. Verify correct Mobile Key and Stream URI", code); shutdownFuture.set(FDv2SourceResult.status( FDv2SourceResult.Status.terminalError(failure), fdv1Fallback)); - EventSource es; synchronized (closeLock) { closed = true; - es = eventSource; - eventSource = null; - } - if (es != null) { - es.close(); } if (requestor != null) { try { @@ -425,45 +485,33 @@ private void handleError(FaultEvent event) { } catch (IOException ignored) { } } - } else { - logger.warn("Stream received HTTP error {}; will retry", code); - streamStarted = System.currentTimeMillis(); - resultQueue.put(FDv2SourceResult.status(FDv2SourceResult.Status.interrupted(failure), fdv1Fallback)); + return -1; } + logger.warn("Stream received HTTP error {}; will retry", code); } else { + // Every transport-level failure, including the server closing the connection, is + // recoverable. LDUtil.logExceptionAtWarnLevel(logger, t, "Stream network error"); - streamStarted = System.currentTimeMillis(); - resultQueue.put(FDv2SourceResult.status( - FDv2SourceResult.Status.interrupted( - new LDFailure("Stream network error", t, - LDFailure.FailureType.NETWORK_FAILURE)), - fdv1Fallback)); + failure = new LDFailure("Stream network error", t, LDFailure.FailureType.NETWORK_FAILURE); } + + long delay = recordFailureAndGetDelay(false); + logger.info("Will reconnect to stream in {} ms", delay); + resultQueue.put(FDv2SourceResult.status(FDv2SourceResult.Status.interrupted(failure), fdv1Fallback)); + return delay; } /** - * Interrupts the current connection so the EventSource reconnects immediately on the - * streaming thread, and resets the diagnostic timer for the new connection attempt. - * {@link EventSource#interrupt()} is safe to call from the streaming thread itself. + * Asks the connection attempt to drop the current connection once the current message has + * been handled, and to reconnect after a backoff. * * @param failed true if the restart is due to an error (for diagnostic recording) */ private void restartStream(boolean failed) { recordStreamInit(failed); - streamStarted = System.currentTimeMillis(); - - EventSource es; - synchronized (closeLock) { - if (closed) { - return; - } - es = eventSource; - } - - if (es != null) { - es.interrupt(); - } protocolHandler.reset(); + restartRequested = true; + restartDelayMillis = recordFailureAndGetDelay(false); } @Override diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/HttpFeatureFlagFetcher.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/HttpFeatureFlagFetcher.java index 782807dc..42934477 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/HttpFeatureFlagFetcher.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/HttpFeatureFlagFetcher.java @@ -122,7 +122,8 @@ public void onResponse(@NonNull Call call, @NonNull final Response response) { logger.error("Received 400 response when fetching flag values. Please check recommended ProGuard settings"); } callback.onError(new LDInvalidResponseCodeFailure("Unexpected response when retrieving Feature Flags: " + response + " using url: " - + request.url() + " with body: " + body, response.code(), true)); + + request.url() + " with body: " + body, response.code(), + LDUtil.isHttpErrorRecoverable(response.code()))); return; } logger.debug(body); diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDInvalidResponseCodeFailure.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDInvalidResponseCodeFailure.java index 6c5b6f38..d1168162 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDInvalidResponseCodeFailure.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDInvalidResponseCodeFailure.java @@ -14,14 +14,14 @@ public class LDInvalidResponseCodeFailure extends LDFailure { private final int responseCode; /** - * Whether or not the failure may be fixed by retrying + * Whether the failure is one that may resolve on its own if retried soon. */ private final boolean retryable; /** * @param message the message * @param responseCode the response code - * @param retryable whether or not retrying may resolve the issue + * @param retryable whether the failure may resolve on its own if retried soon */ public LDInvalidResponseCodeFailure(String message, int responseCode, boolean retryable) { super(message, FailureType.UNEXPECTED_RESPONSE_CODE); @@ -33,7 +33,7 @@ public LDInvalidResponseCodeFailure(String message, int responseCode, boolean re * @param message the message * @param cause the cause of the failure * @param responseCode the response code - * @param retryable whether or not retrying may resolve the issue + * @param retryable whether the failure may resolve on its own if retried soon */ public LDInvalidResponseCodeFailure(String message, Throwable cause, int responseCode, boolean retryable) { super(message, cause, FailureType.UNEXPECTED_RESPONSE_CODE); @@ -42,7 +42,8 @@ public LDInvalidResponseCodeFailure(String message, Throwable cause, int respons } /** - * @return true if retrying may resolve the issue + * @return true if the failure may resolve on its own if retried soon. An authentication + * failure, for example, is unlikely to. */ public boolean isRetryable() { return retryable; diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDUtil.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDUtil.java index ae8ed1f5..71a56df1 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDUtil.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/LDUtil.java @@ -212,9 +212,11 @@ public void updateHeaders(Map headers) { } /** - * Tests whether an HTTP error status represents a condition that might resolve on its own if we retry. + * Classifies an HTTP error status as either a {@code normal} failure, which may resolve on its + * own if retried soon, or an {@code unexpected} failure, which is not expected to. + * * @param statusCode the HTTP status - * @return true if retrying makes sense; false if it should be considered a permanent failure + * @return true if the status is {@code normal}, or false if it is {@code unexpected} */ static boolean isHttpErrorRecoverable(int statusCode) { if (statusCode >= 400 && statusCode < 500) { @@ -230,6 +232,18 @@ static boolean isHttpErrorRecoverable(int statusCode) { return true; } + /** + * Classifies a data source failure as either a {@code normal} failure, which may resolve on + * its own if retried soon, or an {@code unexpected} failure, which is not expected to. + * + * @param failure the failure reported by a data source, or null + * @return true if the failure is {@code unexpected}, or false if it is {@code normal} + */ + static boolean isUnexpectedFailure(Throwable failure) { + return failure instanceof LDInvalidResponseCodeFailure && + !isHttpErrorRecoverable(((LDInvalidResponseCodeFailure) failure).getResponseCode()); + } + static void logExceptionAtErrorLevel(LDLogger logger, Throwable ex, String msgFormat, Object... msgArgs) { logException(logger, ex, true, msgFormat, msgArgs); } diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingDataSource.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingDataSource.java index a68697d2..86ae8be6 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingDataSource.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingDataSource.java @@ -16,6 +16,10 @@ * streaming with Components.pollingDataSource(), or 2. streaming is enabled, but the application is * in the background so we do polling instead. The logic for this is in * ComponentsImpl.PollingDataSourceBuilderImpl and ComponentsImpl.StreamingDataSourceBuilderImpl. + *

+ * Polls happen at the poll interval, except that a failure that is not expected to resolve on + * its own, such as HTTP 401, is followed by a much longer wait. No failure stops the data source + * from polling again. */ final class PollingDataSource implements DataSource { private final LDContext context; @@ -29,6 +33,10 @@ final class PollingDataSource implements DataSource { private final LDLogger logger; final AtomicReference> currentPollTask = new AtomicReference<>(); // visible for testing + // Guarded by the lock on this instance. + private final PollingRetryState retryState; + private boolean running = false; + /** * @param context that this data source will fetch data for * @param dataSourceUpdateSink to send data to @@ -52,6 +60,28 @@ final class PollingDataSource implements DataSource { PlatformState platformState, TaskExecutor taskExecutor, LDLogger logger + ) { + this(context, dataSourceUpdateSink, initialDelayMillis, pollIntervalMillis, maxNumberOfPolls, + fetcher, platformState, taskExecutor, new PollingRetryState(pollIntervalMillis), logger); + } + + /** + * This constructor allows tests to supply a {@link PollingRetryState} with short delays. See the + * other constructor for the remaining parameters. + * + * @param retryState the retry state that decides the wait after a failed poll + */ + PollingDataSource( + LDContext context, + DataSourceUpdateSink dataSourceUpdateSink, + long initialDelayMillis, + long pollIntervalMillis, + long maxNumberOfPolls, + FeatureFetcher fetcher, + PlatformState platformState, + TaskExecutor taskExecutor, + PollingRetryState retryState, + LDLogger logger ) { this.context = context; this.dataSourceUpdateSink = dataSourceUpdateSink; @@ -61,6 +91,7 @@ final class PollingDataSource implements DataSource { this.fetcher = fetcher; this.platformState = platformState; this.taskExecutor = taskExecutor; + this.retryState = retryState; this.logger = logger; } @@ -73,16 +104,19 @@ public void start(final Callback resultCallback) { return; } - Runnable pollRunnable = () -> poll(resultCallback); + synchronized (this) { + running = true; + } logger.debug("Scheduling polling task with interval of {}ms, starting after {}ms, with number of polls {}", pollIntervalMillis, initialDelayMillis, numberOfPollsRemaining); - ScheduledFuture task = taskExecutor.startRepeatingTask(pollRunnable, - initialDelayMillis, pollIntervalMillis); - currentPollTask.set(task); + schedulePoll(initialDelayMillis, resultCallback); } @Override public void stop(Callback completionCallback) { + synchronized (this) { + running = false; + } ScheduledFuture task = currentPollTask.getAndSet(null); if (task != null) { task.cancel(true); @@ -90,18 +124,62 @@ public void stop(Callback completionCallback) { completionCallback.onSuccess(null); } - private void poll(Callback resultCallback) { - // poll if there are polls remaining - if (numberOfPollsRemaining > 0) { + /** + * Schedules the next poll, unless the data source has been stopped or has used up its polls. + */ + private synchronized void schedulePoll(long delayMillis, Callback resultCallback) { + if (!running || numberOfPollsRemaining <= 0) { + return; + } + currentPollTask.set(taskExecutor.scheduleTask(() -> poll(resultCallback), delayMillis)); + } + + private void poll(final Callback resultCallback) { + synchronized (this) { + if (!running || numberOfPollsRemaining <= 0) { + return; + } numberOfPollsRemaining--; - ConnectivityManager.fetchAndSetData(fetcher, context, dataSourceUpdateSink, - resultCallback, logger); - } else { - // terminate if we have no polls remaining - ScheduledFuture task = currentPollTask.getAndSet(null); - if (task != null) { - task.cancel(true); + } + + Callback pollCallback = new Callback() { + @Override + public void onSuccess(Boolean result) { + long delay; + synchronized (PollingDataSource.this) { + retryState.recordSuccess(); + delay = retryState.nextDelayMillis(); + } + resultCallback.onSuccess(result); + schedulePoll(delay, resultCallback); } + + @Override + public void onError(Throwable error) { + boolean unexpected = LDUtil.isUnexpectedFailure(error); + long delay; + synchronized (PollingDataSource.this) { + retryState.recordFailure(unexpected); + delay = retryState.nextDelayMillis(); + } + if (unexpected) { + logger.error("Received HTTP error {} from polling request. This is not expected to resolve on its own; verify correct Mobile Key and Polling URI. Will retry in {} ms.", + ((LDInvalidResponseCodeFailure) error).getResponseCode(), delay); + } else { + logger.warn("Polling request failed. Will retry in {} ms.", delay); + } + resultCallback.onError(error); + schedulePoll(delay, resultCallback); + } + }; + + try { + ConnectivityManager.fetchAndSetData(fetcher, context, dataSourceUpdateSink, pollCallback, logger); + } catch (RuntimeException e) { + // If the fetcher throws instead of using its callback, the caller is still owed a + // result and the next poll. + LDUtil.logExceptionAtErrorLevel(logger, e, "Unexpected exception while polling for flags"); + pollCallback.onError(new LDFailure("Exception while fetching flags", e, LDFailure.FailureType.UNKNOWN_ERROR)); } } } diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingRetryState.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingRetryState.java new file mode 100644 index 00000000..13070303 --- /dev/null +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/PollingRetryState.java @@ -0,0 +1,90 @@ +package com.launchdarkly.sdk.android; + +import androidx.annotation.NonNull; + +import java.util.Random; + +/** + * Computes the wait before a polling component polls again. The backoff clears after enough + * consecutive successful polls. + *

+ * This class is not thread-safe. The owning component must serialize access to it. + */ +final class PollingRetryState { + /** + * The backoff clears after this many consecutive successful polls. + */ + static final int RESET_THRESHOLD_SUCCESSES = 2; + + private final long pollIntervalMillis; + private final RetryRegime extended; + private final Random random; + + private boolean inExtendedRegime = false; + private int attempts = 0; + private int consecutiveSuccesses = 0; + + /** + * @param pollIntervalMillis the configured poll interval, which must be positive + */ + PollingRetryState(long pollIntervalMillis) { + this(pollIntervalMillis, RetryRegime.EXTENDED, new Random()); + } + + /** + * Creates a retry state with an explicit extended regime. + * + * @param pollIntervalMillis the configured poll interval, which must be positive + * @param extended the regime entered by an {@code unexpected} failure, raised to the + * poll interval where smaller + * @param random source of jitter + */ + PollingRetryState( + long pollIntervalMillis, + @NonNull RetryRegime extended, + @NonNull Random random + ) { + this.pollIntervalMillis = Math.max(1, pollIntervalMillis); + this.extended = extended.atLeast(this.pollIntervalMillis); + this.random = random; + } + + /** + * Records a successful poll. The backoff clears after {@link #RESET_THRESHOLD_SUCCESSES} + * consecutive successes. + */ + void recordSuccess() { + consecutiveSuccesses++; + if (consecutiveSuccesses >= RESET_THRESHOLD_SUCCESSES) { + inExtendedRegime = false; + attempts = 0; + } + } + + /** + * Records a failed poll. Call this before {@link #nextDelayMillis()}. + * + * @param unexpected true if the failure is {@code unexpected}, false if {@code normal} + */ + void recordFailure(boolean unexpected) { + consecutiveSuccesses = 0; + if (unexpected && !inExtendedRegime) { + inExtendedRegime = true; + attempts = 1; + return; + } + attempts++; + } + + /** + * @return the wait in milliseconds before the next poll + */ + long nextDelayMillis() { + // A normal failure does not grow the delay, and a success shows the service works. Only + // an unresolved unexpected failure waits longer than the poll interval. + if (!inExtendedRegime || consecutiveSuccesses > 0) { + return pollIntervalMillis; + } + return Math.max(pollIntervalMillis, extended.jitteredDelayMillis(attempts, random)); + } +} diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/RetryRegime.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/RetryRegime.java new file mode 100644 index 00000000..e07f8d2e --- /dev/null +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/RetryRegime.java @@ -0,0 +1,72 @@ +package com.launchdarkly.sdk.android; + +import androidx.annotation.NonNull; + +import java.util.Random; + +/** + * The delay bounds for one retry regime. + */ +final class RetryRegime { + /** + * The base delay for the first attempt after an {@code unexpected} failure. + */ + static final long EXTENDED_INITIAL_DELAY_MILLIS = 5L * 60_000L; + /** + * The largest delay after an {@code unexpected} failure. + */ + static final long EXTENDED_MAX_DELAY_MILLIS = 60L * 60_000L; + /** + * The regime entered after an {@code unexpected} failure. + */ + static final RetryRegime EXTENDED = + new RetryRegime(EXTENDED_INITIAL_DELAY_MILLIS, EXTENDED_MAX_DELAY_MILLIS); + + // Caps the doubling so that initialDelayMillis * 2^exponent cannot overflow. + private static final int MAX_EXPONENT = 30; + + /** + * The base delay for the first attempt in this regime. + */ + final long initialDelayMillis; + /** + * The largest delay this regime produces. + */ + final long maxDelayMillis; + + /** + * @param initialDelayMillis the base delay for the first attempt, floored at zero + * @param maxDelayMillis the largest delay, raised to the initial delay if smaller + */ + RetryRegime(long initialDelayMillis, long maxDelayMillis) { + this.initialDelayMillis = Math.max(0, initialDelayMillis); + this.maxDelayMillis = Math.max(this.initialDelayMillis, maxDelayMillis); + } + + /** + * @param floorMillis the smallest value allowed for either bound + * @return a regime with the same bounds, except that neither is below the floor + */ + @NonNull + RetryRegime atLeast(long floorMillis) { + return new RetryRegime( + Math.max(initialDelayMillis, floorMillis), Math.max(maxDelayMillis, floorMillis)); + } + + /** + * The exponential backoff for an attempt, less a random amount of up to half of it. + * + * @param attempts the number of failures so far, starting at 1 + * @param random source of jitter + * @return the wait in milliseconds + */ + long jitteredDelayMillis(int attempts, @NonNull Random random) { + int exponent = Math.min(Math.max(0, attempts - 1), MAX_EXPONENT); + long base = initialDelayMillis * (1L << exponent); + if (base < 0 || base > maxDelayMillis) { // negative means the multiplication overflowed + base = maxDelayMillis; + } + long jitter = base > 1 ? (long) (random.nextDouble() * (base / 2.0)) : 0; + return base - jitter; + } +} diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SourceManager.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SourceManager.java index bbbd8adb..67e75f5a 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SourceManager.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SourceManager.java @@ -1,6 +1,7 @@ package com.launchdarkly.sdk.android; import androidx.annotation.NonNull; +import androidx.annotation.Nullable; import com.launchdarkly.sdk.android.subsystems.Initializer; import com.launchdarkly.sdk.android.subsystems.Synchronizer; @@ -8,18 +9,26 @@ import java.io.Closeable; import java.io.IOException; import java.util.List; +import java.util.concurrent.Future; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; /** * Manages the state of synchronizers and initializers: tracks which is active, - * advances through the lists (with optional block state for synchronizers), + * advances through the lists (skipping synchronizers that are blocked or backing off), * and closes the previous source when switching. *

+ * A synchronizer that reports an unexpected error is put into a backoff, and its slot is skipped + * until the backoff ends. No synchronizer is ever permanently removed. + *

* Package-private for internal use by FDv2DataSource. */ final class SourceManager implements Closeable { private final List synchronizerFactories; private final List> initializers; + private final ScheduledExecutorService executor; private final Object activeSourceLock = new Object(); private Closeable activeSource; @@ -31,12 +40,19 @@ final class SourceManager implements Closeable { private SynchronizerFactoryWithState currentSynchronizerFactory; + // Handed out by nextAvailableSynchronizer() while every usable slot is backing off, and + // completed by the first backoff to end or by close(). + @Nullable + private LDAwaitFuture pendingNext; + SourceManager( @NonNull List synchronizerFactories, - @NonNull List> initializers + @NonNull List> initializers, + @NonNull ScheduledExecutorService executor ) { this.synchronizerFactories = synchronizerFactories; this.initializers = initializers; + this.executor = executor; } /** @@ -49,7 +65,7 @@ void resetSourceIndex() { } } - /** True if any synchronizer is marked as FDv1 fallback (Android: not used yet). */ + /** True if any synchronizer is marked as FDv1 fallback. */ boolean hasFDv1Fallback() { for (SynchronizerFactoryWithState s : synchronizerFactories) { if (s.isFDv1Fallback()) { @@ -61,8 +77,8 @@ boolean hasFDv1Fallback() { /** * Block all non-FDv1 synchronizers, unblock the FDv1 fallback, and reset the - * synchronizer index so the next {@link #getNextAvailableSynchronizerAndSetActive()} - * picks the now-unblocked FDv1 slot. + * synchronizer index so the next {@link #nextAvailableSynchronizer()} picks the now-unblocked + * FDv1 slot. */ void fdv1Fallback() { synchronized (activeSourceLock) { @@ -70,6 +86,7 @@ void fdv1Fallback() { if (s.isFDv1Fallback()) { s.unblock(); } else { + s.cancelPendingUnblock(); s.block(); } } @@ -100,7 +117,7 @@ private SynchronizerFactoryWithState getNextAvailableSynchronizer() { * and return it. Returns null if shutdown or no available synchronizers. * Skips synchronizers whose factory returns null from build(). */ - Synchronizer getNextAvailableSynchronizerAndSetActive() { + private Synchronizer getNextAvailableSynchronizerAndSetActive() { synchronized (activeSourceLock) { if (isShutdown) { currentSynchronizerFactory = null; @@ -131,6 +148,34 @@ Synchronizer getNextAvailableSynchronizerAndSetActive() { } } + /** + * Selects the synchronizer to run next, builds it, and makes it the active source in place of + * the previous one. If every usable slot is backing off, the returned future completes when + * the first backoff ends. It completes with null once this manager is closed or no + * synchronizer is left to try. + * + * @return a future for the next synchronizer + */ + @NonNull + Future nextAvailableSynchronizer() { + synchronized (activeSourceLock) { + Synchronizer synchronizer = getNextAvailableSynchronizerAndSetActive(); + if (synchronizer != null || !hasBackingOffSynchronizers()) { + return completed(synchronizer); + } + if (pendingNext == null) { + pendingNext = new LDAwaitFuture<>(); + } + return pendingNext; + } + } + + private static LDAwaitFuture completed(@Nullable T value) { + LDAwaitFuture future = new LDAwaitFuture<>(); + future.set(value); + return future; + } + boolean hasAvailableSources() { return hasInitializers() || getAvailableSynchronizerCount() > 0; } @@ -151,12 +196,97 @@ private FDv2DataSource.DataSourceFactory getNextInitializer() { return initializers.get(initializerIndex); } - /** Block the current synchronizer so it will not be returned again (e.g. after TERMINAL_ERROR). */ - void blockCurrentSynchronizer() { + /** + * Puts the current synchronizer's slot into backoff. The slot is skipped by + * {@link #nextAvailableSynchronizer()} until the backoff ends. + * + * @param nowMillis the current time in milliseconds, on the same clock as every other call on + * this manager + * @return how long the slot stays in backoff, in milliseconds, or -1 if there is no current + * synchronizer + */ + long backOffCurrentSynchronizer(long nowMillis) { + synchronized (activeSourceLock) { + final SynchronizerFactoryWithState slot = currentSynchronizerFactory; + if (slot == null || isShutdown) { + return -1; + } + long delayMillis = slot.startBackoff(nowMillis); + ScheduledFuture unblock = executor.schedule(new Runnable() { + @Override + public void run() { + endBackoff(slot); + } + }, delayMillis, TimeUnit.MILLISECONDS); + slot.setPendingUnblock(unblock); + return delayMillis; + } + } + + /** + * Records healthy operation by the current synchronizer, which counts toward resetting its + * slot's backoff. + * + * @param nowMillis the current time in milliseconds, on the same clock as every other call on + * this manager + */ + void recordCurrentSynchronizerHealthy(long nowMillis) { synchronized (activeSourceLock) { if (currentSynchronizerFactory != null) { - currentSynchronizerFactory.block(); + currentSynchronizerFactory.recordHealthy(nowMillis); + } + } + } + + /** + * Ends a slot's backoff. If a caller is waiting in {@link #nextAvailableSynchronizer()}, the + * next synchronizer is selected and handed to it. + */ + void endBackoff(@NonNull SynchronizerFactoryWithState slot) { + LDAwaitFuture waiting = null; + Synchronizer synchronizer = null; + synchronized (activeSourceLock) { + if (isShutdown || !slot.endBackoff()) { + return; + } + if (pendingNext != null) { + synchronizer = getNextAvailableSynchronizerAndSetActive(); + if (synchronizer != null || !hasBackingOffSynchronizers()) { + waiting = pendingNext; + pendingNext = null; + } + } + } + if (waiting != null) { + waiting.set(synchronizer); + } + } + + private boolean hasBackingOffSynchronizers() { + synchronized (activeSourceLock) { + if (isShutdown) { + return false; + } + for (SynchronizerFactoryWithState s : synchronizerFactories) { + if (s.getState() == SynchronizerFactoryWithState.State.BackingOff) { + return true; + } + } + return false; + } + } + + /** + * @return true if a synchronizer earlier in the list than the current one is available + */ + boolean hasAvailableSynchronizerBeforeCurrent() { + synchronized (activeSourceLock) { + for (int i = 0; i < synchronizerIndex && i < synchronizerFactories.size(); i++) { + if (synchronizerFactories.get(i).getState() == SynchronizerFactoryWithState.State.Available) { + return true; + } } + return false; } } @@ -193,11 +323,14 @@ Initializer getNextInitializerAndSetActive() { } } - /** True if the current synchronizer is the first available one (prime). */ + /** + * True if the current synchronizer is the prime one, meaning that no synchronizer before it + * in the list is available or backing off. + */ boolean isPrimeSynchronizer() { synchronized (activeSourceLock) { for (int i = 0; i < synchronizerFactories.size(); i++) { - if (synchronizerFactories.get(i).getState() == SynchronizerFactoryWithState.State.Available) { + if (synchronizerFactories.get(i).getState() != SynchronizerFactoryWithState.State.Blocked) { return synchronizerIndex == i; } } @@ -217,14 +350,38 @@ int getAvailableSynchronizerCount() { } } + /** + * @return the number of synchronizers that are available or backing off + */ + int getUsableSynchronizerCount() { + synchronized (activeSourceLock) { + int count = 0; + for (SynchronizerFactoryWithState s : synchronizerFactories) { + if (s.getState() != SynchronizerFactoryWithState.State.Blocked) { + count++; + } + } + return count; + } + } + @Override public void close() { + LDAwaitFuture waiting; synchronized (activeSourceLock) { isShutdown = true; if (activeSource != null) { safeClose(activeSource); activeSource = null; } + for (SynchronizerFactoryWithState s : synchronizerFactories) { + s.cancelPendingUnblock(); + } + waiting = pendingNext; + pendingNext = null; + } + if (waiting != null) { + waiting.set(null); } } diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingDataSource.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingDataSource.java index 366254b2..32bbd124 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingDataSource.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingDataSource.java @@ -8,7 +8,7 @@ import com.launchdarkly.eventsource.EventSource; import com.launchdarkly.eventsource.HttpConnectStrategy; import com.launchdarkly.eventsource.MessageEvent; -import com.launchdarkly.eventsource.RetryDelayStrategy; +import com.launchdarkly.eventsource.StreamException; import com.launchdarkly.eventsource.StreamHttpErrorException; import com.launchdarkly.eventsource.background.BackgroundEventHandler; import com.launchdarkly.eventsource.background.BackgroundEventSource; @@ -27,6 +27,7 @@ import com.launchdarkly.sdk.json.SerializationException; import java.net.URI; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import okhttp3.RequestBody; @@ -39,6 +40,10 @@ *

* The SDK uses this implementation if streaming is enabled (as it is by default) and the * application is the foreground. The logic for this is in ComponentsImpl.StreamingDataSourceBuilderImpl. + *

+ * Reconnection is managed by this class rather than by the EventSource library. Every stream + * failure is followed by a reconnect after a backoff. Failures that are unlikely to resolve on + * their own, such as HTTP 401, wait much longer. */ final class StreamingDataSource implements DataSource { private static final String METHOD_REPORT = "REPORT"; @@ -48,13 +53,10 @@ final class StreamingDataSource implements DataSource { private static final String PATCH = "patch"; private static final String DELETE = "delete"; - private static final long MAX_RECONNECT_TIME_MS = 300_000; // 5 minutes - private static final long READ_TIMEOUT_MS = 300_000; // 5 minutes is the standard read timeout used for all LaunchDarkly stream connections, based on // an expectation that the server will send heartbeats at a shorter interval than that - private BackgroundEventSource es; private final LDContext context; private final HttpProperties httpProperties; private final boolean evaluationReasons; @@ -64,14 +66,20 @@ final class StreamingDataSource implements DataSource { private final DataSourceUpdateSink dataSourceUpdateSink; private final FeatureFetcher fetcher; private final boolean streamEvenInBackground; - private volatile boolean running = false; - // volatile because it is written from the EventSource background thread (onError) - // and read from start(), which may be invoked on a different thread. - private volatile boolean connection401Error = false; private final DiagnosticStore diagnosticStore; - private long eventSourceStarted; + private final TaskExecutor taskExecutor; private final LDLogger logger; + // The following fields are guarded by the lock on this instance. + private final StreamingRetryState retryState; + private BackgroundEventSource es; + private BackgroundEventHandler handler; + private ScheduledFuture pendingReconnect; + private boolean running = false; + + // Written when a connection attempt begins, read by the EventSource callbacks for diagnostics. + private volatile long eventSourceStarted; + StreamingDataSource( @NonNull ClientContext clientContext, @NonNull LDContext context, @@ -79,6 +87,22 @@ final class StreamingDataSource implements DataSource { @NonNull FeatureFetcher fetcher, int initialReconnectDelayMillis, boolean streamEvenInBackground + ) { + this(clientContext, context, dataSourceUpdateSink, fetcher, initialReconnectDelayMillis, + streamEvenInBackground, new StreamingRetryState(initialReconnectDelayMillis)); + } + + /** + * This constructor allows tests to supply a {@link StreamingRetryState} with short delays. + */ + StreamingDataSource( + @NonNull ClientContext clientContext, + @NonNull LDContext context, + @NonNull DataSourceUpdateSink dataSourceUpdateSink, + @NonNull FeatureFetcher fetcher, + int initialReconnectDelayMillis, + boolean streamEvenInBackground, + @NonNull StreamingRetryState retryState ) { this.context = context; this.dataSourceUpdateSink = dataSourceUpdateSink; @@ -90,111 +114,146 @@ final class StreamingDataSource implements DataSource { this.initialReconnectDelayMillis = initialReconnectDelayMillis; this.streamEvenInBackground = streamEvenInBackground; this.diagnosticStore = ClientContextImpl.get(clientContext).getDiagnosticStore(); + this.taskExecutor = ClientContextImpl.get(clientContext).getTaskExecutor(); + this.retryState = retryState; this.logger = clientContext.getBaseLogger(); } public void start(@NonNull Callback resultCallback) { - if (!running && !connection401Error) { - logger.debug("Starting."); + synchronized (this) { + if (running) { + return; + } + running = true; + handler = makeHandler(resultCallback); + } + logger.debug("Starting."); + connect(); + } - BackgroundEventHandler handler = new BackgroundEventHandler() { - @Override - public void onOpen() { - logger.info("Started LaunchDarkly EventStream"); - if (diagnosticStore != null) { - diagnosticStore.recordStreamInit(eventSourceStarted, (int) (System.currentTimeMillis() - eventSourceStarted), false); - } + private BackgroundEventHandler makeHandler(@NonNull Callback resultCallback) { + return new BackgroundEventHandler() { + @Override + public void onOpen() { + logger.info("Started LaunchDarkly EventStream"); + if (diagnosticStore != null) { + diagnosticStore.recordStreamInit(eventSourceStarted, (int) (System.currentTimeMillis() - eventSourceStarted), false); } + } - @Override - public void onClosed() { - logger.info("Closed LaunchDarkly EventStream"); - } + @Override + public void onClosed() { + logger.info("Closed LaunchDarkly EventStream"); + } - @Override - public void onMessage(final String name, MessageEvent event) { - final String eventData = event.getData(); - logger.debug("onMessage: {}: {}", name, eventData); - handle(name, eventData, resultCallback); + @Override + public void onMessage(final String name, MessageEvent event) { + final String eventData = event.getData(); + logger.debug("onMessage: {}: {}", name, eventData); + synchronized (StreamingDataSource.this) { + retryState.recordSuccess(System.currentTimeMillis()); } + handle(name, eventData, resultCallback); + } + + @Override + public void onComment(String comment) { + // intentionally empty + } - @Override - public void onComment(String comment) { - // intentionally empty + @Override + public void onError(Throwable t) { + LDUtil.logExceptionAtErrorLevel(logger, t, + "Encountered EventStream error connecting to URI: {}", + getUri(context)); + + LDFailure failure; + boolean unexpected = false; + int code = 0; + if (t instanceof StreamHttpErrorException) { + if (diagnosticStore != null) { + diagnosticStore.recordStreamInit(eventSourceStarted, (int) (System.currentTimeMillis() - eventSourceStarted), true); + } + code = ((StreamHttpErrorException) t).getCode(); + boolean recoverable = LDUtil.isHttpErrorRecoverable(code); + unexpected = !recoverable; + failure = new LDInvalidResponseCodeFailure("Unexpected Response Code From Stream Connection", t, code, recoverable); + } else { + failure = new LDFailure("Network error in stream connection", t, LDFailure.FailureType.NETWORK_FAILURE); } - @Override - public void onError(Throwable t) { - LDUtil.logExceptionAtErrorLevel(logger, t, - "Encountered EventStream error connecting to URI: {}", - getUri(context)); - if (t instanceof StreamHttpErrorException) { - if (diagnosticStore != null) { - diagnosticStore.recordStreamInit(eventSourceStarted, (int) (System.currentTimeMillis() - eventSourceStarted), true); - } - int code = ((StreamHttpErrorException) t).getCode(); - if (!LDUtil.isHttpErrorRecoverable(code)) { - logger.error("Encountered non-retriable error: {}. Aborting connection to stream. Verify correct Mobile Key and Stream URI", code); - running = false; - // Set the connection401Error guard before notifying the callback. A consumer - // may react to the error by synchronously calling start() again, and the - // guard must already be set so that retry is a no-op. - if (code == 401) { - connection401Error = true; - } - resultCallback.onError(new LDInvalidResponseCodeFailure("Unexpected Response Code From Stream Connection", t, code, false)); - if (code == 401) { - dataSourceUpdateSink.shutDown(); - } - stop(null); + // Only a StreamException means the connection has ended. Anything else was thrown + // while processing an event on a stream that is still open. + if (t instanceof StreamException) { + long delay = scheduleReconnectAfterFailure(unexpected); + if (delay >= 0) { + if (unexpected) { + logger.error("Encountered HTTP error {} from stream. This is not expected to resolve on its own; verify correct Mobile Key and Stream URI. Will retry in {} ms.", code, delay); } else { - eventSourceStarted = System.currentTimeMillis(); - resultCallback.onError(new LDInvalidResponseCodeFailure("Unexpected Response Code From Stream Connection", t, code, true)); + logger.warn("Will retry stream connection in {} ms.", delay); } - } else { - resultCallback.onError(new LDFailure("Network error in stream connection", t, LDFailure.FailureType.NETWORK_FAILURE)); } } - }; - - HttpConnectStrategy connectStrategy = ConnectStrategy.http(getUri(context)) - .clientBuilderActions(clientBuilder -> { - httpProperties.applyToHttpClientBuilder(clientBuilder); - clientBuilder.readTimeout(READ_TIMEOUT_MS, TimeUnit.MILLISECONDS); - }) - .requestTransformer(input -> - input.newBuilder() - .headers( - input.headers().newBuilder().addAll(httpProperties.toHeadersBuilder().build()).build() - ).build()); - - if (useReport) { - connectStrategy = connectStrategy.methodAndBody(METHOD_REPORT, getRequestBody(context)); + resultCallback.onError(failure); } + }; + } - EventSource.Builder esBuilder = new EventSource.Builder(connectStrategy) - .retryDelay(initialReconnectDelayMillis, TimeUnit.MILLISECONDS) - .retryDelayStrategy(RetryDelayStrategy.defaultStrategy() - .maxDelay(MAX_RECONNECT_TIME_MS, TimeUnit.MILLISECONDS)); - - eventSourceStarted = System.currentTimeMillis(); - es = new BackgroundEventSource.Builder(handler, esBuilder) - // The stream thread asks this handler, before it reconnects, whether an - // error ends the stream. The onError callback above runs on a different - // thread, so a stop() from there can arrive after a fast reconnect has - // already opened a new connection. A decision made here cannot lose that race. - .connectionErrorHandler(t -> { - if (t instanceof StreamHttpErrorException && - !LDUtil.isHttpErrorRecoverable(((StreamHttpErrorException) t).getCode())) { - return ConnectionErrorHandler.Action.SHUTDOWN; - } - return ConnectionErrorHandler.Action.PROCEED; - }) - .build(); - es.start(); + /** + * Opens a new stream connection if this data source is still running. + */ + private synchronized void connect() { + if (!running) { + return; + } + pendingReconnect = null; + + HttpConnectStrategy connectStrategy = ConnectStrategy.http(getUri(context)) + .clientBuilderActions(clientBuilder -> { + httpProperties.applyToHttpClientBuilder(clientBuilder); + clientBuilder.readTimeout(READ_TIMEOUT_MS, TimeUnit.MILLISECONDS); + }) + .requestTransformer(input -> + input.newBuilder() + .headers( + input.headers().newBuilder().addAll(httpProperties.toHeadersBuilder().build()).build() + ).build()); + + if (useReport) { + connectStrategy = connectStrategy.methodAndBody(METHOD_REPORT, getRequestBody(context)); + } - running = true; + EventSource.Builder esBuilder = new EventSource.Builder(connectStrategy); + + eventSourceStarted = System.currentTimeMillis(); + // Reconnects run only after the previous stream has ended, so the previous event source + // is not closed here. + es = new BackgroundEventSource.Builder(handler, esBuilder) + // End the stream on every error. This class schedules its own reconnect from + // onError, so the library must never reconnect on its own. + .connectionErrorHandler(t -> ConnectionErrorHandler.Action.SHUTDOWN) + .build(); + es.start(); + } + + /** + * Records a stream failure and schedules the next connection attempt. + * + * @param unexpected whether the failure is classified as {@code unexpected} + * @return the delay before the next attempt in milliseconds, or -1 if the data source has been + * stopped and no reconnect was scheduled + */ + private synchronized long scheduleReconnectAfterFailure(boolean unexpected) { + if (!running) { + return -1; + } + retryState.recordFailure(unexpected, System.currentTimeMillis()); + long delay = retryState.nextDelayMillis(); + if (pendingReconnect != null) { + pendingReconnect.cancel(false); } + pendingReconnect = taskExecutor.scheduleTask(this::connect, delay); + return delay; } @NonNull @@ -283,12 +342,22 @@ public boolean needsRefresh(boolean newInBackground, LDContext newEvaluationCont (newInBackground && !streamEvenInBackground); } - private synchronized void stopSync() { - if (es != null) { - es.close(); + private void stopSync() { + BackgroundEventSource esToClose; + synchronized (this) { + running = false; + if (pendingReconnect != null) { + pendingReconnect.cancel(false); + pendingReconnect = null; + } + esToClose = es; + es = null; + } + // Closing waits for the EventSource's threads to finish, so it is done outside the lock + // in case one of those threads is in a callback that needs the lock. + if (esToClose != null) { + esToClose.close(); } - running = false; - es = null; logger.debug("Stopped."); } diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingRetryState.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingRetryState.java new file mode 100644 index 00000000..44a4ce1a --- /dev/null +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/StreamingRetryState.java @@ -0,0 +1,108 @@ +package com.launchdarkly.sdk.android; + +import androidx.annotation.NonNull; + +import java.util.Random; + +/** + * Computes the wait before a streaming connection is retried. The backoff clears after enough + * continuous healthy operation. + *

+ * This class is not thread-safe. The owning component must serialize access to it. + */ +final class StreamingRetryState { + /** + * The largest delay after a {@code normal} failure. + */ + static final long NORMAL_MAX_DELAY_MILLIS = 30_000L; + /** + * The backoff clears after this much continuous healthy operation. + */ + static final long RESET_THRESHOLD_MILLIS = 60_000L; + + private final RetryRegime normal; + private final RetryRegime extended; + private final long resetThresholdMillis; + private final Random random; + + private boolean inExtendedRegime = false; + private int attempts = 0; + private boolean healthy = false; + private long healthySinceMillis = 0; + + /** + * @param initialReconnectDelayMillis the delay before the first reconnect after a + * {@code normal} failure, which may be zero + */ + StreamingRetryState(long initialReconnectDelayMillis) { + this(new RetryRegime(initialReconnectDelayMillis, NORMAL_MAX_DELAY_MILLIS), + RetryRegime.EXTENDED.atLeast(initialReconnectDelayMillis), + RESET_THRESHOLD_MILLIS, + new Random()); + } + + /** + * Creates a retry state with explicit regimes. + * + * @param normal the regime for {@code normal} failures + * @param extended the regime entered by an {@code unexpected} failure + * @param resetThresholdMillis how long healthy operation must last before the backoff clears + * @param random source of jitter + */ + StreamingRetryState( + @NonNull RetryRegime normal, + @NonNull RetryRegime extended, + long resetThresholdMillis, + @NonNull Random random + ) { + this.normal = normal; + this.extended = extended; + this.resetThresholdMillis = resetThresholdMillis; + this.random = random; + } + + /** + * Records healthy operation, such as a payload received on the stream. Healthy operation is + * measured from the first call after a failure. + * + * @param nowMillis the current time in milliseconds. Only differences between times matter, + * so any clock will do as long as every call on this instance uses the same + * one. + */ + void recordSuccess(long nowMillis) { + if (!healthy) { + healthy = true; + healthySinceMillis = nowMillis; + } + } + + /** + * Records a failure. Call this before {@link #nextDelayMillis()}. + * + * @param unexpected true if the failure is {@code unexpected}, false if {@code normal} + * @param nowMillis the current time in milliseconds. Only differences between times matter, + * so any clock will do as long as every call on this instance uses the same + * one. + */ + void recordFailure(boolean unexpected, long nowMillis) { + if (healthy && nowMillis - healthySinceMillis >= resetThresholdMillis) { + // The connection worked for long enough that this failure starts the backoff over. + inExtendedRegime = false; + attempts = 0; + } + healthy = false; + if (unexpected && !inExtendedRegime) { + inExtendedRegime = true; + attempts = 1; + return; + } + attempts++; + } + + /** + * @return the wait in milliseconds before the next connection attempt + */ + long nextDelayMillis() { + return (inExtendedRegime ? extended : normal).jitteredDelayMillis(attempts, random); + } +} diff --git a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SynchronizerFactoryWithState.java b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SynchronizerFactoryWithState.java index 8afcb658..9248c93a 100644 --- a/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SynchronizerFactoryWithState.java +++ b/launchdarkly-android-client-sdk/src/main/java/com/launchdarkly/sdk/android/SynchronizerFactoryWithState.java @@ -1,35 +1,56 @@ package com.launchdarkly.sdk.android; import androidx.annotation.NonNull; +import androidx.annotation.Nullable; import com.launchdarkly.sdk.android.subsystems.Synchronizer; +import java.util.Random; +import java.util.concurrent.ScheduledFuture; + /** - * Wraps a synchronizer factory with availability state (available/blocked). - * Used by {@link SourceManager} to skip synchronizers that have been blocked (e.g. after TERMINAL_ERROR). + * Wraps a synchronizer factory with availability state. A synchronizer is not usable while it + * is waiting out a backoff after an unexpected error, or while it is blocked. *

- * Package-private for internal use by FDv2DataSource. + * Package-private for internal use by FDv2DataSource. Callers synchronize on the + * {@link SourceManager}'s lock. */ final class SynchronizerFactoryWithState { enum State { /** This synchronizer is available to use. */ Available, - /** This synchronizer is no longer available (e.g. after TERMINAL_ERROR). */ + /** + * This synchronizer reported an unexpected error and is waiting out a backoff before it + * may be used again. + */ + BackingOff, + /** This synchronizer is not available until something unblocks it. */ Blocked } private final FDv2DataSource.DataSourceFactory factory; private State state = State.Available; private final boolean isFDv1Fallback; + // Backoff after unexpected errors from this slot's synchronizers. + private final StreamingRetryState retryState; + @Nullable + private ScheduledFuture pendingUnblock; - SynchronizerFactoryWithState(@NonNull FDv2DataSource.DataSourceFactory factory) { - this(factory, false); - } - - SynchronizerFactoryWithState(@NonNull FDv2DataSource.DataSourceFactory factory, boolean isFDv1Fallback) { + /** + * @param factory builds this slot's synchronizer + * @param isFDv1Fallback true if this slot holds the FDv1 fallback synchronizer + * @param backoff the delay bounds for the backoff after an unexpected error + */ + SynchronizerFactoryWithState( + @NonNull FDv2DataSource.DataSourceFactory factory, + boolean isFDv1Fallback, + @NonNull RetryRegime backoff + ) { this.factory = factory; this.isFDv1Fallback = isFDv1Fallback; + this.retryState = new StreamingRetryState( + backoff, backoff, StreamingRetryState.RESET_THRESHOLD_MILLIS, new Random()); } State getState() { @@ -51,4 +72,52 @@ Synchronizer build() { boolean isFDv1Fallback() { return isFDv1Fallback; } + + /** + * Records healthy operation by this slot's synchronizer. + * + * @param nowMillis the current time in milliseconds, on the same clock as every other call on + * this instance + */ + void recordHealthy(long nowMillis) { + retryState.recordSuccess(nowMillis); + } + + /** + * Records an unexpected error from this slot's synchronizer and puts the slot into backoff. + * + * @param nowMillis the current time in milliseconds, on the same clock as every other call on + * this instance + * @return how long the slot stays in backoff, in milliseconds + */ + long startBackoff(long nowMillis) { + retryState.recordFailure(true, nowMillis); + state = State.BackingOff; + return retryState.nextDelayMillis(); + } + + /** + * Ends this slot's backoff, if it is in one. + * + * @return true if the slot became available + */ + boolean endBackoff() { + pendingUnblock = null; + if (state != State.BackingOff) { + return false; + } + state = State.Available; + return true; + } + + void setPendingUnblock(@Nullable ScheduledFuture pendingUnblock) { + this.pendingUnblock = pendingUnblock; + } + + void cancelPendingUnblock() { + if (pendingUnblock != null) { + pendingUnblock.cancel(false); + pendingUnblock = null; + } + } } diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceConditionsTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceConditionsTest.java index 7acade43..c560f5a2 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceConditionsTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceConditionsTest.java @@ -22,6 +22,7 @@ import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.TimeoutException; import static org.junit.Assert.assertEquals; @@ -196,6 +197,17 @@ public void recovery_informDoesNothing() throws Exception { assertEquals(ConditionType.RECOVERY, type); } + @Test + public void recovery_gateClosed_rearmsTimerUntilGateOpens() throws Exception { + // The gate stays closed for the first two firings and opens on the third, so the future + // completes only after the timer has been re-armed twice. + AtomicInteger firings = new AtomicInteger(0); + RecoveryCondition condition = new RecoveryCondition(executor, 0, () -> firings.incrementAndGet() >= 3); + + assertEquals(ConditionType.RECOVERY, condition.getFuture().get(500, TimeUnit.MILLISECONDS)); + assertEquals(3, firings.get()); + } + // ==== Conditions (wrapper) ==== @Test diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceTest.java index 33ecfc61..8d0316d1 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2DataSourceTest.java @@ -4,6 +4,7 @@ import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -104,6 +105,30 @@ private FDv2DataSource buildDataSource( recoveryTimeoutSeconds); } + // A synchronizer backoff short enough to wait out in a test. + private static final RetryRegime SHORT_BACKOFF = new RetryRegime(100, 100); + + /** + * Builds a data source with the given backoff after a synchronizer's unexpected error. + */ + private FDv2DataSource buildDataSource( + MockComponents.MockDataSourceUpdateSink sink, + List> initializers, + List> synchronizers, + RetryRegime synchronizerBackoff) { + return new FDv2DataSource( + CONTEXT, + initializers, + synchronizers, + null, + sink, + executor, + logging.logger, + FDv2DataSourceConditions.DEFAULT_FALLBACK_TIMEOUT_SECONDS, + FDv2DataSourceConditions.DEFAULT_RECOVERY_TIMEOUT_SECONDS, + synchronizerBackoff); + } + /** Starts the data source and returns a callback that will receive the start result. */ private AwaitableCallback startDataSource(FDv2DataSource dataSource) { AwaitableCallback cb = new AwaitableCallback<>(); @@ -607,26 +632,36 @@ public void terminalErrorBlocksSynchronizer() throws Exception { } @Test - public void allThreeSynchronizersFailReportsExhaustion() throws Exception { + public void allSynchronizersFailingWithUnexpectedErrorsAreRetriedAfterBackoff() throws Exception { MockComponents.MockDataSourceUpdateSink sink = new MockComponents.MockDataSourceUpdateSink(); + AtomicInteger firstBuilds = new AtomicInteger(0); + AtomicInteger secondBuilds = new AtomicInteger(0); + // Both synchronizers fail with an unexpected error the first time they are built, and + // deliver data the second time. FDv2DataSource dataSource = buildDataSource(sink, Collections.emptyList(), Arrays.asList( - () -> new MockQueuedSynchronizer(terminalError()), - () -> new MockQueuedSynchronizer(terminalError()), - () -> new MockQueuedSynchronizer(terminalError()))); - + () -> firstBuilds.incrementAndGet() == 1 + ? new MockQueuedSynchronizer(terminalError()) + : new MockQueuedSynchronizer(FDv2SourceResult.changeSet(makeChangeSet(false), false)), + () -> secondBuilds.incrementAndGet() == 1 + ? new MockQueuedSynchronizer(terminalError()) + : new MockQueuedSynchronizer(FDv2SourceResult.changeSet(makeChangeSet(false), false))), + SHORT_BACKOFF); AwaitableCallback startCallback = startDataSource(dataSource); - awaitExpectingError(startCallback); - List statuses = sink.awaitStatuses(4, AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); - assertEquals(4, statuses.size()); - assertEquals(DataSourceState.INTERRUPTED, statuses.get(0)); - assertEquals(DataSourceState.INTERRUPTED, statuses.get(1)); - assertEquals(DataSourceState.INTERRUPTED, statuses.get(2)); - assertEquals(DataSourceState.OFF, statuses.get(3)); - assertNotNull(sink.getLastError()); + // Each failure is reported as an interruption, and the data source waits for a backoff + // to end rather than reporting OFF. + assertEquals(Arrays.asList(DataSourceState.INTERRUPTED, DataSourceState.INTERRUPTED), + sink.awaitStatuses(2, AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)); + + // Once the backoffs end, whichever synchronizer returns first is tried again and + // initialization completes. The two backoffs have independent jitter, so either may be + // the one that is rebuilt. + assertTrue(startCallback.await(AWAIT_TIMEOUT_SECONDS * 1000)); + assertEquals(3, firstBuilds.get() + secondBuilds.get()); + stopDataSource(dataSource); } @Test @@ -650,20 +685,6 @@ public void blockedSynchronizerSkippedInRotation() throws Exception { stopDataSource(dataSource); } - @Test - public void allSynchronizersBlockedReturnsNullAndExits() throws Exception { - MockComponents.MockDataSourceUpdateSink sink = new MockComponents.MockDataSourceUpdateSink(); - - FDv2DataSource dataSource = buildDataSource(sink, - Collections.emptyList(), - Arrays.asList( - () -> new MockQueuedSynchronizer(terminalError()), - () -> new MockQueuedSynchronizer(terminalError()))); - - AwaitableCallback startCallback = startDataSource(dataSource); - awaitExpectingError(startCallback); - } - @Test public void recoveryResetsToFirstAvailableSynchronizer() throws Exception { MockComponents.MockDataSourceUpdateSink sink = new MockComponents.MockDataSourceUpdateSink(); @@ -1466,17 +1487,15 @@ public void statusIncludesErrorInfoOnFailure() throws Exception { Collections.singletonList(() -> new MockQueuedSynchronizer( FDv2SourceResult.status(FDv2SourceResult.Status.terminalError(terminalErr), false)))); - AwaitableCallback startCallback = startDataSource(dataSource); - awaitExpectingError(startCallback); + startDataSource(dataSource); + // The unexpected error is reported as an interruption that carries the error, and the + // data source then waits out the synchronizer's backoff rather than going OFF. DataSourceState first = sink.awaitStatus(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); assertEquals(DataSourceState.INTERRUPTED, first); - - DataSourceState second = sink.awaitStatus(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); - assertEquals(DataSourceState.OFF, second); - - assertEquals(DataSourceState.OFF, sink.getLastState()); - assertNotNull(sink.getLastError()); + assertEquals(DataSourceState.INTERRUPTED, sink.getLastState()); + assertSame(terminalErr, sink.getLastError()); + stopDataSource(dataSource); } @Test @@ -1500,25 +1519,33 @@ public void statusRemainsValidDuringSynchronizerOperation() throws Exception { } @Test - public void statusTransitionsFromValidToOffWhenAllSynchronizersFail() throws Exception { + public void statusStaysInterruptedWhileTheOnlySynchronizerBacksOffThenReturnsToValid() throws Exception { MockComponents.MockDataSourceUpdateSink sink = new MockComponents.MockDataSourceUpdateSink(); RuntimeException err = new RuntimeException("server error"); + AtomicInteger builds = new AtomicInteger(0); + // The synchronizer delivers data and then fails with an unexpected error. When it is + // built again it delivers data. FDv2DataSource dataSource = buildDataSource(sink, Collections.emptyList(), - Collections.singletonList(() -> new MockQueuedSynchronizer( - FDv2SourceResult.changeSet(makeChangeSet(false), false), - FDv2SourceResult.status(FDv2SourceResult.Status.terminalError(err), false)))); + Collections.singletonList(() -> builds.incrementAndGet() == 1 + ? new MockQueuedSynchronizer( + FDv2SourceResult.changeSet(makeChangeSet(false), false), + FDv2SourceResult.status(FDv2SourceResult.Status.terminalError(err), false)) + : new MockQueuedSynchronizer(FDv2SourceResult.changeSet(makeChangeSet(false), false))), + SHORT_BACKOFF); AwaitableCallback startCallback = startDataSource(dataSource); - assertTrue(startCallback.await(AWAIT_TIMEOUT_SECONDS * 1000)); // changeset arrives first, so start succeeds - + assertTrue(startCallback.await(AWAIT_TIMEOUT_SECONDS * 1000)); assertEquals(DataSourceState.VALID, sink.awaitStatus(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)); assertEquals(DataSourceState.INTERRUPTED, sink.awaitStatus(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)); - assertEquals(DataSourceState.OFF, sink.awaitStatus(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)); - - assertEquals(DataSourceState.OFF, sink.getLastState()); assertNotNull(sink.getLastError()); + + // Once the backoff ends the synchronizer is tried again and the status returns to VALID, + // never having reached OFF. + assertEquals(DataSourceState.VALID, sink.awaitStatus(AWAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)); + assertEquals(2, builds.get()); + stopDataSource(dataSource); } @Test @@ -2079,29 +2106,28 @@ public void orchestrationLogging_fdv1Fallback_logsInfo() throws Exception { } @Test - public void orchestrationLogging_permanentFailure_logsWarn() throws Exception { + public void orchestrationLogging_unexpectedError_logsWarn() throws Exception { MockComponents.MockDataSourceUpdateSink sink = new MockComponents.MockDataSourceUpdateSink(); FDv2DataSource dataSource = buildDataSource(sink, Collections.emptyList(), Collections.singletonList(() -> new MockQueuedSynchronizer(terminalError()))); - AwaitableCallback startCallback = startDataSource(dataSource); - awaitExpectingError(startCallback); + startDataSource(dataSource); awaitLogContains(logging, - "Synchronizer 'MockQueuedSynchronizer' permanently failed and will not be used again until application restart."); + "Synchronizer 'MockQueuedSynchronizer' reported an unexpected error and will not be tried again for"); + stopDataSource(dataSource); } @Test - public void orchestrationLogging_allSynchronizersExhausted_logsWarn() throws Exception { + public void orchestrationLogging_allSynchronizersBackingOff_logsInfo() throws Exception { MockComponents.MockDataSourceUpdateSink sink = new MockComponents.MockDataSourceUpdateSink(); FDv2DataSource dataSource = buildDataSource(sink, Collections.emptyList(), Arrays.asList( - () -> new MockQueuedSynchronizer(terminalError()), () -> new MockQueuedSynchronizer(terminalError()), () -> new MockQueuedSynchronizer(terminalError()))); - AwaitableCallback startCallback = startDataSource(dataSource); - awaitExpectingError(startCallback); - awaitLogContains(logging, "No more synchronizers available."); + startDataSource(dataSource); + awaitLogContains(logging, "All synchronizers are waiting out a backoff after unexpected errors"); + stopDataSource(dataSource); } @Test diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizerTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizerTest.java index 9aa3891d..f4218d93 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizerTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/FDv2StreamingSynchronizerTest.java @@ -21,9 +21,8 @@ import java.io.IOException; import java.net.URI; import java.util.HashMap; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.Future; +import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -38,7 +37,7 @@ public class FDv2StreamingSynchronizerTest { @Rule public Timeout globalTimeout = Timeout.seconds(10); - private final ExecutorService executor = Executors.newCachedThreadPool(); + private final ScheduledThreadPoolExecutor executor = new ScheduledThreadPoolExecutor(4); @After public void tearDown() { @@ -99,6 +98,13 @@ private FDv2StreamingSynchronizer makeSynchronizer( httpProperties(), executor, LOGGER, null); } + private FDv2StreamingSynchronizer makeSynchronizer(URI streamBaseUri, int initialReconnectDelayMillis) { + return new FDv2StreamingSynchronizer( + CONTEXT, EMPTY_SELECTOR_SOURCE, streamBaseUri, STREAM_PATH, + null, initialReconnectDelayMillis, false, false, + httpProperties(), executor, LOGGER, null); + } + private static DiagnosticStore basicDiagnosticStore() { return new DiagnosticStore(new DiagnosticStore.SdkDiagnosticParams( "mobile-key", "android-client-sdk", "1.0.0", "Android", null, null, null)); @@ -337,6 +343,49 @@ public void httpNonRecoverableError() throws Exception { } } + // ---- backoff after failures ---- + + @Test + public void recoverableErrorIsReportedAndTheStreamReconnects() throws Exception { + String serverIntent = makeEvent("server-intent", "{\"payloads\":[{\"id\":\"payload-1\",\"target\":100,\"intentCode\":\"xfer-full\",\"reason\":\"payload-missing\"}]}"); + String payloadTransferred = makeEvent("payload-transferred", "{\"state\":\"(p:payload-1:100)\",\"version\":100}"); + + try (HttpServer server = HttpServer.start(Handlers.sequential( + Handlers.status(503), + Handlers.status(503), + Handlers.all( + Handlers.SSE.start(), + Handlers.SSE.event(serverIntent), + Handlers.SSE.event(payloadTransferred), + Handlers.SSE.leaveOpen())))) { + + FDv2StreamingSynchronizer sync = makeSynchronizer(server.getUri(), 1); + + // Each 503 is reported as an interruption, and the synchronizer reconnects on its own + // until the stream delivers data. + assertEquals(SourceSignal.INTERRUPTED, sync.next().get(5, TimeUnit.SECONDS).getStatus().getState()); + assertEquals(SourceSignal.INTERRUPTED, sync.next().get(5, TimeUnit.SECONDS).getStatus().getState()); + assertEquals(SourceResultType.CHANGE_SET, sync.next().get(5, TimeUnit.SECONDS).getResultType()); + + sync.close(); + } + } + + @Test + public void closeDuringBackoffCancelsReconnect() throws Exception { + try (HttpServer server = HttpServer.start(Handlers.status(503))) { + // A long reconnect delay keeps the scheduled attempt pending until close() runs. + executor.setRemoveOnCancelPolicy(true); + FDv2StreamingSynchronizer sync = makeSynchronizer(server.getUri(), 60_000); + assertEquals(SourceSignal.INTERRUPTED, sync.next().get(5, TimeUnit.SECONDS).getStatus().getState()); + + // Closing cancels the scheduled attempt and reports shutdown. + sync.close(); + assertTrue(executor.getQueue().isEmpty()); + assertEquals(SourceSignal.SHUTDOWN, sync.next().get(1, TimeUnit.SECONDS).getStatus().getState()); + } + } + @Test public void httpRecoverableError() throws Exception { try (HttpServer server = HttpServer.start(Handlers.status(503))) { diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/LDUtilTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/LDUtilTest.java index 01b1d411..1bfce914 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/LDUtilTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/LDUtilTest.java @@ -27,4 +27,36 @@ public void testSanitizeSpaces() { Assert.assertEquals("--hello--", LDUtil.sanitizeSpaces(" hello ")); Assert.assertEquals("world", LDUtil.sanitizeSpaces("world")); } + + @Test + public void isHttpErrorRecoverableClassifiesStatusCodes() { + // 400, 408, 429, and 5xx are normal. Every other 4xx is unexpected. + Assert.assertTrue(LDUtil.isHttpErrorRecoverable(400)); + Assert.assertTrue(LDUtil.isHttpErrorRecoverable(408)); + Assert.assertTrue(LDUtil.isHttpErrorRecoverable(429)); + Assert.assertTrue(LDUtil.isHttpErrorRecoverable(500)); + Assert.assertTrue(LDUtil.isHttpErrorRecoverable(503)); + Assert.assertTrue(LDUtil.isHttpErrorRecoverable(302)); + + Assert.assertFalse(LDUtil.isHttpErrorRecoverable(401)); + Assert.assertFalse(LDUtil.isHttpErrorRecoverable(403)); + Assert.assertFalse(LDUtil.isHttpErrorRecoverable(404)); + Assert.assertFalse(LDUtil.isHttpErrorRecoverable(405)); + } + + @Test + public void isUnexpectedFailureIsTrueOnlyForUnexpectedHttpStatuses() { + Assert.assertTrue(LDUtil.isUnexpectedFailure(new LDInvalidResponseCodeFailure("x", 401, false))); + Assert.assertTrue(LDUtil.isUnexpectedFailure(new LDInvalidResponseCodeFailure("x", 403, false))); + // Classification is by status code, not by the retryable flag the failure was built with. + Assert.assertTrue(LDUtil.isUnexpectedFailure(new LDInvalidResponseCodeFailure("x", 401, true))); + + Assert.assertFalse(LDUtil.isUnexpectedFailure(new LDInvalidResponseCodeFailure("x", 429, true))); + Assert.assertFalse(LDUtil.isUnexpectedFailure(new LDInvalidResponseCodeFailure("x", 500, true))); + // Malformed bodies and transport errors are normal failures. + Assert.assertFalse(LDUtil.isUnexpectedFailure(new LDFailure("x", LDFailure.FailureType.INVALID_RESPONSE_BODY))); + Assert.assertFalse(LDUtil.isUnexpectedFailure(new LDFailure("x", LDFailure.FailureType.NETWORK_FAILURE))); + Assert.assertFalse(LDUtil.isUnexpectedFailure(new RuntimeException("x"))); + Assert.assertFalse(LDUtil.isUnexpectedFailure(null)); + } } diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingDataSourceTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingDataSourceTest.java index d16fdbc0..82af9d66 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingDataSourceTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingDataSourceTest.java @@ -22,7 +22,11 @@ import org.junit.Test; import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; import java.util.Map; +import java.util.Random; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledFuture; @@ -271,8 +275,6 @@ public void terminatesAfterMaxNumberOfPolls() throws Exception { try { ds.start(LDUtil.noOpCallback()); - ScheduledFuture pollTask = ds.currentPollTask.get(); - assertFalse(pollTask.isCancelled()); LDContext context1 = requireValue(fetcher.receivedContexts, 500, TimeUnit.MILLISECONDS); @@ -280,12 +282,192 @@ public void terminatesAfterMaxNumberOfPolls() throws Exception { // if a third request is sent, this will fail here requireNoMoreValues(fetcher.receivedContexts, 200, TimeUnit.MILLISECONDS); - assertTrue(pollTask.isCancelled()); + ScheduledFuture pollTask = ds.currentPollTask.get(); + assertTrue("no further poll should be pending", pollTask == null || pollTask.isDone()); } finally { ds.stop(LDUtil.noOpCallback()); } } + // --- backoff after failures --- + // + // These tests drive the data source's timers with a ManualTaskExecutor and a + // PollingRetryState without jitter. The mock fetcher answers synchronously, so every poll and + // its outcome happen inside runPendingTasks() and the scheduled delays can be asserted exactly. + + private static final long POLL_INTERVAL_MILLIS = 30_000; + private static final long EXTENDED_DELAY_MILLIS = 300_000; + + private final ManualTaskExecutor manualTaskExecutor = new ManualTaskExecutor(); + + private static LDInvalidResponseCodeFailure httpFailure(int status) { + return new LDInvalidResponseCodeFailure("test failure", status, LDUtil.isHttpErrorRecoverable(status)); + } + + private static PollingRetryState retryStateWithoutJitter() { + return new PollingRetryState(POLL_INTERVAL_MILLIS, + new RetryRegime(EXTENDED_DELAY_MILLIS, EXTENDED_DELAY_MILLIS * 4), + new Random() { + @Override + public double nextDouble() { + return 0; + } + }); + } + + private PollingDataSource makePollingDataSource(long maxNumberOfPolls) { + ClientContextImpl clientContext = makeClientContext(false, null); + return new PollingDataSource( + clientContext.getEvaluationContext(), + clientContext.getDataSourceUpdateSink(), + 0, + POLL_INTERVAL_MILLIS, + maxNumberOfPolls, + clientContext.getFetcher(), + clientContext.getPlatformState(), + manualTaskExecutor, + retryStateWithoutJitter(), + clientContext.getBaseLogger() + ); + } + + private static class TrackingCallback implements Callback { + final List successes = new ArrayList<>(); + final List errors = new ArrayList<>(); + + @Override + public void onSuccess(Boolean result) { + successes.add(result != null ? result : false); + } + + @Override + public void onError(Throwable error) { + errors.add(error); + } + } + + @Test + public void normalErrorIsPolledAgainAfterPollInterval() { + PollingDataSource ds = makePollingDataSource(Long.MAX_VALUE); + fetcher.setupErrorResponse(httpFailure(500)); + fetcher.setupSuccessResponse("{}"); + TrackingCallback callback = new TrackingCallback(); + + // The first poll fails with a 500. + ds.start(callback); + manualTaskExecutor.runPendingTasks(); + assertEquals(1, callback.errors.size()); + + // The next poll is scheduled for the regular interval and succeeds. + assertEquals(Collections.singletonList(POLL_INTERVAL_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + assertEquals(1, callback.successes.size()); + assertEquals(2, fetcher.receivedContexts.size()); + } + + @Test + public void unexpectedErrorIsPolledAgainAfterExtendedDelay() { + PollingDataSource ds = makePollingDataSource(Long.MAX_VALUE); + fetcher.setupErrorResponse(httpFailure(401)); + fetcher.setupSuccessResponse("{}"); + TrackingCallback callback = new TrackingCallback(); + + // The first poll fails with a 401, which is reported without shutting the SDK down. + ds.start(callback); + manualTaskExecutor.runPendingTasks(); + assertEquals(401, ((LDInvalidResponseCodeFailure) callback.errors.get(0)).getResponseCode()); + assertFalse(dataSourceUpdateSink.shutDownCalled); + + // The next poll is scheduled for the extended delay rather than the poll interval, and + // succeeds once that delay has passed. + assertEquals(Collections.singletonList(EXTENDED_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + assertEquals(1, callback.successes.size()); + assertEquals(2, fetcher.receivedContexts.size()); + } + + @Test + public void sustainedUnexpectedErrorsKeepPollingWithGrowingDelay() { + PollingDataSource ds = makePollingDataSource(Long.MAX_VALUE); + fetcher.setupErrorResponse(httpFailure(401)); + fetcher.setupErrorResponse(httpFailure(403)); + fetcher.setupErrorResponse(httpFailure(405)); + fetcher.setupSuccessResponse("{}"); + TrackingCallback callback = new TrackingCallback(); + + // Each failure schedules the next poll with double the delay, and the data source never + // gives up. + ds.start(callback); + manualTaskExecutor.runPendingTasks(); + for (long expectedDelay : new long[] {EXTENDED_DELAY_MILLIS, EXTENDED_DELAY_MILLIS * 2, EXTENDED_DELAY_MILLIS * 4}) { + assertEquals(Collections.singletonList(expectedDelay), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + } + + // The fourth poll succeeded. + assertEquals(3, callback.errors.size()); + assertEquals(1, callback.successes.size()); + } + + @Test + public void twoConsecutiveSuccessfulPollsResetBackoff() { + PollingDataSource ds = makePollingDataSource(Long.MAX_VALUE); + fetcher.setupErrorResponse(httpFailure(401)); + fetcher.setupSuccessResponse("{}"); + fetcher.setupSuccessResponse("{}"); + fetcher.setupErrorResponse(httpFailure(500)); + fetcher.setupSuccessResponse("{}"); + TrackingCallback callback = new TrackingCallback(); + + // The 401 moves the data source to the extended delay. + ds.start(callback); + manualTaskExecutor.runPendingTasks(); + assertEquals(Collections.singletonList(EXTENDED_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + + // One success returns to the regular interval. A second one clears the backoff, so the + // 500 that follows is retried at the regular interval instead of a doubled extended delay. + manualTaskExecutor.runPendingTasks(); + assertEquals(Collections.singletonList(POLL_INTERVAL_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + assertEquals(Collections.singletonList(POLL_INTERVAL_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + assertEquals(2, callback.errors.size()); + assertEquals(Collections.singletonList(POLL_INTERVAL_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + } + + @Test + public void oneShotPollIsNotRetriedAfterFailure() { + PollingDataSource ds = makePollingDataSource(1); + fetcher.setupErrorResponse(httpFailure(500)); + fetcher.setupSuccessResponse("{}"); + TrackingCallback callback = new TrackingCallback(); + + // The single poll fails. + ds.start(callback); + manualTaskExecutor.runPendingTasks(); + assertEquals(1, callback.errors.size()); + + // A one-shot data source is done after its one poll, so nothing further is scheduled. + assertTrue(manualTaskExecutor.pendingDelaysMillis().isEmpty()); + } + + @Test + public void stopCancelsPendingPoll() { + PollingDataSource ds = makePollingDataSource(Long.MAX_VALUE); + fetcher.setupErrorResponse(httpFailure(401)); + fetcher.setupSuccessResponse("{}"); + TrackingCallback callback = new TrackingCallback(); + + // The 401 schedules a poll for the extended delay. + ds.start(callback); + manualTaskExecutor.runPendingTasks(); + assertEquals(Collections.singletonList(EXTENDED_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + + // Stopping cancels it. + ds.stop(LDUtil.noOpCallback()); + assertTrue(manualTaskExecutor.pendingDelaysMillis().isEmpty()); + } + private class MockFetcher implements FeatureFetcher { BlockingQueue receivedContexts = new LinkedBlockingQueue<>(); BlockingQueue responses = new LinkedBlockingQueue<>(); diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingRetryStateTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingRetryStateTest.java new file mode 100644 index 00000000..be5293b1 --- /dev/null +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/PollingRetryStateTest.java @@ -0,0 +1,127 @@ +package com.launchdarkly.sdk.android; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import org.junit.Test; + +import java.util.Random; + +/** + * Unit tests for {@link PollingRetryState}. + */ +public class PollingRetryStateTest { + private static final long SECOND = 1_000L; + private static final long MINUTE = 60 * SECOND; + + /** A random source that always yields the same fraction, so jitter is predictable. */ + private static Random fixedRandom(final double fraction) { + return new Random() { + @Override + public double nextDouble() { + return fraction; + } + }; + } + + private static final Random NO_JITTER = fixedRandom(0); + private static final Random MAX_JITTER = fixedRandom(0.999_999); + + private static PollingRetryState polling(long pollIntervalMillis, Random random) { + return new PollingRetryState(pollIntervalMillis, RetryRegime.EXTENDED, random); + } + + @Test + public void defaultResetThresholdIsTwoSuccesses() { + assertEquals(2, PollingRetryState.RESET_THRESHOLD_SUCCESSES); + } + + @Test + public void delayBeforeAnyFailureIsThePollInterval() { + assertEquals(10 * SECOND, polling(10 * SECOND, NO_JITTER).nextDelayMillis()); + } + + @Test + public void normalFailuresWaitThePollInterval() { + PollingRetryState state = polling(10 * SECOND, MAX_JITTER); + for (int i = 0; i < 5; i++) { + state.recordFailure(false); + assertEquals(10 * SECOND, state.nextDelayMillis()); + } + } + + @Test + public void unexpectedFailureWaitsTheExtendedInitialDelayThenDoubles() { + PollingRetryState state = polling(10 * SECOND, NO_JITTER); + state.recordFailure(true); + assertEquals(5 * MINUTE, state.nextDelayMillis()); + state.recordFailure(false); + assertEquals(10 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void extendedDelayIsFlooredAtThePollInterval() { + PollingRetryState state = polling(10 * MINUTE, NO_JITTER); + state.recordFailure(true); + assertEquals(10 * MINUTE, state.nextDelayMillis()); + state.recordFailure(false); + assertEquals(20 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void extendedCeilingIsAtLeastThePollInterval() { + PollingRetryState state = polling(90 * MINUTE, NO_JITTER); + state.recordFailure(true); + state.recordFailure(false); + state.recordFailure(false); + assertEquals(90 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void successReturnsToPollIntervalWithoutClearingTheBackoff() { + PollingRetryState state = polling(10 * SECOND, NO_JITTER); + state.recordFailure(true); + state.recordSuccess(); + assertEquals(10 * SECOND, state.nextDelayMillis()); + + state.recordFailure(false); + assertEquals(10 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void twoConsecutiveSuccessesClearTheBackoff() { + PollingRetryState state = polling(10 * SECOND, NO_JITTER); + state.recordFailure(true); + state.recordSuccess(); + state.recordSuccess(); + + state.recordFailure(false); + assertEquals(10 * SECOND, state.nextDelayMillis()); + state.recordFailure(true); + assertEquals(5 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void failureRestartsTheSuccessCount() { + PollingRetryState state = polling(10 * SECOND, NO_JITTER); + state.recordFailure(true); + state.recordSuccess(); + state.recordFailure(false); + state.recordSuccess(); + state.recordFailure(false); + assertEquals(20 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void defaultExtendedRegimeIsFiveMinutesWithJitter() { + PollingRetryState state = new PollingRetryState(30 * SECOND); + assertEquals(30 * SECOND, state.nextDelayMillis()); + state.recordFailure(false); + assertEquals(30 * SECOND, state.nextDelayMillis()); + + state.recordFailure(true); + long wait = state.nextDelayMillis(); + assertTrue("wait " + wait, wait > RetryRegime.EXTENDED_INITIAL_DELAY_MILLIS / 2 + && wait <= RetryRegime.EXTENDED_INITIAL_DELAY_MILLIS); + } +} diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/RetryRegimeTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/RetryRegimeTest.java new file mode 100644 index 00000000..9fe2535f --- /dev/null +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/RetryRegimeTest.java @@ -0,0 +1,82 @@ +package com.launchdarkly.sdk.android; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import org.junit.Test; + +import java.util.Random; + +/** + * Unit tests for {@link RetryRegime}. + */ +public class RetryRegimeTest { + private static final long SECOND = 1_000L; + private static final long MINUTE = 60 * SECOND; + + /** A random source that always yields the same fraction, so jitter is predictable. */ + private static Random fixedRandom(final double fraction) { + return new Random() { + @Override + public double nextDouble() { + return fraction; + } + }; + } + + private static final Random NO_JITTER = fixedRandom(0); + private static final Random MAX_JITTER = fixedRandom(0.999_999); + + @Test + public void extendedRegimeIsFiveMinutesToOneHour() { + assertEquals(5 * MINUTE, RetryRegime.EXTENDED.initialDelayMillis); + assertEquals(60 * MINUTE, RetryRegime.EXTENDED.maxDelayMillis); + } + + @Test + public void delaysDoubleFromTheInitialDelayUpToTheCeiling() { + RetryRegime regime = new RetryRegime(SECOND, 30 * SECOND); + long[] expected = { + SECOND, 2 * SECOND, 4 * SECOND, 8 * SECOND, 16 * SECOND, 30 * SECOND, 30 * SECOND}; + for (int i = 0; i < expected.length; i++) { + int attempts = i + 1; + long wait = regime.jitteredDelayMillis(attempts, NO_JITTER); + assertEquals("attempt " + attempts, expected[i], wait); + } + } + + @Test + public void jitterSubtractsUpToHalfOfTheDelay() { + RetryRegime regime = new RetryRegime(SECOND, 30 * SECOND); + // Jitter is chosen from [0, T/2), so the smallest wait is just above T/2. + long minWait = regime.jitteredDelayMillis(1, MAX_JITTER); + assertTrue("wait " + minWait, minWait > SECOND / 2 && minWait <= SECOND); + } + + @Test + public void ceilingIsRaisedToTheInitialDelay() { + RetryRegime regime = new RetryRegime(10 * MINUTE, 8 * MINUTE); + assertEquals(10 * MINUTE, regime.maxDelayMillis); + assertEquals(10 * MINUTE, regime.jitteredDelayMillis(2, NO_JITTER)); + } + + @Test + public void atLeastRaisesBothBoundsToTheFloor() { + RetryRegime raised = RetryRegime.EXTENDED.atLeast(90 * MINUTE); + assertEquals(90 * MINUTE, raised.initialDelayMillis); + assertEquals(90 * MINUTE, raised.maxDelayMillis); + + RetryRegime unchanged = RetryRegime.EXTENDED.atLeast(SECOND); + assertEquals(5 * MINUTE, unchanged.initialDelayMillis); + assertEquals(60 * MINUTE, unchanged.maxDelayMillis); + } + + @Test + public void manyAttemptsStayAtTheCeilingWithoutOverflow() { + RetryRegime regime = new RetryRegime(SECOND, 30 * SECOND); + for (int attempts = 1; attempts <= 200; attempts++) { + long wait = regime.jitteredDelayMillis(attempts, NO_JITTER); + assertTrue("wait " + wait, wait > 0 && wait <= 30 * SECOND); + } + } +} diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/SourceManagerTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/SourceManagerTest.java new file mode 100644 index 00000000..5b7121ea --- /dev/null +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/SourceManagerTest.java @@ -0,0 +1,172 @@ +package com.launchdarkly.sdk.android; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + +import androidx.annotation.NonNull; + +import com.launchdarkly.sdk.android.subsystems.FDv2SourceResult; +import com.launchdarkly.sdk.android.subsystems.Synchronizer; + +import org.junit.After; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.Timeout; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.Future; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +/** + * Unit tests for {@link SourceManager}'s handling of synchronizers that report unexpected + * errors. Such a synchronizer is put aside for a backoff and then becomes available again. + */ +public class SourceManagerTest { + + @Rule + public Timeout globalTimeout = Timeout.seconds(5); + + // A backoff long enough that no scheduled return runs during a test, with room to double. + private static final RetryRegime LONG_BACKOFF = new RetryRegime(60_000, 240_000); + // A backoff short enough to wait out. + private static final RetryRegime SHORT_BACKOFF = new RetryRegime(50, 50); + + private final ScheduledThreadPoolExecutor executor = new ScheduledThreadPoolExecutor(1); + private final List slots = new ArrayList<>(); + + @After + public void tearDown() { + executor.shutdownNow(); + } + + /** A synchronizer that never produces a result. Only its name matters here. */ + private static final class NamedSynchronizer implements Synchronizer { + private final String name; + + NamedSynchronizer(String name) { + this.name = name; + } + + @Override + @NonNull + public Future next() { + return new LDAwaitFuture<>(); + } + + @Override + public void close() {} + + @Override + @NonNull + public String name() { + return name; + } + } + + private SourceManager manager(RetryRegime backoff, String... names) { + for (final String name : names) { + slots.add(new SynchronizerFactoryWithState(() -> new NamedSynchronizer(name), false, backoff)); + } + return new SourceManager(slots, Collections.emptyList(), executor); + } + + /** Asserts that a delay is the first backoff of the regime: its initial delay, less jitter. */ + private static void assertFirstBackoff(RetryRegime regime, long delayMillis) { + assertTrue("delay " + delayMillis, + delayMillis > regime.initialDelayMillis / 2 && delayMillis <= regime.initialDelayMillis); + } + + /** Selects the next synchronizer, asserting that one is known right away. */ + private static String nextNow(SourceManager manager) throws Exception { + Future next = manager.nextAvailableSynchronizer(); + assertTrue(next.isDone()); + Synchronizer synchronizer = next.get(); + return synchronizer == null ? null : synchronizer.name(); + } + + @Test + public void backingOffTheCurrentSlotSkipsItInFavorOfTheNext() throws Exception { + SourceManager manager = manager(LONG_BACKOFF, "a", "b"); + assertEquals("a", nextNow(manager)); + + assertFirstBackoff(LONG_BACKOFF, manager.backOffCurrentSynchronizer(0)); + + assertEquals("b", nextNow(manager)); + assertEquals("b", nextNow(manager)); + } + + @Test + public void slotReturnsWhenItsBackoffEnds() throws Exception { + SourceManager manager = manager(SHORT_BACKOFF, "a"); + assertEquals("a", nextNow(manager)); + manager.backOffCurrentSynchronizer(0); + + // With every slot backing off, the next synchronizer is supplied when the backoff ends. + Future next = manager.nextAvailableSynchronizer(); + assertEquals("a", next.get(2, TimeUnit.SECONDS).name()); + assertEquals("a", nextNow(manager)); + } + + @Test + public void aBackingOffSlotStillOutranksTheCurrentOneForRecovery() throws Exception { + SourceManager manager = manager(LONG_BACKOFF, "a", "b"); + nextNow(manager); + manager.backOffCurrentSynchronizer(0); + assertEquals("b", nextNow(manager)); + + // While "a" is backing off, "b" is not prime, but there is nothing to recover to yet. + assertFalse(manager.isPrimeSynchronizer()); + assertFalse(manager.hasAvailableSynchronizerBeforeCurrent()); + + // Once "a" is available again, recovery to it is possible. + manager.endBackoff(slots.get(0)); + assertTrue(manager.hasAvailableSynchronizerBeforeCurrent()); + } + + @Test + public void repeatedUnexpectedErrorsDoubleTheBackoffUntilHealthyOperationResetsIt() throws Exception { + SourceManager manager = manager(LONG_BACKOFF, "a"); + nextNow(manager); + + // A second unexpected error without any healthy operation in between doubles the wait. + assertFirstBackoff(LONG_BACKOFF, manager.backOffCurrentSynchronizer(0)); + manager.endBackoff(slots.get(0)); + nextNow(manager); + long secondDelay = manager.backOffCurrentSynchronizer(0); + assertTrue("delay " + secondDelay, secondDelay > LONG_BACKOFF.initialDelayMillis); + manager.endBackoff(slots.get(0)); + + // Healthy operation for the reset threshold before the next error starts the backoff over. + nextNow(manager); + long healthyAt = 10_000; + manager.recordCurrentSynchronizerHealthy(healthyAt); + long thirdDelay = manager.backOffCurrentSynchronizer( + healthyAt + StreamingRetryState.RESET_THRESHOLD_MILLIS); + assertFirstBackoff(LONG_BACKOFF, thirdDelay); + } + + @Test + public void closeCompletesTheWaitWithNullAndCancelsPendingBackoffs() throws Exception { + executor.setRemoveOnCancelPolicy(true); + SourceManager manager = manager(LONG_BACKOFF, "a"); + nextNow(manager); + manager.backOffCurrentSynchronizer(0); + Future next = manager.nextAvailableSynchronizer(); + assertFalse(next.isDone()); + + manager.close(); + assertNull(next.get(1, TimeUnit.SECONDS)); + assertTrue(executor.getQueue().isEmpty()); + assertNull(nextNow(manager)); + } + + @Test + public void noSynchronizersYieldsNullRightAway() throws Exception { + assertNull(nextNow(manager(LONG_BACKOFF))); + } +} diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingDataSourceTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingDataSourceTest.java index aa1bbef8..74de79c1 100644 --- a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingDataSourceTest.java +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingDataSourceTest.java @@ -27,8 +27,10 @@ import java.io.IOException; import java.net.URI; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Random; import java.util.concurrent.BlockingQueue; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; @@ -165,6 +167,43 @@ private StreamingDataSource makeStreamingDataSource( .build(clientContext); } + // The tests of backoff behavior below drive the data source's timers with a + // ManualTaskExecutor and a StreamingRetryState without jitter, so the delay chosen for each + // reconnect can be asserted exactly instead of waited for. + private static final long NORMAL_DELAY_MILLIS = 1; + private static final long EXTENDED_DELAY_MILLIS = 300_000; + // A healthy-operation threshold no test reaches. + private static final long NEVER_RESET_MILLIS = 60_000; + + private final ManualTaskExecutor manualTaskExecutor = new ManualTaskExecutor(); + + private static StreamingRetryState retryStateWithoutJitter(long healthyResetThresholdMillis) { + return new StreamingRetryState( + new RetryRegime(NORMAL_DELAY_MILLIS, NORMAL_DELAY_MILLIS), + new RetryRegime(EXTENDED_DELAY_MILLIS, EXTENDED_DELAY_MILLIS * 4), + healthyResetThresholdMillis, + new Random() { + @Override + public double nextDouble() { + return 0; + } + }); + } + + private StreamingDataSource makeStreamingDataSource(URI streamBaseUri, StreamingRetryState retryState) { + LDConfig config = new LDConfig.Builder(AutoEnvAttributes.Disabled) + .serviceEndpoints(Components.serviceEndpoints().streaming(streamBaseUri)) + .build(); + ClientContext baseClientContext = ClientContextImpl.fromConfig( + config, MOBILE_KEY, "", perEnvironmentData, + makeFeatureFetcher(), CONTEXT, + logging.logger, platformState, environmentReporter, manualTaskExecutor); + ClientContext clientContext = ClientContextImpl.forDataSource( + baseClientContext, dataSourceUpdateSink, CONTEXT, false, false); + return new StreamingDataSource(clientContext, CONTEXT, dataSourceUpdateSink, + makeFeatureFetcher(), 1, false, retryState); + } + private static String makeSseEvent(String type, String data) { return "event: " + type + "\ndata: " + data; } @@ -636,7 +675,7 @@ public void startSendsRequestWithReportAndReasons() throws Exception { // --- start(): error handling verified via HttpServer --- @Test - public void startWithHttp401ShutsDownSink() throws Exception { + public void startWithHttp401DoesNotShutDownSink() throws Exception { try (HttpServer server = HttpServer.start(Handlers.status(401))) { StreamingDataSource sds = makeStreamingDataSource( server.getUri(), false, false); @@ -649,13 +688,9 @@ public void startWithHttp401ShutsDownSink() throws Exception { LDInvalidResponseCodeFailure failure = (LDInvalidResponseCodeFailure) error; assertEquals(401, failure.getResponseCode()); assertFalse(failure.isRetryable()); - // shutDown() runs on the EventSource background thread, so poll briefly - // to allow that thread to complete before asserting. - long deadline = System.currentTimeMillis() + 1000; - while (!dataSourceUpdateSink.shutDownCalled && System.currentTimeMillis() < deadline) { - Thread.sleep(10); - } - assertTrue(dataSourceUpdateSink.shutDownCalled); + // Reporting the error is the last thing the data source does with it, so the SDK + // has not been shut down and will not be. + assertFalse(dataSourceUpdateSink.shutDownCalled); } } @@ -770,91 +805,148 @@ public void startWithNetworkErrorReportsNetworkFailure() throws Exception { assertFalse(dataSourceUpdateSink.shutDownCalled); } + // --- start(): backoff after failures --- + @Test - public void startWithHttp401PreventsSubsequentStart() throws Exception { - try (HttpServer server = HttpServer.start(Handlers.status(401))) { - StreamingDataSource sds = makeStreamingDataSource( - server.getUri(), false, false); - TrackingCallback callback1 = new TrackingCallback(); - startDataSource(sds, callback1); + public void unexpectedErrorSchedulesReconnectAfterExtendedDelay() throws Exception { + String putEvent = makeSseEvent("put", VALID_PUT_JSON); - assertNotNull(callback1.awaitError()); + try (HttpServer server = HttpServer.start(Handlers.sequential( + Handlers.status(401), + Handlers.all( + Handlers.SSE.start(), + Handlers.SSE.event(putEvent), + Handlers.SSE.leaveOpen())))) { - // Second start should be a no-op due to connection401Error flag - TrackingCallback callback2 = new TrackingCallback(); - startDataSource(sds, callback2); + StreamingDataSource sds = makeStreamingDataSource(server.getUri(), + retryStateWithoutJitter(NEVER_RESET_MILLIS)); + TrackingCallback callback = new TrackingCallback(); + startDataSource(sds, callback); - assertNull("Second start should not produce a callback", - callback2.errors.poll(500, TimeUnit.MILLISECONDS)); - assertNull(callback2.successes.poll(200, TimeUnit.MILLISECONDS)); + // The 401 is reported, and the reconnect is scheduled for the extended delay. + assertNotNull(callback.awaitError()); + server.getRecorder().requireRequest(); + assertEquals(Collections.singletonList(EXTENDED_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + + // Once that delay has passed, the data source reconnects and receives data. + manualTaskExecutor.runPendingTasks(); + assertNotNull(callback.awaitSuccess()); + server.getRecorder().requireRequest(); + assertFalse(dataSourceUpdateSink.shutDownCalled); } } - // --- start(): no reconnect after an unrecoverable HTTP error --- - @Test - public void unrecoverableErrorOnInitialConnectDoesNotReconnect() throws Exception { + public void unexpectedErrorOnReconnectSchedulesExtendedDelay() throws Exception { String putEvent = makeSseEvent("put", VALID_PUT_JSON); - // A second request would get a working stream. The SDK must not make it. try (HttpServer server = HttpServer.start(Handlers.sequential( - Handlers.status(401), + Handlers.all( + Handlers.SSE.start(), + Handlers.SSE.event(putEvent)), + Handlers.status(403), Handlers.all( Handlers.SSE.start(), Handlers.SSE.event(putEvent), Handlers.SSE.leaveOpen())))) { - StreamingDataSource sds = makeStreamingDataSource( - server.getUri(), dataSourceUpdateSink, false, false, 1); + StreamingDataSource sds = makeStreamingDataSource(server.getUri(), + retryStateWithoutJitter(NEVER_RESET_MILLIS)); TrackingCallback callback = new TrackingCallback(); - sds.start(callback); + startDataSource(sds, callback); + assertNotNull(callback.awaitSuccess()); + // The server ending the first stream is a normal failure, so the reconnect is + // scheduled for the normal delay. Throwable error = callback.awaitError(); - assertNotNull(error); - assertFalse(((LDInvalidResponseCodeFailure) error).isRetryable()); + assertEquals(LDFailure.FailureType.NETWORK_FAILURE, ((LDFailure) error).getFailureType()); + assertEquals(Collections.singletonList(NORMAL_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); - server.getRecorder().requireRequest(); - server.getRecorder().requireNoRequests(500, TimeUnit.MILLISECONDS); - assertNull("no stream data expected after the error", - callback.successes.poll(100, TimeUnit.MILLISECONDS)); + // The 403 on that reconnect moves the data source to the extended delay, after which + // it connects again. + error = callback.awaitError(); + assertEquals(403, ((LDInvalidResponseCodeFailure) error).getResponseCode()); + assertEquals(Collections.singletonList(EXTENDED_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + assertNotNull(callback.awaitSuccess()); + } + } + + @Test + public void sustainedUnexpectedErrorsKeepReconnectingWithGrowingDelay() throws Exception { + try (HttpServer server = HttpServer.start(Handlers.status(401))) { + StreamingDataSource sds = makeStreamingDataSource(server.getUri(), + retryStateWithoutJitter(NEVER_RESET_MILLIS)); + TrackingCallback callback = new TrackingCallback(); + startDataSource(sds, callback); + + // Every attempt fails. Each one is reported and the next is scheduled with double + // the delay, and the data source never gives up. + for (long expectedDelay : new long[] {EXTENDED_DELAY_MILLIS, EXTENDED_DELAY_MILLIS * 2, EXTENDED_DELAY_MILLIS * 4}) { + assertNotNull(callback.awaitError()); + server.getRecorder().requireRequest(); + assertEquals(Collections.singletonList(expectedDelay), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + } + assertNotNull(callback.awaitError()); + assertFalse(dataSourceUpdateSink.shutDownCalled); } } @Test - public void unrecoverableErrorOnReconnectDoesNotReconnectAgain() throws Exception { + public void healthyStreamResetsBackoffToNormalDelay() throws Exception { String putEvent = makeSseEvent("put", VALID_PUT_JSON); + // The data source measures healthy operation on its own clock, so the server must hold + // the stream open for a moment after the message before ending it. This is the shortest + // margin that reliably exceeds the 1 ms threshold used here. + long healthyMarginMillis = 20; try (HttpServer server = HttpServer.start(Handlers.sequential( - // The first connection delivers data, then the server ends the stream. + Handlers.status(401), Handlers.all( Handlers.SSE.start(), - Handlers.SSE.event(putEvent)), - // The reconnect gets an unrecoverable status. - Handlers.status(403), - // A third request would get a working stream. The SDK must not make it. + Handlers.SSE.event(putEvent), + Handlers.delay(healthyMarginMillis)), Handlers.all( Handlers.SSE.start(), Handlers.SSE.event(putEvent), Handlers.SSE.leaveOpen())))) { - StreamingDataSource sds = makeStreamingDataSource( - server.getUri(), dataSourceUpdateSink, false, false, 1); + StreamingDataSource sds = makeStreamingDataSource(server.getUri(), retryStateWithoutJitter(1)); TrackingCallback callback = new TrackingCallback(); - sds.start(callback); + startDataSource(sds, callback); + // The 401 puts the data source in the extended regime. The reconnect then delivers + // data and is ended by the server. + assertNotNull(callback.awaitError()); + manualTaskExecutor.runPendingTasks(); assertNotNull(callback.awaitSuccess()); + assertNotNull(callback.awaitError()); - // The dropped stream reports a network failure first; the 403 follows. - Throwable error = callback.awaitError(); - while (error != null && !(error instanceof LDInvalidResponseCodeFailure)) { - error = callback.awaitError(); - } - assertNotNull(error); - assertEquals(403, ((LDInvalidResponseCodeFailure) error).getResponseCode()); + // Having been healthy for longer than the threshold, the data source is back to the + // normal delay rather than doubling the extended one. + assertEquals(Collections.singletonList(NORMAL_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + manualTaskExecutor.runPendingTasks(); + assertNotNull(callback.awaitSuccess()); + } + } - server.getRecorder().requireRequest(); - server.getRecorder().requireRequest(); - server.getRecorder().requireNoRequests(500, TimeUnit.MILLISECONDS); + @Test + public void stopCancelsPendingReconnect() throws Exception { + try (HttpServer server = HttpServer.start(Handlers.status(401))) { + StreamingDataSource sds = makeStreamingDataSource(server.getUri(), + retryStateWithoutJitter(NEVER_RESET_MILLIS)); + TrackingCallback callback = new TrackingCallback(); + startDataSource(sds, callback); + assertNotNull(callback.awaitError()); + assertEquals(Collections.singletonList(EXTENDED_DELAY_MILLIS), manualTaskExecutor.pendingDelaysMillis()); + + // Stopping cancels the scheduled reconnect. + AwaitableCallback stopped = new AwaitableCallback<>(); + sds.stop(stopped); + stopped.await(STOP_TIMEOUT_MILLIS); + assertTrue(manualTaskExecutor.pendingDelaysMillis().isEmpty()); } } diff --git a/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingRetryStateTest.java b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingRetryStateTest.java new file mode 100644 index 00000000..9f501ae9 --- /dev/null +++ b/launchdarkly-android-client-sdk/src/test/java/com/launchdarkly/sdk/android/StreamingRetryStateTest.java @@ -0,0 +1,167 @@ +package com.launchdarkly.sdk.android; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import org.junit.Test; + +import java.util.Random; + +/** + * Unit tests for {@link StreamingRetryState}. + */ +public class StreamingRetryStateTest { + private static final long SECOND = 1_000L; + private static final long MINUTE = 60 * SECOND; + private static final long THRESHOLD = StreamingRetryState.RESET_THRESHOLD_MILLIS; + + private static final Random NO_JITTER = new Random() { + @Override + public double nextDouble() { + return 0; + } + }; + + private static StreamingRetryState streaming() { + return new StreamingRetryState( + new RetryRegime(SECOND, StreamingRetryState.NORMAL_MAX_DELAY_MILLIS), + RetryRegime.EXTENDED, + THRESHOLD, + NO_JITTER); + } + + @Test + public void defaultsAreThirtySecondNormalCeilingAndOneMinuteResetThreshold() { + assertEquals(30 * SECOND, StreamingRetryState.NORMAL_MAX_DELAY_MILLIS); + assertEquals(60 * SECOND, StreamingRetryState.RESET_THRESHOLD_MILLIS); + } + + // ---- regimes ---- + + @Test + public void normalFailuresBackOffInTheNormalRegime() { + StreamingRetryState state = streaming(); + state.recordFailure(false, 0); + assertEquals(SECOND, state.nextDelayMillis()); + state.recordFailure(false, 0); + assertEquals(2 * SECOND, state.nextDelayMillis()); + state.recordFailure(false, 0); + assertEquals(4 * SECOND, state.nextDelayMillis()); + } + + @Test + public void unexpectedFailureMovesToExtendedRegimeStartingAtExtendedInitialDelay() { + StreamingRetryState state = streaming(); + state.recordFailure(false, 0); + state.recordFailure(false, 0); + state.recordFailure(true, 0); + assertEquals(5 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void normalFailureAfterUnexpectedStaysInExtendedRegime() { + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + state.recordFailure(false, 0); + assertEquals(10 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void repeatedUnexpectedFailuresKeepCountingInExtendedRegime() { + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + state.recordFailure(true, 0); + assertEquals(10 * MINUTE, state.nextDelayMillis()); + } + + // ---- healthy operation ---- + + @Test + public void healthyForThresholdResetsWhenTheNextFailureIsRecorded() { + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + + long connectedAt = 10 * SECOND; + state.recordSuccess(connectedAt); + state.recordFailure(false, connectedAt + THRESHOLD); + + assertEquals(SECOND, state.nextDelayMillis()); + } + + @Test + public void healthyForLessThanThresholdDoesNotReset() { + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + + long connectedAt = 10 * SECOND; + state.recordSuccess(connectedAt); + state.recordFailure(false, connectedAt + THRESHOLD - 1); + + assertEquals(10 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void healthyOperationIsMeasuredFromTheFirstMessageOnTheConnection() { + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + + state.recordSuccess(10 * SECOND); + state.recordSuccess(65 * SECOND); // must not move the marker + state.recordFailure(false, 70 * SECOND); // 60 seconds since the first message + + assertEquals(SECOND, state.nextDelayMillis()); + } + + @Test + public void failureRestartsHealthyOperationMeasurement() { + // The marker from a previous connection does not count toward the next one. + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + + state.recordSuccess(10 * SECOND); + state.recordFailure(false, 20 * SECOND); // healthy for 10 seconds only + state.recordFailure(false, 90 * SECOND); // no message on this connection + + assertEquals(20 * MINUTE, state.nextDelayMillis()); + } + + @Test + public void healthyOperationStartingAtTimeZeroCounts() { + StreamingRetryState state = streaming(); + state.recordFailure(true, 0); + state.recordSuccess(0); + state.recordFailure(false, THRESHOLD); + + assertEquals(SECOND, state.nextDelayMillis()); + } + + // ---- configured initial reconnect delay ---- + + @Test + public void firstNormalRetryWaitsTheConfiguredInitialReconnectDelay() { + StreamingRetryState state = new StreamingRetryState(250); + state.recordFailure(false, 0); + long wait = state.nextDelayMillis(); + assertTrue("wait " + wait, wait > 125 && wait <= 250); + } + + @Test + public void zeroInitialReconnectDelayReconnectsImmediatelyAfterNormalFailure() { + StreamingRetryState state = new StreamingRetryState(0); + state.recordFailure(false, 0); + assertEquals(0, state.nextDelayMillis()); + + state.recordFailure(true, 0); + long wait = state.nextDelayMillis(); + assertTrue("wait " + wait, wait > RetryRegime.EXTENDED_INITIAL_DELAY_MILLIS / 2 + && wait <= RetryRegime.EXTENDED_INITIAL_DELAY_MILLIS); + } + + @Test + public void extendedRegimeIsNeverBelowTheInitialReconnectDelay() { + StreamingRetryState state = new StreamingRetryState(10 * MINUTE); + state.recordFailure(true, 0); + long wait = state.nextDelayMillis(); + assertTrue("wait " + wait, wait > 5 * MINUTE && wait <= 10 * MINUTE); + } +} diff --git a/shared-test-code/src/main/java/com/launchdarkly/sdk/android/ManualTaskExecutor.java b/shared-test-code/src/main/java/com/launchdarkly/sdk/android/ManualTaskExecutor.java index 0f48f4a2..2bc54340 100644 --- a/shared-test-code/src/main/java/com/launchdarkly/sdk/android/ManualTaskExecutor.java +++ b/shared-test-code/src/main/java/com/launchdarkly/sdk/android/ManualTaskExecutor.java @@ -12,7 +12,10 @@ *

* This avoids {@code Thread.sleep}-based timing, which is flaky on loaded CI runners. Cancelled * tasks (e.g. when a debounce timer is reset) are never run, and {@link #cancelledCount()} lets - * tests assert how many times a task was cancelled/rescheduled. + * tests assert how many times a task was cancelled/rescheduled. {@link #pendingDelaysMillis()} + * lets tests assert the delay each pending task was scheduled with. + *

+ * Tasks may be scheduled from any thread. */ public final class ManualTaskExecutor implements TaskExecutor { private final List pending = new ArrayList<>(); @@ -21,17 +24,35 @@ public final class ManualTaskExecutor implements TaskExecutor { /** * @return the number of scheduled tasks that have been cancelled */ - public int cancelledCount() { + public synchronized int cancelledCount() { return cancelledCount; } + /** + * @return the delay each pending, non-cancelled task was scheduled with, in milliseconds, in + * the order the tasks were scheduled + */ + public synchronized List pendingDelaysMillis() { + List delays = new ArrayList<>(); + for (ManualScheduledFuture task : pending) { + if (!task.cancelled) { + delays.add(task.delayMillis); + } + } + return delays; + } + /** * Runs every pending, non-cancelled task that has been scheduled via - * {@link #scheduleTask(Runnable, long)} and clears the pending queue. + * {@link #scheduleTask(Runnable, long)} and clears the pending queue. A task scheduled while + * this runs is left pending for the next call. */ public void runPendingTasks() { - List toRun = new ArrayList<>(pending); - pending.clear(); + List toRun; + synchronized (this) { + toRun = new ArrayList<>(pending); + pending.clear(); + } for (ManualScheduledFuture task : toRun) { if (!task.cancelled) { task.action.run(); @@ -45,37 +66,41 @@ public void executeOnMainThread(Runnable action) { } @Override - public ScheduledFuture scheduleTask(Runnable action, long delayMillis) { - ManualScheduledFuture future = new ManualScheduledFuture(action); + public synchronized ScheduledFuture scheduleTask(Runnable action, long delayMillis) { + ManualScheduledFuture future = new ManualScheduledFuture(action, delayMillis); pending.add(future); return future; } @Override - public ScheduledFuture startRepeatingTask(Runnable action, long initialDelayMillis, long intervalMillis) { - ManualScheduledFuture future = new ManualScheduledFuture(action); + public synchronized ScheduledFuture startRepeatingTask(Runnable action, long initialDelayMillis, long intervalMillis) { + ManualScheduledFuture future = new ManualScheduledFuture(action, initialDelayMillis); pending.add(future); return future; } @Override - public void close() { + public synchronized void close() { pending.clear(); } private final class ManualScheduledFuture implements ScheduledFuture { private final Runnable action; - private boolean cancelled = false; + private final long delayMillis; + private volatile boolean cancelled = false; - ManualScheduledFuture(Runnable action) { + ManualScheduledFuture(Runnable action, long delayMillis) { this.action = action; + this.delayMillis = delayMillis; } @Override public boolean cancel(boolean mayInterruptIfRunning) { - if (!cancelled) { - cancelled = true; - cancelledCount++; + synchronized (ManualTaskExecutor.this) { + if (!cancelled) { + cancelled = true; + cancelledCount++; + } } return true; } @@ -102,7 +127,7 @@ public Object get(long timeout, TimeUnit unit) { @Override public long getDelay(TimeUnit unit) { - return 0; + return unit.convert(delayMillis, TimeUnit.MILLISECONDS); } @Override