diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java index 97df4e29..8290f3ec 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamProcessor.java @@ -49,6 +49,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; +import okhttp3.ConnectionPool; import okhttp3.Headers; /** @@ -191,11 +192,17 @@ public Future start() { // smaller one there because we don't expect long delays within any *non*-streaming response that the // LD client gets. A read timeout on the stream will result in the connection being cycled, so we set // this to be slightly more than the expected interval between heartbeat signals. + // + // 3. The connection pool keeps no idle connections, so each reconnect opens a new connection. After a + // stream read timeout, OkHttp can treat a dead HTTP/2 connection as healthy for up to 1 second. A + // reconnect in that time uses the dead connection and stalls until the next read timeout. HttpConnectStrategy eventSourceHttpConfig = ConnectStrategy.http(this.streamUri) .headers(headers) .clientBuilderActions(clientBuilder -> { httpProperties.applyToHttpClientBuilder(clientBuilder); + // Set the connection pool after httpProperties, which sets its own pool (see note 3 above). + clientBuilder.connectionPool(new ConnectionPool(0, 1, TimeUnit.SECONDS)); }) // Set readTimeout last, to ensure that this hard-coded value overrides any other read // timeout that might have been set by httpProperties (see comment about readTimeout above). diff --git a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamingSynchronizerImpl.java b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamingSynchronizerImpl.java index b5d80286..78c7383e 100644 --- a/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamingSynchronizerImpl.java +++ b/lib/sdk/server/src/main/java/com/launchdarkly/sdk/server/StreamingSynchronizerImpl.java @@ -31,6 +31,7 @@ import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.KeyedItems; import com.launchdarkly.sdk.server.subsystems.SerializationException; import com.google.gson.stream.JsonReader; +import okhttp3.ConnectionPool; import okhttp3.Headers; import org.jetbrains.annotations.NotNull; @@ -109,6 +110,11 @@ private void startStream() { .headers(headers) .clientBuilderActions(clientBuilder -> { httpProperties.applyToHttpClientBuilder(clientBuilder); + // Keep no idle connections, so each reconnect opens a new connection. After a + // stream read timeout, OkHttp can treat a dead HTTP/2 connection as healthy for + // up to 1 second. A reconnect in that time uses the dead connection and stalls + // until the next read timeout. Set this after httpProperties, which sets its own pool. + clientBuilder.connectionPool(new ConnectionPool(0, 1, TimeUnit.SECONDS)); // Add interceptor to inject selector and filter query parameters on each request clientBuilder.addInterceptor(chain -> { okhttp3.Request originalRequest = chain.request(); diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java index 1db21b28..12b84710 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamProcessorTest.java @@ -422,6 +422,36 @@ public void streamWillReconnectAfterGeneralIOException() throws Exception { } } + @Test + public void streamReconnectOpensNewConnection() throws Exception { + // The first stream ends normally. This leaves a connection that the HTTP client could reuse. + Semaphore closeFirstStream = new Semaphore(0); + Handler streamHandler = Handlers.sequential( + closableStreamResponse(EMPTY_DATA_EVENT, closeFirstStream), + streamResponse(EMPTY_DATA_EVENT) + ); + AtomicInteger connectionCount = new AtomicInteger(); + + try (HttpServer server = HttpServer.start(streamHandler)) { + TcpHandler forwardToServer = TcpHandlers.forwardToPort(server.getPort()); + TcpHandler countConnections = socket -> { + connectionCount.incrementAndGet(); + forwardToServer.apply(socket); + }; + try (TcpServer countingServer = TcpServer.start(countConnections)) { + try (StreamProcessor sp = createStreamProcessor(null, countingServer.getHttpUri())) { + startAndWait(sp); + server.getRecorder().requireRequest(); + + closeFirstStream.release(); + server.getRecorder().requireRequest(); + + assertThat(connectionCount.get(), equalTo(2)); + } + } + } + } + @Test public void streamInitDiagnosticRecordedOnOpen() throws Exception { DiagnosticStore acc = basicDiagnosticStore(); diff --git a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamingSynchronizerImplTest.java b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamingSynchronizerImplTest.java index c5b7ca12..5d9dbb2f 100644 --- a/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamingSynchronizerImplTest.java +++ b/lib/sdk/server/src/test/java/com/launchdarkly/sdk/server/StreamingSynchronizerImplTest.java @@ -12,12 +12,16 @@ import com.launchdarkly.testhelpers.httptest.Handlers; import com.launchdarkly.testhelpers.httptest.HttpServer; import com.launchdarkly.testhelpers.httptest.RequestInfo; +import com.launchdarkly.testhelpers.tcptest.TcpHandler; +import com.launchdarkly.testhelpers.tcptest.TcpHandlers; +import com.launchdarkly.testhelpers.tcptest.TcpServer; import org.junit.Test; import java.net.URI; import java.time.Duration; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import static com.launchdarkly.sdk.server.ComponentsImpl.toHttpProperties; import static com.launchdarkly.sdk.server.TestComponents.basicDiagnosticStore; @@ -358,6 +362,57 @@ public void goodbyeEventInResponse() throws Exception { } } + @Test + public void reconnectOpensNewConnection() 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}"); + + // The first stream ends normally. This leaves a connection that the HTTP client could reuse. + try (HttpServer server = HttpServer.start(Handlers.sequential( + Handlers.all( + Handlers.SSE.start(), + Handlers.SSE.event(serverIntent), + Handlers.SSE.event(payloadTransferred)), + Handlers.all( + Handlers.SSE.start(), + Handlers.SSE.event(serverIntent), + Handlers.SSE.event(payloadTransferred), + Handlers.SSE.leaveOpen())))) { + AtomicInteger connectionCount = new AtomicInteger(); + TcpHandler forwardToServer = TcpHandlers.forwardToPort(server.getPort()); + TcpHandler countConnections = socket -> { + connectionCount.incrementAndGet(); + forwardToServer.apply(socket); + }; + + try (TcpServer countingServer = TcpServer.start(countConnections)) { + HttpProperties httpProperties = toHttpProperties(clientContext("sdk-key", baseConfig().build()).getHttp()); + + StreamingSynchronizerImpl synchronizer = new StreamingSynchronizerImpl( + httpProperties, + countingServer.getHttpUri(), + "/stream", + testLogger, + mockSelectorSource(), + null, + Duration.ofMillis(100), + Thread.NORM_PRIORITY, + null + ); + + FDv2SourceResult result = synchronizer.next().get(5, TimeUnit.SECONDS); + assertEquals(SourceResultType.CHANGE_SET, result.getResultType()); + + server.getRecorder().requireRequest(); + server.getRecorder().requireRequest(); + + assertEquals(2, connectionCount.get()); + + synchronizer.close(); + } + } + } + @Test public void heartbeatEvent() throws Exception { String heartbeatEvent = makeEvent("heartbeat", "{}");