This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch quick-fix/nats-jetstream-it-stream-collision in repository https://gitbox.apache.org/repos/asf/camel.git
commit a264bb94c9bc3a6aea3ba3aeb1bc91db0448e27e Author: Claus Ibsen <[email protected]> AuthorDate: Mon Aug 24 17:33:07 2026 +0200 chore: fix NATS JetStream IT tests colliding on shared stream/durable name NatsJetstreamConsumerAckPolicyNoneIT, NatsJetstreamConsumerMaxDeliverIT, and NatsJetstreamConsumerRedeliveryIT all used the identical JetStream stream name (mystream2), subject (mytopic2), and durable consumer name (camel2). NatsITSupport does not tear down streams/consumers between IT classes, and NatsConsumer.setupJetStreamConsumer() never reconciles an already-existing durable consumer's config with a new subscribe request, so whichever of these three tests ran first bound its own ackPolicy/ maxDeliver config to the shared durable consumer server-side. The next test(s) in the suite then silently bound to that stale, mismatched consumer and received zero messages. This is why NatsJetstreamConsumerMaxDeliverIT and NatsJetstreamConsumerRedeliveryIT were previously guarded with @DisabledIfSystemProperty(named = "ci.env.name", ..., "Flaky on GitHub Actions"). That guard was removed without addressing the underlying collision, most likely because it was verified by running each test class individually rather than as part of the full suite. Give each test its own unique stream/subject/durable name, matching the convention already used by the other JetStream IT tests in this package (-manualack, -manualack-nak, -pull, etc). Test-only change, no production code touched. Co-authored-by: Claude <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../nats/jetstream/NatsJetstreamConsumerAckPolicyNoneIT.java | 6 +++--- .../component/nats/jetstream/NatsJetstreamConsumerMaxDeliverIT.java | 4 ++-- .../component/nats/jetstream/NatsJetstreamConsumerRedeliveryIT.java | 6 +++--- 3 files changed, 8 insertions(+), 8 deletions(-) diff --git a/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerAckPolicyNoneIT.java b/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerAckPolicyNoneIT.java index 11a408588ead..a265f88d7931 100644 --- a/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerAckPolicyNoneIT.java +++ b/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerAckPolicyNoneIT.java @@ -40,12 +40,12 @@ public class NatsJetstreamConsumerAckPolicyNoneIT extends NatsITSupport { @Test public void testConsumer() throws Exception { mockResultEndpoint.expectedBodiesReceived("Hello World 3"); - mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2"); + mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2-ackpolicynone"); mockResultEndpoint.expectedHeaderReceived("counter", 3); mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_DELIVERY_COUNTER, 1); mockInputEndpoint.expectedBodiesReceived("Hello World 1", "Hello World 2", "Hello World 3"); - mockInputEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2"); + mockInputEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2-ackpolicynone"); mockInputEndpoint.message(0).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(1); mockInputEndpoint.message(1).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(1); mockInputEndpoint.message(2).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(1); @@ -65,7 +65,7 @@ public class NatsJetstreamConsumerAckPolicyNoneIT extends NatsITSupport { @Override public void configure() { String uri - = "nats:mytopic2?jetstreamEnabled=true&jetstreamName=mystream2&jetstreamAsync=false&durableName=camel2&pullSubscription=false&ackPolicy=none"; + = "nats:mytopic2-ackpolicynone?jetstreamEnabled=true&jetstreamName=mystream2-ackpolicynone&jetstreamAsync=false&durableName=camel2-ackpolicynone&pullSubscription=false&ackPolicy=none"; from("direct:send") // when running full test suite then send can fail due to nats server setup/teardown diff --git a/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerMaxDeliverIT.java b/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerMaxDeliverIT.java index 53470c80b6fa..cc78a43432a2 100644 --- a/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerMaxDeliverIT.java +++ b/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerMaxDeliverIT.java @@ -42,7 +42,7 @@ public class NatsJetstreamConsumerMaxDeliverIT extends NatsITSupport { mockResultEndpoint.expectedMessageCount(0); mockInputEndpoint.expectedMessageCount(3); - mockInputEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2"); + mockInputEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2-maxdeliver"); mockInputEndpoint.message(0).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(1); mockInputEndpoint.message(1).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(2); mockInputEndpoint.message(2).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(3); @@ -60,7 +60,7 @@ public class NatsJetstreamConsumerMaxDeliverIT extends NatsITSupport { @Override public void configure() { String uri - = "nats:mytopic2?jetstreamEnabled=true&jetstreamName=mystream2&jetstreamAsync=false&durableName=camel2&pullSubscription=false&nackWait=10&maxDeliver=3"; + = "nats:mytopic2-maxdeliver?jetstreamEnabled=true&jetstreamName=mystream2-maxdeliver&jetstreamAsync=false&durableName=camel2-maxdeliver&pullSubscription=false&nackWait=10&maxDeliver=3"; from("direct:send") // when running full test suite then send can fail due to nats server setup/teardown diff --git a/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerRedeliveryIT.java b/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerRedeliveryIT.java index 3690e35f8ed8..cbcbf5371500 100644 --- a/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerRedeliveryIT.java +++ b/components/camel-nats/src/test/java/org/apache/camel/component/nats/jetstream/NatsJetstreamConsumerRedeliveryIT.java @@ -40,12 +40,12 @@ public class NatsJetstreamConsumerRedeliveryIT extends NatsITSupport { @Test public void testConsumer() throws Exception { mockResultEndpoint.expectedBodiesReceived("Hello World"); - mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2"); + mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2-redelivery"); mockResultEndpoint.expectedHeaderReceived("counter", 3); mockResultEndpoint.expectedHeaderReceived(NatsConstants.NATS_DELIVERY_COUNTER, 3); mockInputEndpoint.expectedMessageCount(3); - mockInputEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2"); + mockInputEndpoint.expectedHeaderReceived(NatsConstants.NATS_SUBJECT, "mytopic2-redelivery"); mockInputEndpoint.message(0).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(1); mockInputEndpoint.message(1).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(2); mockInputEndpoint.message(2).header(NatsConstants.NATS_DELIVERY_COUNTER).isEqualTo(3); @@ -63,7 +63,7 @@ public class NatsJetstreamConsumerRedeliveryIT extends NatsITSupport { @Override public void configure() { String uri - = "nats:mytopic2?jetstreamEnabled=true&jetstreamName=mystream2&jetstreamAsync=false&durableName=camel2&pullSubscription=false&nackWait=10"; + = "nats:mytopic2-redelivery?jetstreamEnabled=true&jetstreamName=mystream2-redelivery&jetstreamAsync=false&durableName=camel2-redelivery&pullSubscription=false&nackWait=10"; from("direct:send") // when running full test suite then send can fail due to nats server setup/teardown
