Skip to content
Open
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 @@ -4,8 +4,10 @@
package com.azure.cosmos;

import com.azure.cosmos.implementation.AvailabilityStrategyContext;
import com.azure.cosmos.implementation.ClientSideRequestStatistics;
import com.azure.cosmos.implementation.ConnectionPolicy;
import com.azure.cosmos.implementation.CrossRegionAvailabilityContextForRxDocumentServiceRequest;
import com.azure.cosmos.implementation.DiagnosticsClientContext;
import com.azure.cosmos.implementation.GlobalEndpointManager;
import com.azure.cosmos.implementation.OperationType;
import com.azure.cosmos.implementation.PartitionKeyRange;
Expand All @@ -15,13 +17,16 @@
import com.azure.cosmos.implementation.RxDocumentServiceRequest;
import com.azure.cosmos.implementation.SerializationDiagnosticsContext;
import com.azure.cosmos.implementation.apachecommons.collections.list.UnmodifiableList;
import com.azure.cosmos.implementation.directconnectivity.StoreResponseDiagnostics;
import com.azure.cosmos.implementation.guava25.collect.ImmutableList;
import com.azure.cosmos.implementation.perPartitionAutomaticFailover.GlobalPartitionEndpointManagerForPerPartitionAutomaticFailover;
import com.azure.cosmos.implementation.perPartitionAutomaticFailover.PartitionLevelAutomaticFailoverInfo;
import com.azure.cosmos.implementation.perPartitionAutomaticFailover.PerPartitionAutomaticFailoverInfoHolder;
import com.azure.cosmos.implementation.perPartitionCircuitBreaker.PerPartitionCircuitBreakerInfoHolder;
import com.azure.cosmos.implementation.routing.RegionalRoutingContext;
import com.azure.cosmos.rx.TestSuiteBase;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.lang3.tuple.Pair;
import org.assertj.core.api.Assertions;
import org.mockito.Mockito;
Expand All @@ -34,6 +39,7 @@
import java.lang.reflect.Field;

import java.net.URI;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
Expand Down Expand Up @@ -202,6 +208,152 @@ public void tryMarkEndpointAsUnavailableForPartitionKeyRange(
}
}

@Test(groups = {"unit"})
public void diagnosticsAreEmptyWithoutDesignatedOverride() throws Exception {
RxDocumentServiceRequest request = constructRxDocumentServiceRequestInstance(
OperationType.Create,
ResourceType.Document,
"dbs/db1/colls/coll1",
"0",
"dbs/db1/colls/coll1",
"AA",
"BB",
EAST_US_URI_CNST);
request.requestContext.regionalRoutingContextToRoute = null;

ClientSideRequestStatistics directStatistics
= new ClientSideRequestStatistics(mockDiagnosticsClientContext());
directStatistics.recordResponse(request, null, this.singleWriteAccountGlobalEndpointManagerMock);

ClientSideRequestStatistics gatewayStatistics
= new ClientSideRequestStatistics(mockDiagnosticsClientContext());
gatewayStatistics.recordGatewayResponse(
request,
Mockito.mock(StoreResponseDiagnostics.class),
this.singleWriteAccountGlobalEndpointManagerMock);

ObjectMapper objectMapper = new ObjectMapper();
JsonNode directJson = objectMapper.readTree(objectMapper.writeValueAsString(directStatistics));
JsonNode gatewayJson = objectMapper.readTree(objectMapper.writeValueAsString(gatewayStatistics));

Assertions.assertThat(directJson.at("/responseStatisticsList/0/ppaf").isObject()).isTrue();
Assertions.assertThat(directJson.at("/responseStatisticsList/0/ppaf").isEmpty()).isTrue();
Assertions.assertThat(gatewayJson.at("/gatewayStatisticsList/0/ppaf").isObject()).isTrue();
Assertions.assertThat(gatewayJson.at("/gatewayStatisticsList/0/ppaf").isEmpty()).isTrue();
}

