This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new e0cb3aa9d649 CAMEL-25090: camel-core - Event notifiers and vault
configuration: fix bugs found in a deep review (#26991)
e0cb3aa9d649 is described below
commit e0cb3aa9d649846efd6de552809fa47295fb23bc
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 23:20:40 2026 +0200
CAMEL-25090: camel-core - Event notifiers and vault configuration: fix bugs
found in a deep review (#26991)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---
.../impl/engine/DefaultManagementStrategy.java | 16 +-
.../impl/event/EventNotifierEdgeCasesTest.java | 179 +++++++++++++++++++++
.../org/apache/camel/main/BaseMainSupport.java | 22 +--
.../camel/main/VaultConfigurationProperties.java | 120 +++++++++++++-
.../apache/camel/main/MainVaultMultipleTest.java | 62 +++++++
.../management/ManagedCamelContextRestartTest.java | 6 +-
.../java/org/apache/camel/support/EventHelper.java | 24 ++-
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 13 ++
8 files changed, 421 insertions(+), 21 deletions(-)
diff --git
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultManagementStrategy.java
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultManagementStrategy.java
index d203f1d0130e..b0acb3cbd950 100644
---
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultManagementStrategy.java
+++
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultManagementStrategy.java
@@ -81,18 +81,21 @@ public class DefaultManagementStrategy extends
ServiceSupport implements Managem
@Override
public void addEventNotifier(EventNotifier eventNotifier) {
+ if (getCamelContext() != null) {
+ // inject camel context if needed
+ CamelContextAware.trySetCamelContext(eventNotifier,
getCamelContext());
+ }
this.eventNotifiers.add(eventNotifier);
// resort after adding
this.eventNotifiers.sort(OrderedComparator.get());
if (isStarted()) {
- // already started
+ // already started so the notifier must be started as well
+ ServiceHelper.startService(eventNotifier);
this.startedEventNotifiers.add(eventNotifier);
// resort after adding
this.startedEventNotifiers.sort(OrderedComparator.get());
}
if (getCamelContext() != null) {
- // inject camel context if needed
- CamelContextAware.trySetCamelContext(eventNotifier,
getCamelContext());
// okay we have an event notifier that accepts exchange events so
its applicable
if (!eventNotifier.isIgnoreExchangeEvents()) {
getCamelContext().getCamelContextExtension().setEventNotificationApplicable(true);
@@ -103,7 +106,12 @@ public class DefaultManagementStrategy extends
ServiceSupport implements Managem
@Override
public boolean removeEventNotifier(EventNotifier eventNotifier) {
startedEventNotifiers.remove(eventNotifier);
- return eventNotifiers.remove(eventNotifier);
+ boolean removed = eventNotifiers.remove(eventNotifier);
+ if (removed && isStarted()) {
+ // the notifier was started when added (or when this was started)
so stop it
+ ServiceHelper.stopService(eventNotifier);
+ }
+ return removed;
}
@Override
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/event/EventNotifierEdgeCasesTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/event/EventNotifierEdgeCasesTest.java
new file mode 100644
index 000000000000..a431703e225c
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/event/EventNotifierEdgeCasesTest.java
@@ -0,0 +1,179 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.impl.event;
+
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.spi.CamelEvent;
+import org.apache.camel.support.EventNotifierSupport;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class EventNotifierEdgeCasesTest extends ContextTestSupport {
+
+ private static class MyNotifier extends EventNotifierSupport {
+ final List<CamelEvent> events = new CopyOnWriteArrayList<>();
+ final AtomicInteger starts = new AtomicInteger();
+ final AtomicInteger stops = new AtomicInteger();
+
+ @Override
+ public void notify(CamelEvent event) {
+ events.add(event);
+ }
+
+ @Override
+ protected void doStart() throws Exception {
+ starts.incrementAndGet();
+ }
+
+ @Override
+ protected void doStop() throws Exception {
+ stops.incrementAndGet();
+ }
+
+ long count(Class<?> type) {
+ return events.stream().filter(type::isInstance).count();
+ }
+ }
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ return new DefaultCamelContext(createCamelRegistry());
+ }
+
+ @Test
+ public void testExceptionFromIsEnabledDoesNotAffectRouting() throws
Exception {
+ MyNotifier notifier = new MyNotifier() {
+ @Override
+ public boolean isEnabled(CamelEvent event) {
+ if (event instanceof CamelEvent.StepStartedEvent || event
instanceof CamelEvent.ExchangeSendingEvent
+ || event instanceof CamelEvent.RouteStartedEvent) {
+ throw new IllegalStateException("Forced");
+ }
+ return true;
+ }
+ };
+ context.getManagementStrategy().addEventNotifier(notifier);
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").step("foo").to("mock:result").end();
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:result").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello World");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testSentEventsWhenSendingEventsAreIgnored() throws Exception {
+ MyNotifier notifier = new MyNotifier();
+ notifier.setIgnoreExchangeSendingEvents(true);
+ assertSentEvents(notifier);
+ assertEquals(0, notifier.count(CamelEvent.ExchangeSendingEvent.class));
+ }
+
+ @Test
+ public void testSentEventsWhenSendingEventsAreNotEnabled() throws
Exception {
+ MyNotifier notifier = new MyNotifier() {
+ @Override
+ public boolean isEnabled(CamelEvent event) {
+ return event instanceof CamelEvent.ExchangeSentEvent;
+ }
+ };
+ assertSentEvents(notifier);
+ }
+
+ private static void assertSentEvents(MyNotifier notifier) throws Exception
{
+ // use a plain camel context without other event notifiers (such as
the one from the test support)
+ try (CamelContext camel = new DefaultCamelContext()) {
+ camel.getManagementStrategy().addEventNotifier(notifier);
+ camel.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").to("mock:result");
+ }
+ });
+ camel.start();
+
+ MockEndpoint mock = camel.getEndpoint("mock:result",
MockEndpoint.class);
+ mock.expectedMessageCount(1);
+ camel.createProducerTemplate().sendBody("direct:start", "Hello
World");
+ mock.assertIsSatisfied();
+
+ // sent to direct:start and mock:result
+ assertEquals(2,
notifier.count(CamelEvent.ExchangeSentEvent.class));
+ }
+ }
+
+ @Test
+ public void testIgnoreRedeliveryEvents() throws Exception {
+ MyNotifier ignoreRedelivery = new MyNotifier();
+ ignoreRedelivery.setIgnoreExchangeRedeliveryEvents(true);
+ MyNotifier ignoreFailed = new MyNotifier();
+ ignoreFailed.setIgnoreExchangeFailedEvents(true);
+ context.getManagementStrategy().addEventNotifier(ignoreRedelivery);
+ context.getManagementStrategy().addEventNotifier(ignoreFailed);
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+
errorHandler(deadLetterChannel("mock:dead").maximumRedeliveries(2).redeliveryDelay(0));
+
+ from("direct:start").throwException(new
IllegalArgumentException("Forced"));
+ }
+ });
+ context.start();
+
+ getMockEndpoint("mock:dead").expectedMessageCount(1);
+ template.sendBody("direct:start", "Hello World");
+ assertMockEndpointsSatisfied();
+
+ assertEquals(0,
ignoreRedelivery.count(CamelEvent.ExchangeRedeliveryEvent.class));
+ assertEquals(2,
ignoreFailed.count(CamelEvent.ExchangeRedeliveryEvent.class));
+ }
+
+ @Test
+ public void testNotifierAddedAfterStartIsStarted() throws Exception {
+ context.start();
+
+ MyNotifier notifier = new MyNotifier();
+ context.getManagementStrategy().addEventNotifier(notifier);
+ assertTrue(notifier.isStarted());
+ assertEquals(1, notifier.starts.get());
+
+
assertTrue(context.getManagementStrategy().removeEventNotifier(notifier));
+ assertFalse(notifier.isStarted());
+ assertEquals(1, notifier.stops.get());
+ }
+}
diff --git
a/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java
b/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java
index 58528d1dc5e7..a213084e7f8b 100644
--- a/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java
+++ b/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java
@@ -2410,39 +2410,41 @@ public abstract class BaseMainSupport extends
BaseService {
if (mainConfigurationProperties.hasVaultConfiguration()) {
camelContext.setVaultConfiguration(mainConfigurationProperties.vault());
}
- VaultConfiguration target = camelContext.getVaultConfiguration();
+ VaultConfiguration root = camelContext.getVaultConfiguration();
// make defensive copy as we mutate the map
Set<String> keys = new LinkedHashSet<>(properties.asMap().keySet());
// set properties per different vault component
for (String key : keys) {
String name = StringHelper.before(key, ".");
+ // each vault is configured on the vault configuration of the
camel context
+ VaultConfiguration target = root;
if ("aws".equalsIgnoreCase(name)) {
- target = target.aws();
+ target = root.aws();
}
if ("gcp".equalsIgnoreCase(name)) {
- target = target.gcp();
+ target = root.gcp();
}
if ("azure".equalsIgnoreCase(name)) {
- target = target.azure();
+ target = root.azure();
}
if ("hashicorp".equalsIgnoreCase(name)) {
- target = target.hashicorp();
+ target = root.hashicorp();
}
if ("kubernetes".equalsIgnoreCase(name)) {
- target = target.kubernetes();
+ target = root.kubernetes();
}
if ("kubernetescm".equalsIgnoreCase(name)) {
- target = target.kubernetesConfigmaps();
+ target = root.kubernetesConfigmaps();
}
if ("springConfig".equalsIgnoreCase(name)) {
- target = target.springConfig();
+ target = root.springConfig();
}
if ("ibm".equalsIgnoreCase(name)) {
- target = target.ibmSecretsManager();
+ target = root.ibmSecretsManager();
}
if ("cyberark".equalsIgnoreCase(name)) {
- target = target.cyberark();
+ target = root.cyberark();
}
// configure all the properties on the vault at once (to ensure
they are configured in right order)
OrderedLocationProperties config =
MainHelper.extractProperties(properties, name + ".");
diff --git
a/core/camel-main/src/main/java/org/apache/camel/main/VaultConfigurationProperties.java
b/core/camel-main/src/main/java/org/apache/camel/main/VaultConfigurationProperties.java
index b6cbe5219072..36f42ea02a80 100644
---
a/core/camel-main/src/main/java/org/apache/camel/main/VaultConfigurationProperties.java
+++
b/core/camel-main/src/main/java/org/apache/camel/main/VaultConfigurationProperties.java
@@ -17,6 +17,15 @@
package org.apache.camel.main;
import org.apache.camel.spi.BootstrapCloseable;
+import org.apache.camel.vault.AwsVaultConfiguration;
+import org.apache.camel.vault.AzureVaultConfiguration;
+import org.apache.camel.vault.CyberArkVaultConfiguration;
+import org.apache.camel.vault.GcpVaultConfiguration;
+import org.apache.camel.vault.HashicorpVaultConfiguration;
+import org.apache.camel.vault.IBMSecretsManagerVaultConfiguration;
+import org.apache.camel.vault.KubernetesConfigMapVaultConfiguration;
+import org.apache.camel.vault.KubernetesVaultConfiguration;
+import org.apache.camel.vault.SpringCloudConfigConfiguration;
import org.apache.camel.vault.VaultConfiguration;
public class VaultConfigurationProperties extends VaultConfiguration
implements BootstrapCloseable {
@@ -75,7 +84,116 @@ public class VaultConfigurationProperties extends
VaultConfiguration implements
// getter and setters
// --------------------------------------------------------------
- // these are inherited from the parent class
+ // these are inherited from the parent class, but must use the
configurations of the fluent builders
+
+ @Override
+ public AwsVaultConfiguration getAwsVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return aws != null ? aws : super.getAwsVaultConfiguration();
+ }
+
+ @Override
+ public void setAwsVaultConfiguration(AwsVaultConfiguration aws) {
+ super.setAwsVaultConfiguration(aws);
+ this.aws = aws instanceof AwsVaultConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public GcpVaultConfiguration getGcpVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return gcp != null ? gcp : super.getGcpVaultConfiguration();
+ }
+
+ @Override
+ public void setGcpVaultConfiguration(GcpVaultConfiguration gcp) {
+ super.setGcpVaultConfiguration(gcp);
+ this.gcp = gcp instanceof GcpVaultConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public AzureVaultConfiguration getAzureVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return azure != null ? azure : super.getAzureVaultConfiguration();
+ }
+
+ @Override
+ public void setAzureVaultConfiguration(AzureVaultConfiguration azure) {
+ super.setAzureVaultConfiguration(azure);
+ this.azure = azure instanceof AzureVaultConfigurationProperties p ? p
: null;
+ }
+
+ @Override
+ public HashicorpVaultConfiguration getHashicorpVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return hashicorp != null ? hashicorp :
super.getHashicorpVaultConfiguration();
+ }
+
+ @Override
+ public void setHashicorpVaultConfiguration(HashicorpVaultConfiguration
hashicorp) {
+ super.setHashicorpVaultConfiguration(hashicorp);
+ this.hashicorp = hashicorp instanceof
HashicorpVaultConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public KubernetesVaultConfiguration getKubernetesVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return kubernetes != null ? kubernetes :
super.getKubernetesVaultConfiguration();
+ }
+
+ @Override
+ public void setKubernetesVaultConfiguration(KubernetesVaultConfiguration
kubernetes) {
+ super.setKubernetesVaultConfiguration(kubernetes);
+ this.kubernetes = kubernetes instanceof
KubernetesVaultConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public KubernetesConfigMapVaultConfiguration
getKubernetesConfigMapVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return kubernetesConfigmaps != null ? kubernetesConfigmaps :
super.getKubernetesConfigMapVaultConfiguration();
+ }
+
+ @Override
+ public void
setKubernetesConfigMapVaultConfiguration(KubernetesConfigMapVaultConfiguration
kubernetesConfigmaps) {
+ super.setKubernetesConfigMapVaultConfiguration(kubernetesConfigmaps);
+ this.kubernetesConfigmaps
+ = kubernetesConfigmaps instanceof
KubernetesConfigmapsVaultConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public IBMSecretsManagerVaultConfiguration
getIBMSecretsManagerVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return ibmSecretsManager != null ? ibmSecretsManager :
super.getIBMSecretsManagerVaultConfiguration();
+ }
+
+ @Override
+ public void
setIBMSecretsManagerVaultConfiguration(IBMSecretsManagerVaultConfiguration
ibmSecretsManager) {
+ super.setIBMSecretsManagerVaultConfiguration(ibmSecretsManager);
+ this.ibmSecretsManager = ibmSecretsManager instanceof
IBMSecretsManagerVaultConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public SpringCloudConfigConfiguration getSpringCloudConfigConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return springConfig != null ? springConfig :
super.getSpringCloudConfigConfiguration();
+ }
+
+ @Override
+ public void
setSpringCloudConfigConfiguration(SpringCloudConfigConfiguration springConfig) {
+ super.setSpringCloudConfigConfiguration(springConfig);
+ this.springConfig = springConfig instanceof
SpringCloudConfigConfigurationProperties p ? p : null;
+ }
+
+ @Override
+ public CyberArkVaultConfiguration getCyberArkVaultConfiguration() {
+ // the configuration from the fluent builder, or else what has been set
+ return cyberark != null ? cyberark :
super.getCyberArkVaultConfiguration();
+ }
+
+ @Override
+ public void setCyberArkVaultConfiguration(CyberArkVaultConfiguration
cyberark) {
+ super.setCyberArkVaultConfiguration(cyberark);
+ this.cyberark = cyberark instanceof
CyberArkVaultConfigurationProperties p ? p : null;
+ }
// fluent builders
// --------------------------------------------------------------
diff --git
a/core/camel-main/src/test/java/org/apache/camel/main/MainVaultMultipleTest.java
b/core/camel-main/src/test/java/org/apache/camel/main/MainVaultMultipleTest.java
new file mode 100644
index 000000000000..4dcffeeb5c8a
--- /dev/null
+++
b/core/camel-main/src/test/java/org/apache/camel/main/MainVaultMultipleTest.java
@@ -0,0 +1,62 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.main;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.vault.VaultConfiguration;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class MainVaultMultipleTest {
+
+ @Test
+ public void testMainMultipleVaults() {
+ Main main = new Main();
+
+ main.addInitialProperty("camel.vault.aws.region", "myRegion");
+ main.addInitialProperty("camel.vault.hashicorp.host", "myHost");
+ main.addInitialProperty("camel.vault.gcp.projectId", "myProject");
+
+ main.start();
+ try {
+ CamelContext context = main.getCamelContext();
+ VaultConfiguration vault = context.getVaultConfiguration();
+
+ Assertions.assertEquals("myRegion", vault.aws().getRegion());
+ Assertions.assertEquals("myHost", vault.hashicorp().getHost());
+ Assertions.assertEquals("myProject", vault.gcp().getProjectId());
+ } finally {
+ main.stop();
+ }
+ }
+
+ @Test
+ public void testMainVaultGetters() {
+ Main main = new Main();
+ main.configure().vault().aws().setRegion("myRegion");
+
+ main.start();
+ try {
+ VaultConfiguration vault =
main.getCamelContext().getVaultConfiguration();
+
+ Assertions.assertNotNull(vault.getAwsVaultConfiguration());
+ Assertions.assertEquals("myRegion",
vault.getAwsVaultConfiguration().getRegion());
+ } finally {
+ main.stop();
+ }
+ }
+}
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextRestartTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextRestartTest.java
index 1efed9c0dc41..5011bdc41e1d 100644
---
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextRestartTest.java
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextRestartTest.java
@@ -89,11 +89,11 @@ public class ManagedCamelContextRestartTest extends
ManagementTestSupport {
new String[] { "java.lang.String", "java.lang.Object" });
assertEquals("Bye World", reply);
- // restart Camel
- assertEquals(0, starts);
+ // restart Camel (the event notifier was started when added as camel
was already started)
+ assertEquals(1, starts);
assertEquals(0, stops);
mbeanServer.invoke(on, "restart", null, null);
- assertEquals(1, starts);
+ assertEquals(2, starts);
assertEquals(1, stops);
status = (String) mbeanServer.getAttribute(on, "State");
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/EventHelper.java
b/core/camel-support/src/main/java/org/apache/camel/support/EventHelper.java
index 0ff073e7752f..2bdf7445f861 100644
--- a/core/camel-support/src/main/java/org/apache/camel/support/EventHelper.java
+++ b/core/camel-support/src/main/java/org/apache/camel/support/EventHelper.java
@@ -947,7 +947,7 @@ public final class EventHelper {
for (int i = 0; i < notifiers.size(); i++) {
EventNotifier notifier = notifiers.get(i);
- if (isDisabledOrIgnored(notifier) ||
notifier.isIgnoreExchangeFailedEvents()) {
+ if (isDisabledOrIgnored(notifier) ||
notifier.isIgnoreExchangeRedeliveryEvents()) {
continue;
}
@@ -964,6 +964,13 @@ public final class EventHelper {
return answer;
}
+ /**
+ * Notifies that the exchange is being sent to the endpoint.
+ *
+ * @return true if the sending event was notified, or if an exchange sent
event should be notified when the exchange
+ * has been sent (the caller should then time the send and call
+ * {@link #notifyExchangeSent(CamelContext, Exchange, Endpoint,
long)})
+ */
public static boolean notifyExchangeSending(CamelContext context, Exchange
exchange, Endpoint endpoint) {
ManagementStrategy management = context.getManagementStrategy();
if (management == null) {
@@ -990,7 +997,12 @@ public final class EventHelper {
// optimise for loop using index access to avoid creating iterator
object
for (int i = 0; i < notifiers.size(); i++) {
EventNotifier notifier = notifiers.get(i);
- if (isDisabledOrIgnored(notifier) ||
notifier.isIgnoreExchangeSendingEvents()) {
+ if (isDisabledOrIgnored(notifier)) {
+ continue;
+ }
+ // the exchange sent event must be notified even if the notifier
does not want the sending event
+ answer |= !notifier.isIgnoreExchangeSentEvents();
+ if (notifier.isIgnoreExchangeSendingEvents()) {
continue;
}
@@ -1567,7 +1579,13 @@ public final class EventHelper {
}
private static boolean doNotifyEvent(EventNotifier notifier, CamelEvent
event) {
- if (!notifier.isEnabled(event)) {
+ // an exception from the notifier (also from isEnabled) must not
affect routing
+ try {
+ if (!notifier.isEnabled(event)) {
+ return false;
+ }
+ } catch (Throwable e) {
+ LOG.warn("Error checking if event {} is enabled. This exception
will be ignored.", event, e);
return false;
}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index c4719e95dcf1..2fff5abd625c 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -3430,6 +3430,19 @@ healthy while no longer receiving any change event.
Deployments that use readine
now see a Debezium route whose engine has died reported as `DOWN`, where it
was previously reported as `UP`.
The engine is still not restarted automatically.
+=== camel-core - event notifier exchange sent and redelivery events
+
+An event notifier now receives the `ExchangeSentEvent` also when it does not
want the `ExchangeSendingEvent`
+(with `ignoreExchangeSendingEvents=true`, or when `isEnabled` only accepts the
sent event). Previously the sent
+event was only emitted when some event notifier had accepted the sending event.
+
+The `ignoreExchangeRedeliveryEvents` option on an event notifier is now used
to ignore the `ExchangeRedeliveryEvent`.
+Previously the redelivery event was ignored by the
`ignoreExchangeFailedEvents` option instead, and
+`ignoreExchangeRedeliveryEvents` had no effect.
+
+An event notifier that is added after `CamelContext` has been started is now
started (and it is stopped when it is
+removed).
+
=== camel-mongodb - a failed exchange is reported, and does not advance the
consumer's position
Both MongoDB consumers ignored the outcome of the route. A failure left on the
exchange was never looked at,