feat(gax): support transparent retries during mTLS certificate rotations - #13995
macastelaz wants to merge 29 commits into
Conversation
- Add CertificateBasedAccess and WorkloadCertificateUtils for SPIFFE and custom certificate loading - Implement RefreshingHttpJsonChannel and ChannelPool mTLS certificate fingerprint tracking and rotation - Enable transparent retries for retryable UnauthenticatedExceptions in ApiResultRetryAlgorithm and AttemptCallable - Add override delegation for getEndpoint, getHttpTransport, and getExecutor to preserve SLF4J MDC logging in Showcase tests
There was a problem hiding this comment.
Code Review
This pull request introduces support for dynamic mTLS certificate rotation across both gRPC and HTTP/JSON transports by enabling thread-safe channel hot-swapping and automatic refreshing upon encountering an UnauthenticatedException. Key additions include the RefreshingHttpJsonChannel and updates to various callables to intercept and retry unauthenticated errors. However, several critical issues were identified during review: a bug in ChannelPool.refresh() that breaks the GFE channel refresh mechanism for non-mTLS connections; a potential resource leak in RefreshingHttpJsonChannel due to a missing cancel override; regressions caused by the removal of Conscrypt security provider configurations; and incomplete exception wrapping in several streaming callables that results in the loss of the original stack trace, cause, and suppressed exceptions of UnauthenticatedException.
5678ad4 to
e3c70b5
Compare
Addresses AI code review findings from https://paste.googleplex.com/6563525517508608: - GrpcCallContext: Prevent transportChannel stale inheritance in merge() and withChannel() - RefreshingHttpJsonChannel: Set shutdownRequested and shutdownInitiated in shutdownNow() so newCall() throws IllegalStateException - AttemptCallable / StreamingCallables: Pass getCause() when rethrowing retryable UnauthenticatedException to prevent double-wrapping - CertificateBasedAccess: Enforce fail-closed security boundary when certificate config is malformed or missing required keys, and fix JSON unescaping order - ChannelPool: Update ReleasingClientCall Javadoc contract - Unit tests: Add cache invalidation test helpers to eliminate Thread.sleep() delays and add comprehensive tests for all addressed edge cases
e3c70b5 to
a2210c6
Compare
Addresses Gemini code review feedback on ReleasingHttpJsonClientCall and ReleasingClientCall: - Tracks wasStarted atomic flag on client calls to detect if start() has been invoked - If cancel() is invoked before start() (or call is discarded unstarted), cancel() immediately releases the ChannelEntry to decrement the active call reference count - Prevents memory/resource leaks of retired channels that are waiting for outstanding calls to drop to 0 - Adds testCancelBeforeStartReleasesChannelEntry unit tests to both RefreshingHttpJsonChannelTest and ChannelPoolTest
…sensitivity Addresses findings from mTLS security deep-dive code review: - Handle non-workload JSON configs (e.g. PKCS#11 /etc/gcloud/certificate_config.json) gracefully in validateAndResolveConfig without throwing IllegalStateException, preventing initialization failures on Google developer environments - Enforce fail-closed security boundary in getWorkloadCertPath() by validating disk file existence when GOOGLE_API_CERTIFICATE_CONFIG is set and throwing IllegalStateException when mTLS is enabled but no valid cert can be resolved - Make GOOGLE_API_USE_MTLS_ENDPOINT policy comparisons case-insensitive in getMtlsEndpointUsagePolicy()
…nd fail-closed getWorkloadCertPath - Adds testUseMtlsEndpointCaseInsensitive to verify getMtlsEndpointUsagePolicy() handles uppercase 'ALWAYS' and 'NEVER' - Adds assertThrows(IllegalStateException.class, cba::getWorkloadCertPath) in testUseMtlsClientCertificateExplicitTrueNoCredentials to verify getWorkloadCertPath() throws IllegalStateException when mTLS is required but no certificate can be resolved
nbayati
left a comment
There was a problem hiding this comment.
Some feedback on the auth side of things.
…PR 13995 review feedback Address review comments from @nbayati: 1. Make auth library (MtlsUtils) single source of truth for mTLS cert discovery and permission rules. 2. Fix GOOGLE_API_USE_CLIENT_CERTIFICATE flag semantics: true permits mTLS, return null/false cleanly if no certs are found (Row 3). Throw IllegalStateException only when cert config exists but referenced cert/key files are missing (Row 2). 3. Separate GKE and GCE workload certificate resolution paths. 4. Centralize SHA-256 certificate fingerprint calculation in MtlsUtils.
826f766 to
1423299
Compare
- Separate GKE (credentialbundle.pem) and GCE (certificates.pem + private_key.pem) workload certificate fallback paths in MtlsUtils. - Restore full Javadoc on MtlsUtils.getWorkloadCertificateConfiguration. - Format MtlsUtils and MtlsUtilsTest with google-java-format. - Fix Java 8 Mockito reflection error in GrpcLoggingInterceptorTest by instantiating GrpcLoggingInterceptor directly. - Isolate DirectPath environment tests in InstantiatingGrpcChannelProviderTest from host environment variables.
1423299 to
be0a495
Compare
| throw new CertificateSourceUnavailableException( | ||
| "Certificate configuration loaded successfully, but does not contain a 'certificate_file' path."); | ||
| "Certificate configuration loaded successfully, but does not contain a 'certificate_file'" | ||
| + " path."); |
There was a problem hiding this comment.
Could this break the ECP flow? Do we need to check that "workload" exists but "certificate_file" does not exist?
There was a problem hiding this comment.
This method is currently only called from google-auth-library-java/oauth2_http/java/com/google/auth/oauth2/IdentityPoolCredentials.java and from what I understand, ECP is not applicable for IdentityPool credentials so throwing here would be acceptable - but let me know if I'm missing something or if you'd like to see other handling here.
…th go/sdk-mtls-by-default-cert-discovery Address PR 13995 review feedback from @nbayati: - Align discovery and error behavior with go/sdk-mtls-by-default-cert-discovery: - Fail closed (IllegalStateException) when GOOGLE_API_CERTIFICATE_CONFIG points to a missing, unreadable, malformed, or missing cert/key configuration. - Safe fallback (return null) when implicit default gcloud config is missing or is an ECP-only configuration without a workload block. - Fail closed with clear source identification if default gcloud config is unreadable, malformed, or points to missing cert/key files. - Replace .exists() with .isFile() && .canRead() checks across config, certificate, and key paths. - Make getGkeWorkloadCertPath and getGceWorkloadCertPath package-private stubs returning null with explanatory comments for phased rollout. - Explicitly identify the resolution source (GOOGLE_API_CERTIFICATE_CONFIG vs default gcloud location) in all error messages. - Update getCertificatePath exception message to reference 'cert_configs.workload.cert_path' rather than legacy 'certificate_file'. - Add comprehensive test coverage in MtlsUtilsTest and CertificateBasedAccessTest.
…P flow in getCertificatePath
…y and channel refresh - Rename MtlsUtils.validateCertAndKeyFiles to checkCertAndKeyFilesReadable. - Move file readability check outside try-catch in MtlsUtils to clearly separate parsing errors from file existence errors. - Remove GKE/GCE placeholder stubs and internal doc references from MtlsUtils. - Simplify MtlsUtils.getCertificateFingerprint using Files.readAllBytes and Guava BaseEncoding. - Defer activeCertFingerprint mutation in ChannelPool until after channel creation succeeds in refreshAll(). - Add unit test in ChannelPoolTest verifying failed refresh attempts do not mutate fingerprint or prevent subsequent retries.
nbayati
left a comment
There was a problem hiding this comment.
A couple of issues with the UnauthenticatedException handling across ServerStreamingAttemptCallable, BidiStreamingCallable, and ClientStreamingCallable:
transportChannel.refresh()is invoked without atry-catch. Ifrefresh()throws an unchecked exception, the terminalonErrorcallback never fires, which will leave stream observers or retrying futures hanging indefinitely. Any refresh failure should be caught and logged so the original error still reaches the observer.- We shouldn't re-wrap the exception with
isRetryable = true:- For client and bidi streaming, GAX has no stream retry mechanism, so marking it retryable is inert internally and misleading to callers.
- For server streaming, setting
isRetryable = truecausesStreamingRetryAlgorithmto attempt to resume the stream. In the "Graceful Certificate Rotation Handling" section of go/sdk-mds-bound-token , we said streaming calls should not be auto-retried mid-stream. Instead, we should trigger the channel refresh so subsequent calls use the new connection, but propagate the original error directly without marking it retryable. Let me know if you don't agree with this though, maybe it's a shortcoming of the original HLD that we need to revisit and update.
…otation retries - Remove unused FileExistenceProvider/FileContentReader and 3-arg constructor from CertificateBasedAccess. - In ServerStreamingAttemptCallable, BidiStreamingCallable, and ClientStreamingCallable, wrap transportChannel.refresh() in try-catch with warning logging and propagate original exception without marking isRetryable=true. - Add getGeneration() to TransportChannel, ChannelPool, and RefreshingHttpJsonChannel; update AttemptCallable to track attemptGeneration so sibling in-flight requests that failed on the stale connection are retried without redundant channel recreation. - Guard ChannelPool.refresh() and refreshAll() against invocation on shut-down pool and synchronize isShutdown state across shutdown methods. - Add delegating protected constructor in ManagedHttpJsonChannel so RefreshingHttpJsonChannel and ManagedHttpJsonInterceptorChannel do not leak unused parent scheduled executors and default HTTP transports. - Only wrap HTTP/JSON channels with RefreshingHttpJsonChannel when workloadCertPath is not null. - Configure Conscrypt security provider prior to calling NetHttpTransport.Builder.trustCertificates in InstantiatingHttpJsonChannelProvider. - Clear stale transportChannel reference in HttpJsonCallContext.withChannel() and merge() when channel changes. - Add comprehensive unit tests across gax, gax-grpc, and gax-httpjson modules.
…ependency analyzer
|
|
||
| // Prevent expanding timeouts | ||
| if (timeout != null && this.timeout != null && this.timeout.compareTo(timeout) <= 0) { | ||
| if (this.timeout != null && (timeout == null || this.timeout.compareTo(timeout) <= 0)) { |
There was a problem hiding this comment.
In withTimeoutDuration, changing timeout != null && this.timeout != null to this.timeout != null && (timeout == null || ...) prevents callers from clearing an existing timeout. When a context already has a timeout set, passing null or Duration.ZERO (which is normalized to null on line 283) now hits timeout == null and returns this with the old timeout still attached.
Let's restore the base-branch check so null and Duration.ZERO can still clear an existing timeout:
// Prevent expanding timeouts
if (timeout != null && this.timeout != null && this.timeout.compareTo(timeout) <= 0) {
return this;
}There was a problem hiding this comment.
Good catch, thanks. I restored the base-branch check, so withTimeoutDuration(null) and withTimeoutDuration(Duration.ZERO) clear an existing timeout again. I also added GrpcCallContextTest#testWithNullOrZeroTimeoutClearsExistingTimeout to cover both cases.
|
|
||
| // Prevent expanding deadlines | ||
| if (timeout != null && this.timeout != null && this.timeout.compareTo(timeout) <= 0) { | ||
| if (this.timeout != null && (timeout == null || this.timeout.compareTo(timeout) <= 0)) { |
There was a problem hiding this comment.
Same issue as GrpcCallContext.java:L287: let's restore if (timeout != null && this.timeout != null && this.timeout.compareTo(timeout) <= 0) here as well so withTimeoutDuration(null) and withTimeoutDuration(Duration.ZERO) can clear an existing timeout.
There was a problem hiding this comment.
Done. I restored the same check here and added HttpJsonCallContextTest#testWithNullOrZeroTimeoutClearsExistingTimeout to cover both the null and Duration.ZERO cases.
|
|
||
| public RefreshingHttpJsonChannel( | ||
| Supplier<ManagedHttpJsonChannel> channelFactory, String workloadCertPath) { | ||
| this(channelFactory.get(), channelFactory, workloadCertPath); |
There was a problem hiding this comment.
Same startup ordering issue as InstantiatingHttpJsonChannelProvider.java:L252: this(channelFactory.get(), ...) creates the initial channel from disk before the delegated constructor initializes CertificateRotationTracker. We should initialize rotationTracker before calling channelFactory.get() so the baseline fingerprint is recorded before the transport reads the certificate.
There was a problem hiding this comment.
Done. The constructor now creates rotationTracker before calling the transport factory for the initial transport, so the baseline fingerprint is recorded before the certificate is read. RefreshingHttpJsonChannelTest#rotationDuringInitialTransportCreation_isDetectedAndRefreshed covers this (see the provider thread for details).
| } | ||
| boolean shouldRetry = transportChannel.getGeneration() > attemptGeneration; | ||
| if (shouldRetry) { | ||
| UnauthenticatedException newEx = |
There was a problem hiding this comment.
nit: Creating newEx replaces the original stack trace with this frame in onErrorImpl. Let's copy the stack trace via newEx.setStackTrace(unauthenticatedException.getStackTrace()) before reassigning t (matching AttemptCallable.java:L124).
There was a problem hiding this comment.
Done. newEx now copies the original stack trace, which matches AttemptCallable. I updated ServerStreamingAttemptCallableTest#testUnauthenticatedRefreshWithGenerationAdvanceRetries to assert that the stack trace is kept.
| @Nullable Throwable previousThrowable, | ||
| @Nullable ResponseT previousResponse, | ||
| TimedAttemptSettings previousSettings) { | ||
| if (previousThrowable instanceof UnauthenticatedException |
There was a problem hiding this comment.
shouldRetry marks any retryable UnauthenticatedException as eligible for retry, while createNextAttempt only builds custom zero-delay settings for the first certificate rotation attempt and returns null if a subsequent attempt also fails with a retryable auth error. Because returning null causes the retry framework to fall back to standard exponential backoff instead of stopping, consecutive rotation-marked auth failures continue retrying:
- On RPCs configured with a total timeout and no maximum attempt cap, the call retries with exponential backoff until the total timeout expires instead of failing after the single rotation retry.
- On RPCs configured with multiple attempts where
UNAUTHENTICATEDis not a retryable status code, subsequent auth failures consume the method's normal retry budget up to the attempt limit.
Can we have createNextAttempt return an explicitly exhausted attempt configuration to stop immediately once the single rotation retry has been used?
There was a problem hiding this comment.
Good catch, thanks. Once the single rotation retry has been used, createNextAttempt now returns exhausted settings (maxAttempts == attemptCount) instead of null. The retry framework therefore stops right away instead of falling back to exponential backoff. This fixes both cases you described.
While doing this I found a related issue in server streaming. StreamingRetryAlgorithm resets attemptCount after a stream makes progress but kept overallAttemptCount. So after any earlier retry on a stream, the rotation-retry condition (overallAttemptCount == attemptCount) could never match again, and with the exhausted settings a later rotation would stop the stream.
On progress, StreamingRetryAlgorithm now computes the next attempt from a fresh baseline and then adds the previous overall count back on. As a result:
- Each progress window gets its own single rotation retry.
- Attempt numbering for tracing still only goes up.
- A repeated rotation failure without progress still stops.
Tests in ApiResultRetryAlgorithmTest:
testSecondRotationFailureStopsWithTotalTimeoutAndNoMaxAttempts: your case 1.testSecondRotationFailureDoesNotConsumeNormalRetryBudget: your case 2.testStreamRotationRetryIsAvailableAgainAfterProgress: the streaming re-arm, and that it stays bounded.testStreamProgressResetKeepsOverallAttemptCountIncreasing: attempt numbering across resets.testStreamProgressResetReturnsNullWhenNotRetryable: the non-retryable path after progress.
| allEntries.removeIf(entry -> entry != newEntry && entry.channel.isTerminated()); | ||
|
|
||
| ChannelEntry oldEntry = activeEntry.getAndSet(newEntry); | ||
| rotationTracker.markRefreshed(currentDiskFingerprint); |
There was a problem hiding this comment.
rotationTracker.markRefreshed(currentDiskFingerprint) runs before generation.incrementAndGet(), leaving a timing window during certificate rotation where a concurrent failing RPC can miss its retry attempt.
When an RPC fails with an authentication error, AttemptCallable and ServerStreamingAttemptCallable check shouldRefresh() and getGeneration() outside refreshLock. If a failing call runs its retry check between these two updates, it sees that the rotation tracker already has the new certificate and skips refreshing. Because it never enters refresh(), it does not wait on the lock and immediately reads the old generation count, causing it to surface the authentication failure instead of retrying on the new channel.
I think incrementing generation right after swapping activeEntry and before marking the tracker refreshed can close this window.
There was a problem hiding this comment.
Agreed, thanks for spotting the window. refresh() now does three steps in order:
- Swap in the new transport.
- Increment
generation. - Call
rotationTracker.markRefreshed(...).
Any failing RPC that sees the new fingerprint therefore also sees the new generation and retries. After the refactor for the allEntries thread there's no activeEntry anymore, but the same order applies to the transport swap, and a comment next to the three statements explains why the order matters.
RefreshingHttpJsonChannelTest#refresh_swapsTransportAndKeepsChannel and #rotationDuringInitialTransportCreation_isDetectedAndRefreshed cover the state after a refresh.
| private final String workloadCertPath; | ||
| private final AtomicReference<ChannelEntry> activeEntry; | ||
| // Keep track of all entries to properly await their termination | ||
| private final ConcurrentLinkedQueue<ChannelEntry> allEntries = new ConcurrentLinkedQueue<>(); |
There was a problem hiding this comment.
IIUC, historically HttpJson's channel is essentially 1:1 with a HttpUrlConnection.
I'm a bit lost, but what does allEntries here hold? I think it would essentially be something like oldChannel (or is it possible to have multiple old channels)?
I'm wondering below as it seems like there is quite a bit of logic pertaining to finding and updating the channel reference. I wonder if this can be simplified towards having just something like an atmoicreference to the old channel and we initialize shutdown to the old channel on refresh. Then we update the activeEntry.
Basically tracked via oldEntry and activeEntry.
There was a problem hiding this comment.
I think we should probably double check this, because if it could be simplified this way, I think all we need to do is to re-create the the underlying HttpTransport and initialize an orderly shutdown for the old HttpTransport (in flight requests should be continue until they can finish).
In my head, I think it's easier to manage and we don't need to have a ReleasingHttpJsonClientCalls
There was a problem hiding this comment.
Thanks, this was a great simplification. On the original question: allEntries held every channel that was still active or retired but not yet terminated. It could hold more than one retired channel, for example when a long-running call on the first channel spans two rotations, so a single oldEntry reference could have lost track of a channel.
I went with your follow-up: only the HttpTransport is recreated. RefreshingHttpJsonChannel is now a thin wrapper around a single ManagedHttpJsonChannel. On rotation, refresh() creates a new mTLS transport and swaps it into the delegate. The delegate keeps its executors and endpoint. allEntries, ChannelEntry ref-counting and the Releasing* call and listener wrappers are all gone.
On shutting down the old transport: each call holds the transport it was created with, so calls already in flight finish on the old transport and new calls use the new one. We don't explicitly shut down the old transport. In google-http-client 2.2.0, NetHttpTransport doesn't override HttpTransport.shutdown(), which is a no-op, so there's nothing to shut down in order. Its idle keep-alive connections expire on their own. shutdown()/shutdownNow() take the same lock as refresh(), so a transport swap can't race with channel shutdown.
One small behavior change: newCall() after shutdownNow() no longer throws. It now behaves the same as a plain ManagedHttpJsonChannel.
Tests in RefreshingHttpJsonChannelTest:
refresh_swapsTransportAndKeepsChannelcallCreatedBeforeRefresh_usesOriginalTransport: an in-flight call stays on the old transport and the next call uses the new one, usingMockHttpService.testConcurrentNewCallDuringRefreshshutdown_waitsForInProgressRefreshtestRefreshDoesNotCreateTransportWhenShutdowntestRefreshFactoryExceptionDoesNotWedgeFingerprinttestChannelDelegationMethodsclose_shutsDownUnderlyingChannel
| private final String workloadCertPath; | ||
| private final AtomicReference<ChannelEntry> activeEntry; | ||
| // Keep track of all entries to properly await their termination | ||
| private final ConcurrentLinkedQueue<ChannelEntry> allEntries = new ConcurrentLinkedQueue<>(); |
There was a problem hiding this comment.
I think we should probably double check this, because if it could be simplified this way, I think all we need to do is to re-create the the underlying HttpTransport and initialize an orderly shutdown for the old HttpTransport (in flight requests should be continue until they can finish).
In my head, I think it's easier to manage and we don't need to have a ReleasingHttpJsonClientCalls
| this(null, true, null, null, true); | ||
| } | ||
|
|
||
| protected ManagedHttpJsonChannel(boolean isDelegatingWrapper) { |
There was a problem hiding this comment.
is isDelegatingWrapper being used?
There was a problem hiding this comment.
The parameter's value isn't read. It only exists to tell this constructor apart from the protected no-arg constructor, which creates a default transport and executor. ManagedHttpJsonInterceptorChannel and RefreshingHttpJsonChannel call super(true), so they don't create a transport and executor that they would never use or shut down, since they pass every call to the channel they wrap.
I made the constructor package-private, so it isn't new public or protected API, and added Javadoc that explains this. If you'd prefer a clearer mechanism, such as a private marker type instead of the boolean, I'm happy to switch.
| SslUtils.initSslContext( | ||
| sslContext, | ||
| null, | ||
| SslUtils.getPkixTrustManagerFactory(), | ||
| mtlsKeyStore, | ||
| "", | ||
| SslUtils.getDefaultKeyManagerFactory()); |
There was a problem hiding this comment.
I believe builder.trustCertificates(null, mtlsKeyStore, ""); already invokes the logic to intialize the SslUtils.
I think the logic should be updated to:
- invoke
builder.trustCertificates(null, mtlsKeyStore, ""); - always intialize
HttpJsonConscryptUtils.configureConscryptSecurityProvider(builder);(this should try to load Conscrypt and configure the ssl socket factory)
This does get a bit tricky since I believe the logic only upgrades to mTLS if we are using the default, so if a user can't comply with PQC but needs to upgrade, we may not. I think we'll need to double check this, but this could be an edge case where we add configurations in the future.
There was a problem hiding this comment.
Thanks for looking closely. I checked this against google-http-client 2.2.0 and I think the extra SSLContext step is needed. Without it, PQC would be silently dropped on the mTLS path:
trustCertificates(null, mtlsKeyStore, "")builds anSSLContextusing the default JDK provider and sets it as the builder'sSSLSocketFactory.configureConscryptSecurityProvider(builder)only sets the security provider and a socket configurator. When the transport is built, an explicitly setSSLSocketFactorytakes precedence over the security provider. So with just steps 1 and 2, sockets come from the JDK context,Conscrypt.isConscrypt(socket)is false, and the PQC named groups are never applied. mTLS still works, but PQC doesn't.
That's why, when Conscrypt is available, the code builds an SSLContext with the Conscrypt provider and the mTLS key material and sets it as the socket factory. configureConscryptSecurityProvider then adds the named-groups configurator. The trustCertificates call stays as the fallback for when Conscrypt can't be loaded. This order gives both the client certificate and PQC. The logic matches what's already on agentic-identities-bound-token; this PR only moved it around.
On the edge case: agreed. Today there's no way to use mTLS while opting out of Conscrypt/PQC. I'd suggest tracking an opt-out setting as a follow-up rather than adding it to this PR. Happy to file an issue for it.
… stack trace - Restore the base-branch guard in GrpcCallContext/HttpJsonCallContext withTimeoutDuration so a null or zero timeout clears an existing one. - Copy the original stack trace onto the wrapped UnauthenticatedException in ServerStreamingAttemptCallable, matching AttemptCallable.
…on and fix refresh ordering - RefreshingHttpJsonChannel now initializes its CertificateRotationTracker before creating the initial channel, so a rotation during startup is detected. The provider lets RefreshingHttpJsonChannel create the initial channel, unwrapping checked IOException/GeneralSecurityException to keep the getTransportChannel() contract. - In refresh(), increment the generation before marking the tracker refreshed so a failing RPC that sees the new fingerprint also sees the new generation. - Catch and log channel factory failures in refresh(), keeping the old channel.
Return explicitly exhausted settings from ApiResultRetryAlgorithm createNextAttempt once the free certificate-rotation retry has been used, so a subsequent rotation-marked UNAUTHENTICATED failure stops instead of falling back to exponential backoff.
On a rotation-triggered refresh, ChannelPool now keeps only successfully recreated channels so no traffic is routed to the old certificate, and marks the certificate refreshed once at least one channel was recreated. Statically sized pools schedule a one-time refill; dynamically sized pools are regrown by resize(). Non-rotation refreshes keep per-slot fallback.
- Use Guava @VisibleForTesting in RefreshingHttpJsonChannel. - Make the delegating-wrapper ManagedHttpJsonChannel constructor package-private and document it; note refresh() is an intentional no-op. - Add Javadocs for getGeneration/refresh/shouldRefresh with {@inheritdoc} on overrides. - Drop the redundant useMtlsClientCertificate() check in InstantiatingGrpcChannelProvider and document the rotation-tracking conditions.
…stead of replacing the channel RefreshingHttpJsonChannel now wraps a single ManagedHttpJsonChannel and replaces its HttpTransport when the workload certificate rotates. Calls capture their transport when created, so in-flight requests complete on the previous transport, and NetHttpTransport holds no pooled resources that need releasing. This removes the per-generation channel tracking and reference counting (ChannelEntry, allEntries, ReleasingHttpJsonClientCall) and keeps a single set of executors for the lifetime of the channel.
| @Nullable HttpTransport createHttpTransport() throws IOException, GeneralSecurityException { | ||
| if (mtlsProvider == null) { | ||
| return null; |
There was a problem hiding this comment.
qq, this looks like a behavior change. Can we double check this?
This looks like it's from the previous impl of createHttpTransport. Maybe lost from an earlier rebase?
There was a problem hiding this comment.
Good catch, thanks. You're right: the PR started before #13853 and the rebase kept the old shape of createHttpTransport(). It also dropped two of the #13853 tests. I restored createHttpTransport() and configureMtls() to the base implementation and restored testCreateHttpTransport_returnsValidTransport and testConfigureConscryptSecurityProvider_returnsConfiguredBuilder.
One intentional difference remains, in a small separate helper used only when the channel is created. If mTLS is enabled but the provider returns no client certificate, we throw IOException("Failed to initialize mTLS HttpTransport") instead of silently using a transport without it. This matches the gRPC provider's createChannelBuilder(). During a rotation, it means the refresh keeps the existing authenticated transport. For callers that don't use mTLS, behavior is the same as the base.
Tests:
InstantiatingHttpJsonChannelProviderTest#getTransportChannel_withMtlsKeyStore_usesMtlsTransport: the normal mTLS path still produces an mTLS transport.InstantiatingHttpJsonChannelProviderTest#refresh_whenKeyStoreUnavailableDuringRotation_keepsCurrentTransport: if the key store is unavailable during a rotation, the current transport is kept, and the next refresh swaps in a new mTLS transport once it's available again.
| if (currentDiskFingerprint.isEmpty()) { | ||
| return; | ||
| } |
There was a problem hiding this comment.
from the javadocs, this looks like it can occur if no workload certificate path is configured or if the file is currently unreadable/empty
I think we have guards in the InstantiatingHttpJsonChannelProvider to ensure that certPath is non-null here as non-null will create the RefreshingHttpJsonChannel, but what could be the cause of a file being unreadable/empty? I assume empty may be the possibility during rotation, but unreadable result in this triggered over and over.
I'm worried that there may be cases where this just retries over and over as the generation doesn't get incremented here.
There was a problem hiding this comment.
It can't loop. The extra attempt is only allowed when the channel's generation moves past the generation the attempt started on (AttemptCallable / ServerStreamingAttemptCallable). If refresh() returns here without swapping the transport, the generation is unchanged and the UNAUTHENTICATED error goes back to the caller unchanged (AttemptCallableTest#testRefreshReturnsWithoutAdvancingGeneration_notMarkedRetryable). Separately, ApiResultRetryAlgorithm allows at most one rotation retry per call.
If the file stays unreadable, we don't even get here: shouldRefresh() returns false for an empty or unreadable file, so refresh() isn't called. This line covers the short window where shouldRefresh() saw a new, readable certificate but the re-read under the lock finds the file empty, typically a rotator that truncates and then rewrites. In that case the call fails once with the original error, and the next auth failure after the write completes refreshes normally. Other causes are the file being deleted, replaced non-atomically, or having its permissions changed.
I also added RefreshingHttpJsonChannelTest#refresh_whenCertificateFileEmpty_keepsTransportAndGeneration to pin this down: the transport and generation are unchanged, shouldRefresh() stays false while the file is empty, and the next refresh after the certificate is written succeeds.
| <test>!InstantiatingGrpcChannelProviderTest#testLogDirectPathMisconfig_AttemptDirectPathNotSetAndAttemptDirectPathXdsSetViaEnv_warns,!InstantiatingGrpcChannelProviderTest#canUseDirectPath_directPathEnvVarNotSet_attemptDirectPathIsTrue,InstantiatingGrpcChannelProviderTest#testLogDirectPathMisconfigWrongCredential</test> | ||
| <!-- <test>!InstantiatingGrpcChannelProviderTest#testLogDirectPathMisconfig_AttemptDirectPathNotSetAndAttemptDirectPathXdsSetViaEnv_warns</test> --> |
There was a problem hiding this comment.
This was to get the gax-grpc unit tests running in CI. The current config ends with InstantiatingGrpcChannelProviderTest#testLogDirectPathMisconfigWrongCredential without a !. Surefire treats that as an inclusion, so only that one test runs for gax-grpc. You can see it in the sdk-platform-java units logs: on a main-based branch gax-grpc runs Tests run: 1, while on this PR it runs the full suite (238), including the new ChannelPoolTest coverage.
I've narrowed the change so it now only drops that stray inclusion and keeps both existing exclusions (including canUseDirectPath_directPathEnvVarNotSet_attemptDirectPathIsTrue). main has the same config, so I'll open a separate issue/PR for it there.
…and narrow surefire config - Restore createHttpTransport()/configureMtls() from the base (lost in a rebase) and the two tests that were dropped; keep failing closed when mTLS is enabled but no client certificate is available, in a separate helper used only for channel creation. - Add provider tests for the mTLS channel path and for keeping the current transport when the key store is unavailable during a rotation. - Add a test for refresh() when the certificate file is empty mid-rotation. - gax-grpc surefire: keep both existing exclusions and only drop the stray inclusion pattern that limited the module to a single test.
| e); | ||
| } | ||
| } | ||
| boolean shouldRetry = channel.getGeneration() > attemptGeneration; |
There was a problem hiding this comment.
qq, do you think we can check this first? If another thread has already refreshed the channels (e.g. acquired the lock first), then we don't need to check the channel via shouldRefresh.
channel.getGeneration() > attemptGeneration should indicate that we need to retry and we can skip the check above. IIUC, the check above is needed for the first call to acquire the lock
There was a problem hiding this comment.
Good call, done. In both AttemptCallable and ServerStreamingAttemptCallable the handler now checks getGeneration() > attemptGeneration first, and only calls shouldRefresh()/refresh() if the generation hasn't moved. As you said, the disk check is still needed for the first failing call to detect the rotation. Retry decisions don't change; this just skips the cert read and hash for calls that fail after another thread has already refreshed. I also moved shouldRefresh() inside the try so a failed fingerprint read can't hide the original UNAUTHENTICATED. Added testGenerationAlreadyAdvanced_skipsShouldRefresh and the streaming equivalent.
| if (workloadCertPath != null && currentDiskFingerprint.isEmpty()) { | ||
| return; | ||
| } | ||
| if (refreshAll() && !currentDiskFingerprint.isEmpty()) { |
There was a problem hiding this comment.
Hmm, I didn't realize refreshAll() increments the generation. refreshSafely gets called every 50 minutes as part of channelpool so I think we will need to distinguish between a rotation (e.g. 401 channel refresh) vs a refresh from a background task
There was a problem hiding this comment.
Good catch. Generation was bumped on every refreshAll(), including the 50-minute preemptive refresh, so a non-rotation UNAUTHENTICATED that overlapped a background refresh got a free retry. The increment is now separate from refreshAll() and only happens when the pool switches to a new certificate: in refresh() after the rotation swap, and in refreshSafely() only when the disk fingerprint differs from the active one and every channel was recreated. That also keeps "a generation bump means every channel has the new certificate" true for the preemptive path. It still happens after the swap and before the new fingerprint is marked active. HttpJson isn't affected since RefreshingHttpJsonChannel has no periodic refresh. Added preemptiveRefresh_withoutRotation_doesNotIncrementGeneration and related tests.
While in here I also fixed a related shutdown issue: shutdown()/shutdownNow() cancelled the refresh future while holding entryWriteLock, which the running refresh itself holds, so shutdown waited for an in-progress refresh instead of interrupting it. The futures are now cancelled before taking the lock; setting isShutdown and shutting down the entries still happen under it. Added shutdown_interruptsInProgressRefresh.
| if (shouldRetry) { | ||
| UnauthenticatedException newEx = | ||
| new UnauthenticatedException( | ||
| unauthenticatedException.getMessage(), | ||
| unauthenticatedException.getCause(), | ||
| unauthenticatedException.getStatusCode(), | ||
| true, // isRetryable = true | ||
| unauthenticatedException.getErrorDetails()); | ||
| newEx.setStackTrace(unauthenticatedException.getStackTrace()); | ||
| for (Throwable suppressed : unauthenticatedException.getSuppressed()) { | ||
| newEx.addSuppressed(suppressed); | ||
| } | ||
| throw newEx; |
There was a problem hiding this comment.
(Flagged from Gemini) I think I may have been wrong on the original behavior of hijacking the isRetryable parameter. Gemini flags that overloads the value when user configures custom status code params.
Suggestion from Gemini:
@NullMarked
public class UnauthenticatedException extends ApiException {
private final boolean channelRefreshed;
// Existing public constructors delegate with channelRefreshed = false ...
private UnauthenticatedException(
String message,
Throwable cause,
StatusCode statusCode,
boolean retryable,
ErrorDetails errorDetails,
boolean channelRefreshed) {
super(message, cause, statusCode, retryable, errorDetails);
this.channelRefreshed = channelRefreshed;
}
boolean isChannelRefreshed() {
return channelRefreshed;
}
UnauthenticatedException withChannelRefreshed() {
UnauthenticatedException newEx =
new UnauthenticatedException(
getMessage(), getCause(), getStatusCode(), true, getErrorDetails(), true);
newEx.setStackTrace(getStackTrace());
for (Throwable suppressed : getSuppressed()) {
newEx.addSuppressed(suppressed);
}
return newEx;
}
}
if shouldRetry, then we will call unauthenticatedException.withChannelRefreshed() to mark that we need to retry and pass this over to an update ApiResultRetryAlgo:
@Override
public @Nullable TimedAttemptSettings createNextAttempt(
@Nullable RetryingContext context,
@Nullable Throwable previousThrowable,
@Nullable ResponseT previousResponse,
TimedAttemptSettings previousSettings) {
if (previousThrowable instanceof UnauthenticatedException
&& ((UnauthenticatedException) previousThrowable).isChannelRefreshed()) {
if (previousSettings.getOverallAttemptCount() == previousSettings.getAttemptCount()) {
RetrySettings globalSettings = previousSettings.getGlobalSettings();
if (globalSettings.getMaxAttempts() == 0
&& globalSettings.getTotalTimeoutDuration().isZero()) {
globalSettings = globalSettings.toBuilder().setMaxAttempts(1).build();
}
return previousSettings.toBuilder()
.setGlobalSettings(globalSettings)
.setRetryDelayDuration(java.time.Duration.ZERO)
.setRandomizedRetryDelayDuration(java.time.Duration.ZERO)
.setAttemptCount(previousSettings.getAttemptCount())
.setOverallAttemptCount(previousSettings.getOverallAttemptCount() + 1)
.build();
}
// The single rotation retry has already been used. Return exhausted settings so the retry
// framework stops, rather than returning null and falling back to exponential backoff.
int exhaustedAttemptCount = previousSettings.getAttemptCount() + 1;
return previousSettings.toBuilder()
.setGlobalSettings(
previousSettings.getGlobalSettings().toBuilder()
.setMaxAttempts(exhaustedAttemptCount)
.build())
.setAttemptCount(exhaustedAttemptCount)
.setOverallAttemptCount(previousSettings.getOverallAttemptCount() + 1)
.build();
}
return null;
}
I'll need to take a look at this tomorrow.
There was a problem hiding this comment.
Thanks, agreed, this was a real problem and not only for mTLS. The exception factories set isRetryable() from the method's retry codes. So when UNAUTHENTICATED is configured as retryable, every 401 already arrived retryable, and the algorithm treated it as a rotation retry: one zero-delay retry and then stop, instead of the configured attempts and backoff. The early return in shouldRetry also overrode per-call retry codes. I went with your channelRefreshed flag, with a few adjustments:
withChannelRefreshed()keeps the originalisRetryable()instead of forcingtrue, so a surfaced error still reflects the configured retry codes.createNextAttemptand bothshouldRetryoverloads key onisChannelRefreshed(), so per-call retry codes are no longer overridden for ordinary 401s.- Once the free retry is used, a flagged failure still returns exhausted settings, as in your snippet.
Server streaming goes through the same algorithm via StreamingRetryAlgorithm, so ServerStreamingAttemptCallable just calls withChannelRefreshed(). Both new methods are package-private, so there's no public API change.
Two notes for transparency. (1) Because of the last bullet, a second rotation within a single call stops that call, even for clients that configured UNAUTHENTICATED as retryable. This keeps the "stop after one rotation retry" behavior from the earlier round. (2) UnauthenticatedException had no explicit serialVersionUID, so the new methods would have changed the computed one. I pinned it to the value from the current releases (2.83–2.87) and made the flag transient, so the serialized form is unchanged. Tests check the UID, deserializing an exception serialized by the released jar, and the existing retry behaviour for configured and unconfigured UNAUTHENTICATED, per-call retry codes, and a null context.
| LOG.fine( | ||
| "Refreshing all channels" | ||
| + (Strings.isNullOrEmpty(activeFingerprint) | ||
| ? "" | ||
| : " with certificate fingerprint: " + activeFingerprint)); |
There was a problem hiding this comment.
(Flagged by Gemini)
Gemini noticed that this is supposed to log the current active fingerprint, but rotationTracker.getActiveCertFingerprint(); returns the old fingerprint as the new fingerprint value doesn't get set until after refreshAll() is called.
Hmm, perhaps we can move this log to after refreshAll() completes (or maybe we hold the diskFingerprint in a static var so that it always has the latest state)
There was a problem hiding this comment.
Agreed, at that point the tracker still holds the old fingerprint. refreshAll now just logs "Refreshing all channels" (as before this PR), and the fingerprint is logged after a successful switch, next to markRefreshed(), using the fingerprint read from disk. No static or extra state. Added refresh_onRotation_logsNewCertificateFingerprint.
- Check the channel generation before reading the certificate from disk, and keep shouldRefresh() inside the try block. - Only advance the ChannelPool generation when every channel switched to a new certificate, so periodic refreshes don't enable rotation retries. - Mark rotation retries with a package-private channelRefreshed flag on UnauthenticatedException instead of reusing isRetryable(), so configured UNAUTHENTICATED retries and per-call retry codes keep their behaviour. Pin serialVersionUID to the released value; the flag is transient. - Log the new certificate fingerprint after a successful switch. - Cancel resize/refresh futures before taking entryWriteLock on shutdown so an in-progress refresh is interrupted.
Description
This PR adds transparent retries for mTLS workload certificate rotation to the gRPC and HTTP/JSON transports. When a request fails with
UNAUTHENTICATEDand the workload certificate on disk has changed, the transport is rebuilt with the new certificate and the request is retried once, without interrupting in-flight RPCs or streams.🚀 Core Features & Architectural Updates
• Rotation detection:
CertificateRotationTrackercompares the SHA-256 fingerprint of the workload certificate on disk with the certificate the transport was built with. The check only runs after an auth failure (never on the request path), and positive results are cached for at most 1 second so a burst of failures doesn't trigger a burst of file reads. Empty or mid-write files are ignored.• gRPC:
ChannelPoolreplaces all of its channels when a rotation is detected. Calls already in flight finish on their old channels, which are shut down once idle. During a rotation refresh, a channel that can't be recreated is dropped rather than kept, so no traffic keeps using the old certificate; a statically sized pool is refilled in the background.• HTTP/JSON:
RefreshingHttpJsonChannelswaps the channel'sHttpTransportfor one built with the new certificate. Calls already created keep using the transport they started with, so in-flight requests aren't interrupted.• Retries:
AttemptCallable(unary) andServerStreamingAttemptCallable(server streaming) refresh the transport after anUNAUTHENTICATEDfailure. If the transport moved to a new certificate during or after the attempt, the failure is flagged for a single immediate retry on the refreshed channel; its configuredisRetryable()value is unchanged.•
ApiResultRetryAlgorithmgives that retry once, immediately, without using a regular attempt.• If the free retry also fails after another rotation, the call stops instead of falling back to backoff retries.
• For server streams, the free retry becomes available again once the stream has made progress, and the stream resumes through the existing resumption strategy.
• Client and bidi streams refresh the transport on
UNAUTHENTICATEDbut don't retry; the next stream uses the new certificate.• Plumbing:
TransportChannelgainsshouldRefresh(),refresh()andgetGeneration(), which default to no-ops.ApiCallContext.getTransportChannel()gives retrying callables access to the channel, andGrpcCallContext/HttpJsonCallContextcarry it throughmerge()andwithChannel().🔒 System Hardening & Bug Fixes
• Outstanding RPC leak (
ChannelPool.ReleasingClientCall): if a call was cancelled beforestart(),start()threw without releasing the channel entry, leaving the channel with a permanently outstanding RPC count so it could never be cleaned up after a refresh. The entry is now released on that path and on cancellation before start.• HTTP/JSON refresh vs. shutdown:
shutdown()/shutdownNow()are serialized withrefresh()so a transport swap can't race with teardown.mTLS is enabled automatically when a workload certificate config is present (
MtlsUtils, auth library):GOOGLE_API_CERTIFICATE_CONFIG, or the default gcloudcertificate_config.json, contains aworkloadcertificate, client certificates are used even ifGOOGLE_API_USE_CLIENT_CERTIFICATEis unset.GOOGLE_API_USE_CLIENT_CERTIFICATE=falsestill turns mTLS off.mTLS misconfiguration now fails closed (
MtlsUtils):IllegalStateExceptionin these cases:GOOGLE_API_CERTIFICATE_CONFIGfile is missing, unreadable or malformed;IOExceptioninstead of falling back to a non-mTLS connection.🧪 Testing
Automated Testing
• Added and updated comprehensive unit-tests reflecting the thread-safety
fixes inside ChannelPoolTest.java and RefreshingHttpJsonChannelTest.java.
• Corrected edge case test configurations to leverage realistic mocked X.509
certificates to properly exercise deep WorkloadCertificateUtils.
getCertificateFingerprint() filesystem caching mechanisms.
Manual Testing
End-to-end tests with real generated GAPIC clients against local mTLS servers, run on this PR's head. Each scenario runs in its own JVM and checks the exact sequence of client certificates and outcomes the server saw.
KeyManagementServiceClientover gRPC and HTTP/JSON, andGrpcBigQueryReadStubfor server-streamingReadRowswith offset-based resumption.UNAUTHENTICATED/ 401 for "revoked" certificates (by CN).GOOGLE_API_CERTIFICATE_CONFIGis replaced by atomic rename, key first, then cert.close()with calls in flight: 70,624/70,624 calls succeeded, no hangs.close()(49,561/49,561 calls succeeded, no hangs). The socket factory was verified for both Conscrypt and plain JDK TLS.GOOGLE_API_USE_CLIENT_CERTIFICATE=false, or no certificate config): a single attempt with no client cert and no refresh, over both gRPC and HTTP/JSON, even when the cert files on disk change.Unknown authType: GENERIChandshake error, which fix(gax-httpjson): use Conscrypt TrustManagerFactory for mTLS SSLContext [blocked on #13995] #14556 fixes. With fix(gax-httpjson): use Conscrypt TrustManagerFactory for mTLS SSLContext [blocked on #13995] #14556 applied on top of this PR, all 22 scenarios pass on JDK 21 with Conscrypt active.