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

jamesnetherton pushed a commit to branch 3.40.x
in repository https://gitbox.apache.org/repos/asf/camel-quarkus.git


The following commit(s) were added to refs/heads/3.40.x by this push:
     new a1088d199b Relates to #9257. Surface Azure Service Bus send failures 
in AzureServiceBusTest
a1088d199b is described below

commit a1088d199be58be1d4eeceb9d140d6d0a0f62789
Author: Jiří Ondrušek <[email protected]>
AuthorDate: Fri Oct 2 12:30:27 2026 +0200

    Relates to #9257. Surface Azure Service Bus send failures in 
AzureServiceBusTest
    
    The send endpoint answered 201 even when the Camel exchange failed, so a
    failed send was indistinguishable from a consumer that never received the
    message. Return 500 with the stack trace and assert it in the test, and put
    the route start / send timeline into the receive timeout error, so the cause
    reaches the CI console even when test output is redirected to a file.
    
    Co-authored-by: Claude Fable 5.1 <[email protected]>
---
 .../servicebus/it/AzureServiceBusResource.java     |  11 +-
 .../azure/servicebus/it/AzureServiceBusTest.java   | 127 +++++++++++----------
 2 files changed, 77 insertions(+), 61 deletions(-)

diff --git 
a/integration-test-groups/azure/azure-servicebus/src/main/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusResource.java
 
b/integration-test-groups/azure/azure-servicebus/src/main/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusResource.java
index 74a3249c2f..d7cd5b24f3 100644
--- 
a/integration-test-groups/azure/azure-servicebus/src/main/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusResource.java
+++ 
b/integration-test-groups/azure/azure-servicebus/src/main/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusResource.java
@@ -16,6 +16,8 @@
  */
 package org.apache.camel.quarkus.component.azure.servicebus.it;
 
+import java.io.PrintWriter;
+import java.io.StringWriter;
 import java.net.URI;
 import java.time.Instant;
 import java.time.OffsetDateTime;
@@ -119,7 +121,7 @@ public class AzureServiceBusResource {
             offsetDateTime = 
OffsetDateTime.ofInstant(Instant.ofEpochMilli(scheduledEnqueueTime), 
ZoneId.systemDefault());
         }
 
-        fluentProducerTemplate.to(directEndpointUri)
+        Exchange result = fluentProducerTemplate.to(directEndpointUri)
                 .withHeader("serviceBusType", serviceBusType)
                 .withHeader("destination", destination)
                 .withHeader("transportType", transportType)
@@ -128,6 +130,13 @@ public class AzureServiceBusResource {
                 .withBody(payload)
                 .send();
 
+        // send() does not throw, the failure is stored on the exchange. 
Report it instead of answering 201
+        if (result.getException() != null) {
+            StringWriter stackTrace = new StringWriter();
+            result.getException().printStackTrace(new PrintWriter(stackTrace));
+            return 
Response.serverError().entity(stackTrace.toString()).build();
+        }
+
         return Response.created(new URI("https://camel.apache.org/";)).build();
     }
 
diff --git 
a/integration-test-groups/azure/azure-servicebus/src/test/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusTest.java
 
b/integration-test-groups/azure/azure-servicebus/src/test/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusTest.java
index f922f7e782..7f0f2a3a9d 100644
--- 
a/integration-test-groups/azure/azure-servicebus/src/test/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusTest.java
+++ 
b/integration-test-groups/azure/azure-servicebus/src/test/java/org/apache/camel/quarkus/component/azure/servicebus/it/AzureServiceBusTest.java
@@ -33,11 +33,16 @@ import io.quarkus.test.common.QuarkusTestResource;
 import io.quarkus.test.junit.QuarkusTest;
 import io.restassured.RestAssured;
 import io.restassured.http.ContentType;
+import io.restassured.response.Response;
+import io.restassured.specification.RequestSpecification;
 import org.apache.camel.quarkus.test.EnabledIf;
 import org.apache.camel.quarkus.test.mock.backend.MockBackendDisabled;
 import org.apache.camel.quarkus.test.support.azure.AzureServiceBusTestResource;
 import org.awaitility.Awaitility;
+import org.awaitility.core.ConditionTimeoutException;
+import org.awaitility.core.ThrowingRunnable;
 import org.jboss.logging.Logger;
