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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
import com.launchdarkly.eventsource.StreamHttpErrorException;
import com.launchdarkly.eventsource.background.BackgroundEventHandler;
import com.launchdarkly.eventsource.background.BackgroundEventSource;
import com.launchdarkly.eventsource.background.ConnectionErrorHandler;
import com.launchdarkly.logging.LDLogger;
import com.launchdarkly.sdk.LDContext;
import com.launchdarkly.sdk.android.subsystems.Callback;
Expand Down Expand Up @@ -177,7 +178,19 @@ public void onError(Throwable t) {
.maxDelay(MAX_RECONNECT_TIME_MS, TimeUnit.MILLISECONDS));

eventSourceStarted = System.currentTimeMillis();
es = new BackgroundEventSource.Builder(handler, esBuilder).build();
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();

running = true;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,14 @@ private StreamingDataSource makeStreamingDataSource(
URI streamBaseUri,
MockComponents.MockDataSourceUpdateSink sink,
boolean evaluationReasons, boolean useReport) {
return makeStreamingDataSource(streamBaseUri, sink, evaluationReasons, useReport, 100);
}

private StreamingDataSource makeStreamingDataSource(
URI streamBaseUri,
MockComponents.MockDataSourceUpdateSink sink,
boolean evaluationReasons, boolean useReport,
int initialReconnectDelayMillis) {
LDConfig.Builder configBuilder = new LDConfig.Builder(AutoEnvAttributes.Disabled)
.serviceEndpoints(Components.serviceEndpoints().streaming(streamBaseUri))
.evaluationReasons(evaluationReasons);
Expand All @@ -129,7 +137,7 @@ private StreamingDataSource makeStreamingDataSource(
ClientContext clientContext = ClientContextImpl.forDataSource(
baseClientContext, sink, CONTEXT, false, false);
return (StreamingDataSource) Components.streamingDataSource()
.initialReconnectDelayMillis(100)
.initialReconnectDelayMillis(initialReconnectDelayMillis)
.build(clientContext);
}

Expand Down Expand Up @@ -758,6 +766,74 @@ public void startWithHttp401PreventsSubsequentStart() throws Exception {
}
}

// --- start(): no reconnect after an unrecoverable HTTP error ---

@Test
public void unrecoverableErrorOnInitialConnectDoesNotReconnect() 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.SSE.leaveOpen())))) {

StreamingDataSource sds = makeStreamingDataSource(
server.getUri(), dataSourceUpdateSink, false, false, 1);
TrackingCallback callback = new TrackingCallback();
sds.start(callback);

Throwable error = callback.awaitError();
assertNotNull(error);
assertFalse(((LDInvalidResponseCodeFailure) error).isRetryable());

server.getRecorder().requireRequest();
server.getRecorder().requireNoRequests(500, TimeUnit.MILLISECONDS);
assertNull("no stream data expected after the error",
callback.successes.poll(100, TimeUnit.MILLISECONDS));
}
}

@Test
public void unrecoverableErrorOnReconnectDoesNotReconnectAgain() throws Exception {
String putEvent = makeSseEvent("put", VALID_PUT_JSON);

try (HttpServer server = HttpServer.start(Handlers.sequential(
// The first connection delivers data, then the server ends the stream.
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.all(
Handlers.SSE.start(),
Handlers.SSE.event(putEvent),
Handlers.SSE.leaveOpen())))) {

StreamingDataSource sds = makeStreamingDataSource(
server.getUri(), dataSourceUpdateSink, false, false, 1);
TrackingCallback callback = new TrackingCallback();
sds.start(callback);

assertNotNull(callback.awaitSuccess());

// 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());

server.getRecorder().requireRequest();
server.getRecorder().requireRequest();
server.getRecorder().requireNoRequests(500, TimeUnit.MILLISECONDS);
}
}

// --- start(): SSE event processing via HttpServer ---

@Test
Expand Down
Loading