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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,29 @@ public interface DataSourceFactory<T> {
@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<DataSourceFactory<Initializer>> initializers,
@NonNull List<DataSourceFactory<Synchronizer>> synchronizers,
@Nullable DataSourceFactory<Synchronizer> fdv1FallbackSynchronizer,
@NonNull DataSourceUpdateSinkV2 dataSourceUpdateSink,
@NonNull ScheduledExecutorService sharedExecutor,
@NonNull LDLogger logger,
long fallbackTimeoutSeconds,
long recoveryTimeoutSeconds,
@NonNull RetryRegime synchronizerBackoff
) {
this.evaluationContext = evaluationContext;
this.dataSourceUpdateSink = dataSourceUpdateSink;
Expand All @@ -140,16 +163,17 @@ public interface DataSourceFactory<T> {

List<SynchronizerFactoryWithState> allSynchronizers = new ArrayList<>();
for (DataSourceFactory<Synchronizer> 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;
Expand Down Expand Up @@ -475,11 +499,20 @@ private List<FDv2DataSourceConditions.Condition> getConditions(int synchronizerC
List<FDv2DataSourceConditions.Condition> 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";
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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.");
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading