allow striping clients in FaultTolerantHttpClient

This commit is contained in:
Ravi Khadiwala
2024-04-02 13:01:56 -05:00
committed by ravi-signal
parent bb0da69c9e
commit 3a1ecb342f
5 changed files with 124 additions and 22 deletions

View File

@@ -76,7 +76,8 @@ public class Cdn3RemoteStorageManagerTest {
new Cdn3StorageManagerConfiguration(
wireMock.url("storage-manager/"),
"clientId",
new SecretString("clientSecret")));
new SecretString("clientSecret"),
2));
wireMock.stubFor(get(urlEqualTo("/cdn2/source/small"))
.willReturn(aResponse()

View File

@@ -22,21 +22,27 @@ import io.github.resilience4j.circuitbreaker.CallNotPermittedException;
import java.io.IOException;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpHeaders;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.time.Duration;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.mockito.Mockito;
import org.whispersystems.textsecuregcm.configuration.CircuitBreakerConfiguration;
import org.whispersystems.textsecuregcm.configuration.RetryConfiguration;
import javax.net.ssl.SSLSession;
class FaultTolerantHttpClientTest {
@@ -126,7 +132,7 @@ class FaultTolerantHttpClientTest {
final HttpClient mockHttpClient = mock(HttpClient.class);
final FaultTolerantHttpClient client = new FaultTolerantHttpClient(
"test",
mockHttpClient,
List.of(mockHttpClient),
retryExecutor,
Duration.ofSeconds(1),
new RetryConfiguration(),
@@ -150,6 +156,53 @@ class FaultTolerantHttpClientTest {
verify(mockHttpClient, times(3)).sendAsync(any(), any());
}
@Test
void testMultipleClients() throws IOException, InterruptedException {
final HttpClient mockHttpClient1 = mock(HttpClient.class);
final HttpClient mockHttpClient2 = mock(HttpClient.class);
final FaultTolerantHttpClient client = new FaultTolerantHttpClient(
"test",
List.of(mockHttpClient1, mockHttpClient2),
retryExecutor,
Duration.ofSeconds(1),
new RetryConfiguration(),
throwable -> throwable instanceof IOException,
new CircuitBreakerConfiguration());
// Just to get a dummy HttpResponse
wireMock.stubFor(get(urlEqualTo("/ping"))
.willReturn(aResponse()
.withHeader("Content-Type", "text/plain")
.withBody("Pong!")));
final HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create("http://localhost:" + wireMock.getPort() + "/ping"))
.GET()
.build();
final HttpResponse response = HttpClient.newHttpClient().send(request, HttpResponse.BodyHandlers.discarding());
final AtomicInteger client1Calls = new AtomicInteger(0);
final AtomicInteger client2Calls = new AtomicInteger(0);
when(mockHttpClient1.sendAsync(any(), any()))
.thenAnswer(args -> {
client1Calls.incrementAndGet();
return CompletableFuture.completedFuture(response);
});
when(mockHttpClient2.sendAsync(any(), any()))
.thenAnswer(args -> {
client2Calls.incrementAndGet();
return CompletableFuture.completedFuture(response);
});
final int numCalls = 100;
for (int i = 0; i < numCalls; i++) {
client.sendAsync(request, HttpResponse.BodyHandlers.ofString()).join();
}
assertThat(client2Calls.get()).isGreaterThan(0);
assertThat(client1Calls.get()).isGreaterThan(0);
assertThat(client1Calls.get() + client2Calls.get()).isEqualTo(numCalls);
}
@Test
void testNetworkFailureCircuitBreaker() throws InterruptedException {
CircuitBreakerConfiguration circuitBreakerConfiguration = new CircuitBreakerConfiguration();