DATAES-504 - Add configuration options for logging and timeouts to ReactiveElasticsearchClient.

We now log HTTP requests and responses with the org.springframework.data.elasticsearch.client.WIRE logger for both, the HighLevelRestClient and our reactive client and associate a logging Id for improved traceability.

Configuration of connection/socket timeouts is now available via the ClientConfiguration.

Along the lines we also aligned entity handling to EntityOperations already common in other modules. EntityOperations centralizes how aspects of entities (versioning, retrieval of index name/index type) are handled.

Original Pull Request: #229
This commit is contained in:
Mark Paluch
2018-11-29 14:26:34 +01:00
committed by Christoph Strobl
parent ba890cb7eb
commit 2dcd1cfbad
15 changed files with 1158 additions and 266 deletions

View File

@@ -19,6 +19,7 @@ import static org.assertj.core.api.Assertions.*;
import static org.mockito.Mockito.*;
import java.net.InetSocketAddress;
import java.time.Duration;
import javax.net.ssl.SSLContext;
@@ -40,7 +41,7 @@ public class ClientConfigurationUnitTests {
assertThat(clientConfiguration.getEndpoints()).containsOnly(InetSocketAddress.createUnresolved("localhost", 9200));
}
@Test // DATAES-488
@Test // DATAES-488, DATAES-504
public void shouldCreateCustomizedConfiguration() {
HttpHeaders headers = new HttpHeaders();
@@ -50,15 +51,17 @@ public class ClientConfigurationUnitTests {
.connectedTo("foo", "bar") //
.usingSsl() //
.withDefaultHeaders(headers) //
.build();
.withConnectTimeout(Duration.ofDays(1)).withSocketTimeout(Duration.ofDays(2)).build();
assertThat(clientConfiguration.getEndpoints()).containsOnly(InetSocketAddress.createUnresolved("foo", 9200),
InetSocketAddress.createUnresolved("bar", 9200));
assertThat(clientConfiguration.useSsl()).isTrue();
assertThat(clientConfiguration.getDefaultHeaders().get("foo")).containsOnly("bar");
assertThat(clientConfiguration.getConnectTimeout()).isEqualTo(Duration.ofDays(1));
assertThat(clientConfiguration.getSocketTimeout()).isEqualTo(Duration.ofDays(2));
}
@Test // DATAES-488
@Test // DATAES-488, DATAES-504
public void shouldCreateSslConfiguration() {
SSLContext sslContext = mock(SSLContext.class);
@@ -72,5 +75,7 @@ public class ClientConfigurationUnitTests {
InetSocketAddress.createUnresolved("bar", 9200));
assertThat(clientConfiguration.useSsl()).isTrue();
assertThat(clientConfiguration.getSslContext()).contains(sslContext);
assertThat(clientConfiguration.getConnectTimeout()).isEqualTo(Duration.ofSeconds(10));
assertThat(clientConfiguration.getSocketTimeout()).isEqualTo(Duration.ofSeconds(5));
}
}

View File

@@ -214,11 +214,13 @@ public class ReactiveMockClientTestsUtils {
Mockito.when(headersUriSpec.uri(any(String.class))).thenReturn(headersUriSpec);
Mockito.when(headersUriSpec.uri(any(), any(Map.class))).thenReturn(headersUriSpec);
Mockito.when(headersUriSpec.headers(any(Consumer.class))).thenReturn(headersUriSpec);
Mockito.when(headersUriSpec.attribute(anyString(), anyString())).thenReturn(headersUriSpec);
RequestBodyUriSpec bodyUriSpec = mock(RequestBodyUriSpec.class);
Mockito.when(webClient.method(any())).thenReturn(bodyUriSpec);
Mockito.when(bodyUriSpec.body(any())).thenReturn(headersUriSpec);
Mockito.when(bodyUriSpec.uri(any(), any(Map.class))).thenReturn(bodyUriSpec);
Mockito.when(bodyUriSpec.attribute(anyString(), anyString())).thenReturn(bodyUriSpec);
Mockito.when(bodyUriSpec.headers(any(Consumer.class))).thenReturn(bodyUriSpec);
ClientResponse response = mock(ClientResponse.class);

View File

@@ -21,25 +21,26 @@ import static org.elasticsearch.index.query.QueryBuilders.*;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.junit.Rule;
import org.springframework.data.elasticsearch.ElasticsearchVersion;
import org.springframework.data.elasticsearch.ElasticsearchVersionRule;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.net.ConnectException;
import java.util.Arrays;
import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.stream.Collectors;
import org.junit.Before;
import org.junit.Rule;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.dao.DataAccessResourceFailureException;
import org.springframework.data.annotation.Id;
import org.springframework.data.elasticsearch.ElasticsearchVersion;
import org.springframework.data.elasticsearch.ElasticsearchVersionRule;
import org.springframework.data.elasticsearch.TestUtils;
import org.springframework.data.elasticsearch.annotations.Document;
import org.springframework.data.elasticsearch.core.query.Criteria;
@@ -53,7 +54,10 @@ import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.util.StringUtils;
/**
* Integration tests for {@link ReactiveElasticsearchTemplate}.
*
* @author Christoph Strobl
* @author Mark Paluch
* @currentRead Golden Fool - Robin Hobb
*/
@RunWith(SpringRunner.class)
@@ -140,7 +144,8 @@ public class ReactiveElasticsearchTemplateTests {
SampleEntity sampleEntity = randomEntity("in another index");
template.save(sampleEntity, ALTERNATE_INDEX).as(StepVerifier::create)//
template.save(sampleEntity, ALTERNATE_INDEX) //
.as(StepVerifier::create)//
.expectNextCount(1)//
.verifyComplete();
@@ -154,12 +159,13 @@ public class ReactiveElasticsearchTemplateTests {
@Test // DATAES-504
public void insertShouldAcceptPlainMapStructureAsSource() {
Map<String, Object> map = Collections.singletonMap("foo", "bar");
Map<String, Object> map = new LinkedHashMap<>(Collections.singletonMap("foo", "bar"));
template.save(map, ALTERNATE_INDEX, "singleton-map") //
.as(StepVerifier::create) //
.expectNextCount(1) //
.verifyComplete();
.consumeNextWith(actual -> {
assertThat(map).containsKey("id");
}).verifyComplete();
}
@Test(expected = IllegalArgumentException.class) // DATAES-504
@@ -434,7 +440,7 @@ public class ReactiveElasticsearchTemplateTests {
@Test // DATAES-504
@ElasticsearchVersion(asOf = "6.5.0")
public void deleteByQueryShouldReturnZeroIfNothingDeleted() {
public void deleteByQueryShouldReturnZeroIfNothingDeleted() throws Exception {
index(randomEntity("test message"));