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 @@ -49,6 +49,7 @@
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;

import okhttp3.ConnectionPool;
import okhttp3.Headers;

/**
Expand Down Expand Up @@ -191,11 +192,17 @@ public Future<Void> 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).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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", "{}");
Expand Down
Loading