This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 44c5faca0f51 CAMEL-25342: Fix flaky test NatsConsumerIT (#27505)
44c5faca0f51 is described below
commit 44c5faca0f51d057317a9966ff0c204beb6c2fec
Author: Guillaume Nodet <[email protected]>
AuthorDate: Wed Oct 7 21:30:23 2026 +0200
CAMEL-25342: Fix flaky test NatsConsumerIT (#27505)
Co-authored-by: Claude Sonnet 4.5 <[email protected]>
---
.../camel/component/nats/integration/NatsConsumerIT.java | 14 +++++++++++++-
1 file changed, 13 insertions(+), 1 deletion(-)
diff --git
a/components/camel-nats/src/test/java/org/apache/camel/component/nats/integration/NatsConsumerIT.java
b/components/camel-nats/src/test/java/org/apache/camel/component/nats/integration/NatsConsumerIT.java
index c768042602ee..6b75bb47349c 100644
---
a/components/camel-nats/src/test/java/org/apache/camel/component/nats/integration/NatsConsumerIT.java
+++
b/components/camel-nats/src/test/java/org/apache/camel/component/nats/integration/NatsConsumerIT.java
@@ -16,12 +16,18 @@
*/
package org.apache.camel.component.nats.integration;
+import java.util.concurrent.TimeUnit;
+
import org.apache.camel.EndpointInject;
+import org.apache.camel.Route;
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.component.mock.MockEndpoint;
import org.apache.camel.component.nats.NatsConstants;
+import org.apache.camel.component.nats.NatsConsumer;
import org.junit.jupiter.api.Test;
+import static org.awaitility.Awaitility.await;
+
public class NatsConsumerIT extends NatsITSupport {
@EndpointInject("mock:result")
@@ -32,9 +38,15 @@ public class NatsConsumerIT extends NatsITSupport {
mockResultEndpoint.expectedBodiesReceived("Hello World");
mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT,
"test");
+ // Wait for the NATS consumer to be subscribed before sending messages,
+ // since core NATS does not persist messages for inactive subscribers
+ await().atMost(10, TimeUnit.SECONDS)
+ .until(() -> context.getRoutes().stream()
+ .map(Route::getConsumer)
+ .anyMatch(c -> c instanceof NatsConsumer nc &&
nc.isActive()));
+
template.sendBody("direct:send", "Hello World");
- mockResultEndpoint.setAssertPeriod(5000);
mockResultEndpoint.assertIsSatisfied();
}