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 f575ae5e3 fix(tencent): coordinate client creation and invalidation
(#2196)
f575ae5e3 is described below
commit f575ae5e37a29e788d9f92b55fb0a3949c0e8bfa
Author: yyqdbngt <[email protected]>
AuthorDate: Tue Aug 18 18:54:30 2026 +0800
fix(tencent): coordinate client creation and invalidation (#2196)
Credential invalidation could finish while a cache miss was still creating
a client, allowing the old client to enter the cache afterward.
Serialize cache-miss creation with invalidation while keeping the normal
cache-hit path lock-free.
Signed-off-by: yyqdbngt <[email protected]>
---
.../provider/tencent/TencentClientFactory.java | 10 +++-
.../provider/tencent/TencentClientFactoryTest.java | 67 ++++++++++++++++++++++
2 files changed, 75 insertions(+), 2 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactory.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactory.java
index f6cd9c7d8..a10384776 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactory.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactory.java
@@ -48,10 +48,16 @@ public class TencentClientFactory {
public TrocketClient client(Long credentialId, String region) {
String key = cacheKey(credentialId, region);
- return clients.computeIfAbsent(key, ignored ->
createClient(credentialId, region));
+ TrocketClient cached = clients.get(key);
+ if (cached != null) {
+ return cached;
+ }
+ synchronized (this) {
+ return clients.computeIfAbsent(key, ignored ->
createClient(credentialId, region));
+ }
}
- public void invalidateCredential(Long credentialId) {
+ public synchronized void invalidateCredential(Long credentialId) {
String prefix = credentialId + "#";
clients.keySet().removeIf(key -> key.startsWith(prefix));
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactoryTest.java
index f4b8ce71f..f3604f4ed 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentClientFactoryTest.java
@@ -16,9 +16,17 @@
*/
package org.apache.rocketmq.studio.provider.tencent;
+import com.tencentcloudapi.trocket.v20230308.TrocketClient;
+import
org.apache.rocketmq.studio.provider.credential.CloudCredentialRepository;
import org.junit.jupiter.api.Test;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
class TencentClientFactoryTest {
@@ -45,4 +53,63 @@ class TencentClientFactoryTest {
assertThat(TencentClientFactory.endpointFor("ap-shanghai-fsi"))
.isEqualTo("trocket.ap-shanghai-fsi.tencentcloudapi.com");
}
+
+ @Test
+ void invalidateCredentialShouldRemoveAnInFlightClientCreationTest() throws
Exception {
+ long credentialId = 1L;
+ BlockingClientFactory factory = new BlockingClientFactory();
+ FutureTask<TrocketClient> firstClientTask = new FutureTask<>(
+ () -> factory.client(credentialId, "ap-shanghai"));
+ Thread creationThread = new Thread(firstClientTask,
"tencent-client-creation-test");
+ creationThread.start();
+ assertThat(factory.creationStarted.await(5,
TimeUnit.SECONDS)).isTrue();
+
+ Thread invalidationThread = new Thread(
+ () -> factory.invalidateCredential(credentialId),
+ "tencent-client-invalidation-test");
+ invalidationThread.start();
+ awaitBlockedOrTerminated(invalidationThread);
+
+ factory.allowCreation.countDown();
+ TrocketClient firstClient = firstClientTask.get(5, TimeUnit.SECONDS);
+ invalidationThread.join(TimeUnit.SECONDS.toMillis(5));
+
+ assertThat(invalidationThread.isAlive()).isFalse();
+ assertThat(factory.client(credentialId,
"ap-shanghai")).isNotSameAs(firstClient);
+ assertThat(factory.creationCount).hasValue(2);
+ }
+
+ private static void awaitBlockedOrTerminated(Thread thread) {
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
+ while (thread.getState() != Thread.State.BLOCKED && thread.isAlive()
+ && System.nanoTime() < deadline) {
+ Thread.onSpinWait();
+ }
+ assertThat(thread.getState() == Thread.State.BLOCKED ||
!thread.isAlive()).isTrue();
+ }
+
+ private static final class BlockingClientFactory extends
TencentClientFactory {
+ private final CountDownLatch creationStarted = new CountDownLatch(1);
+ private final CountDownLatch allowCreation = new CountDownLatch(1);
+ private final AtomicInteger creationCount = new AtomicInteger();
+
+ private BlockingClientFactory() {
+ super(mock(CloudCredentialRepository.class));
+ }
+
+ @Override
+ protected TrocketClient createClient(Long credentialId, String region)
{
+ TrocketClient client = mock(TrocketClient.class);
+ if (creationCount.incrementAndGet() == 1) {
+ creationStarted.countDown();
+ try {
+ assertThat(allowCreation.await(5,
TimeUnit.SECONDS)).isTrue();
+ } catch (InterruptedException exception) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Client creation
interrupted", exception);
+ }
+ }
+ return client;
+ }
+ }
}