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 017499a2..366254b2 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 @@ -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; @@ -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; 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 c4652d17..f3ec18d2 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 @@ -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); @@ -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); } @@ -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