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);