From b7e05872dac81b3a4f8e9a5621f92dabe112cb82 Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Wed, 12 Aug 2026 19:22:52 -0400 Subject: [PATCH 1/8] Improve PPCB failback diagnostics --- ...PartitionCircuitBreakerInfoHolderTest.java | 188 +++++++++++++++++ .../PpcbFailbackLoggingTest.java | 138 +++++++++++++ .../ClientSideRequestStatistics.java | 10 +- ...tManagerForPerPartitionCircuitBreaker.java | 190 ++++++++++++++---- .../PerPartitionCircuitBreakerInfoHolder.java | 26 ++- 5 files changed, 502 insertions(+), 50 deletions(-) create mode 100644 sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolderTest.java create mode 100644 sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolderTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolderTest.java new file mode 100644 index 000000000000..4b9d350a3996 --- /dev/null +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolderTest.java @@ -0,0 +1,188 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.azure.cosmos.implementation.perPartitionCircuitBreaker; + +import com.azure.cosmos.implementation.ClientSideRequestStatistics; +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; +import com.azure.cosmos.implementation.ResourceType; +import com.azure.cosmos.implementation.RxDocumentServiceRequest; +import com.azure.cosmos.implementation.apachecommons.collections.list.UnmodifiableList; +import com.azure.cosmos.implementation.directconnectivity.StoreResponseDiagnostics; +import com.azure.cosmos.implementation.perPartitionAutomaticFailover.PerPartitionAutomaticFailoverInfoHolder; +import com.azure.cosmos.implementation.routing.RegionalRoutingContext; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.module.SimpleModule; +import org.mockito.Mockito; +import org.testng.annotations.Test; + +import java.time.Instant; +import java.net.URI; +import java.util.Arrays; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.doReturn; + +public class PerPartitionCircuitBreakerInfoHolderTest { + + @Test(groups = {"unit"}) + public void storesImmutableStateSnapshot() { + LocationSpecificHealthContext healthContext = createHealthContext(LocationHealthStatus.Unavailable); + Map currentState = new LinkedHashMap<>(); + currentState.put("eastus", healthContext); + + PerPartitionCircuitBreakerInfoHolder holder = new PerPartitionCircuitBreakerInfoHolder(); + holder.setPerPartitionCircuitBreakerInfoHolder(currentState); + currentState.clear(); + + assertThat(holder.getPerPartitionCircuitBreakerInfoHolder()) + .containsOnlyKeys("eastus") + .containsValue(healthContext); + assertThatThrownBy(() -> holder.getPerPartitionCircuitBreakerInfoHolder().clear()) + .isInstanceOf(UnsupportedOperationException.class); + } + + @Test(groups = {"unit"}) + public void initializedEmptyStateIsSerialized() throws Exception { + PerPartitionCircuitBreakerInfoHolder holder = new PerPartitionCircuitBreakerInfoHolder(); + holder.setPerPartitionCircuitBreakerInfoHolder(Collections.emptyMap()); + + ObjectMapper objectMapper = new ObjectMapper(); + objectMapper.registerModule(new SimpleModule().addSerializer( + PerPartitionCircuitBreakerInfoHolder.class, + new PerPartitionCircuitBreakerInfoHolder.PerPartitionCircuitBreakerInfoHolderSerializer())); + + assertThat(objectMapper.writeValueAsString(holder)) + .isEqualTo("{\"stateByRegion\":{}}"); + assertThat(objectMapper.writeValueAsString(PerPartitionCircuitBreakerInfoHolder.EMPTY)) + .isEqualTo("null"); + } + + @Test(groups = {"unit"}) + public void responseStatisticsRetainStateAtRecordTime() { + DiagnosticsClientContext diagnosticsClientContext = Mockito.mock(DiagnosticsClientContext.class); + PerPartitionCircuitBreakerInfoHolder holder = new PerPartitionCircuitBreakerInfoHolder(); + holder.setPerPartitionCircuitBreakerInfoHolder(Collections.singletonMap( + "eastus", + createHealthContext(LocationHealthStatus.Unavailable))); + + RxDocumentServiceRequest request = RxDocumentServiceRequest.create( + diagnosticsClientContext, + OperationType.Read, + ResourceType.Document); + request.requestContext.setCrossRegionAvailabilityContext( + new CrossRegionAvailabilityContextForRxDocumentServiceRequest( + null, + null, + null, + new AtomicBoolean(false), + holder, + new PerPartitionAutomaticFailoverInfoHolder())); + + ClientSideRequestStatistics statistics = new ClientSideRequestStatistics(diagnosticsClientContext); + statistics.recordResponse(request, null, null); + holder.setPerPartitionCircuitBreakerInfoHolder(Collections.singletonMap( + "westus", + createHealthContext(LocationHealthStatus.Healthy))); + + PerPartitionCircuitBreakerInfoHolder recordedHolder = statistics.getResponseStatisticsList() + .iterator() + .next() + .getPerPartitionCircuitBreakerInfoHolder(); + assertThat(recordedHolder.getPerPartitionCircuitBreakerInfoHolder()).containsOnlyKeys("eastus"); + } + + @Test(groups = {"unit"}) + public void gatewayStatisticsRetainStateAtRecordTime() throws Exception { + DiagnosticsClientContext diagnosticsClientContext = Mockito.mock(DiagnosticsClientContext.class); + PerPartitionCircuitBreakerInfoHolder holder = new PerPartitionCircuitBreakerInfoHolder(); + holder.setPerPartitionCircuitBreakerInfoHolder(Collections.singletonMap( + "eastus", + createHealthContext(LocationHealthStatus.Unavailable))); + RxDocumentServiceRequest request = createRequest(diagnosticsClientContext, holder); + + ClientSideRequestStatistics statistics = new ClientSideRequestStatistics(diagnosticsClientContext); + statistics.recordGatewayResponse(request, Mockito.mock(StoreResponseDiagnostics.class), null); + holder.setPerPartitionCircuitBreakerInfoHolder(Collections.singletonMap( + "westus", + createHealthContext(LocationHealthStatus.Healthy))); + + PerPartitionCircuitBreakerInfoHolder recordedHolder = statistics.getGatewayStatisticsList() + .get(0) + .getPerPartitionCircuitBreakerInfoHolder(); + assertThat(recordedHolder.getPerPartitionCircuitBreakerInfoHolder()).containsOnlyKeys("eastus"); + assertThat(new ObjectMapper().writeValueAsString(statistics)) + .contains("\"ppcb\":{\"stateByRegion\":{\"eastus\":"); + } + + @Test(groups = {"unit"}) + public void routingLookupInitializesEmptyStateWhenNoCircuitExists() throws Exception { + DiagnosticsClientContext diagnosticsClientContext = Mockito.mock(DiagnosticsClientContext.class); + PerPartitionCircuitBreakerInfoHolder holder = new PerPartitionCircuitBreakerInfoHolder(); + RxDocumentServiceRequest request = createRequest(diagnosticsClientContext, holder); + request.setResourceId("collectionRid"); + PartitionKeyRange partitionKeyRange = new PartitionKeyRange("0", "AA", "BB"); + request.requestContext.resolvedPartitionKeyRange = partitionKeyRange; + request.requestContext.resolvedPartitionKeyRangeForCircuitBreaker = partitionKeyRange; + + RegionalRoutingContext eastUs = new RegionalRoutingContext(URI.create("https://eastus.documents.azure.com")); + RegionalRoutingContext westUs = new RegionalRoutingContext(URI.create("https://westus.documents.azure.com")); + GlobalEndpointManager globalEndpointManager = Mockito.mock(GlobalEndpointManager.class); + doReturn(false).when(globalEndpointManager).canUseMultipleWriteLocations(request); + doReturn(UnmodifiableList.unmodifiableList(Arrays.asList(eastUs, westUs))) + .when(globalEndpointManager) + .getApplicableReadRegionalRoutingContexts(Collections.emptyList()); + + GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker manager + = new GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(globalEndpointManager); + manager.resetCircuitBreakerConfig(PartitionLevelCircuitBreakerConfig.fromJsonString( + "{\"isPartitionLevelCircuitBreakerEnabled\":true," + + "\"consecutiveExceptionCountToleratedForReads\":10," + + "\"consecutiveExceptionCountToleratedForWrites\":5}")); + + assertThat(manager.getUnavailableRegionsForPartitionKeyRange(request, "collectionRid", partitionKeyRange)) + .isEmpty(); + assertThat(holder.isInitialized()).isTrue(); + assertThat(holder.getPerPartitionCircuitBreakerInfoHolder()).isEmpty(); + + ClientSideRequestStatistics statistics = new ClientSideRequestStatistics(diagnosticsClientContext); + statistics.recordResponse(request, null, null); + assertThat(new ObjectMapper().writeValueAsString(statistics)) + .contains("\"ppcb\":{\"stateByRegion\":{}}"); + } + + private static RxDocumentServiceRequest createRequest( + DiagnosticsClientContext diagnosticsClientContext, + PerPartitionCircuitBreakerInfoHolder holder) { + + RxDocumentServiceRequest request = RxDocumentServiceRequest.create( + diagnosticsClientContext, + OperationType.Read, + ResourceType.Document); + request.requestContext.setCrossRegionAvailabilityContext( + new CrossRegionAvailabilityContextForRxDocumentServiceRequest( + null, + null, + null, + new AtomicBoolean(false), + holder, + new PerPartitionAutomaticFailoverInfoHolder())); + return request; + } + + private static LocationSpecificHealthContext createHealthContext(LocationHealthStatus healthStatus) { + return new LocationSpecificHealthContext.Builder() + .withLocationHealthStatus(healthStatus) + .withUnavailableSince(Instant.EPOCH) + .build(); + } +} \ No newline at end of file diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java new file mode 100644 index 000000000000..60a369488bac --- /dev/null +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java @@ -0,0 +1,138 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.azure.cosmos.implementation.perPartitionCircuitBreaker; + +import com.azure.cosmos.implementation.GlobalEndpointManager; +import com.azure.cosmos.implementation.OperationType; +import com.azure.cosmos.implementation.PartitionKeyRange; +import com.azure.cosmos.implementation.PartitionKeyRangeWrapper; +import com.azure.cosmos.implementation.routing.RegionalRoutingContext; +import org.mockito.Mockito; +import org.slf4j.Logger; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; +import reactor.core.publisher.Flux; + +import java.net.URI; +import java.time.Duration; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.contains; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +public class PpcbFailbackLoggingTest { + + private static final PartitionKeyRangeWrapper PARTITION = new PartitionKeyRangeWrapper( + new PartitionKeyRange("0", "AA", "BB"), + "collectionRid"); + private static final RegionalRoutingContext REGION = new RegionalRoutingContext( + URI.create("https://contoso-east-us.documents.azure.com")); + + private GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker manager; + private Logger logger; + + @BeforeMethod(groups = {"unit"}) + public void setup() { + GlobalEndpointManager globalEndpointManager = Mockito.mock(GlobalEndpointManager.class); + doReturn("eastus").when(globalEndpointManager).getRegionName( + REGION.getGatewayRegionalEndpoint(), + OperationType.Read); + this.logger = Mockito.mock(Logger.class); + doReturn(true).when(this.logger).isWarnEnabled(); + doReturn(true).when(this.logger).isDebugEnabled(); + this.manager = new GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker( + globalEndpointManager, + this.logger); + } + + @Test(groups = {"unit"}) + public void repeatedFailuresAreSampledByPartitionRegionStageAndReason() { + RuntimeException failure = new RuntimeException("failure"); + + for (int occurrence = 0; occurrence < 10; occurrence++) { + this.manager.logFailbackFailure(PARTITION, REGION, "OPEN_CONNECTION_TASK", failure); + } + + verify(this.logger, times(2)).warn( + contains("collectionResourceId: collectionRid, partitionKeyRangeId: 0, region: eastus, stage: OPEN_CONNECTION_TASK, reason: RuntimeException"), + same(failure)); + verify(this.logger, times(8)).debug( + contains("collectionResourceId: collectionRid, partitionKeyRangeId: 0, region: eastus, stage: OPEN_CONNECTION_TASK, reason: RuntimeException"), + same(failure)); + } + + @Test(groups = {"unit"}) + public void changedFailureReasonUsesTheSameCounter() { + RuntimeException firstFailure = new RuntimeException("first"); + IllegalStateException changedFailure = new IllegalStateException("changed"); + + this.manager.logFailbackFailure(PARTITION, REGION, "OPEN_CONNECTION_TASK", firstFailure); + this.manager.logFailbackFailure(PARTITION, REGION, "OPEN_CONNECTION_TASK", firstFailure); + this.manager.logFailbackFailure(PARTITION, REGION, "OPEN_CONNECTION_TASK", changedFailure); + + verify(this.logger).warn(contains("reason: RuntimeException"), same(firstFailure)); + verify(this.logger).debug(contains("reason: RuntimeException"), same(firstFailure)); + verify(this.logger).debug(contains("reason: IllegalStateException"), same(changedFailure)); + } + + @Test(groups = {"unit"}) + public void differentStagesUseTheSameCounter() { + RuntimeException failure = new RuntimeException("failure"); + + this.manager.logFailbackFailure(PARTITION, REGION, "OPEN_CONNECTION_TASK", failure); + this.manager.logFailbackFailure(PARTITION, REGION, "RECOVERY_PIPELINE", failure); + + verify(this.logger).warn(contains("stage: OPEN_CONNECTION_TASK"), same(failure)); + verify(this.logger).debug(contains("stage: RECOVERY_PIPELINE"), same(failure)); + } + + @Test(groups = {"unit"}) + public void streamFailureWithoutPartitionIdentityIsStillLogged() { + RuntimeException failure = new RuntimeException("failure"); + + this.manager.logFailbackFailure(null, null, "RECOVERY_STREAM", failure); + + verify(this.logger).warn( + contains("collectionResourceId: , partitionKeyRangeId: , region: , stage: RECOVERY_STREAM, reason: RuntimeException"), + same(failure)); + } + + @Test(groups = {"unit"}) + public void manyPartitionsUseConstantSamplingState() { + RuntimeException failure = new RuntimeException("failure"); + + for (int rangeId = 0; rangeId < 100; rangeId++) { + this.manager.logFailbackFailure( + new PartitionKeyRangeWrapper( + new PartitionKeyRange(String.valueOf(rangeId), "AA", "BB"), + "collectionRid"), + REGION, + "OPEN_CONNECTION_TASK", + failure); + } + + verify(this.logger, times(11)).warn(contains("reason: RuntimeException"), same(failure)); + verify(this.logger, times(89)).debug(contains("reason: RuntimeException"), same(failure)); + } + + @Test(groups = {"unit"}) + public void unexpectedStreamFailureIsLoggedAndRetried() { + RuntimeException failure = new RuntimeException("failure"); + AtomicInteger subscriptions = new AtomicInteger(); + Flux recoveryWork = Flux.defer(() -> subscriptions.incrementAndGet() == 1 + ? Flux.error(failure) + : Flux.just("recovered")); + + Object result = this.manager.keepFailbackRecoveryAlive(recoveryWork) + .blockFirst(Duration.ofSeconds(1)); + + assertThat(result).isEqualTo("recovered"); + assertThat(subscriptions.get()).isEqualTo(2); + verify(this.logger).warn(contains("stage: RECOVERY_STREAM"), same(failure)); + } +} \ No newline at end of file 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 aa1f7974848b..b557bbe6b1cf 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 @@ -12,6 +12,7 @@ import com.azure.cosmos.implementation.routing.RegionalRoutingContext; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.core.JsonGenerator; import com.fasterxml.jackson.databind.SerializerProvider; import com.fasterxml.jackson.databind.annotation.JsonSerialize; @@ -174,7 +175,8 @@ public void recordResponse(RxDocumentServiceRequest request, StoreResultDiagnost this.approximateInsertionCountInBloomFilter = request.requestContext.getApproximateBloomFilterInsertionCount(); storeResponseStatistics.sessionTokenEvaluationResults = request.requestContext.getSessionTokenEvaluationResults(); - storeResponseStatistics.perPartitionCircuitBreakerInfoHolder = request.requestContext.getPerPartitionCircuitBreakerInfoHolder(); + storeResponseStatistics.perPartitionCircuitBreakerInfoHolder + = request.requestContext.getPerPartitionCircuitBreakerInfoHolder().snapshot(); storeResponseStatistics.perPartitionAutomaticFailoverInfoHolder = request.requestContext.getPerPartitionFailoverContextHolder(); if (request.requestContext.getCrossRegionAvailabilityContext() != null) { @@ -268,7 +270,8 @@ public void recordGatewayResponse( if (rxDocumentServiceRequest.requestContext != null) { gatewayStatistics.sessionTokenEvaluationResults = rxDocumentServiceRequest.requestContext.getSessionTokenEvaluationResults(); - gatewayStatistics.perPartitionCircuitBreakerInfoHolder = rxDocumentServiceRequest.requestContext.getPerPartitionCircuitBreakerInfoHolder(); + gatewayStatistics.perPartitionCircuitBreakerInfoHolder + = rxDocumentServiceRequest.requestContext.getPerPartitionCircuitBreakerInfoHolder().snapshot(); gatewayStatistics.perPartitionAutomaticFailoverInfoHolder = rxDocumentServiceRequest.requestContext.getPerPartitionFailoverContextHolder(); gatewayStatistics.isHubRegionProcessingOnly = "false"; @@ -742,6 +745,7 @@ public static class StoreResponseStatistics { private Set sessionTokenEvaluationResults; @JsonSerialize(using = PerPartitionCircuitBreakerInfoHolder.PerPartitionCircuitBreakerInfoHolderSerializer.class) + @JsonProperty("ppcb") private PerPartitionCircuitBreakerInfoHolder perPartitionCircuitBreakerInfoHolder; @JsonSerialize(using = PerPartitionAutomaticFailoverInfoHolder.PerPartitionFailoverInfoHolderSerializer.class) @@ -1113,7 +1117,7 @@ public void serialize(GatewayStatistics gatewayStatistics, } this.writeNonEmptyStringSetField(jsonGenerator, "sessionTokenEvaluationResults", gatewayStatistics.getSessionTokenEvaluationResults()); - this.writeNonNullObjectField(jsonGenerator, "perPartitionCircuitBreakerInfoHolder", gatewayStatistics.getPerPartitionCircuitBreakerInfoHolder()); + this.writeNonNullObjectField(jsonGenerator, "ppcb", gatewayStatistics.getPerPartitionCircuitBreakerInfoHolder()); this.writeNonNullObjectField(jsonGenerator, "perPartitionAutomaticFailoverInfoHolder", gatewayStatistics.getPerPartitionFailoverInfoHolder()); this.writeNonNullStringField(jsonGenerator, "requestTCG", gatewayStatistics.getRequestThroughputControlGroupName()); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java index cf5ffe7aa874..97a582c3483a 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java @@ -38,6 +38,7 @@ import java.util.PriorityQueue; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import static com.azure.cosmos.implementation.guava25.base.Preconditions.checkNotNull; @@ -57,10 +58,21 @@ public class GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker impleme private final AtomicBoolean isClosed = new AtomicBoolean(false); private final AtomicBoolean isPartitionRecoveryTaskRunning = new AtomicBoolean(false); private final AtomicReference partitionRecoveryDisposable = new AtomicReference<>(); + private final Logger failbackLogger; + private final AtomicInteger failbackFailureLogCount; public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpointManager globalEndpointManager) { + this(globalEndpointManager, logger); + } + + GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker( + GlobalEndpointManager globalEndpointManager, + Logger failbackLogger) { + this.partitionKeyRangeToLocationSpecificUnavailabilityInfo = new ConcurrentHashMap<>(); this.globalEndpointManager = globalEndpointManager; + this.failbackLogger = checkNotNull(failbackLogger, "Argument 'failbackLogger' cannot be null!"); + this.failbackFailureLogCount = new AtomicInteger(); PartitionLevelCircuitBreakerConfig partitionLevelCircuitBreakerConfig = Configs.getPartitionLevelCircuitBreakerConfig(); this.consecutiveExceptionBasedCircuitBreaker = new ConsecutiveExceptionBasedCircuitBreaker(partitionLevelCircuitBreakerConfig); @@ -74,9 +86,15 @@ public void init() { if (this.consecutiveExceptionBasedCircuitBreaker.isPartitionLevelCircuitBreakerEnabled() && this.isPartitionRecoveryTaskRunning.compareAndSet(false, true)) { - this.partitionRecoveryDisposable.set(this.updateStaleLocationInfo() + Disposable recoveryDisposable = this.updateStaleLocationInfo() .subscribeOn(CosmosSchedulers.PARTITION_AVAILABILITY_CHECK_BOUNDED_ELASTIC) - .subscribe()); + .subscribe(); + this.partitionRecoveryDisposable.set(recoveryDisposable); + + if (this.isClosed.get() + && this.partitionRecoveryDisposable.compareAndSet(recoveryDisposable, null)) { + recoveryDisposable.dispose(); + } } } @@ -109,7 +127,8 @@ public void handleLocationExceptionForPartitionKeyRange( // so we skip circuit breaking in this case if cancellation kick in ungracefully (e.g. user cancelled the request, or end-to-end timeout on the operation before routing decision is made) // if the exception is not due to a cancellation, then we should have enough information to decide if we should circuit break or not // so we proceed with circuit breaking in this case - if (resolvedPartitionKeyRangeForCircuitBreaker == null && isCancellationException) { + if (resolvedPartitionKeyRangeForCircuitBreaker == null + && isCancellationException) { logger.warn("Skipping circuit breaking for operation as partitionKeyRange information isn't available for an e2e timeout cancelled request with operationType: " + request.getOperationType() + " and collectionResourceId: " + @@ -146,9 +165,10 @@ public void handleLocationExceptionForPartitionKeyRange( isFailoverPossible.set( partitionLevelLocationUnavailabilityInfoAsVal.areLocationsAvailableForPartitionKeyRange(applicableRegionalRoutingContexts)); + } - request.requestContext.setPerPartitionCircuitBreakerInfoHolder(partitionLevelLocationUnavailabilityInfoAsVal.regionToLocationSpecificHealthContext); + this.publishSnapshot(request, partitionLevelLocationUnavailabilityInfoAsVal); return partitionLevelLocationUnavailabilityInfoAsVal; }); @@ -207,9 +227,10 @@ public void handleLocationSuccessForPartitionKeyRange(RxDocumentServiceRequest r partitionKeyRangeToFailoverInfoAsVal.handleSuccess( partitionKeyRangeWrapper, succeededRegionalRoutingContext, - request.isReadOnlyRequest()); + request.isReadOnlyRequest(), + false); - request.requestContext.setPerPartitionCircuitBreakerInfoHolder(partitionKeyRangeToFailoverInfoAsVal.regionToLocationSpecificHealthContext); + this.publishSnapshot(request, partitionKeyRangeToFailoverInfoAsVal); return partitionKeyRangeToFailoverInfoAsVal; }); } catch (Exception e) { @@ -236,6 +257,7 @@ public List getUnavailableRegionsForPartitionKeyRange( this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.get(partitionKeyRangeWrapper); List unavailableRegions = new ArrayList<>(); + this.publishSnapshot(request, partitionLevelLocationUnavailabilityInfoSnapshot); if (partitionLevelLocationUnavailabilityInfoSnapshot != null) { Map locationEndpointToFailureMetricsForPartition = @@ -278,18 +300,26 @@ public List getUnavailableRegionsForPartitionKeyRange( } } + private void publishSnapshot( + RxDocumentServiceRequest request, + PartitionLevelLocationUnavailabilityInfo info) { + + request.requestContext.setPerPartitionCircuitBreakerInfoHolder( + info == null ? Collections.emptyMap() : info.regionToLocationSpecificHealthContext); + } + private Flux updateStaleLocationInfo() { - return Mono.just(1) + Flux recoveryWork = Mono.just(1) .delayElement(Duration.ofSeconds(Configs.getStalePartitionUnavailabilityRefreshIntervalInSeconds())) .repeat(() -> !this.isClosed.get()) .flatMap(ignore -> Flux.fromIterable(this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.entrySet()), 1, 1) .flatMap(partitionKeyRangeWrapperToPartitionKeyRangeWrapperPair -> { logger.debug("Background updateStaleLocationInfo kicking in..."); + PartitionKeyRangeWrapper partitionKeyRangeWrapper + = partitionKeyRangeWrapperToPartitionKeyRangeWrapperPair.getKey(); try { - PartitionKeyRangeWrapper partitionKeyRangeWrapper = partitionKeyRangeWrapperToPartitionKeyRangeWrapperPair.getKey(); - PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.get(partitionKeyRangeWrapper); if (partitionLevelLocationUnavailabilityInfo != null) { @@ -320,7 +350,11 @@ private Flux updateStaleLocationInfo() { return Mono.empty(); } } catch (Exception e) { - logger.warn("An exception was thrown trying to recover an Unavailable partitionKeyRange!", e); + this.logFailbackFailure( + partitionKeyRangeWrapper, + null, + "SCAN_UNAVAILABLE_PARTITIONS", + e); return Flux.empty(); } }, 1, 1) @@ -353,52 +387,118 @@ private Flux updateStaleLocationInfo() { + partitionKeyRangeWrapper.getCollectionResourceId() + " has succeeded..."); - partitionLevelLocationUnavailabilityInfo.locationEndpointToLocationSpecificContextForPartition.compute(locationWithStaleUnavailabilityInfo, (locationWithStaleUnavailabilityInfoAsKey, locationSpecificContextAsVal) -> { - - if (locationSpecificContextAsVal != null) { - locationSpecificContextAsVal = GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker - .this.locationSpecificHealthContextTransitionHandler.handleSuccess( - locationSpecificContextAsVal, - partitionKeyRangeWrapper, - this.regionalRoutingContextToRegion.getOrDefault(locationWithStaleUnavailabilityInfoAsKey, StringUtils.EMPTY), - false, - true); - } - return locationSpecificContextAsVal; - }); + partitionLevelLocationUnavailabilityInfo.handleSuccess( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + true, + true); + }) .onErrorResume(throwable -> { - logger.debug("An exception was thrown trying to recover an Unavailable partition key range!", throwable); + this.logFailbackFailure( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + "OPEN_CONNECTION_TASK", + throwable); return Mono.empty(); }); + } else { + this.logFailbackFailure( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + "RESOLVE_GATEWAY_ADDRESS_CACHE", + new IllegalStateException("GatewayAddressCache is not available.")); } } else { - partitionLevelLocationUnavailabilityInfo.locationEndpointToLocationSpecificContextForPartition.compute(locationWithStaleUnavailabilityInfo, (locationWithStaleUnavailabilityInfoAsKey, locationSpecificContextAsVal) -> { + partitionLevelLocationUnavailabilityInfo.handleSuccess( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + true, + true); - if (locationSpecificContextAsVal != null) { - locationSpecificContextAsVal = GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker - .this.locationSpecificHealthContextTransitionHandler.handleSuccess( - locationSpecificContextAsVal, - partitionKeyRangeWrapper, - this.regionalRoutingContextToRegion.getOrDefault(locationWithStaleUnavailabilityInfoAsKey, StringUtils.EMPTY), - false, - true); - } - return locationSpecificContextAsVal; - }); } } } catch (Exception e) { - logger.debug("An exception was thrown trying to recover an Unavailable partition key range!", e); + PartitionKeyRangeWrapper partitionKeyRangeWrapper = locationToLocationSpecificHealthContextPair.getLeft(); + RegionalRoutingContext locationWithStaleUnavailabilityInfo + = locationToLocationSpecificHealthContextPair.getRight().getLeft(); + this.logFailbackFailure( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + "RECOVERY_PIPELINE", + e); return Flux.empty(); } return Flux.empty(); - }, 1, 1) - .onErrorResume(throwable -> { - logger.warn("An exception : was thrown trying to recover an Unavailable partitionKeyRange!, fail-back flow won't be executed!", throwable); - return Flux.empty(); - }); + }, 1, 1); + + return this.keepFailbackRecoveryAlive(recoveryWork); + } + + Flux keepFailbackRecoveryAlive(Flux recoveryWork) { + return recoveryWork + .doOnError(throwable -> this.logFailbackFailure( + null, + null, + "RECOVERY_STREAM", + throwable)) + .retry(); + } + + void logFailbackFailure( + PartitionKeyRangeWrapper partitionKeyRangeWrapper, + RegionalRoutingContext regionalRoutingContext, + String stage, + Throwable throwable) { + + String region = this.resolveRegionName(regionalRoutingContext); + String reason = throwable == null ? "UNKNOWN" : throwable.getClass().getSimpleName(); + String collectionResourceId = partitionKeyRangeWrapper == null + ? StringUtils.EMPTY + : partitionKeyRangeWrapper.getCollectionResourceId(); + String partitionKeyRangeId = partitionKeyRangeWrapper == null + || partitionKeyRangeWrapper.getPartitionKeyRange() == null + ? StringUtils.EMPTY + : partitionKeyRangeWrapper.getPartitionKeyRange().getId(); + String message = "PPCB failback failed for collectionResourceId: " + + collectionResourceId + + ", partitionKeyRangeId: " + + partitionKeyRangeId + + ", region: " + + region + + ", stage: " + + stage + + ", reason: " + + reason; + + if (this.shouldLogFailbackFailureAtWarn()) { + this.failbackLogger.warn(message, throwable); + } else { + this.failbackLogger.debug(message, throwable); + } + + } + + private boolean shouldLogFailbackFailureAtWarn() { + int count = this.failbackFailureLogCount.updateAndGet( + current -> current == Integer.MAX_VALUE ? 1 : current + 1); + return count == 1 || count % 10 == 0; + } + + private String resolveRegionName(RegionalRoutingContext regionalRoutingContext) { + if (regionalRoutingContext == null) { + return StringUtils.EMPTY; + } + + String region = this.regionalRoutingContextToRegion.get(regionalRoutingContext); + if (!StringUtils.isEmpty(region)) { + return region; + } + + return this.globalEndpointManager.getRegionName( + regionalRoutingContext.getGatewayRegionalEndpoint(), + OperationType.Read); } public boolean isPerPartitionLevelCircuitBreakingApplicable(RxDocumentServiceRequest request) { @@ -451,6 +551,7 @@ public void setGlobalAddressResolver(GlobalAddressResolver globalAddressResolver @Override public void close() { this.isClosed.set(true); + this.failbackFailureLogCount.set(0); Disposable disposable = this.partitionRecoveryDisposable.getAndSet(null); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); @@ -520,7 +621,8 @@ private boolean handleException( private void handleSuccess( PartitionKeyRangeWrapper partitionKeyRangeWrapper, RegionalRoutingContext succeededLocation, - boolean isReadOnlyRequest) { + boolean isReadOnlyRequest, + boolean forceStatusChange) { this.locationEndpointToLocationSpecificContextForPartition.compute(succeededLocation, (locationAsKey, locationSpecificContextAsVal) -> { @@ -543,7 +645,7 @@ private void handleSuccess( locationSpecificContextAsVal, partitionKeyRangeWrapper, GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.this.regionalRoutingContextToRegion.getOrDefault(succeededLocation, StringUtils.EMPTY), - false, + forceStatusChange, isReadOnlyRequest); // used only for building diagnostics - so creating a lookup for URI and region name diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolder.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolder.java index 8fed008204a0..f422a8ddde41 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolder.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PerPartitionCircuitBreakerInfoHolder.java @@ -6,25 +6,45 @@ import com.azure.cosmos.implementation.Utils; 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.util.Collections; +import java.util.LinkedHashMap; import java.util.Map; +@JsonSerialize(using = PerPartitionCircuitBreakerInfoHolder.PerPartitionCircuitBreakerInfoHolderSerializer.class) public class PerPartitionCircuitBreakerInfoHolder implements Serializable { public static final PerPartitionCircuitBreakerInfoHolder EMPTY = new PerPartitionCircuitBreakerInfoHolder(); private final Utils.ValueHolder> perPartitionCircuitBreakerInfoHolder = new Utils.ValueHolder>(); + private boolean initialized; public synchronized void setPerPartitionCircuitBreakerInfoHolder(final Map locationSpecificHealthContext) { - this.perPartitionCircuitBreakerInfoHolder.v = locationSpecificHealthContext; + this.initialized = true; + this.perPartitionCircuitBreakerInfoHolder.v = locationSpecificHealthContext == null + ? Collections.emptyMap() + : Collections.unmodifiableMap(new LinkedHashMap<>(locationSpecificHealthContext)); } public synchronized Map getPerPartitionCircuitBreakerInfoHolder() { return perPartitionCircuitBreakerInfoHolder.v; } + public synchronized PerPartitionCircuitBreakerInfoHolder snapshot() { + PerPartitionCircuitBreakerInfoHolder snapshot = new PerPartitionCircuitBreakerInfoHolder(); + if (this.initialized) { + snapshot.setPerPartitionCircuitBreakerInfoHolder(this.perPartitionCircuitBreakerInfoHolder.v); + } + return snapshot; + } + + synchronized boolean isInitialized() { + return this.initialized; + } + public static class PerPartitionCircuitBreakerInfoHolderSerializer extends com.fasterxml.jackson.databind.JsonSerializer { @Override @@ -32,10 +52,10 @@ public void serialize(PerPartitionCircuitBreakerInfoHolder value, JsonGenerator Map locationToLocationSpecificHealthContext = value.getPerPartitionCircuitBreakerInfoHolder(); - if (locationToLocationSpecificHealthContext != null && !locationToLocationSpecificHealthContext.isEmpty()) { + if (value.isInitialized()) { gen.writeStartObject(); - gen.writePOJOField("locSpecificHealthCtx", locationToLocationSpecificHealthContext); + gen.writePOJOField("stateByRegion", locationToLocationSpecificHealthContext); gen.writeEndObject(); } else { From ba23e94f4309793ee2958aa5e6d04e38b5c99834 Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Wed, 12 Aug 2026 19:46:54 -0400 Subject: [PATCH 2/8] Validate PPCB state in diagnostics E2E tests --- .../PerPartitionCircuitBreakerE2ETests.java | 60 +++++++++++++++++++ 1 file changed, 60 insertions(+) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java index 8e1ad9f6c37f..13c95b48653d 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java @@ -58,6 +58,7 @@ import com.azure.cosmos.test.faultinjection.FaultInjectionRuleBuilder; import com.azure.cosmos.test.faultinjection.FaultInjectionServerErrorResult; import com.azure.cosmos.test.faultinjection.FaultInjectionServerErrorType; +import com.fasterxml.jackson.databind.JsonNode; import org.testng.SkipException; import org.testng.annotations.AfterClass; import org.testng.annotations.AfterMethod; @@ -3807,6 +3808,7 @@ private void execute( testId, executeDataPlaneOperation, operationInvocationParamsWrapper); + List ppcbStateByRegionNodes = getPpcbStateByRegionNodes(response); ConsecutiveExceptionBasedCircuitBreaker consecutiveExceptionBasedCircuitBreaker = globalPartitionEndpointManagerForPerPartitionCircuitBreaker.getConsecutiveExceptionBasedCircuitBreaker(); @@ -3832,6 +3834,9 @@ private void execute( if (executionCountAfterCircuitBreakingThresholdBreached > 1) { validateResponseInAbsenceOfFailures.accept(response); + assertPpcbHealthStatus( + ppcbStateByRegionNodes, + LocationHealthStatus.Unavailable); } if (response.cosmosItemResponse != null) { @@ -3898,6 +3903,10 @@ private void execute( executeDataPlaneOperation, operationInvocationParamsWrapper); validateResponseInAbsenceOfFailures.accept(response); + assertPpcbHealthStatus( + getPpcbStateByRegionNodes(response), + LocationHealthStatus.HealthyTentative, + LocationHealthStatus.Healthy); if (response.cosmosItemResponse != null) { assertThat(response.cosmosItemResponse).isNotNull(); @@ -3958,6 +3967,57 @@ private static CosmosDiagnosticsContext getDiagnosticsContext(ResponseWrapper return null; } + private static List getPpcbStateByRegionNodes(ResponseWrapper response) { + CosmosDiagnosticsContext diagnosticsContext = getDiagnosticsContext(response); + assertThat(diagnosticsContext).isNotNull(); + + try { + JsonNode diagnostics = Utils.getSimpleObjectMapper().readTree(diagnosticsContext.toJson()); + List stateByRegionNodes = new ArrayList<>(); + for (JsonNode ppcbNode : diagnostics.findValues("ppcb")) { + JsonNode stateByRegion = ppcbNode.get("stateByRegion"); + if (stateByRegion != null && stateByRegion.isObject()) { + stateByRegionNodes.add(stateByRegion); + } + } + + assertThat(stateByRegionNodes) + .as("Expected every PPCB-enabled data-plane operation to include ppcb.stateByRegion. Diagnostics: %s", diagnostics) + .isNotEmpty(); + return stateByRegionNodes; + } catch (Exception e) { + throw new AssertionError("Failed to parse CosmosDiagnostics for PPCB state", e); + } + } + + private static void assertPpcbHealthStatus( + List stateByRegionNodes, + LocationHealthStatus... expectedStatuses) { + + List actualStatuses = new ArrayList<>(); + for (JsonNode stateByRegion : stateByRegionNodes) { + Iterator regionStates = stateByRegion.elements(); + while (regionStates.hasNext()) { + JsonNode healthStatus = regionStates.next().get("locationHealthStatus"); + if (healthStatus != null) { + actualStatuses.add(healthStatus.asText()); + } + } + } + + boolean expectedStatusFound = false; + for (LocationHealthStatus expectedStatus : expectedStatuses) { + if (actualStatuses.contains(expectedStatus.toString())) { + expectedStatusFound = true; + break; + } + } + + assertThat(expectedStatusFound) + .as("Expected PPCB health status to be one of %s but found %s", Arrays.toString(expectedStatuses), actualStatuses) + .isTrue(); + } + private ResponseWrapper executeDataPlaneOperationWithTransient4041002Retry( String testId, Function> executeDataPlaneOperation, From 688cd1154d056527905bb43421905c1ac9888ecf Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Thu, 13 Aug 2026 10:50:02 -0400 Subject: [PATCH 3/8] Limit PPCB diagnostics assertion to data-plane requests --- .../com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java index 13c95b48653d..0660ab0d717a 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/PerPartitionCircuitBreakerE2ETests.java @@ -3808,7 +3808,6 @@ private void execute( testId, executeDataPlaneOperation, operationInvocationParamsWrapper); - List ppcbStateByRegionNodes = getPpcbStateByRegionNodes(response); ConsecutiveExceptionBasedCircuitBreaker consecutiveExceptionBasedCircuitBreaker = globalPartitionEndpointManagerForPerPartitionCircuitBreaker.getConsecutiveExceptionBasedCircuitBreaker(); @@ -3835,7 +3834,7 @@ private void execute( if (executionCountAfterCircuitBreakingThresholdBreached > 1) { validateResponseInAbsenceOfFailures.accept(response); assertPpcbHealthStatus( - ppcbStateByRegionNodes, + getPpcbStateByRegionNodes(response), LocationHealthStatus.Unavailable); } From 84c1091e07c14eb9206560d3e03f53a9aa469922 Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Thu, 13 Aug 2026 11:30:55 -0400 Subject: [PATCH 4/8] Log PPCB failback backlog progress --- .../PpcbFailbackLoggingTest.java | 72 +++++++++-- ...tManagerForPerPartitionCircuitBreaker.java | 117 +++++++++++------- 2 files changed, 131 insertions(+), 58 deletions(-) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java index 60a369488bac..304758cf8aac 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java @@ -13,9 +13,13 @@ import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; import reactor.core.publisher.Flux; +import reactor.core.scheduler.Scheduler; +import reactor.core.scheduler.Schedulers; import java.net.URI; import java.time.Duration; +import java.util.HashSet; +import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; import static org.assertj.core.api.Assertions.assertThat; @@ -116,23 +120,67 @@ public void manyPartitionsUseConstantSamplingState() { failure); } - verify(this.logger, times(11)).warn(contains("reason: RuntimeException"), same(failure)); - verify(this.logger, times(89)).debug(contains("reason: RuntimeException"), same(failure)); + verify(this.logger, times(11)).warn(contains("reason: RuntimeException"), same(failure)); + verify(this.logger, times(89)).debug(contains("reason: RuntimeException"), same(failure)); } - @Test(groups = {"unit"}) - public void unexpectedStreamFailureIsLoggedAndRetried() { - RuntimeException failure = new RuntimeException("failure"); - AtomicInteger subscriptions = new AtomicInteger(); - Flux recoveryWork = Flux.defer(() -> subscriptions.incrementAndGet() == 1 - ? Flux.error(failure) - : Flux.just("recovered")); + @Test(groups = {"unit"}) + public void unexpectedStreamFailureIsLoggedAndRetried() { + RuntimeException failure = new RuntimeException("failure"); + AtomicInteger subscriptions = new AtomicInteger(); + Flux recoveryWork = Flux.defer(() -> subscriptions.incrementAndGet() == 1 + ? Flux.error(failure) + : Flux.just("recovered")); + + Object result = this.manager.keepFailbackRecoveryAlive(recoveryWork) + .blockFirst(Duration.ofSeconds(1)); + + assertThat(result).isEqualTo("recovered"); + assertThat(subscriptions.get()).isEqualTo(2); + verify(this.logger).warn(contains("stage: RECOVERY_STREAM"), same(failure)); + } + + @Test(groups = {"unit"}) + public void recoverySurvivesRepeatedFailuresOnOneThreadScheduler() { + Scheduler scheduler = Schedulers.newBoundedElastic(1, 1, "ppcb-resilience-test"); + AtomicInteger subscriptions = new AtomicInteger(); + Set threadNames = new HashSet<>(); + + try { + Flux recoveryWork = Flux.defer(() -> { + threadNames.add(Thread.currentThread().getName()); + return subscriptions.incrementAndGet() <= 3 + ? Flux.error(new RuntimeException("failure")) + : Flux.just("recovered"); + }).subscribeOn(scheduler); Object result = this.manager.keepFailbackRecoveryAlive(recoveryWork) - .blockFirst(Duration.ofSeconds(1)); + .blockFirst(Duration.ofSeconds(5)); assertThat(result).isEqualTo("recovered"); - assertThat(subscriptions.get()).isEqualTo(2); - verify(this.logger).warn(contains("stage: RECOVERY_STREAM"), same(failure)); + assertThat(subscriptions.get()).isEqualTo(4); + assertThat(threadNames).hasSize(1); + assertThat(threadNames.iterator().next()).startsWith("ppcb-resilience-test"); + verify(this.logger, times(1)).warn(contains("stage: RECOVERY_STREAM"), Mockito.any(RuntimeException.class)); + verify(this.logger, times(2)).debug(contains("stage: RECOVERY_STREAM"), Mockito.any(RuntimeException.class)); + } finally { + scheduler.dispose(); + } + } + + @Test(groups = {"unit"}) + public void failbackBacklogProgressIsSampledAndCompletionIsLogged() { + for (int scan = 0; scan < 10; scan++) { + this.manager.logFailbackBacklog(7); } + this.manager.logFailbackBacklog(0); + this.manager.logFailbackBacklog(0); + + verify(this.logger, times(2)).info(contains( + "PPCB failback backlog: unavailablePartitionRegionCount: 7")); + verify(this.logger, times(8)).debug(contains( + "PPCB failback backlog: unavailablePartitionRegionCount: 7")); + verify(this.logger).info(contains( + "PPCB failback backlog: unavailablePartitionRegionCount: 0")); + } } \ No newline at end of file diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java index 97a582c3483a..f35c86761df9 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java @@ -60,6 +60,7 @@ public class GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker impleme private final AtomicReference partitionRecoveryDisposable = new AtomicReference<>(); private final Logger failbackLogger; private final AtomicInteger failbackFailureLogCount; + private final AtomicInteger failbackBacklogScanCount; public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpointManager globalEndpointManager) { this(globalEndpointManager, logger); @@ -73,6 +74,7 @@ public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpoin this.globalEndpointManager = globalEndpointManager; this.failbackLogger = checkNotNull(failbackLogger, "Argument 'failbackLogger' cannot be null!"); this.failbackFailureLogCount = new AtomicInteger(); + this.failbackBacklogScanCount = new AtomicInteger(); PartitionLevelCircuitBreakerConfig partitionLevelCircuitBreakerConfig = Configs.getPartitionLevelCircuitBreakerConfig(); this.consecutiveExceptionBasedCircuitBreaker = new ConsecutiveExceptionBasedCircuitBreaker(partitionLevelCircuitBreakerConfig); @@ -312,52 +314,7 @@ private Flux updateStaleLocationInfo() { Flux recoveryWork = Mono.just(1) .delayElement(Duration.ofSeconds(Configs.getStalePartitionUnavailabilityRefreshIntervalInSeconds())) .repeat(() -> !this.isClosed.get()) - .flatMap(ignore -> Flux.fromIterable(this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.entrySet()), 1, 1) - .flatMap(partitionKeyRangeWrapperToPartitionKeyRangeWrapperPair -> { - - logger.debug("Background updateStaleLocationInfo kicking in..."); - PartitionKeyRangeWrapper partitionKeyRangeWrapper - = partitionKeyRangeWrapperToPartitionKeyRangeWrapperPair.getKey(); - - try { - PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.get(partitionKeyRangeWrapper); - - if (partitionLevelLocationUnavailabilityInfo != null) { - - List>> locationToLocationSpecificHealthContextList = new ArrayList<>(); - - for (Map.Entry locationToLocationLevelMetrics : partitionLevelLocationUnavailabilityInfo.locationEndpointToLocationSpecificContextForPartition.entrySet()) { - - RegionalRoutingContext locationWithStaleUnavailabilityInfo = locationToLocationLevelMetrics.getKey(); - LocationSpecificHealthContext locationSpecificHealthContext = locationToLocationLevelMetrics.getValue(); - - if (!locationSpecificHealthContext.isRegionAvailableToProcessRequests()) { - locationToLocationSpecificHealthContextList.add( - Pair.of( - partitionKeyRangeWrapper, - Pair.of( - locationWithStaleUnavailabilityInfo, - locationSpecificHealthContext))); - } - } - - if (locationToLocationSpecificHealthContextList.isEmpty()) { - return Flux.empty(); - } else { - return Flux.fromIterable(locationToLocationSpecificHealthContextList); - } - } else { - return Mono.empty(); - } - } catch (Exception e) { - this.logFailbackFailure( - partitionKeyRangeWrapper, - null, - "SCAN_UNAVAILABLE_PARTITIONS", - e); - return Flux.empty(); - } - }, 1, 1) + .flatMap(ignore -> this.getFailbackBacklog(), 1, 1) .flatMap(locationToLocationSpecificHealthContextPair -> { try { @@ -436,6 +393,73 @@ private Flux updateStaleLocationInfo() { return this.keepFailbackRecoveryAlive(recoveryWork); } + private Flux>> getFailbackBacklog() { + List>> failbackBacklog + = new ArrayList<>(); + + for (Map.Entry entry + : this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.entrySet()) { + + PartitionKeyRangeWrapper partitionKeyRangeWrapper = entry.getKey(); + + try { + PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = entry.getValue(); + + if (partitionLevelLocationUnavailabilityInfo != null) { + + for (Map.Entry locationToLocationLevelMetrics + : partitionLevelLocationUnavailabilityInfo.locationEndpointToLocationSpecificContextForPartition.entrySet()) { + + RegionalRoutingContext locationWithStaleUnavailabilityInfo = locationToLocationLevelMetrics.getKey(); + LocationSpecificHealthContext locationSpecificHealthContext = locationToLocationLevelMetrics.getValue(); + + if (!locationSpecificHealthContext.isRegionAvailableToProcessRequests()) { + failbackBacklog.add( + Pair.of( + partitionKeyRangeWrapper, + Pair.of( + locationWithStaleUnavailabilityInfo, + locationSpecificHealthContext))); + } + } + } + } catch (Exception e) { + this.logFailbackFailure( + partitionKeyRangeWrapper, + null, + "SCAN_UNAVAILABLE_PARTITIONS", + e); + } + } + + this.logFailbackBacklog(failbackBacklog.size()); + return Flux.fromIterable(failbackBacklog); + } + + void logFailbackBacklog(int unavailablePartitionRegionCount) { + int previousScanCount = unavailablePartitionRegionCount == 0 + ? this.failbackBacklogScanCount.getAndSet(0) + : this.failbackBacklogScanCount.updateAndGet( + current -> current == Integer.MAX_VALUE ? 1 : current + 1); + + if (unavailablePartitionRegionCount == 0 && previousScanCount == 0) { + return; + } + + String message = "PPCB failback backlog: unavailablePartitionRegionCount: " + + unavailablePartitionRegionCount + + ", trackedPartitionKeyRangeCount: " + + this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.size() + + ", consecutiveBacklogScanCount: " + + (unavailablePartitionRegionCount == 0 ? 0 : previousScanCount); + + if (unavailablePartitionRegionCount == 0 || previousScanCount == 1 || previousScanCount % 10 == 0) { + this.failbackLogger.info(message); + } else { + this.failbackLogger.debug(message); + } + } + Flux keepFailbackRecoveryAlive(Flux recoveryWork) { return recoveryWork .doOnError(throwable -> this.logFailbackFailure( @@ -552,6 +576,7 @@ public void setGlobalAddressResolver(GlobalAddressResolver globalAddressResolver public void close() { this.isClosed.set(true); this.failbackFailureLogCount.set(0); + this.failbackBacklogScanCount.set(0); Disposable disposable = this.partitionRecoveryDisposable.getAndSet(null); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); From 55cfe30f871fcc7b0bc1d7a81b33ecd46ffe2184 Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Thu, 13 Aug 2026 12:04:52 -0400 Subject: [PATCH 5/8] Add PPCB failback remaining meter --- .../PpcbFailbackLoggingTest.java | 51 ++++++++++++++++++ sdk/cosmos/azure-cosmos/CHANGELOG.md | 1 + .../com/azure/cosmos/CosmosAsyncClient.java | 12 +++++ ...tManagerForPerPartitionCircuitBreaker.java | 52 +++++++++++++++++++ .../azure/cosmos/models/CosmosMetricName.java | 8 +++ 5 files changed, 124 insertions(+) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java index 304758cf8aac..6fd70f5c2b31 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java @@ -8,6 +8,10 @@ import com.azure.cosmos.implementation.PartitionKeyRange; import com.azure.cosmos.implementation.PartitionKeyRangeWrapper; import com.azure.cosmos.implementation.routing.RegionalRoutingContext; +import com.azure.cosmos.models.CosmosMetricName; +import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.Tag; +import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import org.mockito.Mockito; import org.slf4j.Logger; import org.testng.annotations.BeforeMethod; @@ -18,7 +22,11 @@ import java.net.URI; import java.time.Duration; +import java.util.Arrays; +import java.util.Collections; import java.util.HashSet; +import java.util.HashMap; +import java.util.Map; import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; @@ -183,4 +191,47 @@ public void failbackBacklogProgressIsSampledAndCompletionIsLogged() { verify(this.logger).info(contains( "PPCB failback backlog: unavailablePartitionRegionCount: 0")); } + + @Test(groups = {"unit"}) + public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { + SimpleMeterRegistry registry = new SimpleMeterRegistry(); + Map> remainingByCollection = new HashMap<>(); + remainingByCollection.put("collectionA", new HashSet<>(Arrays.asList("0", "1", "2"))); + remainingByCollection.put("collectionB", new HashSet<>(Arrays.asList("3", "4"))); + + try { + assertThat(CosmosMetricName.fromString("cosmos.client.ppcb.failback.remaining")) + .isSameAs(CosmosMetricName.PPCB_FAILBACK_REMAINING); + this.manager.registerFailbackRemainingMeter( + registry, + Tag.of("ClientCorrelationId", "client1")); + this.manager.recordFailbackRemainingByCollection(remainingByCollection); + + assertThat(getFailbackRemainingGauge(registry, "collectionA").value()).isEqualTo(3); + assertThat(getFailbackRemainingGauge(registry, "collectionB").value()).isEqualTo(2); + + Map> updatedRemainingByCollection = new HashMap<>(); + updatedRemainingByCollection.put("collectionA", Collections.emptySet()); + updatedRemainingByCollection.put("collectionB", Collections.singleton("4")); + this.manager.recordFailbackRemainingByCollection(updatedRemainingByCollection); + + assertThat(getFailbackRemainingGauge(registry, "collectionA").value()).isZero(); + assertThat(getFailbackRemainingGauge(registry, "collectionB").value()).isEqualTo(1); + + this.manager.close(); + assertThat(registry.find(CosmosMetricName.PPCB_FAILBACK_REMAINING.toString()) + .gauges()).isEmpty(); + } finally { + this.manager.close(); + registry.close(); + } + } + + private static Gauge getFailbackRemainingGauge(SimpleMeterRegistry registry, String collectionRid) { + return registry + .get(CosmosMetricName.PPCB_FAILBACK_REMAINING.toString()) + .tag("ClientCorrelationId", "client1") + .tag("CollectionRid", collectionRid) + .gauge(); + } } \ No newline at end of file diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index 4105b1622f13..232915fb2a12 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -5,6 +5,7 @@ #### Features Added * Enabled Gateway V2 (thin-client) data-plane routing by default for `Cosmos(Async)Client` instances configured with `gatewayMode` and HTTP/2, gated by an HTTP/2 connectivity probe with automatic fallback to Gateway V1. - See [PR 49437](https://github.com/Azure/azure-sdk-for-java/pull/49437) * Added support for QueryPlan and Execute Stored Procedure requests to be routed to Gateway V2. - See [PR 47759](https://github.com/Azure/azure-sdk-for-java/pull/47759) +* Added the `cosmos.client.ppcb.failback.remaining` gauge (`CosmosMetricName.PPCB_FAILBACK_REMAINING`) to report partition key ranges remaining to fail back per collection. - See [PR 50122](https://github.com/Azure/azure-sdk-for-java/pull/50122) #### Breaking Changes diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java index a5db3fbfd167..b5cb00ee6214 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java @@ -220,6 +220,8 @@ private static ImplementationBridgeHelpers.FeedResponseHelper.FeedResponseAccess .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.SYSTEM_CPU); CosmosMeterOptions memoryMeterOptions = clientTelemetryConfigAccessor() .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.SYSTEM_MEMORY_FREE); + CosmosMeterOptions failbackRemainingMeterOptions = clientTelemetryConfigAccessor() + .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.PPCB_FAILBACK_REMAINING); if (clientMetricRegistrySnapshot != null) { ClientTelemetryMetrics.add(clientMetricRegistrySnapshot, cpuMeterOptions, memoryMeterOptions); @@ -236,6 +238,13 @@ private static ImplementationBridgeHelpers.FeedResponseHelper.FeedResponseAccess effectiveTelemetryConfig, this.accountTagValue ); + if (failbackRemainingMeterOptions.isEnabled()) { + this.asyncDocumentClient + .getGlobalPartitionEndpointManagerForCircuitBreaker() + .registerFailbackRemainingMeter( + this.clientMetricRegistrySnapshot, + this.clientCorrelationTag); + } clientTelemetryConfigAccessor().addDiagnosticsHandler( effectiveTelemetryConfig, @@ -563,6 +572,9 @@ public CosmosAsyncDatabase getDatabase(String id) { @Override public void close() { if (this.clientMetricRegistrySnapshot != null) { + this.asyncDocumentClient + .getGlobalPartitionEndpointManagerForCircuitBreaker() + .removeFailbackRemainingMeter(); ClientTelemetryMetrics.remove(this.clientMetricRegistrySnapshot); } asyncDocumentClient.close(); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java index f35c86761df9..881f161a3d22 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java @@ -21,6 +21,11 @@ import com.azure.cosmos.implementation.directconnectivity.GatewayAddressCache; import com.azure.cosmos.implementation.directconnectivity.GlobalAddressResolver; import com.azure.cosmos.implementation.routing.RegionalRoutingContext; +import com.azure.cosmos.models.CosmosMetricName; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.MultiGauge; +import io.micrometer.core.instrument.Tag; +import io.micrometer.core.instrument.Tags; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import reactor.core.Disposable; @@ -33,9 +38,11 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.PriorityQueue; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -46,6 +53,7 @@ public class GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker implements AutoCloseable { private static final Logger logger = LoggerFactory.getLogger(GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.class); + private static final String COLLECTION_RID_TAG_NAME = "CollectionRid"; private static final Map EMPTY_MAP = new HashMap<>(); private static final String BASE_EXCEPTION_MESSAGE = "FAILED IN Per-Partition Circuit Breaker: "; @@ -61,6 +69,7 @@ public class GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker impleme private final Logger failbackLogger; private final AtomicInteger failbackFailureLogCount; private final AtomicInteger failbackBacklogScanCount; + private final AtomicReference failbackRemainingGauge; public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpointManager globalEndpointManager) { this(globalEndpointManager, logger); @@ -75,6 +84,7 @@ public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpoin this.failbackLogger = checkNotNull(failbackLogger, "Argument 'failbackLogger' cannot be null!"); this.failbackFailureLogCount = new AtomicInteger(); this.failbackBacklogScanCount = new AtomicInteger(); + this.failbackRemainingGauge = new AtomicReference<>(); PartitionLevelCircuitBreakerConfig partitionLevelCircuitBreakerConfig = Configs.getPartitionLevelCircuitBreakerConfig(); this.consecutiveExceptionBasedCircuitBreaker = new ConsecutiveExceptionBasedCircuitBreaker(partitionLevelCircuitBreakerConfig); @@ -396,11 +406,16 @@ private Flux updateStaleLocationInfo() { private Flux>> getFailbackBacklog() { List>> failbackBacklog = new ArrayList<>(); + Map> remainingPartitionRangeIdsByCollection = new HashMap<>(); for (Map.Entry entry : this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.entrySet()) { PartitionKeyRangeWrapper partitionKeyRangeWrapper = entry.getKey(); + Set remainingPartitionRangeIds = remainingPartitionRangeIdsByCollection + .computeIfAbsent( + partitionKeyRangeWrapper.getCollectionResourceId(), + ignored -> new HashSet<>()); try { PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = entry.getValue(); @@ -420,6 +435,7 @@ private Flux> remainingPartitionRangeIdsByCollection) { + MultiGauge gauge = this.failbackRemainingGauge.get(); + if (gauge == null) { + return; + } + + List> rows = new ArrayList<>(); + remainingPartitionRangeIdsByCollection.forEach((collectionRid, partitionRangeIds) -> + rows.add(MultiGauge.Row.of( + Tags.of(COLLECTION_RID_TAG_NAME, collectionRid), + partitionRangeIds.size()))); + gauge.register(rows, true); + } + + public void removeFailbackRemainingMeter() { + MultiGauge gauge = this.failbackRemainingGauge.getAndSet(null); + if (gauge != null) { + gauge.register(Collections.emptyList(), true); + } + } + void logFailbackBacklog(int unavailablePartitionRegionCount) { int previousScanCount = unavailablePartitionRegionCount == 0 ? this.failbackBacklogScanCount.getAndSet(0) @@ -577,6 +628,7 @@ public void close() { this.isClosed.set(true); this.failbackFailureLogCount.set(0); this.failbackBacklogScanCount.set(0); + this.removeFailbackRemainingMeter(); Disposable disposable = this.partitionRecoveryDisposable.getAndSet(null); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java index 54dd12c73b21..1cc5cc4557ec 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java @@ -76,6 +76,13 @@ private CosmosMetricName(String name, CosmosMetricCategory metricCategory) { nameOf("op.maxItemCount"), CosmosMetricCategory.OPERATION_DETAILS); + /** + * Number of PPCB partition key ranges remaining to fail back, by collection (Gauge) + */ + public static final CosmosMetricName PPCB_FAILBACK_REMAINING = new CosmosMetricName( + nameOf("ppcb.failback.remaining"), + CosmosMetricCategory.OPERATION_SUMMARY); + /** * Number of requests (Counter) * NOTE: No percentiles or histogram supported @@ -465,6 +472,7 @@ private static Map createMeterNameMap() { map.put(nameOf("op.maxitemcount"), CosmosMetricName.OPERATION_DETAILS_MAX_ITEM_COUNT); map.put(nameOf("op.actualitemcount"), CosmosMetricName.OPERATION_DETAILS_ACTUAL_ITEM_COUNT); map.put(nameOf("op.regionscontacted"), CosmosMetricName.OPERATION_DETAILS_REGIONS_CONTACTED); + map.put(nameOf("ppcb.failback.remaining"), CosmosMetricName.PPCB_FAILBACK_REMAINING); map.put(nameOf("req.rntbd.requests"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_REQUESTS); map.put(nameOf("req.rntbd.latency"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_LATENCY); map.put(nameOf("req.rntbd.backendlatency"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_BACKEND_LATENCY); From 7bd8549c571c8e33e8d429207f3596742c14273a Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Thu, 13 Aug 2026 12:11:21 -0400 Subject: [PATCH 6/8] Clarify PPCB failback meter name --- .../PpcbFailbackLoggingTest.java | 8 ++++---- sdk/cosmos/azure-cosmos/CHANGELOG.md | 2 +- .../src/main/java/com/azure/cosmos/CosmosAsyncClient.java | 2 +- ...itionEndpointManagerForPerPartitionCircuitBreaker.java | 2 +- .../java/com/azure/cosmos/models/CosmosMetricName.java | 8 +++++--- 5 files changed, 12 insertions(+), 10 deletions(-) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java index 6fd70f5c2b31..549ada4b7348 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java @@ -200,8 +200,8 @@ public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { remainingByCollection.put("collectionB", new HashSet<>(Arrays.asList("3", "4"))); try { - assertThat(CosmosMetricName.fromString("cosmos.client.ppcb.failback.remaining")) - .isSameAs(CosmosMetricName.PPCB_FAILBACK_REMAINING); + assertThat(CosmosMetricName.fromString("cosmos.client.ppcb.failback.pendingPartitionCount")) + .isSameAs(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT); this.manager.registerFailbackRemainingMeter( registry, Tag.of("ClientCorrelationId", "client1")); @@ -219,7 +219,7 @@ public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { assertThat(getFailbackRemainingGauge(registry, "collectionB").value()).isEqualTo(1); this.manager.close(); - assertThat(registry.find(CosmosMetricName.PPCB_FAILBACK_REMAINING.toString()) + assertThat(registry.find(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT.toString()) .gauges()).isEmpty(); } finally { this.manager.close(); @@ -229,7 +229,7 @@ public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { private static Gauge getFailbackRemainingGauge(SimpleMeterRegistry registry, String collectionRid) { return registry - .get(CosmosMetricName.PPCB_FAILBACK_REMAINING.toString()) + .get(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT.toString()) .tag("ClientCorrelationId", "client1") .tag("CollectionRid", collectionRid) .gauge(); diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index 232915fb2a12..dfb143407ccf 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -5,7 +5,7 @@ #### Features Added * Enabled Gateway V2 (thin-client) data-plane routing by default for `Cosmos(Async)Client` instances configured with `gatewayMode` and HTTP/2, gated by an HTTP/2 connectivity probe with automatic fallback to Gateway V1. - See [PR 49437](https://github.com/Azure/azure-sdk-for-java/pull/49437) * Added support for QueryPlan and Execute Stored Procedure requests to be routed to Gateway V2. - See [PR 47759](https://github.com/Azure/azure-sdk-for-java/pull/47759) -* Added the `cosmos.client.ppcb.failback.remaining` gauge (`CosmosMetricName.PPCB_FAILBACK_REMAINING`) to report partition key ranges remaining to fail back per collection. - See [PR 50122](https://github.com/Azure/azure-sdk-for-java/pull/50122) +* Added the `cosmos.client.ppcb.failback.pendingPartitionCount` gauge (`CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT`) to report partition key ranges remaining to fail back per collection. - See [PR 50122](https://github.com/Azure/azure-sdk-for-java/pull/50122) #### Breaking Changes diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java index b5cb00ee6214..17de2664803c 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java @@ -221,7 +221,7 @@ private static ImplementationBridgeHelpers.FeedResponseHelper.FeedResponseAccess CosmosMeterOptions memoryMeterOptions = clientTelemetryConfigAccessor() .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.SYSTEM_MEMORY_FREE); CosmosMeterOptions failbackRemainingMeterOptions = clientTelemetryConfigAccessor() - .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.PPCB_FAILBACK_REMAINING); + .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT); if (clientMetricRegistrySnapshot != null) { ClientTelemetryMetrics.add(clientMetricRegistrySnapshot, cpuMeterOptions, memoryMeterOptions); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java index 881f161a3d22..3a1a5420d412 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java @@ -459,7 +459,7 @@ public synchronized void registerFailbackRemainingMeter(MeterRegistry meterRegis } this.failbackRemainingGauge.set( - MultiGauge.builder(CosmosMetricName.PPCB_FAILBACK_REMAINING.toString()) + MultiGauge.builder(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT.toString()) .description("Number of PPCB partition key ranges remaining to fail back") .baseUnit("partitions") .tags(Collections.singletonList(clientCorrelationTag)) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java index 1cc5cc4557ec..589b8a4c9943 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java @@ -79,8 +79,8 @@ private CosmosMetricName(String name, CosmosMetricCategory metricCategory) { /** * Number of PPCB partition key ranges remaining to fail back, by collection (Gauge) */ - public static final CosmosMetricName PPCB_FAILBACK_REMAINING = new CosmosMetricName( - nameOf("ppcb.failback.remaining"), + public static final CosmosMetricName PPCB_FAILBACK_PENDING_PARTITION_COUNT = new CosmosMetricName( + nameOf("ppcb.failback.pendingPartitionCount"), CosmosMetricCategory.OPERATION_SUMMARY); /** @@ -472,7 +472,9 @@ private static Map createMeterNameMap() { map.put(nameOf("op.maxitemcount"), CosmosMetricName.OPERATION_DETAILS_MAX_ITEM_COUNT); map.put(nameOf("op.actualitemcount"), CosmosMetricName.OPERATION_DETAILS_ACTUAL_ITEM_COUNT); map.put(nameOf("op.regionscontacted"), CosmosMetricName.OPERATION_DETAILS_REGIONS_CONTACTED); - map.put(nameOf("ppcb.failback.remaining"), CosmosMetricName.PPCB_FAILBACK_REMAINING); + map.put( + nameOf("ppcb.failback.pendingpartitioncount"), + CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT); map.put(nameOf("req.rntbd.requests"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_REQUESTS); map.put(nameOf("req.rntbd.latency"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_LATENCY); map.put(nameOf("req.rntbd.backendlatency"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_BACKEND_LATENCY); From 2683720e7778c9ba5ea7f3fb3aebfd644e00ac56 Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Thu, 13 Aug 2026 12:21:56 -0400 Subject: [PATCH 7/8] Reduce PPCB failback meter allocations --- .../PpcbFailbackLoggingTest.java | 15 ++++++------ ...tManagerForPerPartitionCircuitBreaker.java | 23 +++++++++++-------- 2 files changed, 20 insertions(+), 18 deletions(-) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java index 549ada4b7348..40db8eee81e4 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java @@ -22,10 +22,9 @@ import java.net.URI; import java.time.Duration; -import java.util.Arrays; import java.util.Collections; -import java.util.HashSet; import java.util.HashMap; +import java.util.HashSet; import java.util.Map; import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; @@ -195,9 +194,9 @@ public void failbackBacklogProgressIsSampledAndCompletionIsLogged() { @Test(groups = {"unit"}) public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { SimpleMeterRegistry registry = new SimpleMeterRegistry(); - Map> remainingByCollection = new HashMap<>(); - remainingByCollection.put("collectionA", new HashSet<>(Arrays.asList("0", "1", "2"))); - remainingByCollection.put("collectionB", new HashSet<>(Arrays.asList("3", "4"))); + Map remainingByCollection = new HashMap<>(); + remainingByCollection.put("collectionA", new AtomicInteger(3)); + remainingByCollection.put("collectionB", new AtomicInteger(2)); try { assertThat(CosmosMetricName.fromString("cosmos.client.ppcb.failback.pendingPartitionCount")) @@ -210,9 +209,9 @@ public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { assertThat(getFailbackRemainingGauge(registry, "collectionA").value()).isEqualTo(3); assertThat(getFailbackRemainingGauge(registry, "collectionB").value()).isEqualTo(2); - Map> updatedRemainingByCollection = new HashMap<>(); - updatedRemainingByCollection.put("collectionA", Collections.emptySet()); - updatedRemainingByCollection.put("collectionB", Collections.singleton("4")); + Map updatedRemainingByCollection = new HashMap<>(); + updatedRemainingByCollection.put("collectionA", new AtomicInteger()); + updatedRemainingByCollection.put("collectionB", new AtomicInteger(1)); this.manager.recordFailbackRemainingByCollection(updatedRemainingByCollection); assertThat(getFailbackRemainingGauge(registry, "collectionA").value()).isZero(); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java index 3a1a5420d412..8582b2d615a2 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java @@ -38,11 +38,9 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; -import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.PriorityQueue; -import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; @@ -406,16 +404,17 @@ private Flux updateStaleLocationInfo() { private Flux>> getFailbackBacklog() { List>> failbackBacklog = new ArrayList<>(); - Map> remainingPartitionRangeIdsByCollection = new HashMap<>(); + Map remainingPartitionCountByCollection = new HashMap<>(); for (Map.Entry entry : this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.entrySet()) { PartitionKeyRangeWrapper partitionKeyRangeWrapper = entry.getKey(); - Set remainingPartitionRangeIds = remainingPartitionRangeIdsByCollection + AtomicInteger remainingPartitionCount = remainingPartitionCountByCollection .computeIfAbsent( partitionKeyRangeWrapper.getCollectionResourceId(), - ignored -> new HashSet<>()); + ignored -> new AtomicInteger()); + boolean isPartitionRangeUnavailable = false; try { PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = entry.getValue(); @@ -429,13 +428,13 @@ private Flux> remainingPartitionRangeIdsByCollection) { + void recordFailbackRemainingByCollection(Map remainingPartitionCountByCollection) { MultiGauge gauge = this.failbackRemainingGauge.get(); if (gauge == null) { return; } List> rows = new ArrayList<>(); - remainingPartitionRangeIdsByCollection.forEach((collectionRid, partitionRangeIds) -> + remainingPartitionCountByCollection.forEach((collectionRid, partitionCount) -> rows.add(MultiGauge.Row.of( Tags.of(COLLECTION_RID_TAG_NAME, collectionRid), - partitionRangeIds.size()))); + partitionCount))); gauge.register(rows, true); } From 01dc822b90533b6e1d35e00ed1509e34d4aae0a4 Mon Sep 17 00:00:00 2001 From: Abhijeet Mohanty Date: Thu, 13 Aug 2026 12:44:58 -0400 Subject: [PATCH 8/8] Track PPCB pending recoveries by collection --- .../PpcbFailbackLoggingTest.java | 41 +++---- sdk/cosmos/azure-cosmos/CHANGELOG.md | 2 +- .../com/azure/cosmos/CosmosAsyncClient.java | 10 +- ...tManagerForPerPartitionCircuitBreaker.java | 103 ++++++++++-------- .../azure/cosmos/models/CosmosMetricName.java | 10 +- 5 files changed, 89 insertions(+), 77 deletions(-) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java index 40db8eee81e4..1e774e9f4e93 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/PpcbFailbackLoggingTest.java @@ -192,33 +192,36 @@ public void failbackBacklogProgressIsSampledAndCompletionIsLogged() { } @Test(groups = {"unit"}) - public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { + public void failbackPendingRecoveryGaugeIsPerCollectionAndReportsZero() { SimpleMeterRegistry registry = new SimpleMeterRegistry(); - Map remainingByCollection = new HashMap<>(); - remainingByCollection.put("collectionA", new AtomicInteger(3)); - remainingByCollection.put("collectionB", new AtomicInteger(2)); + Map pendingRecoveryByCollection = new HashMap<>(); + pendingRecoveryByCollection.put("collectionA", new AtomicInteger(3)); + pendingRecoveryByCollection.put("collectionB", new AtomicInteger(2)); try { - assertThat(CosmosMetricName.fromString("cosmos.client.ppcb.failback.pendingPartitionCount")) - .isSameAs(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT); - this.manager.registerFailbackRemainingMeter( + assertThat(CosmosMetricName.fromString("cosmos.client.ppcb.failback.pendingRecoveryCount")) + .isSameAs(CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT); + this.manager.registerFailbackPendingRecoveryMeter( registry, Tag.of("ClientCorrelationId", "client1")); - this.manager.recordFailbackRemainingByCollection(remainingByCollection); + this.manager.recordFailbackPendingRecoveryByCollection(pendingRecoveryByCollection); - assertThat(getFailbackRemainingGauge(registry, "collectionA").value()).isEqualTo(3); - assertThat(getFailbackRemainingGauge(registry, "collectionB").value()).isEqualTo(2); + assertThat(getPendingRecoveryGauge(registry, "collectionA").value()).isEqualTo(3); + assertThat(getPendingRecoveryGauge(registry, "collectionB").value()).isEqualTo(2); - Map updatedRemainingByCollection = new HashMap<>(); - updatedRemainingByCollection.put("collectionA", new AtomicInteger()); - updatedRemainingByCollection.put("collectionB", new AtomicInteger(1)); - this.manager.recordFailbackRemainingByCollection(updatedRemainingByCollection); + pendingRecoveryByCollection.get("collectionA").decrementAndGet(); + assertThat(getPendingRecoveryGauge(registry, "collectionA").value()).isEqualTo(2); - assertThat(getFailbackRemainingGauge(registry, "collectionA").value()).isZero(); - assertThat(getFailbackRemainingGauge(registry, "collectionB").value()).isEqualTo(1); + Map updatedPendingRecoveryByCollection = new HashMap<>(); + updatedPendingRecoveryByCollection.put("collectionA", new AtomicInteger()); + updatedPendingRecoveryByCollection.put("collectionB", new AtomicInteger(1)); + this.manager.recordFailbackPendingRecoveryByCollection(updatedPendingRecoveryByCollection); + + assertThat(getPendingRecoveryGauge(registry, "collectionA").value()).isZero(); + assertThat(getPendingRecoveryGauge(registry, "collectionB").value()).isEqualTo(1); this.manager.close(); - assertThat(registry.find(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT.toString()) + assertThat(registry.find(CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT.toString()) .gauges()).isEmpty(); } finally { this.manager.close(); @@ -226,9 +229,9 @@ public void failbackRemainingGaugeIsPerCollectionAndReportsZero() { } } - private static Gauge getFailbackRemainingGauge(SimpleMeterRegistry registry, String collectionRid) { + private static Gauge getPendingRecoveryGauge(SimpleMeterRegistry registry, String collectionRid) { return registry - .get(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT.toString()) + .get(CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT.toString()) .tag("ClientCorrelationId", "client1") .tag("CollectionRid", collectionRid) .gauge(); diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index dfb143407ccf..59ed34744ec2 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -5,7 +5,7 @@ #### Features Added * Enabled Gateway V2 (thin-client) data-plane routing by default for `Cosmos(Async)Client` instances configured with `gatewayMode` and HTTP/2, gated by an HTTP/2 connectivity probe with automatic fallback to Gateway V1. - See [PR 49437](https://github.com/Azure/azure-sdk-for-java/pull/49437) * Added support for QueryPlan and Execute Stored Procedure requests to be routed to Gateway V2. - See [PR 47759](https://github.com/Azure/azure-sdk-for-java/pull/47759) -* Added the `cosmos.client.ppcb.failback.pendingPartitionCount` gauge (`CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT`) to report partition key ranges remaining to fail back per collection. - See [PR 50122](https://github.com/Azure/azure-sdk-for-java/pull/50122) +* Added the `cosmos.client.ppcb.failback.pendingRecoveryCount` gauge (`CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT`) to report partition range-region recovery actions pending failback per collection. - See [PR 50122](https://github.com/Azure/azure-sdk-for-java/pull/50122) #### Breaking Changes diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java index 17de2664803c..4e1382cc8fc3 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/CosmosAsyncClient.java @@ -220,8 +220,8 @@ private static ImplementationBridgeHelpers.FeedResponseHelper.FeedResponseAccess .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.SYSTEM_CPU); CosmosMeterOptions memoryMeterOptions = clientTelemetryConfigAccessor() .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.SYSTEM_MEMORY_FREE); - CosmosMeterOptions failbackRemainingMeterOptions = clientTelemetryConfigAccessor() - .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT); + CosmosMeterOptions failbackPendingRecoveryMeterOptions = clientTelemetryConfigAccessor() + .getMeterOptions(effectiveTelemetryConfig, CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT); if (clientMetricRegistrySnapshot != null) { ClientTelemetryMetrics.add(clientMetricRegistrySnapshot, cpuMeterOptions, memoryMeterOptions); @@ -238,10 +238,10 @@ private static ImplementationBridgeHelpers.FeedResponseHelper.FeedResponseAccess effectiveTelemetryConfig, this.accountTagValue ); - if (failbackRemainingMeterOptions.isEnabled()) { + if (failbackPendingRecoveryMeterOptions.isEnabled()) { this.asyncDocumentClient .getGlobalPartitionEndpointManagerForCircuitBreaker() - .registerFailbackRemainingMeter( + .registerFailbackPendingRecoveryMeter( this.clientMetricRegistrySnapshot, this.clientCorrelationTag); } @@ -574,7 +574,7 @@ public void close() { if (this.clientMetricRegistrySnapshot != null) { this.asyncDocumentClient .getGlobalPartitionEndpointManagerForCircuitBreaker() - .removeFailbackRemainingMeter(); + .removeFailbackPendingRecoveryMeter(); ClientTelemetryMetrics.remove(this.clientMetricRegistrySnapshot); } asyncDocumentClient.close(); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java index 8582b2d615a2..e07f925048c3 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/perPartitionCircuitBreaker/GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.java @@ -67,7 +67,7 @@ public class GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker impleme private final Logger failbackLogger; private final AtomicInteger failbackFailureLogCount; private final AtomicInteger failbackBacklogScanCount; - private final AtomicReference failbackRemainingGauge; + private final AtomicReference failbackPendingRecoveryGauge; public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpointManager globalEndpointManager) { this(globalEndpointManager, logger); @@ -82,7 +82,7 @@ public GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker(GlobalEndpoin this.failbackLogger = checkNotNull(failbackLogger, "Argument 'failbackLogger' cannot be null!"); this.failbackFailureLogCount = new AtomicInteger(); this.failbackBacklogScanCount = new AtomicInteger(); - this.failbackRemainingGauge = new AtomicReference<>(); + this.failbackPendingRecoveryGauge = new AtomicReference<>(); PartitionLevelCircuitBreakerConfig partitionLevelCircuitBreakerConfig = Configs.getPartitionLevelCircuitBreakerConfig(); this.consecutiveExceptionBasedCircuitBreaker = new ConsecutiveExceptionBasedCircuitBreaker(partitionLevelCircuitBreakerConfig); @@ -323,11 +323,12 @@ private Flux updateStaleLocationInfo() { .delayElement(Duration.ofSeconds(Configs.getStalePartitionUnavailabilityRefreshIntervalInSeconds())) .repeat(() -> !this.isClosed.get()) .flatMap(ignore -> this.getFailbackBacklog(), 1, 1) - .flatMap(locationToLocationSpecificHealthContextPair -> { + .flatMap(failbackRecoveryCandidate -> { try { - PartitionKeyRangeWrapper partitionKeyRangeWrapper = locationToLocationSpecificHealthContextPair.getLeft(); - RegionalRoutingContext locationWithStaleUnavailabilityInfo = locationToLocationSpecificHealthContextPair.getRight().getLeft(); + PartitionKeyRangeWrapper partitionKeyRangeWrapper = failbackRecoveryCandidate.getLeft(); + RegionalRoutingContext locationWithStaleUnavailabilityInfo = failbackRecoveryCandidate.getRight().getLeft(); + AtomicInteger pendingRecoveryCount = failbackRecoveryCandidate.getRight().getRight(); PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.get(partitionKeyRangeWrapper); @@ -352,11 +353,13 @@ private Flux updateStaleLocationInfo() { + partitionKeyRangeWrapper.getCollectionResourceId() + " has succeeded..."); - partitionLevelLocationUnavailabilityInfo.handleSuccess( - partitionKeyRangeWrapper, - locationWithStaleUnavailabilityInfo, - true, - true); + if (partitionLevelLocationUnavailabilityInfo.handleSuccess( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + true, + true)) { + this.markFailbackRecoveryCompleted(pendingRecoveryCount); + } }) .onErrorResume(throwable -> { @@ -375,18 +378,20 @@ private Flux updateStaleLocationInfo() { new IllegalStateException("GatewayAddressCache is not available.")); } } else { - partitionLevelLocationUnavailabilityInfo.handleSuccess( - partitionKeyRangeWrapper, - locationWithStaleUnavailabilityInfo, - true, - true); + if (partitionLevelLocationUnavailabilityInfo.handleSuccess( + partitionKeyRangeWrapper, + locationWithStaleUnavailabilityInfo, + true, + true)) { + this.markFailbackRecoveryCompleted(pendingRecoveryCount); + } } } } catch (Exception e) { - PartitionKeyRangeWrapper partitionKeyRangeWrapper = locationToLocationSpecificHealthContextPair.getLeft(); + PartitionKeyRangeWrapper partitionKeyRangeWrapper = failbackRecoveryCandidate.getLeft(); RegionalRoutingContext locationWithStaleUnavailabilityInfo - = locationToLocationSpecificHealthContextPair.getRight().getLeft(); + = failbackRecoveryCandidate.getRight().getLeft(); this.logFailbackFailure( partitionKeyRangeWrapper, locationWithStaleUnavailabilityInfo, @@ -401,20 +406,19 @@ private Flux updateStaleLocationInfo() { return this.keepFailbackRecoveryAlive(recoveryWork); } - private Flux>> getFailbackBacklog() { - List>> failbackBacklog + private Flux>> getFailbackBacklog() { + List>> failbackBacklog = new ArrayList<>(); - Map remainingPartitionCountByCollection = new HashMap<>(); + Map pendingRecoveryCountByCollection = new HashMap<>(); for (Map.Entry entry : this.partitionKeyRangeToLocationSpecificUnavailabilityInfo.entrySet()) { PartitionKeyRangeWrapper partitionKeyRangeWrapper = entry.getKey(); - AtomicInteger remainingPartitionCount = remainingPartitionCountByCollection + AtomicInteger pendingRecoveryCount = pendingRecoveryCountByCollection .computeIfAbsent( partitionKeyRangeWrapper.getCollectionResourceId(), ignored -> new AtomicInteger()); - boolean isPartitionRangeUnavailable = false; try { PartitionLevelLocationUnavailabilityInfo partitionLevelLocationUnavailabilityInfo = entry.getValue(); @@ -428,13 +432,10 @@ private Flux Math.max(0, current - 1)); + } + + public synchronized void registerFailbackPendingRecoveryMeter(MeterRegistry meterRegistry, Tag clientCorrelationTag) { + if (meterRegistry == null || clientCorrelationTag == null || this.failbackPendingRecoveryGauge.get() != null) { return; } - this.failbackRemainingGauge.set( - MultiGauge.builder(CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT.toString()) - .description("Number of PPCB partition key ranges remaining to fail back") - .baseUnit("partitions") + this.failbackPendingRecoveryGauge.set( + MultiGauge.builder(CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT.toString()) + .description("Number of PPCB partition range-region recovery actions pending failback") + .baseUnit("recoveries") .tags(Collections.singletonList(clientCorrelationTag)) .register(meterRegistry)); } - void recordFailbackRemainingByCollection(Map remainingPartitionCountByCollection) { - MultiGauge gauge = this.failbackRemainingGauge.get(); + void recordFailbackPendingRecoveryByCollection(Map pendingRecoveryCountByCollection) { + MultiGauge gauge = this.failbackPendingRecoveryGauge.get(); if (gauge == null) { return; } List> rows = new ArrayList<>(); - remainingPartitionCountByCollection.forEach((collectionRid, partitionCount) -> + pendingRecoveryCountByCollection.forEach((collectionRid, recoveryCount) -> rows.add(MultiGauge.Row.of( Tags.of(COLLECTION_RID_TAG_NAME, collectionRid), - partitionCount))); + recoveryCount))); gauge.register(rows, true); } - public void removeFailbackRemainingMeter() { - MultiGauge gauge = this.failbackRemainingGauge.getAndSet(null); + public void removeFailbackPendingRecoveryMeter() { + MultiGauge gauge = this.failbackPendingRecoveryGauge.getAndSet(null); if (gauge != null) { gauge.register(Collections.emptyList(), true); } @@ -631,7 +633,7 @@ public void close() { this.isClosed.set(true); this.failbackFailureLogCount.set(0); this.failbackBacklogScanCount.set(0); - this.removeFailbackRemainingMeter(); + this.removeFailbackPendingRecoveryMeter(); Disposable disposable = this.partitionRecoveryDisposable.getAndSet(null); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); @@ -698,12 +700,13 @@ private boolean handleException( return isExceptionThresholdBreached.get(); } - private void handleSuccess( + private boolean handleSuccess( PartitionKeyRangeWrapper partitionKeyRangeWrapper, RegionalRoutingContext succeededLocation, boolean isReadOnlyRequest, boolean forceStatusChange) { + AtomicBoolean transitionedFromUnavailable = new AtomicBoolean(false); this.locationEndpointToLocationSpecificContextForPartition.compute(succeededLocation, (locationAsKey, locationSpecificContextAsVal) -> { LocationSpecificHealthContext locationSpecificHealthContextAfterTransition; @@ -721,12 +724,16 @@ private void handleSuccess( .build(); } + LocationHealthStatus previousHealthStatus = locationSpecificContextAsVal.getLocationHealthStatus(); locationSpecificHealthContextAfterTransition = this.locationSpecificHealthContextTransitionHandler.handleSuccess( locationSpecificContextAsVal, partitionKeyRangeWrapper, GlobalPartitionEndpointManagerForPerPartitionCircuitBreaker.this.regionalRoutingContextToRegion.getOrDefault(succeededLocation, StringUtils.EMPTY), forceStatusChange, isReadOnlyRequest); + transitionedFromUnavailable.set( + previousHealthStatus == LocationHealthStatus.Unavailable + && locationSpecificHealthContextAfterTransition.getLocationHealthStatus() != LocationHealthStatus.Unavailable); // used only for building diagnostics - so creating a lookup for URI and region name @@ -743,6 +750,8 @@ private void handleSuccess( return locationSpecificHealthContextAfterTransition; }); + + return transitionedFromUnavailable.get(); } public boolean areLocationsAvailableForPartitionKeyRange(List availableLocationsAtAccountLevel) { diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java index 589b8a4c9943..8f284354b0c5 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/models/CosmosMetricName.java @@ -77,10 +77,10 @@ private CosmosMetricName(String name, CosmosMetricCategory metricCategory) { CosmosMetricCategory.OPERATION_DETAILS); /** - * Number of PPCB partition key ranges remaining to fail back, by collection (Gauge) + * Number of PPCB partition range-region recovery actions pending failback, by collection (Gauge) */ - public static final CosmosMetricName PPCB_FAILBACK_PENDING_PARTITION_COUNT = new CosmosMetricName( - nameOf("ppcb.failback.pendingPartitionCount"), + public static final CosmosMetricName PPCB_FAILBACK_PENDING_RECOVERY_COUNT = new CosmosMetricName( + nameOf("ppcb.failback.pendingRecoveryCount"), CosmosMetricCategory.OPERATION_SUMMARY); /** @@ -473,8 +473,8 @@ private static Map createMeterNameMap() { map.put(nameOf("op.actualitemcount"), CosmosMetricName.OPERATION_DETAILS_ACTUAL_ITEM_COUNT); map.put(nameOf("op.regionscontacted"), CosmosMetricName.OPERATION_DETAILS_REGIONS_CONTACTED); map.put( - nameOf("ppcb.failback.pendingpartitioncount"), - CosmosMetricName.PPCB_FAILBACK_PENDING_PARTITION_COUNT); + nameOf("ppcb.failback.pendingrecoverycount"), + CosmosMetricName.PPCB_FAILBACK_PENDING_RECOVERY_COUNT); map.put(nameOf("req.rntbd.requests"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_REQUESTS); map.put(nameOf("req.rntbd.latency"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_LATENCY); map.put(nameOf("req.rntbd.backendlatency"), CosmosMetricName.REQUEST_SUMMARY_DIRECT_BACKEND_LATENCY);