This is an automated email from the ASF dual-hosted git repository.

lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git


The following commit(s) were added to refs/heads/rocketmq-studio by this push:
     new c99b9ad50 fix(aliyun): prevent stale clients after credential rotation 
(#5252)
c99b9ad50 is described below

commit c99b9ad50b79467546997e51f3401d57ceebea5f
Author: zmuxuny <[email protected]>
AuthorDate: Fri Oct 9 00:03:41 2026 -0700

    fix(aliyun): prevent stale clients after credential rotation (#5252)
    
    Rotating an Aliyun credential evicts the cached OpenAPI clients, but 
`invalidateCredential` could
    finish while a cache miss was still building a client from the pre-rotation 
secret - and that client
    was published into the map after the eviction, so subsequent calls kept 
using a secret that no longer
    exists. `ConcurrentHashMap.computeIfAbsent` only locks the bin for its own 
key, so the
    `entrySet().removeIf` in `invalidateCredential` never sees a creation that 
is still in flight.
    
    Both paths now take the factory monitor: a lock-free `get` fast path, then 
`computeIfAbsent` under
    the same lock that `invalidateCredential` holds. This is the shape 
`TencentClientFactory` already
    uses, so the two vendor factories agree.
    
    #5319 carried the same fix and was closed as a duplicate in favour of this 
PR, which has the earlier
    lineage (#5052).
---
 .../provider/alibaba/AliyunClientFactory.java      | 10 +++-
 .../provider/alibaba/AliyunClientFactoryTest.java  | 55 ++++++++++++++++++++++
 2 files changed, 63 insertions(+), 2 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactory.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactory.java
index 76234f3eb..2e4295192 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactory.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactory.java
@@ -61,13 +61,19 @@ public class AliyunClientFactory {
 
     public AsyncClient client(Long credentialId, String region) {
         String key = cacheKey(credentialId, region);
-        return clients.computeIfAbsent(key, ignored -> 
createClient(credentialId, region));
+        AsyncClient cached = clients.get(key);
+        if (cached != null) {
+            return cached;
+        }
+        synchronized (this) {
+            return clients.computeIfAbsent(key, ignored -> 
createClient(credentialId, region));
+        }
     }
 
     /**
      * Releases all clients created with a credential so the next call 
observes rotated secrets.
      */
-    public void invalidateCredential(Long credentialId) {
+    public synchronized void invalidateCredential(Long credentialId) {
         String prefix = credentialId + "#";
         clients.entrySet().removeIf(entry -> {
             if (!entry.getKey().startsWith(prefix)) {
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactoryTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactoryTest.java
index b75338fa8..506481c4a 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactoryTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunClientFactoryTest.java
@@ -38,6 +38,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.FutureTask;
 import java.util.concurrent.TimeUnit;
 
 import static org.assertj.core.api.Assertions.assertThat;
@@ -116,6 +117,60 @@ class AliyunClientFactoryTest {
         assertThat(factory.client(CREDENTIAL_ID, REGION)).isSameAs(second);
     }
 
+    @Test
+    void invalidateCredentialShouldEvictClientCreatedDuringRotationTest() 
throws Exception {
+        AsyncClient oldClient = Mockito.mock(AsyncClient.class);
+        AsyncClient replacement = Mockito.mock(AsyncClient.class);
+        CountDownLatch creationStarted = new CountDownLatch(1);
+        CountDownLatch allowCreation = new CountDownLatch(1);
+        AtomicInteger creations = new AtomicInteger();
+        factory = new AliyunClientFactory(credentialRepository) {
+            @Override
+            protected AsyncClient createClient(Long credentialId, String 
region) {
+                if (creations.incrementAndGet() == 1) {
+                    creationStarted.countDown();
+                    try {
+                        if (!allowCreation.await(5, TimeUnit.SECONDS)) {
+                            throw new IllegalStateException("Timed out waiting 
to create client");
+                        }
+                    } catch (InterruptedException exception) {
+                        Thread.currentThread().interrupt();
+                        throw new IllegalStateException("Client creation 
interrupted", exception);
+                    }
+                    return oldClient;
+                }
+                return replacement;
+            }
+        };
+
+        FutureTask<AsyncClient> creation = new FutureTask<>(() -> 
factory.client(CREDENTIAL_ID, REGION));
+        Thread creationThread = new Thread(creation, 
"aliyun-client-creation-test");
+        creationThread.start();
+        try {
+            assertThat(creationStarted.await(5, TimeUnit.SECONDS)).isTrue();
+            FutureTask<Void> invalidation = new FutureTask<>(() -> {
+                factory.invalidateCredential(CREDENTIAL_ID);
+                return null;
+            });
+            Thread invalidationThread = new Thread(invalidation, 
"aliyun-client-invalidation-test");
+            invalidationThread.start();
+            long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
+            while (!invalidation.isDone() && invalidationThread.getState() != 
Thread.State.BLOCKED
+                    && System.nanoTime() < deadline) {
+                Thread.onSpinWait();
+            }
+            assertThat(invalidation.isDone() || invalidationThread.getState() 
== Thread.State.BLOCKED).isTrue();
+
+            allowCreation.countDown();
+            assertThat(creation.get(5, TimeUnit.SECONDS)).isSameAs(oldClient);
+            invalidation.get(5, TimeUnit.SECONDS);
+            assertThat(factory.client(CREDENTIAL_ID, 
REGION)).isSameAs(replacement);
+            Mockito.verify(oldClient).close();
+        } finally {
+            allowCreation.countDown();
+        }
+    }
+
     @Test
     void callShouldThrow504OnTimeoutTest() {
         factory.setCallTimeoutSeconds(1L);

Reply via email to