@Test(groups = {"unit"})
public void designatedOverrideDiagnosticsContainRegionAndStableSince() throws Exception {
String collectionResourceId = "dbs/db1/colls/coll1";
GlobalPartitionEndpointManagerForPerPartitionAutomaticFailover manager
= new GlobalPartitionEndpointManagerForPerPartitionAutomaticFailover(
this.singleWriteAccountGlobalEndpointManagerMock,
true);
RxDocumentServiceRequest failoverRequest = constructRxDocumentServiceRequestInstance(
OperationType.Create,
ResourceType.Document,
collectionResourceId,
"0",
collectionResourceId,
"AA",
"BB",
EAST_US_URI_CNST);

Mockito.when(this.singleWriteAccountGlobalEndpointManagerMock.getRegionName(
EAST_US_2_URI_CNST,
OperationType.Read)).thenReturn(EAST_US_2_CNST);

Mockito.clearInvocations(this.singleWriteAccountGlobalEndpointManagerMock);
Instant beforeDesignation = Instant.now();
Assertions.assertThat(manager.tryMarkEndpointAsUnavailableForPartitionKeyRange(failoverRequest, false))
.isTrue();
Instant afterDesignation = Instant.now();

java.lang.reflect.Method getDiagnosticsSnapshot
= PerPartitionAutomaticFailoverInfoHolder.class.getDeclaredMethod("getDiagnosticsSnapshot");
getDiagnosticsSnapshot.setAccessible(true);
Object firstDiagnosticsSnapshot = getDiagnosticsSnapshot.invoke(
failoverRequest.requestContext.getPerPartitionFailoverContextHolder());

ObjectMapper objectMapper = new ObjectMapper();
JsonNode firstSnapshot = objectMapper.readTree(objectMapper.writeValueAsString(
failoverRequest.requestContext.getPerPartitionFailoverContextHolder()));
Instant designatedSince = Instant.parse(firstSnapshot.get("since").asText());

Assertions.assertThat(firstSnapshot.get("currWriteRegion").asText()).isEqualTo(EAST_US_2_CNST);
Assertions.assertThat(designatedSince).isBetween(beforeDesignation, afterDesignation);

RxDocumentServiceRequest reuseRequest = constructRxDocumentServiceRequestInstance(
OperationType.Create,
ResourceType.Document,
collectionResourceId,
"0",
collectionResourceId,
"AA",
"BB",
EAST_US_URI_CNST);
Assertions.assertThat(manager.tryAddPartitionLevelLocationOverride(reuseRequest)).isTrue();

Assertions.assertThat(getDiagnosticsSnapshot.invoke(reuseRequest.requestContext.getPerPartitionFailoverContextHolder()))
.isSameAs(firstDiagnosticsSnapshot);
JsonNode reusedSnapshot = objectMapper.readTree(objectMapper.writeValueAsString(
reuseRequest.requestContext.getPerPartitionFailoverContextHolder()));
Assertions.assertThat(reusedSnapshot.get("currWriteRegion").asText()).isEqualTo(EAST_US_2_CNST);
Assertions.assertThat(reusedSnapshot.get("since").asText()).isEqualTo(firstSnapshot.get("since").asText());
Mockito.verify(this.singleWriteAccountGlobalEndpointManagerMock, Mockito.times(1))
.getRegionName(EAST_US_2_URI_CNST, OperationType.Read);
}

@Test(groups = {"unit"})
public void responseStatisticsRetainDesignatedOverrideAtRecordTime() throws Exception {
String collectionResourceId = "dbs/db1/colls/coll1";
GlobalPartitionEndpointManagerForPerPartitionAutomaticFailover manager
= new GlobalPartitionEndpointManagerForPerPartitionAutomaticFailover(
this.singleWriteAccountGlobalEndpointManagerMock,
true);
RxDocumentServiceRequest request = constructRxDocumentServiceRequestInstance(
OperationType.Create,
ResourceType.Document,
collectionResourceId,
"0",
collectionResourceId,
"AA",
"BB",
EAST_US_URI_CNST);

Mockito.when(this.singleWriteAccountGlobalEndpointManagerMock.getRegionName(
EAST_US_2_URI_CNST,
OperationType.Read)).thenReturn(EAST_US_2_CNST);
Assertions.assertThat(manager.tryMarkEndpointAsUnavailableForPartitionKeyRange(request, false)).isTrue();
request.requestContext.regionalRoutingContextToRoute = null;

ClientSideRequestStatistics directStatistics
= new ClientSideRequestStatistics(mockDiagnosticsClientContext());
directStatistics.recordResponse(request, null, this.singleWriteAccountGlobalEndpointManagerMock);
request.requestContext.setPerPartitionAutomaticFailoverInfoHolder(null);

Assertions.assertThat(manager.tryAddPartitionLevelLocationOverride(request)).isTrue();
request.requestContext.regionalRoutingContextToRoute = null;
ClientSideRequestStatistics gatewayStatistics
= new ClientSideRequestStatistics(mockDiagnosticsClientContext());
gatewayStatistics.recordGatewayResponse(
request,
Mockito.mock(StoreResponseDiagnostics.class),
this.singleWriteAccountGlobalEndpointManagerMock);
request.requestContext.setPerPartitionAutomaticFailoverInfoHolder(null);

ObjectMapper objectMapper = new ObjectMapper();
JsonNode directJson = objectMapper.readTree(objectMapper.writeValueAsString(directStatistics));
JsonNode gatewayJson = objectMapper.readTree(objectMapper.writeValueAsString(gatewayStatistics));

assertPopulatedPpaf(directJson.at("/responseStatisticsList/0/ppaf"));
assertPopulatedPpaf(gatewayJson.at("/gatewayStatisticsList/0/ppaf"));
Assertions.assertThat(directJson.at("/responseStatisticsList/0")
.has("perPartitionAutomaticFailoverInfoHolder")).isFalse();
Assertions.assertThat(gatewayJson.at("/gatewayStatisticsList/0")
.has("perPartitionAutomaticFailoverInfoHolder")).isFalse();
}

@Test(groups = {"unit"})
public void allRegionUnavailableHandlingWithMultiThreading() {

Expand Down Expand Up @@ -430,6 +582,12 @@ private RxDocumentServiceRequest constructRxDocumentServiceRequestInstance(
return request;
}

private static void assertPopulatedPpaf(JsonNode ppaf) {
Assertions.assertThat(ppaf.isObject()).isTrue();
Assertions.assertThat(ppaf.get("currWriteRegion").asText()).isEqualTo(EAST_US_2_CNST);
Assertions.assertThat(Instant.parse(ppaf.get("since").asText())).isNotNull();
}

private static URI createUrl(String url) {
try {
return new URI(url);
Expand Down
Loading
Loading