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

davsclaus pushed a commit to branch fix/CAMEL-24158-backport-4.14.x
in repository https://gitbox.apache.org/repos/asf/camel.git

commit fbbbd173e44d6007c0ee402e62af1d9cd1c1ad60
Author: Claus Ibsen <[email protected]>
AuthorDate: Sat Jul 18 10:11:54 2026 +0200

    CAMEL-24158: camel-azure-eventhubs - Add HeaderFilterStrategy to prevent 
Camel headers leaking to AMQP properties
    
    Backport of header filter strategy fix from PR #24884 to 4.14.x.
    EventHubsComponent now extends HeaderFilterStrategyComponent so that
    the default DefaultHeaderFilterStrategy filters out Camel-internal
    headers before they are put on EventData application properties.
    Also guards against setting a null PARTITION_KEY header on the consumer 
side.
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../eventhubs/EventHubsComponentConfigurer.java    |  6 ++++
 .../component/azure/eventhubs/azure-eventhubs.json | 11 ++++---
 .../azure/eventhubs/EventHubsComponent.java        |  4 +--
 .../azure/eventhubs/EventHubsConsumer.java         |  4 ++-
 .../azure/eventhubs/EventHubsProducer.java         |  5 ++-
 .../operations/EventHubsProducerOperations.java    | 36 +++++++++++++++-------
 .../operations/EventHubsProducerOperationsIT.java  |  9 ++++--
 7 files changed, 52 insertions(+), 23 deletions(-)

diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/generated/java/org/apache/camel/component/azure/eventhubs/EventHubsComponentConfigurer.java
 
b/components/camel-azure/camel-azure-eventhubs/src/generated/java/org/apache/camel/component/azure/eventhubs/EventHubsComponentConfigurer.java
index 4cf39976fd6c..809c16d4fac4 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/generated/java/org/apache/camel/component/azure/eventhubs/EventHubsComponentConfigurer.java
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/generated/java/org/apache/camel/component/azure/eventhubs/EventHubsComponentConfigurer.java
@@ -61,6 +61,8 @@ public class EventHubsComponentConfigurer extends 
PropertyConfigurerSupport impl
         case "credentialType": 
getOrCreateConfiguration(target).setCredentialType(property(camelContext, 
org.apache.camel.component.azure.eventhubs.CredentialType.class, value)); 
return true;
         case "eventposition":
         case "eventPosition": 
getOrCreateConfiguration(target).setEventPosition(property(camelContext, 
java.util.Map.class, value)); return true;
+        case "headerfilterstrategy":
+        case "headerFilterStrategy": 
target.setHeaderFilterStrategy(property(camelContext, 
org.apache.camel.spi.HeaderFilterStrategy.class, value)); return true;
         case "lazystartproducer":
         case "lazyStartProducer": 
target.setLazyStartProducer(property(camelContext, boolean.class, value)); 
return true;
         case "partitionid":
@@ -120,6 +122,8 @@ public class EventHubsComponentConfigurer extends 
PropertyConfigurerSupport impl
         case "credentialType": return 
org.apache.camel.component.azure.eventhubs.CredentialType.class;
         case "eventposition":
         case "eventPosition": return java.util.Map.class;
+        case "headerfilterstrategy":
+        case "headerFilterStrategy": return 
org.apache.camel.spi.HeaderFilterStrategy.class;
         case "lazystartproducer":
         case "lazyStartProducer": return boolean.class;
         case "partitionid":
@@ -175,6 +179,8 @@ public class EventHubsComponentConfigurer extends 
PropertyConfigurerSupport impl
         case "credentialType": return 
getOrCreateConfiguration(target).getCredentialType();
         case "eventposition":
         case "eventPosition": return 
getOrCreateConfiguration(target).getEventPosition();
+        case "headerfilterstrategy":
+        case "headerFilterStrategy": return target.getHeaderFilterStrategy();
         case "lazystartproducer":
         case "lazyStartProducer": return target.isLazyStartProducer();
         case "partitionid":
diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/generated/resources/META-INF/org/apache/camel/component/azure/eventhubs/azure-eventhubs.json
 
b/components/camel-azure/camel-azure-eventhubs/src/generated/resources/META-INF/org/apache/camel/component/azure/eventhubs/azure-eventhubs.json
index 5dc538cdf023..2708df0f2347 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/generated/resources/META-INF/org/apache/camel/component/azure/eventhubs/azure-eventhubs.json
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/generated/resources/META-INF/org/apache/camel/component/azure/eventhubs/azure-eventhubs.json
@@ -43,11 +43,12 @@
     "partitionKey": { "index": 16, "kind": "property", "displayName": 
"Partition Key", "group": "producer", "label": "producer", "required": false, 
"type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": false, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "Sets a hashing key to be 
provided for the batch of events, which instructs the Event Hubs  [...]
     "producerAsyncClient": { "index": 17, "kind": "property", "displayName": 
"Producer Async Client", "group": "producer", "label": "producer", "required": 
false, "type": "object", "javaType": 
"com.azure.messaging.eventhubs.EventHubProducerAsyncClient", "deprecated": 
false, "deprecationNote": "", "autowired": true, "secret": false, 
"configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "Sets the Eve [...]
     "autowiredEnabled": { "index": 18, "kind": "property", "displayName": 
"Autowired Enabled", "group": "advanced", "label": "advanced", "required": 
false, "type": "boolean", "javaType": "boolean", "deprecated": false, 
"autowired": false, "secret": false, "defaultValue": true, "description": 
"Whether autowiring is enabled. This is used for automatic autowiring options 
(the option must be marked as autowired) by looking up in the registry to find 
if there is a single instance of matching  [...]
-    "connectionString": { "index": 19, "kind": "property", "displayName": 
"Connection String", "group": "security", "label": "security", "required": 
false, "type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": true, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "Instead of supplying 
namespace, sharedAccessKey, sharedAccessName, etc. you can sup [...]
-    "credentialType": { "index": 20, "kind": "property", "displayName": 
"Credential Type", "group": "security", "label": "security", "required": false, 
"type": "enum", "javaType": 
"org.apache.camel.component.azure.eventhubs.CredentialType", "enum": [ 
"AZURE_IDENTITY", "CONNECTION_STRING", "TOKEN_CREDENTIAL" ], "deprecated": 
false, "autowired": false, "secret": false, "defaultValue": 
"CONNECTION_STRING", "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfigurat [...]
-    "sharedAccessKey": { "index": 21, "kind": "property", "displayName": 
"Shared Access Key", "group": "security", "label": "security", "required": 
false, "type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": true, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "The generated value for 
the SharedAccessName." },
-    "sharedAccessName": { "index": 22, "kind": "property", "displayName": 
"Shared Access Name", "group": "security", "label": "security", "required": 
false, "type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": false, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "The name you chose for 
your EventHubs SAS keys." },
-    "tokenCredential": { "index": 23, "kind": "property", "displayName": 
"Token Credential", "group": "security", "label": "security", "required": 
false, "type": "object", "javaType": 
"com.azure.core.credential.TokenCredential", "deprecated": false, 
"deprecationNote": "", "autowired": true, "secret": true, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "Provide custom 
authentication credenti [...]
+    "headerFilterStrategy": { "index": 19, "kind": "property", "displayName": 
"Header Filter Strategy", "group": "filter", "label": "filter", "required": 
false, "type": "object", "javaType": 
"org.apache.camel.spi.HeaderFilterStrategy", "deprecated": false, "autowired": 
false, "secret": false, "description": "To use a custom 
org.apache.camel.spi.HeaderFilterStrategy to filter header to and from Camel 
message." },
+    "connectionString": { "index": 20, "kind": "property", "displayName": 
"Connection String", "group": "security", "label": "security", "required": 
false, "type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": true, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "Instead of supplying 
namespace, sharedAccessKey, sharedAccessName, etc. you can sup [...]
+    "credentialType": { "index": 21, "kind": "property", "displayName": 
"Credential Type", "group": "security", "label": "security", "required": false, 
"type": "enum", "javaType": 
"org.apache.camel.component.azure.eventhubs.CredentialType", "enum": [ 
"AZURE_IDENTITY", "CONNECTION_STRING", "TOKEN_CREDENTIAL" ], "deprecated": 
false, "autowired": false, "secret": false, "defaultValue": 
"CONNECTION_STRING", "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfigurat [...]
+    "sharedAccessKey": { "index": 22, "kind": "property", "displayName": 
"Shared Access Key", "group": "security", "label": "security", "required": 
false, "type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": true, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "The generated value for 
the SharedAccessName." },
+    "sharedAccessName": { "index": 23, "kind": "property", "displayName": 
"Shared Access Name", "group": "security", "label": "security", "required": 
false, "type": "string", "javaType": "java.lang.String", "deprecated": false, 
"autowired": false, "secret": false, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "The name you chose for 
your EventHubs SAS keys." },
+    "tokenCredential": { "index": 24, "kind": "property", "displayName": 
"Token Credential", "group": "security", "label": "security", "required": 
false, "type": "object", "javaType": 
"com.azure.core.credential.TokenCredential", "deprecated": false, 
"deprecationNote": "", "autowired": true, "secret": true, "configurationClass": 
"org.apache.camel.component.azure.eventhubs.EventHubsConfiguration", 
"configurationField": "configuration", "description": "Provide custom 
authentication credenti [...]
   },
   "headers": {
     "CamelAzureEventHubsPartitionKey": { "index": 0, "kind": "header", 
"displayName": "", "group": "common", "label": "", "required": false, 
"javaType": "String", "deprecated": false, "deprecationNote": "", "autowired": 
false, "secret": false, "description": "(producer) Overrides the hashing key to 
be provided for the batch of events, which instructs the Event Hubs service to 
map this key to a specific partition. (consumer) It sets the partition hashing 
key if it was set when originally  [...]
diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java
 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java
index 04fd5b3de54e..8a18befff120 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsComponent.java
@@ -22,7 +22,7 @@ import com.azure.identity.DefaultAzureCredential;
 import org.apache.camel.Endpoint;
 import org.apache.camel.spi.Metadata;
 import org.apache.camel.spi.annotations.Component;
-import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.HeaderFilterStrategyComponent;
 import org.apache.camel.util.ObjectHelper;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -31,7 +31,7 @@ import org.slf4j.LoggerFactory;
  * Azure EventHubs component
  */
 @Component("azure-eventhubs")
-public class EventHubsComponent extends DefaultComponent {
+public class EventHubsComponent extends HeaderFilterStrategyComponent {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(EventHubsComponent.class);
 
diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java
 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java
index af425ddee6cc..7e8378e8c3b8 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsConsumer.java
@@ -95,7 +95,9 @@ public class EventHubsConsumer extends DefaultConsumer {
         message.setBody(eventContext.getEventData().getBody());
         // set headers
         message.setHeader(EventHubsConstants.PARTITION_ID, 
eventContext.getPartitionContext().getPartitionId());
-        message.setHeader(EventHubsConstants.PARTITION_KEY, 
eventContext.getEventData().getPartitionKey());
+        if (eventContext.getEventData().getPartitionKey() != null) {
+            message.setHeader(EventHubsConstants.PARTITION_KEY, 
eventContext.getEventData().getPartitionKey());
+        }
         message.setHeader(EventHubsConstants.OFFSET, 
eventContext.getEventData().getOffset());
         message.setHeader(EventHubsConstants.ENQUEUED_TIME, 
eventContext.getEventData().getEnqueuedTime());
         message.setHeader(EventHubsConstants.SEQUENCE_NUMBER, 
eventContext.getEventData().getSequenceNumber());
diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java
 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java
index 861348cf52b5..a10f9348a2db 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/EventHubsProducer.java
@@ -22,6 +22,7 @@ import org.apache.camel.Endpoint;
 import org.apache.camel.Exchange;
 import 
org.apache.camel.component.azure.eventhubs.client.EventHubsClientFactory;
 import 
org.apache.camel.component.azure.eventhubs.operations.EventHubsProducerOperations;
+import org.apache.camel.spi.HeaderFilterStrategy;
 import org.apache.camel.support.DefaultAsyncProducer;
 
 public class EventHubsProducer extends DefaultAsyncProducer {
@@ -45,7 +46,9 @@ public class EventHubsProducer extends DefaultAsyncProducer {
         }
 
         // create our operations
-        producerOperations = new 
EventHubsProducerOperations(producerAsyncClient, getConfiguration());
+        HeaderFilterStrategy headerFilterStrategy
+                = ((EventHubsComponent) 
getEndpoint().getComponent()).getHeaderFilterStrategy();
+        producerOperations = new 
EventHubsProducerOperations(producerAsyncClient, getConfiguration(), 
headerFilterStrategy);
     }
 
     @Override
diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/operations/EventHubsProducerOperations.java
 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/operations/EventHubsProducerOperations.java
index 23498d29d0e0..8d70509b23b8 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/operations/EventHubsProducerOperations.java
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/main/java/org/apache/camel/component/azure/eventhubs/operations/EventHubsProducerOperations.java
@@ -30,6 +30,7 @@ import org.apache.camel.Message;
 import org.apache.camel.TypeConverter;
 import org.apache.camel.component.azure.eventhubs.EventHubsConfiguration;
 import 
org.apache.camel.component.azure.eventhubs.EventHubsConfigurationOptionsProxy;
+import org.apache.camel.spi.HeaderFilterStrategy;
 import org.apache.camel.util.ObjectHelper;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -41,12 +42,15 @@ public class EventHubsProducerOperations {
 
     private final EventHubProducerAsyncClient producerAsyncClient;
     private final EventHubsConfigurationOptionsProxy configurationOptionsProxy;
+    private final HeaderFilterStrategy headerFilterStrategy;
 
     public EventHubsProducerOperations(final EventHubProducerAsyncClient 
producerAsyncClient,
-                                       final EventHubsConfiguration 
configuration) {
+                                       final EventHubsConfiguration 
configuration,
+                                       final HeaderFilterStrategy 
headerFilterStrategy) {
         ObjectHelper.notNull(producerAsyncClient, "client cannot be null");
 
         this.producerAsyncClient = producerAsyncClient;
+        this.headerFilterStrategy = headerFilterStrategy;
         configurationOptionsProxy = new 
EventHubsConfigurationOptionsProxy(configuration);
     }
 
@@ -107,7 +111,7 @@ public class EventHubsProducerOperations {
         // check if our exchange is list or contain some values
         if (exchange.getIn().getBody() instanceof Iterable) {
             return createEventDataFromIterable((Iterable<Object>) 
exchange.getIn().getBody(),
-                    exchange.getContext().getTypeConverter(), 
exchange.getIn().getHeaders());
+                    exchange.getContext().getTypeConverter(), 
exchange.getIn().getHeaders(), exchange);
         }
 
         // we have only a single event here
@@ -115,16 +119,17 @@ public class EventHubsProducerOperations {
     }
 
     private Iterable<EventData> createEventDataFromIterable(
-            final Iterable<Object> inputData, final TypeConverter converter, 
Map<String, Object> headers) {
+            final Iterable<Object> inputData, final TypeConverter converter,
+            Map<String, Object> headers, final Exchange exchange) {
         final List<EventData> finalEventData = new LinkedList<>();
 
         inputData.forEach(data -> {
-            if (data instanceof Exchange) {
-                finalEventData.add(createEventDataFromExchange((Exchange) 
data));
-            } else if (data instanceof Message) {
-                finalEventData.add(createEventDataFromMessage((Message) data));
+            if (data instanceof Exchange ex) {
+                finalEventData.add(createEventDataFromExchange(ex));
+            } else if (data instanceof Message msg) {
+                finalEventData.add(createEventDataFromMessage(msg));
             } else {
-                finalEventData.add(createEventDataFromObject(data, converter, 
headers));
+                finalEventData.add(createEventDataFromObject(data, converter, 
headers, exchange));
             }
         });
 
@@ -137,11 +142,12 @@ public class EventHubsProducerOperations {
 
     private EventData createEventDataFromMessage(final Message message) {
         return createEventDataFromObject(message.getBody(), 
message.getExchange().getContext().getTypeConverter(),
-                message.getHeaders());
+                message.getHeaders(), message.getExchange());
     }
 
     private EventData createEventDataFromObject(
-            final Object inputData, final TypeConverter converter, Map<String, 
Object> headers) {
+            final Object inputData, final TypeConverter converter,
+            Map<String, Object> headers, final Exchange exchange) {
         final byte[] data = converter.convertTo(byte[].class, inputData);
 
         if (ObjectHelper.isEmpty(data)) {
@@ -152,7 +158,15 @@ public class EventHubsProducerOperations {
         }
 
         EventData eventData = new EventData(data);
-        eventData.getProperties().putAll(headers);
+        if (headerFilterStrategy != null) {
+            headers.forEach((key, value) -> {
+                if (!headerFilterStrategy.applyFilterToCamelHeaders(key, 
value, exchange)) {
+                    eventData.getProperties().put(key, value);
+                }
+            });
+        } else {
+            eventData.getProperties().putAll(headers);
+        }
 
         return eventData;
     }
diff --git 
a/components/camel-azure/camel-azure-eventhubs/src/test/java/org/apache/camel/component/azure/eventhubs/integration/operations/EventHubsProducerOperationsIT.java
 
b/components/camel-azure/camel-azure-eventhubs/src/test/java/org/apache/camel/component/azure/eventhubs/integration/operations/EventHubsProducerOperationsIT.java
index a152acdcd9b1..71ea6d119ae1 100644
--- 
a/components/camel-azure/camel-azure-eventhubs/src/test/java/org/apache/camel/component/azure/eventhubs/integration/operations/EventHubsProducerOperationsIT.java
+++ 
b/components/camel-azure/camel-azure-eventhubs/src/test/java/org/apache/camel/component/azure/eventhubs/integration/operations/EventHubsProducerOperationsIT.java
@@ -64,7 +64,8 @@ class EventHubsProducerOperationsIT extends CamelTestSupport {
 
     @Test
     public void testSendEventWithSpecificPartition() {
-        final EventHubsProducerOperations operations = new 
EventHubsProducerOperations(producerAsyncClient, configuration);
+        final EventHubsProducerOperations operations
+                = new EventHubsProducerOperations(producerAsyncClient, 
configuration, null);
         final String firstPartition = 
producerAsyncClient.getPartitionIds().blockLast();
         final Exchange exchange = new DefaultExchange(context);
 
@@ -96,7 +97,8 @@ class EventHubsProducerOperationsIT extends CamelTestSupport {
 
     @Test
     public void testIterableExchangesSendEventsWithSpecificPartition() {
-        final EventHubsProducerOperations operations = new 
EventHubsProducerOperations(producerAsyncClient, configuration);
+        final EventHubsProducerOperations operations
+                = new EventHubsProducerOperations(producerAsyncClient, 
configuration, null);
         final String firstPartition = 
producerAsyncClient.getPartitionIds().blockLast();
 
         final Exchange exchange1 = new DefaultExchange(context);
@@ -144,7 +146,8 @@ class EventHubsProducerOperationsIT extends 
CamelTestSupport {
 
     @Test
     public void testIterableStringSendEventsWithSpecificPartition() {
-        final EventHubsProducerOperations operations = new 
EventHubsProducerOperations(producerAsyncClient, configuration);
+        final EventHubsProducerOperations operations
+                = new EventHubsProducerOperations(producerAsyncClient, 
configuration, null);
         final String firstPartition = 
producerAsyncClient.getPartitionIds().blockLast();
 
         final List<String> messages = new LinkedList<>();

Reply via email to