+import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Assumptions;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.Test;
@@ -55,6 +60,9 @@ class AzureServiceBusTest {
 
     private static final Logger LOG = 
Logger.getLogger(AzureServiceBusTest.class);
 
+    // Kept at 1 minute on purpose while #9257 is investigated, a longer 
window would hide late deliveries
+    private static final Duration MESSAGE_RECEIVE_TIMEOUT = 
Duration.ofMinutes(1);
+
     @BeforeAll
     public static void beforeAll() {
         //this is not necessary for mocked testing
@@ -100,22 +108,16 @@ class AzureServiceBusTest {
         final String messageBody = UUID.randomUUID().toString();
 
         try {
-            RestAssured.given()
-                    .post("/azure-servicebus/route/" + consumerRouteId + 
"/start")
-                    .then()
-                    .statusCode(204);
+            final Instant routeStarted = startRoute(consumerRouteId);
 
-            RestAssured.given()
+            final Sent sent = send(RestAssured.given()
                     .contentType(ContentType.TEXT)
                     .queryParam("serviceBusType", destinationType)
                     .queryParam("payloadType", payloadType)
                     .queryParam("transportType", transportType.name())
-                    .body(messageBody)
-                    .post("/azure-servicebus/send/message/" + destination)
-                    .then()
-                    .statusCode(201);
+                    .body(messageBody), destination);
 
-            Awaitility.await().pollInterval(1, TimeUnit.SECONDS).atMost(1, 
TimeUnit.MINUTES).untilAsserted(() -> {
+            awaitReceived(routeStarted, sent, () -> {
                 RestAssured.given()
                         .queryParam("endpointUri", mockEndpointUri)
                         .get("/azure-servicebus/receive/messages")
@@ -145,22 +147,16 @@ class AzureServiceBusTest {
             messages.add("cq-azure-servicebus-test-" + 
UUID.randomUUID().toString());
         }
         try {
-            RestAssured.given()
-                    .post("/azure-servicebus/route/" + consumerRouteId + 
"/start")
-                    .then()
-                    .statusCode(204);
+            final Instant routeStarted = startRoute(consumerRouteId);
 
-            RestAssured.given()
+            final Sent sent = send(RestAssured.given()
                     .contentType(ContentType.JSON)
                     .queryParam("serviceBusType", "queue")
                     .queryParam("payloadType", List.class.getSimpleName())
                     .queryParam("transportType", AmqpTransportType.AMQP.name())
-                    .body(messages)
-                    .post("/azure-servicebus/send/message/" + destination)
-                    .then()
-                    .statusCode(201);
+                    .body(messages), destination);
 
-            Awaitility.await().pollInterval(1, TimeUnit.SECONDS).atMost(1, 
TimeUnit.MINUTES).untilAsserted(() -> {
+            awaitReceived(routeStarted, sent, () -> {
                 RestAssured.given()
                         .queryParam("endpointUri", mockEndpointUri)
                         .get("/azure-servicebus/receive/messages")
@@ -186,21 +182,15 @@ class AzureServiceBusTest {
     void produceConsumeWithCustomClients() {
         final String messageBody = UUID.randomUUID().toString();
         try {
-            RestAssured.given()
-                    
.post("/azure-servicebus/route/servicebus-queue-consumer-custom-processor/start")
-                    .then()
-                    .statusCode(204);
+            final Instant routeStarted = 
startRoute("servicebus-queue-consumer-custom-processor");
 
-            RestAssured.given()
+            final Sent sent = send(RestAssured.given()
                     .contentType(ContentType.TEXT)
                     .queryParam("directEndpointUri", 
"direct:send-message-custom-client")
                     .queryParam("payloadType", String.class.getSimpleName())
-                    .body(messageBody)
-                    .post("/azure-servicebus/send/message/" + 
AzureServiceBusHelper.getDestination("queue"))
-                    .then()
-                    .statusCode(201);
+                    .body(messageBody), 
AzureServiceBusHelper.getDestination("queue"));
 
-            Awaitility.await().pollInterval(1, TimeUnit.SECONDS).atMost(1, 
TimeUnit.MINUTES).untilAsserted(() -> {
+            awaitReceived(routeStarted, sent, () -> {
                 RestAssured.given()
                         .queryParam("endpointUri", 
AzureServiceBusProducers.MOCK_ENDPOINT_URI)
                         .get("/azure-servicebus/receive/messages")
@@ -223,21 +213,15 @@ class AzureServiceBusTest {
     void tokenCredentialAuthentication() {
         final String messageBody = UUID.randomUUID().toString();
         try {
-            RestAssured.given()
-                    
.post("/azure-servicebus/route/servicebus-queue-consumer-token-credential/start")
-                    .then()
-                    .statusCode(204);
+            final Instant routeStarted = 
startRoute("servicebus-queue-consumer-token-credential");
 
-            RestAssured.given()
+            final Sent sent = send(RestAssured.given()
                     .contentType(ContentType.TEXT)
                     .queryParam("directEndpointUri", "direct:token-credential")
                     .queryParam("payloadType", String.class.getSimpleName())
-                    .body(messageBody)
-                    .post("/azure-servicebus/send/message/" + 
AzureServiceBusHelper.getDestination("queue"))
-                    .then()
-                    .statusCode(201);
+                    .body(messageBody), 
AzureServiceBusHelper.getDestination("queue"));
 
-            Awaitility.await().pollInterval(1, TimeUnit.SECONDS).atMost(1, 
TimeUnit.MINUTES).untilAsserted(() -> {
+            awaitReceived(routeStarted, sent, () -> {
                 RestAssured.given()
                         .queryParam("endpointUri", 
"mock:servicebus-token-credential-results")
                         .get("/azure-servicebus/receive/messages")
@@ -259,25 +243,19 @@ class AzureServiceBusTest {
     void scheduled() {
         final String messageBody = UUID.randomUUID().toString();
         try {
-            RestAssured.given()
-                    
.post("/azure-servicebus/route/servicebus-queue-scheduled-consumer/start")
-                    .then()
-                    .statusCode(204);
+            final Instant routeStarted = 
startRoute("servicebus-queue-scheduled-consumer");
 
             // Schedule message for 15 seconds in the future
             long scheduledEnqueueTime = Instant.now()
                     .plus(Duration.of(15, ChronoUnit.SECONDS))
                     .toEpochMilli();
 
-            RestAssured.given()
+            final Sent sent = send(RestAssured.given()
                     .contentType(ContentType.TEXT)
                     .queryParam("directEndpointUri", "direct:scheduled")
                     .queryParam("payloadType", String.class.getSimpleName())
                     .queryParam("scheduledEnqueueTime", scheduledEnqueueTime)
-                    .body(messageBody)
-                    .post("/azure-servicebus/send/message/" + 
AzureServiceBusHelper.getDestination("queue"))
-                    .then()
-                    .statusCode(201);
+                    .body(messageBody), 
AzureServiceBusHelper.getDestination("queue"));
 
             //we are checking that there is no message in the next 10 seconds.
             //we should rather avoid checking the message near the scheduled 
time, because process can be delayed and receives the message
@@ -298,7 +276,7 @@ class AzureServiceBusTest {
             }
 
             // Message should be enqueued and eventually consumed
-            Awaitility.await().pollInterval(1, TimeUnit.SECONDS).atMost(1, 
TimeUnit.MINUTES).untilAsserted(() -> {
+            awaitReceived(routeStarted, sent, () -> {
                 RestAssured.given()
                         .queryParam("endpointUri", 
"mock:servicebus-queue-scheduled-consumer-results")
                         .get("/azure-servicebus/receive/messages")
@@ -323,21 +301,15 @@ class AzureServiceBusTest {
 
         final String messageBody = UUID.randomUUID().toString();
         try {
-            RestAssured.given()
-                    
.post("/azure-servicebus/route/servicebus-queue-consumer-azure-identity/start")
-                    .then()
-                    .statusCode(204);
+            final Instant routeStarted = 
startRoute("servicebus-queue-consumer-azure-identity");
 
-            RestAssured.given()
+            final Sent sent = send(RestAssured.given()
                     .contentType(ContentType.TEXT)
                     .queryParam("directEndpointUri", "direct:azure-identity")
                     .queryParam("payloadType", String.class.getSimpleName())
-                    .body(messageBody)
-                    .post("/azure-servicebus/send/message/" + 
AzureServiceBusHelper.getDestination("queue"))
-                    .then()
-                    .statusCode(201);
+                    .body(messageBody), 
AzureServiceBusHelper.getDestination("queue"));
 
-            Awaitility.await().pollInterval(1, TimeUnit.SECONDS).atMost(1, 
TimeUnit.MINUTES).untilAsserted(() -> {
+            awaitReceived(routeStarted, sent, () -> {
                 RestAssured.given()
                         .queryParam("endpointUri", 
"mock:servicebus-azure-identity-results")
                         .get("/azure-servicebus/receive/messages")
@@ -363,6 +335,41 @@ class AzureServiceBusTest {
                 .body(is("true"));
     }
 
+    private static Instant startRoute(String routeId) {
+        RestAssured.given()
+                .post("/azure-servicebus/route/" + routeId + "/start")
+                .then()
+                .statusCode(204);
+        return Instant.now();
+    }
+
+    // Fails with the server side stack trace when the send did not succeed, 
so the cause reaches the CI console
+    private static Sent send(RequestSpecification request, String destination) 
{
+        final Instant start = Instant.now();
+        final Response response = 
request.post("/azure-servicebus/send/message/" + destination);
+        final long millis = Duration.between(start, Instant.now()).toMillis();
+        Assertions.assertEquals(201, response.statusCode(),
+                () -> "Sending the message failed after " + millis + " ms:\n" 
+ response.asString());
+        return new Sent(Instant.now(), millis);
+    }
+
+    // Puts the route start / send timeline into the timeout error, the only 
text that reaches a CI console when
+    // test output is redirected to a file
+    private static void awaitReceived(Instant routeStarted, Sent sent, 
ThrowingRunnable assertion) {
+        try {
+            Awaitility.await().pollInterval(1, 
TimeUnit.SECONDS).atMost(MESSAGE_RECEIVE_TIMEOUT).untilAsserted(assertion);
+        } catch (ConditionTimeoutException e) {
+            throw new AssertionError(
+                    "No message received within %s. The consumer route was 
started %d ms before the send completed, the send took %d ms. %s"
+                            .formatted(MESSAGE_RECEIVE_TIMEOUT, 
Duration.between(routeStarted, sent.at()).toMillis(),
+                                    sent.millis(), e.getMessage()),
+                    e);
+        }
+    }
+
+    private record Sent(Instant at, long millis) {
+    }
+
     static Stream<Arguments> produceConsumeOptions() {
         String destinationTypes = "queue";
 

Reply via email to