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