diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/GlobalPartitionEndpointManagerForPPAFUnitTests.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/GlobalPartitionEndpointManagerForPPAFUnitTests.java index c2ef4713edb8..0465c46e952b 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/GlobalPartitionEndpointManagerForPPAFUnitTests.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/GlobalPartitionEndpointManagerForPPAFUnitTests.java @@ -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; @@ -15,6 +17,7 @@ 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; @@ -22,6 +25,8 @@ 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; @@ -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; @@ -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() { @@ -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); diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionAutomaticFailoverE2ETests.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionAutomaticFailoverE2ETests.java index 73120d7db31f..afcd26df167d 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionAutomaticFailoverE2ETests.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionAutomaticFailoverE2ETests.java @@ -3,6 +3,7 @@ package com.azure.cosmos; +import com.azure.cosmos.implementation.ClientSideRequestStatistics; import com.azure.cosmos.implementation.Configs; import com.azure.cosmos.implementation.ConnectionPolicy; import com.azure.cosmos.implementation.DatabaseAccount; @@ -63,6 +64,7 @@ import com.azure.cosmos.test.faultinjection.FaultInjectionServerErrorResult; import com.azure.cosmos.test.faultinjection.FaultInjectionServerErrorType; import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBufAllocator; @@ -85,6 +87,7 @@ import java.net.URI; import java.net.URISyntaxException; import java.time.Duration; +import java.time.Instant; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -264,6 +267,7 @@ public class PerPartitionAutomaticFailoverE2ETests extends TestSuiteBase { assertThat(cosmosDiagnosticsValueHolder.v).isNotNull(); CosmosDiagnostics cosmosDiagnostics = cosmosDiagnosticsValueHolder.v; + getPpafBookmarks(responseWrapper); assertThat(cosmosDiagnostics.getDiagnosticsContext()).isNotNull(); assertThat(cosmosDiagnostics.getDiagnosticsContext().getContactedRegionNames()).isNotNull(); @@ -1388,6 +1392,11 @@ public void testPpafWithWriteFailoverWithEligibleErrorStatusCodes( ResponseWrapper responseAfterFailover = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseAfterFailover, expectedResponseCharacteristicsAfterFailover); + JsonNode bookmarkAfterFailover = assertPpafOverride(responseAfterFailover, preferredRegions.get(1)); + ResponseWrapper responseWithReusedOverride = dataPlaneOperation.apply(operationInvocationParamsWrapper); + this.validateExpectedResponseCharacteristics.accept(responseWithReusedOverride, expectedResponseCharacteristicsAfterFailover); + assertThat(assertPpafOverride(responseWithReusedOverride, preferredRegions.get(1))) + .isEqualTo(bookmarkAfterFailover); } catch (Exception e) { Assertions.fail("The test ran into an exception {}", e); } finally { @@ -1502,6 +1511,11 @@ public void testPpafWithWriteFailoverWithEligibleErrorStatusCodes( ResponseWrapper responseAfterFailover = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseAfterFailover, expectedResponseCharacteristicsAfterFailover); + JsonNode bookmarkAfterFailover = assertPpafOverride(responseAfterFailover, preferredRegions.get(1)); + ResponseWrapper responseWithReusedOverride = dataPlaneOperation.apply(operationInvocationParamsWrapper); + this.validateExpectedResponseCharacteristics.accept(responseWithReusedOverride, expectedResponseCharacteristicsAfterFailover); + assertThat(assertPpafOverride(responseWithReusedOverride, preferredRegions.get(1))) + .isEqualTo(bookmarkAfterFailover); } catch (Exception e) { Assertions.fail("The test ran into an exception {}", e); } finally { @@ -1648,18 +1662,22 @@ public void testPpafWithWriteFailoverWithEligibleErrorStatusCodesWithPpafDynamic globalEndpointManager.refreshLocationAsync(null, true).block(); ResponseWrapper responseWithPpafDisabled = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseWithPpafDisabled, expectedResponseCharacteristicsWhenPpafIsDisabled); + assertThat(getPpafBookmarks(responseWithPpafDisabled)).allMatch(JsonNode::isEmpty); // Phase 2: PPAF enabled -> expect characteristics provided for ENABLED ppafEnabledRef.set(Boolean.TRUE); globalEndpointManager.refreshLocationAsync(null, true).block(); ResponseWrapper responseWithPpafEnabled = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseWithPpafEnabled, expectedResponseCharacteristicsWhenPpafIsEnabled); + JsonNode enabledBookmark = assertPpafOverride(responseWithPpafEnabled, preferredRegions.get(1)); // Phase 3: PPAF disabled again -> confirm behavior reverts ppafEnabledRef.set(Boolean.FALSE); globalEndpointManager.refreshLocationAsync(null, true).block(); responseWithPpafDisabled = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseWithPpafDisabled, expectedResponseCharacteristicsWhenPpafIsDisabled); + assertThat(getPpafBookmarks(responseWithPpafDisabled)).allMatch(JsonNode::isEmpty); + assertThat(assertPpafOverride(responseWithPpafEnabled, preferredRegions.get(1))).isEqualTo(enabledBookmark); } catch (Exception e) { Assertions.fail("The test ran into an exception {}", e); } finally { @@ -1757,18 +1775,22 @@ public void testPpafWithWriteFailoverWithEligibleErrorStatusCodesWithPpafDynamic globalEndpointManager.refreshLocationAsync(null, true).block(); ResponseWrapper responseWithPpafDisabled = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseWithPpafDisabled, expectedResponseCharacteristicsWhenPpafIsDisabled); + assertThat(getPpafBookmarks(responseWithPpafDisabled)).allMatch(JsonNode::isEmpty); // Phase 2: PPAF enabled -> expect characteristics provided for ENABLED ppafEnabledRef.set(Boolean.TRUE); globalEndpointManager.refreshLocationAsync(null, true).block(); ResponseWrapper responseWithPpafEnabled = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseWithPpafEnabled, expectedResponseCharacteristicsWhenPpafIsEnabled); + JsonNode enabledBookmark = assertPpafOverride(responseWithPpafEnabled, preferredRegions.get(1)); // Phase 3: PPAF disabled again -> confirm behavior reverts ppafEnabledRef.set(Boolean.FALSE); globalEndpointManager.refreshLocationAsync(null, true).block(); responseWithPpafDisabled = dataPlaneOperation.apply(operationInvocationParamsWrapper); this.validateExpectedResponseCharacteristics.accept(responseWithPpafDisabled, expectedResponseCharacteristicsWhenPpafIsDisabled); + assertThat(getPpafBookmarks(responseWithPpafDisabled)).allMatch(JsonNode::isEmpty); + assertThat(assertPpafOverride(responseWithPpafEnabled, preferredRegions.get(1))).isEqualTo(enabledBookmark); } catch (Exception e) { Assertions.fail("The test ran into an exception {}", e); } finally { @@ -2136,6 +2158,46 @@ private void runHedgingPhasesForNonWrite( this.validateExpectedResponseCharacteristics.accept(postWindow, expectedAfterWindow); } + private static List getPpafBookmarks(ResponseWrapper response) { + CosmosDiagnostics diagnostics = extractDiagnostics(response); + assertThat(diagnostics).isNotNull(); + List bookmarks = new ArrayList<>(); + for (ClientSideRequestStatistics statistics : diagnostics.getClientSideRequestStatistics()) { + JsonNode statisticsJson = OBJECT_MAPPER.valueToTree(statistics); + for (String statisticsList : Arrays.asList("responseStatisticsList", "gatewayStatisticsList")) { + for (JsonNode attempt : statisticsJson.path(statisticsList)) { + assertThat(attempt.has("perPartitionAutomaticFailoverInfoHolder")).isFalse(); + JsonNode bookmark = attempt.path("ppaf"); + assertThat(bookmark.isObject()).as("ppaf must be an object in %s", attempt).isTrue(); + if (!bookmark.isEmpty()) { + assertThat(bookmark.size()).isEqualTo(2); + assertThat(bookmark.path("currWriteRegion").isTextual()).isTrue(); + assertThat(bookmark.path("currWriteRegion").asText()).isNotBlank(); + assertThat(bookmark.path("since").isTextual()).isTrue(); + assertThat(Instant.parse(bookmark.path("since").asText())).isBeforeOrEqualTo(Instant.now()); + } + bookmarks.add(bookmark); + } + } + } + assertThat(bookmarks).as("PPAF attempt diagnostics must be present").isNotEmpty(); + return bookmarks; + } + + private static JsonNode assertPpafOverride(ResponseWrapper response, String expectedRegion) { + List populatedBookmarks = new ArrayList<>(); + for (JsonNode bookmark : getPpafBookmarks(response)) { + if (!bookmark.isEmpty()) { + assertThat(bookmark.path("currWriteRegion").asText()).isEqualToIgnoringCase(expectedRegion); + populatedBookmarks.add(bookmark); + } + } + assertThat(populatedBookmarks).as("A designated PPAF override must be recorded").isNotEmpty(); + JsonNode firstBookmark = populatedBookmarks.get(0); + assertThat(populatedBookmarks).allMatch(firstBookmark::equals); + return firstBookmark; + } + private static CosmosDiagnostics extractDiagnostics(ResponseWrapper response) { if (response.cosmosItemResponse != null) { return response.cosmosItemResponse.getDiagnostics(); diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index d99a6b198f3d..ddeb06c7bc47 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -11,6 +11,7 @@ #### Bugs Fixed #### Other Changes +* Added a compact `ppaf` bookmark to each data-plane attempt in `CosmosDiagnostics`, containing the current per-partition write region and the time it was designated, or an empty object when no override is active. ### 4.82.0 (2026-08-26) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/ClientSideRequestStatistics.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/ClientSideRequestStatistics.java index b557bbe6b1cf..c044b43b4492 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/ClientSideRequestStatistics.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/ClientSideRequestStatistics.java @@ -177,7 +177,8 @@ public void recordResponse(RxDocumentServiceRequest request, StoreResultDiagnost storeResponseStatistics.sessionTokenEvaluationResults = request.requestContext.getSessionTokenEvaluationResults(); storeResponseStatistics.perPartitionCircuitBreakerInfoHolder = request.requestContext.getPerPartitionCircuitBreakerInfoHolder().snapshot(); - storeResponseStatistics.perPartitionAutomaticFailoverInfoHolder = request.requestContext.getPerPartitionFailoverContextHolder(); + storeResponseStatistics.perPartitionAutomaticFailoverInfoHolder + = request.requestContext.getPerPartitionFailoverContextHolder().snapshot(); if (request.requestContext.getCrossRegionAvailabilityContext() != null) { CrossRegionAvailabilityContextForRxDocumentServiceRequest crossRegionAvailabilityContextForRequest @@ -272,7 +273,8 @@ public void recordGatewayResponse( gatewayStatistics.sessionTokenEvaluationResults = rxDocumentServiceRequest.requestContext.getSessionTokenEvaluationResults(); gatewayStatistics.perPartitionCircuitBreakerInfoHolder = rxDocumentServiceRequest.requestContext.getPerPartitionCircuitBreakerInfoHolder().snapshot(); - gatewayStatistics.perPartitionAutomaticFailoverInfoHolder = rxDocumentServiceRequest.requestContext.getPerPartitionFailoverContextHolder(); + gatewayStatistics.perPartitionAutomaticFailoverInfoHolder + = rxDocumentServiceRequest.requestContext.getPerPartitionFailoverContextHolder().snapshot(); gatewayStatistics.isHubRegionProcessingOnly = "false"; CrossRegionAvailabilityContextForRxDocumentServiceRequest crossRegionAvailabilityContextForRequest @@ -749,7 +751,9 @@ public static class StoreResponseStatistics { private PerPartitionCircuitBreakerInfoHolder perPartitionCircuitBreakerInfoHolder; @JsonSerialize(using = PerPartitionAutomaticFailoverInfoHolder.PerPartitionFailoverInfoHolderSerializer.class) - private PerPartitionAutomaticFailoverInfoHolder perPartitionAutomaticFailoverInfoHolder; + @JsonProperty("ppaf") + private PerPartitionAutomaticFailoverInfoHolder perPartitionAutomaticFailoverInfoHolder + = PerPartitionAutomaticFailoverInfoHolder.EMPTY; @JsonSerialize private String isHubRegionProcessingOnly; @@ -969,7 +973,8 @@ public static class GatewayStatistics { private List faultInjectionEvaluationResults; private Set sessionTokenEvaluationResults; private PerPartitionCircuitBreakerInfoHolder perPartitionCircuitBreakerInfoHolder; - private PerPartitionAutomaticFailoverInfoHolder perPartitionAutomaticFailoverInfoHolder; + private PerPartitionAutomaticFailoverInfoHolder perPartitionAutomaticFailoverInfoHolder + = PerPartitionAutomaticFailoverInfoHolder.EMPTY; private String endpoint; private String requestThroughputControlGroupName; private String requestThroughputControlGroupConfig; @@ -1118,7 +1123,7 @@ public void serialize(GatewayStatistics gatewayStatistics, this.writeNonEmptyStringSetField(jsonGenerator, "sessionTokenEvaluationResults", gatewayStatistics.getSessionTokenEvaluationResults()); this.writeNonNullObjectField(jsonGenerator, "ppcb", gatewayStatistics.getPerPartitionCircuitBreakerInfoHolder()); - this.writeNonNullObjectField(jsonGenerator, "perPartitionAutomaticFailoverInfoHolder", gatewayStatistics.getPerPartitionFailoverInfoHolder()); + jsonGenerator.writeObjectField("ppaf", gatewayStatistics.getPerPartitionFailoverInfoHolder()); this.writeNonNullStringField(jsonGenerator, "requestTCG", gatewayStatistics.getRequestThroughputControlGroupName()); this.writeNonNullStringField(jsonGenerator, "requestTCGConfig", gatewayStatistics.getRequestThroughputControlGroupConfig()); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PartitionLevelAutomaticFailoverInfo.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PartitionLevelAutomaticFailoverInfo.java index b59c136b9723..85aefcf99319 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PartitionLevelAutomaticFailoverInfo.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PartitionLevelAutomaticFailoverInfo.java @@ -6,18 +6,13 @@ import com.azure.cosmos.implementation.GlobalEndpointManager; import com.azure.cosmos.implementation.OperationType; import com.azure.cosmos.implementation.routing.RegionalRoutingContext; -import com.fasterxml.jackson.core.JsonGenerator; -import com.fasterxml.jackson.databind.SerializerProvider; -import com.fasterxml.jackson.databind.annotation.JsonSerialize; -import java.io.IOException; import java.io.Serializable; -import java.net.URI; +import java.time.Instant; import java.util.List; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; -@JsonSerialize(using = PartitionLevelAutomaticFailoverInfo.PartitionLevelFailoverInfoSerializer.class) public class PartitionLevelAutomaticFailoverInfo implements Serializable { // Set of URIs which have seen 503s (specific to document writes) or 403/3s @@ -25,6 +20,8 @@ public class PartitionLevelAutomaticFailoverInfo implements Serializable { // The current URI corresponds to the regional endpoint to use as an override private RegionalRoutingContext current; + private volatile PerPartitionAutomaticFailoverDiagnostics diagnosticsSnapshot + = PerPartitionAutomaticFailoverDiagnostics.EMPTY; private final GlobalEndpointManager globalEndpointManager; PartitionLevelAutomaticFailoverInfo(RegionalRoutingContext current, GlobalEndpointManager globalEndpointManager) { @@ -52,6 +49,13 @@ synchronized boolean tryMoveToNextLocation( this.failedRegionalRoutingContexts.add(failedRegionalRoutingContext); this.current = regionalRoutingContext; + Instant currentWriteRegionSince = Instant.now(); + String currentWriteRegion = this.globalEndpointManager.getRegionName( + regionalRoutingContext.getGatewayRegionalEndpoint(), + OperationType.Read); + this.diagnosticsSnapshot = new PerPartitionAutomaticFailoverDiagnostics( + currentWriteRegion, + currentWriteRegionSince); return true; } @@ -59,42 +63,11 @@ synchronized boolean tryMoveToNextLocation( return false; } - public RegionalRoutingContext getCurrent() { + public synchronized RegionalRoutingContext getCurrent() { return this.current; } - static class PartitionLevelFailoverInfoSerializer extends com.fasterxml.jackson.databind.JsonSerializer { - - @Override - public void serialize(PartitionLevelAutomaticFailoverInfo value, JsonGenerator gen, SerializerProvider serializers) throws IOException { - - gen.writeStartObject(); - - if (!value.failedRegionalRoutingContexts.isEmpty()) { - - StringBuilder sb = new StringBuilder("["); - - for (RegionalRoutingContext location : value.failedRegionalRoutingContexts) { - - URI gatewayRegionalEndpoint = location.getGatewayRegionalEndpoint(); - - sb.append(value.globalEndpointManager.getRegionName(gatewayRegionalEndpoint, OperationType.Read)).append(","); - } - - sb.deleteCharAt(sb.length() - 1); - sb.append("]"); - - gen.writePOJOField("failedRegions", sb.toString()); - } else { - gen.writePOJOField("failedRegions", "[]"); - } - - if (value.current != null) { - URI gatewayRegionalEndpoint = value.current.getGatewayRegionalEndpoint(); - gen.writePOJOField("overrideRegion", value.globalEndpointManager.getRegionName(gatewayRegionalEndpoint, OperationType.Read)); - } - - gen.writeEndObject(); - } + PerPartitionAutomaticFailoverDiagnostics snapshot() { + return this.diagnosticsSnapshot; } } diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverDiagnostics.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverDiagnostics.java new file mode 100644 index 000000000000..3e35f1fc6f8e --- /dev/null +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverDiagnostics.java @@ -0,0 +1,30 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.azure.cosmos.implementation.perPartitionAutomaticFailover; + +import java.io.Serializable; +import java.time.Instant; + +final class PerPartitionAutomaticFailoverDiagnostics implements Serializable { + private static final long serialVersionUID = 1L; + + static final PerPartitionAutomaticFailoverDiagnostics EMPTY + = new PerPartitionAutomaticFailoverDiagnostics(null, null); + + private final String currentWriteRegion; + private final Instant since; + + PerPartitionAutomaticFailoverDiagnostics(String currentWriteRegion, Instant since) { + this.currentWriteRegion = currentWriteRegion; + this.since = since; + } + + String getCurrentWriteRegion() { + return this.currentWriteRegion; + } + + Instant getSince() { + return this.since; + } +} \ No newline at end of file diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverInfoHolder.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverInfoHolder.java index 0c4d91680732..587a8e653fac 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverInfoHolder.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionAutomaticFailover/PerPartitionAutomaticFailoverInfoHolder.java @@ -3,6 +3,7 @@ package com.azure.cosmos.implementation.perPartitionAutomaticFailover; +import com.azure.cosmos.implementation.DiagnosticsInstantSerializer; import com.azure.cosmos.implementation.Utils; import com.fasterxml.jackson.core.JsonGenerator; import com.fasterxml.jackson.databind.SerializerProvider; @@ -14,34 +15,52 @@ @JsonSerialize(using = PerPartitionAutomaticFailoverInfoHolder.PerPartitionFailoverInfoHolderSerializer.class) public class PerPartitionAutomaticFailoverInfoHolder implements Serializable { - public static final PerPartitionAutomaticFailoverInfoHolder EMPTY = new PerPartitionAutomaticFailoverInfoHolder(); + public static final PerPartitionAutomaticFailoverInfoHolder EMPTY + = new PerPartitionAutomaticFailoverInfoHolder(PerPartitionAutomaticFailoverDiagnostics.EMPTY); - private final Utils.ValueHolder partitionLevelFailoverInfoValueHolder = new Utils.ValueHolder<>(); + private final Utils.ValueHolder diagnosticsSnapshotValueHolder + = new Utils.ValueHolder<>(); - public synchronized PartitionLevelAutomaticFailoverInfo getPartitionLevelFailoverInfo() { - return partitionLevelFailoverInfoValueHolder.v; + public PerPartitionAutomaticFailoverInfoHolder() { + this(PerPartitionAutomaticFailoverDiagnostics.EMPTY); } - public synchronized void setPartitionLevelFailoverInfo(PartitionLevelAutomaticFailoverInfo partitionLevelAutomaticFailoverInfo) { - this.partitionLevelFailoverInfoValueHolder.v = partitionLevelAutomaticFailoverInfo; + private PerPartitionAutomaticFailoverInfoHolder(PerPartitionAutomaticFailoverDiagnostics diagnosticsSnapshot) { + this.diagnosticsSnapshotValueHolder.v = diagnosticsSnapshot; } - public static class PerPartitionFailoverInfoHolderSerializer extends com.fasterxml.jackson.databind.JsonSerializer { + synchronized PerPartitionAutomaticFailoverDiagnostics getDiagnosticsSnapshot() { + return this.diagnosticsSnapshotValueHolder.v; + } - @Override - public void serialize(PerPartitionAutomaticFailoverInfoHolder value, JsonGenerator gen, SerializerProvider serializers) throws IOException { + public synchronized void setPartitionLevelFailoverInfo(PartitionLevelAutomaticFailoverInfo partitionLevelAutomaticFailoverInfo) { + if (this == EMPTY) { + return; + } - PartitionLevelAutomaticFailoverInfo partitionLevelAutomaticFailoverInfo = value.getPartitionLevelFailoverInfo(); + this.diagnosticsSnapshotValueHolder.v = partitionLevelAutomaticFailoverInfo == null + ? PerPartitionAutomaticFailoverDiagnostics.EMPTY + : partitionLevelAutomaticFailoverInfo.snapshot(); + } - if (partitionLevelAutomaticFailoverInfo != null) { - gen.writeStartObject(); + public synchronized PerPartitionAutomaticFailoverInfoHolder snapshot() { + PerPartitionAutomaticFailoverDiagnostics snapshot = this.diagnosticsSnapshotValueHolder.v; + return snapshot == PerPartitionAutomaticFailoverDiagnostics.EMPTY + ? EMPTY + : new PerPartitionAutomaticFailoverInfoHolder(snapshot); + } - gen.writeObjectField("perPartitionAutomaticFailoverCtx", value.getPartitionLevelFailoverInfo()); + public static class PerPartitionFailoverInfoHolderSerializer extends com.fasterxml.jackson.databind.JsonSerializer { - gen.writeEndObject(); - } else { - gen.writeNull(); + @Override + public void serialize(PerPartitionAutomaticFailoverInfoHolder value, JsonGenerator gen, SerializerProvider serializers) throws IOException { + PerPartitionAutomaticFailoverDiagnostics snapshot = value.getDiagnosticsSnapshot(); + gen.writeStartObject(); + if (snapshot != PerPartitionAutomaticFailoverDiagnostics.EMPTY) { + gen.writeStringField("currWriteRegion", snapshot.getCurrentWriteRegion()); + gen.writeStringField("since", DiagnosticsInstantSerializer.fromInstant(snapshot.getSince())); } + gen.writeEndObject(); } } }