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<>();
