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";