This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-24159-servicebus in repository https://gitbox.apache.org/repos/asf/camel.git
commit c38b05527168eb5ecaa660844324a22f30b600fa Author: Claus Ibsen <[email protected]> AuthorDate: Fri Jul 17 19:47:17 2026 +0200 CAMEL-24159: camel-azure-servicebus - fix medium-severity bugs from code review Co-Authored-By: Claude Opus 4.6 <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../azure/servicebus/ServiceBusComponent.java | 8 +++--- .../azure/servicebus/ServiceBusConsumer.java | 24 ++++++++++++++---- .../azure/servicebus/ServiceBusProducer.java | 19 ++++++++++---- .../ServicebusCloudEventDataTypeTransformer.java | 2 +- .../azure/servicebus/ServiceBusConsumerTest.java | 29 +++++++++++++++++++--- 5 files changed, 64 insertions(+), 18 deletions(-) diff --git a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusComponent.java b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusComponent.java index 42a9f930b3f0..9dd252b7f289 100644 --- a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusComponent.java +++ b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusComponent.java @@ -62,9 +62,11 @@ public class ServiceBusComponent extends DefaultComponent { endpoint.getConfiguration().setCredentialType(CredentialType.CONNECTION_STRING); } } else { - boolean azure = endpoint.getConfiguration().getTokenCredential() instanceof DefaultAzureCredential; - endpoint.getConfiguration() - .setCredentialType(azure ? CredentialType.AZURE_IDENTITY : CredentialType.TOKEN_CREDENTIAL); + if (endpoint.getConfiguration().getCredentialType() == null) { + boolean azure = endpoint.getConfiguration().getTokenCredential() instanceof DefaultAzureCredential; + endpoint.getConfiguration() + .setCredentialType(azure ? CredentialType.AZURE_IDENTITY : CredentialType.TOKEN_CREDENTIAL); + } } return endpoint; diff --git a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumer.java b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumer.java index ffe1c1af1c83..97677316e5b2 100644 --- a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumer.java +++ b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumer.java @@ -55,6 +55,7 @@ public class ServiceBusConsumer extends DefaultConsumer implements ShutdownAware private static final int LOCK_RENEW_INTERVAL_SECONDS = 10; private ServiceBusProcessorClient client; + private boolean clientCloseable; private ServiceBusReceiverAsyncClient renewalClient; private LockRenewer lockRenewer; private final AtomicInteger pendingExchanges = new AtomicInteger(); @@ -75,8 +76,15 @@ public class ServiceBusConsumer extends DefaultConsumer implements ShutdownAware LOG.debug("Creating connection to Azure ServiceBus"); + // close any previously created client (e.g. on route restart after doStop) + if (clientCloseable && client != null) { + client.close(); + client = null; + } + client = getConfiguration().getProcessorClient(); if (client == null) { + clientCloseable = true; // create client as per sessions if (Boolean.FALSE.equals(getConfiguration().isSessionEnabled())) { client = getEndpoint().getServiceBusClientFactory().createServiceBusProcessorClient(getConfiguration(), @@ -85,6 +93,8 @@ public class ServiceBusConsumer extends DefaultConsumer implements ShutdownAware client = getEndpoint().getServiceBusClientFactory().createServiceBusSessionProcessorClient(getConfiguration(), this::processMessage, this::processError); } + } else { + clientCloseable = false; } client.start(); @@ -144,9 +154,13 @@ public class ServiceBusConsumer extends DefaultConsumer implements ShutdownAware renewalClient = null; } if (client != null) { - // stop accepting new messages but keep the connection open - // so that in-flight exchanges can still complete/abandon messages - client.stop(); + if (clientCloseable) { + client.close(); + client = null; + } else { + // user-provided client: stop accepting messages but don't close + client.stop(); + } } // shutdown camel consumer @@ -164,9 +178,9 @@ public class ServiceBusConsumer extends DefaultConsumer implements ShutdownAware renewalClient.close(); renewalClient = null; } - if (client != null) { - // close the client after all in-flight exchanges have completed + if (clientCloseable && client != null) { client.close(); + client = null; } super.doShutdown(); diff --git a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusProducer.java b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusProducer.java index 48794e7e5171..9cc2a4967f1d 100644 --- a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusProducer.java +++ b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/ServiceBusProducer.java @@ -43,6 +43,7 @@ public class ServiceBusProducer extends DefaultProducer { private final Map<ServiceBusProducerOperationDefinition, Consumer<Exchange>> operationsToExecute = new EnumMap<>(ServiceBusProducerOperationDefinition.class); private ServiceBusSenderClient client; + private boolean clientCloseable; private ServiceBusConfigurationOptionsProxy configurationOptionsProxy; private ServiceBusSenderOperations serviceBusSenderOperations; @@ -67,9 +68,13 @@ public class ServiceBusProducer extends DefaultProducer { super.doStart(); // create the senderClient - client = getConfiguration().getSenderClient() != null - ? getConfiguration().getSenderClient() - : getEndpoint().getServiceBusClientFactory().createServiceBusSenderClient(getConfiguration()); + if (getConfiguration().getSenderClient() != null) { + client = getConfiguration().getSenderClient(); + clientCloseable = false; + } else { + client = getEndpoint().getServiceBusClientFactory().createServiceBusSenderClient(getConfiguration()); + clientCloseable = true; + } // create the operations serviceBusSenderOperations = new ServiceBusSenderOperations(client); @@ -99,9 +104,9 @@ public class ServiceBusProducer extends DefaultProducer { @Override protected void doStop() throws Exception { - if (client != null) { - // shutdown client + if (clientCloseable && client != null) { client.close(); + client = null; } super.doStop(); @@ -128,6 +133,8 @@ public class ServiceBusProducer extends DefaultProducer { = exchange.getMessage().getHeader(ServiceBusConstants.APPLICATION_PROPERTIES, Map.class); if (applicationProperties == null) { applicationProperties = new HashMap<>(); + } else { + applicationProperties = new HashMap<>(applicationProperties); } propagateHeaders(exchange, applicationProperties); final String correlationId = exchange.getMessage().getHeader(ServiceBusConstants.CORRELATION_ID, String.class); @@ -159,6 +166,8 @@ public class ServiceBusProducer extends DefaultProducer { = exchange.getMessage().getHeader(ServiceBusConstants.APPLICATION_PROPERTIES, Map.class); if (applicationProperties == null) { applicationProperties = new HashMap<>(); + } else { + applicationProperties = new HashMap<>(applicationProperties); } propagateHeaders(exchange, applicationProperties); final String correlationId = exchange.getMessage().getHeader(ServiceBusConstants.CORRELATION_ID, String.class); diff --git a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/transform/ServicebusCloudEventDataTypeTransformer.java b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/transform/ServicebusCloudEventDataTypeTransformer.java index b8a3e0e397ba..fa02481c5323 100644 --- a/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/transform/ServicebusCloudEventDataTypeTransformer.java +++ b/components/camel-azure/camel-azure-servicebus/src/main/java/org/apache/camel/component/azure/servicebus/transform/ServicebusCloudEventDataTypeTransformer.java @@ -53,7 +53,7 @@ public class ServicebusCloudEventDataTypeTransformer extends Transformer { headers.put(CloudEvent.CAMEL_CLOUD_EVENT_TIME, cloudEvent.getEventTime(message.getExchange())); if (message.getHeaders().containsKey(ServiceBusConstants.CONTENT_TYPE)) { headers.put(CloudEvent.CAMEL_CLOUD_EVENT_CONTENT_TYPE, - message.getHeaders().containsKey(ServiceBusConstants.CONTENT_TYPE)); + message.getHeader(ServiceBusConstants.CONTENT_TYPE)); } else { headers.put(CloudEvent.CAMEL_CLOUD_EVENT_CONTENT_TYPE, CloudEvent.APPLICATION_OCTET_STREAM_MIME_TYPE); } diff --git a/components/camel-azure/camel-azure-servicebus/src/test/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumerTest.java b/components/camel-azure/camel-azure-servicebus/src/test/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumerTest.java index 7153a3325c1f..d63ee50844e1 100644 --- a/components/camel-azure/camel-azure-servicebus/src/test/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumerTest.java +++ b/components/camel-azure/camel-azure-servicebus/src/test/java/org/apache/camel/component/azure/servicebus/ServiceBusConsumerTest.java @@ -612,8 +612,7 @@ public class ServiceBusConsumerTest { consumer.doStop(); verify(renewalClient).close(); - verify(client).stop(); - verify(client, never()).close(); + verify(client).close(); } @Test @@ -703,7 +702,18 @@ public class ServiceBusConsumerTest { } @Test - void doStopStopsClientWithoutClosing() throws Exception { + void doStopClosesInternallyCreatedClient() throws Exception { + ServiceBusConsumer consumer = new ServiceBusConsumer(endpoint, processor); + consumer.doStart(); + + consumer.doStop(); + + verify(client).close(); + } + + @Test + void doStopStopsUserProvidedClientWithoutClosing() throws Exception { + when(configuration.getProcessorClient()).thenReturn(client); ServiceBusConsumer consumer = new ServiceBusConsumer(endpoint, processor); consumer.doStart(); @@ -714,7 +724,7 @@ public class ServiceBusConsumerTest { } @Test - void doShutdownClosesClient() throws Exception { + void doShutdownClosesInternallyCreatedClient() throws Exception { ServiceBusConsumer consumer = new ServiceBusConsumer(endpoint, processor); consumer.doStart(); @@ -723,6 +733,17 @@ public class ServiceBusConsumerTest { verify(client).close(); } + @Test + void doShutdownDoesNotCloseUserProvidedClient() throws Exception { + when(configuration.getProcessorClient()).thenReturn(client); + ServiceBusConsumer consumer = new ServiceBusConsumer(endpoint, processor); + consumer.doStart(); + + consumer.doShutdown(); + + verify(client, never()).close(); + } + private void configureMockMessage() { when(message.getApplicationProperties()).thenReturn(new HashMap<>()); when(message.getBody()).thenReturn(BinaryData.fromBytes(MESSAGE_BODY.getBytes()));
