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

Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git


The following commit(s) were added to refs/heads/master by this push:
     new de5e7dd196 fix: replace RocketMQ e2e sleeps with Awaitility (#6817) 
(#6978)
de5e7dd196 is described below

commit de5e7dd19695112b14696e93d1800cbd89bd3eec
Author: Limbo <[email protected]>
AuthorDate: Mon Sep 28 10:57:53 2026 +0800

    fix: replace RocketMQ e2e sleeps with Awaitility (#6817) (#6978)
    
    Co-authored-by: aias00 <[email protected]>
    Co-authored-by: Liming Deng <[email protected]>
---
 shenyu-e2e/pom.xml                                     |  7 +++++++
 .../shenyu-e2e-case-logging-rocketmq/pom.xml           |  5 +++++
 .../testcase/logging/rocketmq/DividePluginCases.java   | 18 +++++++++++-------
 3 files changed, 23 insertions(+), 7 deletions(-)

diff --git a/shenyu-e2e/pom.xml b/shenyu-e2e/pom.xml
index 59074ad53c..6be2b25708 100644
--- a/shenyu-e2e/pom.xml
+++ b/shenyu-e2e/pom.xml
@@ -42,6 +42,7 @@
     <properties>
         <java.version>17</java.version>
         <junit.version>5.8.2</junit.version>
+        <awaitility.version>4.0.3</awaitility.version>
         <assertj.version>3.27.7</assertj.version>
         <hamcrest.version>1.3</hamcrest.version>
         <jsonassert.version>1.5.0</jsonassert.version>
@@ -196,6 +197,12 @@
                 <scope>import</scope>
             </dependency>
 
+            <dependency>
+                <groupId>org.awaitility</groupId>
+                <artifactId>awaitility</artifactId>
+                <version>${awaitility.version}</version>
+            </dependency>
+
             <dependency>
                 <groupId>com.fasterxml.jackson</groupId>
                 <artifactId>jackson-bom</artifactId>
diff --git 
a/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/pom.xml 
b/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/pom.xml
index f9e8eef963..6f3bd92a54 100644
--- a/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/pom.xml
+++ b/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/pom.xml
@@ -32,5 +32,10 @@
             <artifactId>rocketmq-client</artifactId>
             <version>4.9.3</version>
         </dependency>
+        <dependency>
+            <groupId>org.awaitility</groupId>
+            <artifactId>awaitility</artifactId>
+            <scope>test</scope>
+        </dependency>
     </dependencies>
 </project>
diff --git 
a/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/src/test/java/org/apache/shenyu/e2e/testcase/logging/rocketmq/DividePluginCases.java
 
b/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/src/test/java/org/apache/shenyu/e2e/testcase/logging/rocketmq/DividePluginCases.java
index bc24271d34..87ce0e4a96 100644
--- 
a/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/src/test/java/org/apache/shenyu/e2e/testcase/logging/rocketmq/DividePluginCases.java
+++ 
b/shenyu-e2e/shenyu-e2e-case/shenyu-e2e-case-logging-rocketmq/src/test/java/org/apache/shenyu/e2e/testcase/logging/rocketmq/DividePluginCases.java
@@ -35,6 +35,7 @@ import org.junit.jupiter.api.Assertions;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.time.Duration;
 import java.util.List;
 import java.util.concurrent.atomic.AtomicBoolean;
 
@@ -42,6 +43,7 @@ import static 
org.apache.shenyu.e2e.engine.scenario.function.HttpCheckers.exists
 import static 
org.apache.shenyu.e2e.template.ResourceDataTemplate.newConditions;
 import static 
org.apache.shenyu.e2e.template.ResourceDataTemplate.newRuleBuilder;
 import static 
org.apache.shenyu.e2e.template.ResourceDataTemplate.newSelectorBuilder;
+import static org.awaitility.Awaitility.await;
 
 public class DividePluginCases implements ShenYuScenarioProvider {
 
@@ -53,6 +55,8 @@ public class DividePluginCases implements 
ShenYuScenarioProvider {
 
     private static final String TEST = "/http/order/findById?id=123";
 
+    private static final Duration LOG_CONSUME_TIMEOUT = Duration.ofSeconds(30);
+
     private static final Logger LOG = 
LoggerFactory.getLogger(DividePluginCases.class);
 
     @Override
@@ -99,10 +103,8 @@ public class DividePluginCases implements 
ShenYuScenarioProvider {
                         ShenYuCaseSpec.builder()
                                 .add(request -> {
                                     AtomicBoolean isLog = new 
AtomicBoolean(false);
+                                    DefaultMQPushConsumer consumer = new 
DefaultMQPushConsumer(CONSUMERGROUP);
                                     try {
-                                        Thread.sleep(1000 * 30);
-                                        request.request(Method.GET, 
"/http/order/findById?id=23");
-                                        DefaultMQPushConsumer consumer = new 
DefaultMQPushConsumer(CONSUMERGROUP);
                                         consumer.setNamesrvAddr(NAMESERVER);
                                         consumer.subscribe(TOPIC, "*");
                                         
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, 
consumeConcurrentlyContext) -> {
@@ -116,14 +118,16 @@ public class DividePluginCases implements 
ShenYuScenarioProvider {
                                             }
                                             return 
ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                                         });
-                                        LOG.info("consumer.start ; 
isLog.get():{}", isLog.get());
                                         consumer.start();
-                                        Thread.sleep(1000 * 30);
+                                        LOG.info("RocketMQ consumer started");
+                                        request.request(Method.GET, 
"/http/order/findById?id=23");
+                                        
await().atMost(LOG_CONSUME_TIMEOUT).untilTrue(isLog);
                                         LOG.info("isLog.get():{}", 
isLog.get());
-                                        Assertions.assertTrue(isLog.get());
                                     } catch (Exception e) {
                                         LOG.error("error", e);
-                                        Assertions.assertTrue(isLog.get());
+                                        Assertions.fail("Failed to consume 
RocketMQ access log", e);
+                                    } finally {
+                                        consumer.shutdown();
                                     }
                                 })
                                 .build()

Reply via email to