Repository navigation
feat(gax): support transparent retries during mTLS certificate rotations - #13995
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
…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.
| try { | ||
| refresh(); | ||
| synchronized (entryWriteLock) { | ||
| String currentDiskFingerprint = rotationTracker.readDiskFingerprint(); |
There was a problem hiding this comment.
readDiskFingerprint looks to read from a file which may not be needed for cases where there is no mtls cert required.
Can we do a split where we check for the non-mtls case first:
- If
workloadCertPath== null, we simply refreshAll() without having to load the certificate - Otherwise, we proceed with the logic to read the cert and then to try to rotate?
There was a problem hiding this comment.
Ah, basically what we have for refresh() down below
There was a problem hiding this comment.
Done. refreshSafely() now matches refresh(): with no workloadCertPath it calls refreshAll() and returns without touching the certificate. Only the mTLS path reads the cert. (readDiskFingerprint() already returned "" without any file I/O when the path was null, so behavior is unchanged, but the split makes that explicit.)
| if (workloadCertPath != null && currentDiskFingerprint.isEmpty()) { | ||
| return; | ||
| } |
There was a problem hiding this comment.
IIUC, if workloadCertPath != null, the only case for currentDiskFingerprint being non-empty is during a rotation right?
Can we add a comment to note this (perhaps a debug log here as well that it'll recreate on the next refresh)
There was a problem hiding this comment.
Yes, assuming you meant empty: with a configured path, an empty fingerprint means the certificate couldn't be read. In practice that's almost always a rotation in progress, where the file is empty or only partly written. A missing or unreadable file looks the same. Added a comment and a FINE log saying the refresh is skipped and the channels will be recreated on the next refresh, plus preemptiveRefresh_whenCertUnreadable_skipsRefreshAndLogs.
| for (Entry e : replacedEntries) { | ||
| if (!finalEntries.contains(e)) { | ||
| e.requestShutdown(); | ||
| } | ||
| } |
There was a problem hiding this comment.
I know this is existing code, but I think it may be possible to clean this up. The looping through and finding a match via contains seems less than ideal, especially given the freq that we need to refresh and that we want to speed this up mid-call
Perhaps we can track newEntries and removedEntries.
- Try to create a new channel, if success then add the new channel to newEntries
- If failed, then if we dropUnfrefeshed, add it to shutdownEntries, otherwise, we add the old, non-refreshed channel to new entries
We should be able to just loop through shutdownEntries without checking if the final exists.
Perhaps to handle the case with createdEntries in finally, we'll need to also add new entries to createdEntries
There was a problem hiding this comment.
Agreed, done as you suggested. refreshAll now makes one pass that builds keptEntries and retiredEntries. A successful recreate keeps the new channel and retires the old one. A failure retires the old channel when dropUnrefreshedChannels is set, and keeps it otherwise. After the swap we shut down retiredEntries directly, so the contains() scan is gone. New channels are still tracked in createdEntries so the finally block can shut them down if an Error aborts before the swap. Every write to entries already happens under entryWriteLock, so a plain set replaces the getAndSet. I also added tests confirming that each channel keeps its slot when old channels are kept, and that a rotation where every channel fails leaves the pool untouched (…allChannelsFail_leavesPoolUnchangedAndDoesNotRefill).
| e.requestShutdown(); | ||
| } | ||
| } | ||
| if (dropUnrefreshedChannels && !allCreated && settings.isStaticSize()) { |
There was a problem hiding this comment.
qq, I'm not sure I understand why we check for isStaticSize here. Most of the GAPICs would be static size == 1, but customers can override this value.
There was a problem hiding this comment.
Oh nvm. I remembered out earlier convo about letting resize() populate this back up. We'll probably incur a bit of queueing if channel size drops, then the potential issue where there is a bunch of consecutive resizes. Let me think about this again.
There was a problem hiding this comment.
I went ahead and removed the isStaticSize check. Any rotation refresh that drops channels now schedules the one-time refill, for dynamic pools as well as static ones. Without the refill, a dynamic pool that drops from N channels to a few only grows back by maxResizeDelta per resize interval, so traffic queues on fewer channels and you get the consecutive-resize pattern you mentioned. The refill avoids both. It runs once on the background executor, takes no retries, and is a no-op if resize() has already grown the pool or the pool has shut down. Happy to revisit if you'd rather leave dynamic pools to resize().
| if (isShutdown) { | ||
| return; | ||
| } | ||
| int targetSize = settings.getInitialChannelCount(); |
There was a problem hiding this comment.
I think the target size should probably be the number of old entries in the channelPool (not the initial channel count as this doesn't account for potential channel/ load change - growth or loss of load).
I also realized that the load calculation is stored within a channel itself, so refreshing a channel ends up losing that data which is not idea, but is existing behavior so can be resolved in the future. Maybe something like a snapshot of load stored in channel pool and not in the channel or something.
There was a problem hiding this comment.
Agreed, the refill now targets the number of channels the pool had before the refresh, not initialChannelCount, so it keeps whatever size resizing had settled on. For a static pool the two are the same. Added …partialFailureAfterResize_refillsToPreRefreshSize.
On the load stats: agreed. Each channel's outstanding-RPC peak is lost when it's recreated, which predates this PR. Keeping a load snapshot at the pool level would add new state, so I'd prefer a follow-up. I can open an issue if that works for you.
| boolean rotated = | ||
| !currentDiskFingerprint.isEmpty() | ||
| && !rotationTracker.isAlreadyActive(currentDiskFingerprint); | ||
| if (refreshAll() && !currentDiskFingerprint.isEmpty()) { |
There was a problem hiding this comment.
nit: Can we change the name to be from rotated to certChanged. I think rotated is a bit unclear as it may pertain to channel or cert.
Additionally, I might be missing this case, but why do we need need to markRefresh on L478 below when rotated is false? IIUC, think we should mark it only when we need to rotate?
If the cert changed, then we should opt refresh based on that: refreshAll(certChanged). But we would only do completeCertificateSwitch if certChanged is true as well.
I think the logic would be
if (refreshAll(certChanged) && certChanged) { completeCertificateSwitch(currentDiskFingerprint) };
refreshAll(false) && false is the normal refresh cycle for mtls with no certChanged so it shouldn't call completeCertificateSwitch
refreshAll(true) && true is to refresh the channels with certChanged
- But if refreshAll fails to refresh all the channels, then refresh() on the next call will try to update it.
There was a problem hiding this comment.
Agreed on all three. Renamed to certChanged. The markRefreshed call in the no-change branch was redundant (the fingerprint already equals the active one), so it's gone. The periodic refresh now calls refreshAll(certChanged) and only completes the switch when the cert changed. A normal mTLS cycle keeps any channel that fails to rebuild, as before. When the cert changed, channels still on the old certificate are dropped (and the pool refilled), the same as the reactive refresh(). If no channel can be rebuilt, nothing changes and the next refresh retries. This extends the drop-on-rotation behavior we agreed for the reactive path (option 3) to the periodic refresh when it detects a changed cert, so it supersedes my earlier note that the periodic refresh always keeps per-slot fallback; that now only applies when the cert is unchanged. Both invariants still hold: the cert is only marked active, and the generation only bumps, once every channel left in the pool uses it. Updated the partial-failure test and added tests for the all-fail and no-change cases.
| + " channels will be recreated on the next refresh"); | ||
| return; | ||
| } | ||
| boolean rotated = !rotationTracker.isAlreadyActive(currentDiskFingerprint); |
There was a problem hiding this comment.
I think we should probably have a check similar to what we have in resize:
// Double-check fingerprint inside the lock
if (rotationTracker.isAlreadyActive(currentDiskFingerprint)) {
LOG.fine(
"Channel pool was already refreshed by a concurrent thread, skipping duplicate"
+ " refresh");
return;
}
There was a problem hiding this comment.
For the periodic refresh, an isAlreadyActive early return would skip it in steady state (the cert on disk is normally the active one), and it still needs to recycle channels for the GFE disconnects. When the cert is unchanged it now just refreshes the channels without switching, and the 'already refreshed by a concurrent thread' case is covered by the pre-lock generation check from the other thread.
| private void refreshSafely() { | ||
| try { | ||
| refresh(); | ||
| synchronized (entryWriteLock) { |
There was a problem hiding this comment.
Double check my on this: synchronized ensures that only one refresh request goes out at once, but only ensure that we don't refresh via the rotationTracker.isAlreadyActive call.
This looks like it will still read the disk X - 1 times, if we have X refresh/ refreshSafely calls queued. I think we can optimize this a bit by checking the generation before the lock (e.g. preLockGenerationCount) and comparing it once it obtains the lock (generation.get()). If the obtained generation lock is after the pre-lock then we know that the channels have been refreshed.
I think we'll need to add this to refresh and refreshSafely since they are independent (50 min interval vs cert refresh)
There was a problem hiding this comment.
Good catch: refresh() reads the fingerprint uncached inside the lock, so callers queued behind a rotation refresh each re-read the cert file. Both refresh() and the periodic refresh now snapshot the generation before taking the lock and return early if it advanced while waiting, since a concurrent refresh already switched every channel to a new certificate. If that refresh couldn't rebuild anything, the generation doesn't move, so the next waiter still retries. Added tests that rotate the cert again while a caller waits on the lock, to confirm it skips without reading the disk.
| /** {@inheritDoc} */ | ||
| @Override | ||
| public void refresh() { | ||
| synchronized (refreshLock) { |
There was a problem hiding this comment.
ah, look at this and I think it has the same sync concern with multiple concurrent requests at once. I think we can reduce the need to read from disk by checking the generation.
There was a problem hiding this comment.
Agreed, same fix applied here: refresh() snapshots the generation before taking refreshLock and returns if a concurrent refresh already swapped the transport, so queued callers no longer re-read the cert. Added a test for it.
| if (newFingerprint != null && !newFingerprint.isEmpty()) { | ||
| this.activeCertFingerprint = newFingerprint; | ||
| this.lastDiskCheck = null; | ||
| } |
There was a problem hiding this comment.
(Gemini flagged)
This should be locked by the diskCheckLock as getOrUpdateDiskFingerprint can overwrite the lastDiskCheck.
--
I think this may end up with a race to update the lastDiskCheck. If markRefreshed and getOrUpdateDiskFingerprint, should we keep this lastDiskCheck from getOrUpdateDiskFingerprint?
There was a problem hiding this comment.
Agreed. markRefreshed now updates the active fingerprint and clears lastDiskCheck while holding diskCheckLock, so an in-progress disk check can no longer write a stale result back after the switch. Previously that could make shouldRefresh() return a stale true for up to the 1s cache TTL. That only caused a no-op refresh(), since retries are decided by the generation. Keeping the result from getOrUpdateDiskFingerprint wouldn't be safe, since it can be the old certificate. Lock order is always the pool/refresh lock first, then diskCheckLock, so there's no deadlock. The trade-off is that markRefreshed may wait for one in-progress cert read, and refresh() already reads the cert under the pool lock anyway. Added a test.
|
|
||
| // Restore the pool to its pre-refresh size right away, rather than leave a dynamically | ||
| // sized | ||
| // pool to grow back over several resize() runs while traffic queues on fewer channels. |
There was a problem hiding this comment.
(Gemini flagged)
Not this line, but gemini says that resize() should also check for isShutdown before resizing the channels
There was a problem hiding this comment.
Good catch. A resize that was already waiting for the lock when shutdown() ran could still expand the pool afterwards and leak channels (cancel(true) doesn't interrupt a thread blocked on a monitor). resizeSafely now returns if the pool is shut down, matching refresh and the refill. Added a test.
lqiu96
left a comment
There was a problem hiding this comment.
Took another round of review on the PR. Changes look fine with me.
I think some small nits:
- Let's see if we can update some of full qualified names to use the short name (org.mockito.Mockito.mock, java.util.concurrent.atomic.AtomicLong)
- CertificateRotationTracker probably should have some tests if possible
|
Thanks! Replaced the fully qualified names added in this PR with imports (except |
2d5f9ed
into
googleapis:agentic-identities-bound-token
…ext (#14556) > [!IMPORTANT] > **This must merge before the `agentic-identities-bound-token` feature branch is merged to `main`.** Without it, HTTP/JSON clients fail by default wherever mTLS is enabled automatically (see below). ## Problem When mTLS is active, `InstantiatingHttpJsonChannelProvider.configureMtls()` builds a Conscrypt `SSLContext` but initializes it with the JDK (SunJSSE) PKIX `TrustManagerFactory`. On TLS 1.3, Conscrypt passes authType `"GENERIC"` to the trust manager. SunJSSE's end-entity checks reject that authType for server certificates that are CA-issued and carry a KeyUsage extension, which includes Google front ends. So every mTLS HTTP/JSON handshake fails with: ``` javax.net.ssl.SSLHandshakeException: Unknown authType: GENERIC Caused by: java.security.cert.CertificateException: Unknown authType: GENERIC ``` The bug is already on `main`, but there it only triggers when `GOOGLE_API_USE_CLIENT_CERTIFICATE=true` is set explicitly. On this feature branch, #13995 enables mTLS automatically whenever a workload certificate config is present, so on Cloud Run (agent identity) **every HTTP/JSON client fails by default**. gRPC isn't affected. The non-mTLS path isn't affected either, because google-http-client already pairs Conscrypt with a provider-matched trust manager there. A second provider mismatch shows up on JDK 26: the mTLS `SSLContext` also used the JDK's default (`SunX509`) `KeyManagerFactory`. Since JDK 26 ([JDK-8359956](https://bugs.openjdk.org/browse/JDK-8359956)), that key manager applies algorithm constraints to the client certificate. Called from a Conscrypt handshake, it rejects valid certificates (e.g. `SHA256withRSA`), so the client sends an empty certificate chain and the server aborts with `certificate_required`. ## Fix Use Conscrypt's own PKIX `TrustManagerFactory` (`TrustManagerFactory.getInstance("PKIX", conscryptProvider)`) for the mTLS `SSLContext`. I called the JDK API directly rather than `SslUtils.getPkixTrustManagerFactory(Provider)`, which only exists in google-http-client 2.2.0+. Likewise, use Conscrypt's own `KeyManagerFactory` (`KeyManagerFactory.getInstance("PKIX", conscryptProvider)`), so the trust manager and key manager both come from the same provider as the `SSLContext`. Conscrypt's trust manager loads the same default trust store as the JDK. Verified on JDK 21: 174 anchors in both, the same set, and both honor `-Djavax.net.ssl.trustStore` overrides identically. ## Testing - New `InstantiatingHttpJsonChannelProviderTls13Test`: a local JDK TLS 1.3 server presents a CA-issued leaf with KeyUsage and requires a client certificate. The test asserts that the transport completes the request and presents its client certificate. - Without the fix: fails with `Unknown authType: GENERIC`. - With the fix: passes. - Skips when Conscrypt native or TLS 1.3 is unavailable. - Full `gax-httpjson` suite on the current `agentic-identities-bound-token` (with #13995 merged): 194/194 pass on JDK 8, 11, 17, 21, 25 and 26. On JDK 26, the Tls13 test fails without the `KeyManagerFactory` change: the client sends an empty certificate chain and times out. - End-to-end with #13995's certificate-rotation support, on JDK 21 with Conscrypt active: 22 scenarios with real KMS and BigQuery Storage GAPIC clients against local mTLS servers, where the certificate rotates on disk. - With this fix: all 22 pass. That includes HTTP/JSON rotation, in-flight calls during refresh, and stress with `close()`. The refreshed transport still uses Conscrypt's socket factory. - Without this fix: every HTTP/JSON mTLS scenario fails with `Unknown authType: GENERIC`. Limiting the test server to TLS 1.2 makes the error go away, which confirms that the TLS 1.3 authType is the trigger. - Logs: https://paste.googleplex.com/5563956459077632 (with fix), https://paste.googleplex.com/5467233782988800 (without fix, JDK 21 set) - Live on Cloud Run (agent identity), combined #13873 + #13995 build, default env: | | gRPC | HTTP/JSON | |---|---|---| | Before | ✅ | ❌ `Unknown authType: GENERIC` | | After | ✅ | ✅ authenticated over mTLS with a bound token |
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 #14556 fixes. With fix(gax-httpjson): use Conscrypt TrustManagerFactory for mTLS SSLContext #14556 applied on top of this PR, all 22 scenarios pass on JDK 21 with Conscrypt active.