This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-25072 in repository https://gitbox.apache.org/repos/asf/camel.git
commit 4066d5a8d73d571b4fff2d49b0353e0440edc461 Author: Claus Ibsen <[email protected]> AuthorDate: Tue Sep 29 12:00:21 2026 +0200 CAMEL-25072: camel-management - EIP and service MBeans: fix the remaining follow-ups from the deep review - listEndpointServices is keyed by the route id as well, so two routes consuming the same service no longer fail with KeyAlreadyExistsException, and each consumer gets the hits and route id of its own route. - The min/max processing time statistics no longer lose an update under concurrency (compare and set loop, without allocation). - The notification types advertised by the event notifier MBean are the types it sends. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../engine/DefaultEndpointServiceRegistry.java | 10 +- .../api/management/mbean/CamelOpenMBeanTypes.java | 2 +- .../mbean/ManagedEndpointServiceRegistry.java | 4 + .../management/mbean/ManagedEventNotifier.java | 34 +++--- .../camel/management/mbean/StatisticMaximum.java | 12 +- .../camel/management/mbean/StatisticMinimum.java | 12 +- .../JmxNotificationEventNotifierTest.java | 25 +++++ .../ManagedEndpointServiceRegistryRuntimeTest.java | 37 ++++++ .../ManagedEndpointServiceRegistryTest.java | 125 +++++++++++++++++++++ .../mbean/StatisticMinimumMaximumTest.java | 95 ++++++++++++++++ .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 7 ++ 11 files changed, 330 insertions(+), 33 deletions(-) diff --git a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultEndpointServiceRegistry.java b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultEndpointServiceRegistry.java index c949d02ed401..5a7fe24c1444 100644 --- a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultEndpointServiceRegistry.java +++ b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultEndpointServiceRegistry.java @@ -83,7 +83,8 @@ public class DefaultEndpointServiceRegistry extends ServiceSupport implements En hosted = dc.isHostedService(); routeId = dc.getRouteId(); } - var stat = findStats(endpoint.getEndpointUri(), dir); + // the stats of the route of the consumer, as several routes can consume from the same endpoint + var stat = findStats(endpoint.getEndpointUri(), dir, routeId); long hits = 0; if (stat.isPresent()) { var s = stat.get(); @@ -92,7 +93,7 @@ public class DefaultEndpointServiceRegistry extends ServiceSupport implements En } if ("out".equals(dir) && stat.isEmpty()) { // no OUT stat, then the endpoint may be used only for IN - stat = findStats(endpoint.getEndpointUri(), "in"); + stat = findStats(endpoint.getEndpointUri(), "in", null); if (stat.isPresent()) { return null; } @@ -115,12 +116,13 @@ public class DefaultEndpointServiceRegistry extends ServiceSupport implements En return size; } - private Optional<RuntimeEndpointRegistry.Statistic> findStats(String uri, String direction) { + private Optional<RuntimeEndpointRegistry.Statistic> findStats(String uri, String direction, String routeId) { if (camelContext.getRuntimeEndpointRegistry() == null) { return Optional.empty(); } return camelContext.getRuntimeEndpointRegistry().getEndpointStatistics().stream() - .filter(s -> uri.equals(s.getUri()) && (direction == null || s.getDirection().equals(direction))) + .filter(s -> uri.equals(s.getUri()) && (direction == null || s.getDirection().equals(direction)) + && (routeId == null || routeId.equals(s.getRouteId()))) .findFirst(); } diff --git a/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java b/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java index c1a686c01abf..5d5a65249025 100644 --- a/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java +++ b/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/CamelOpenMBeanTypes.java @@ -34,7 +34,7 @@ public final class CamelOpenMBeanTypes { CompositeType ct = listEndpointServicesCompositeType(); return new TabularType( "listEndpointServices", "Lists all the endpoint services in the registry", ct, - new String[] { "component", "dir", "serviceUrl", "endpointUri" }); + new String[] { "component", "dir", "serviceUrl", "endpointUri", "routeId" }); } public static CompositeType listEndpointServicesCompositeType() throws OpenDataException { diff --git a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEndpointServiceRegistry.java b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEndpointServiceRegistry.java index 56428490709f..b933dfe8aeab 100644 --- a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEndpointServiceRegistry.java +++ b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEndpointServiceRegistry.java @@ -88,6 +88,10 @@ public class ManagedEndpointServiceRegistry extends ManagedService implements Ma m.forEach((k, v) -> sj.add(k + "=" + v)); metadata = sj.toString(); } + if (answer.containsKey(new Object[] { component, dir, serviceUrl, endpointUri, routeId })) { + // endpoints of the same route that only differ in a secret are the same uri when sanitized + continue; + } CompositeData data = new CompositeDataSupport( ct, diff --git a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEventNotifier.java b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEventNotifier.java index 1ccfe2c12dfc..e8383d672990 100644 --- a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEventNotifier.java +++ b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedEventNotifier.java @@ -20,10 +20,12 @@ import java.util.ArrayList; import java.util.List; import javax.management.MBeanNotificationInfo; +import javax.management.Notification; import javax.management.NotificationBroadcasterSupport; import org.apache.camel.CamelContext; import org.apache.camel.api.management.JmxNotificationBroadcasterAware; +import org.apache.camel.spi.CamelEvent; import org.apache.camel.spi.EventNotifier; import org.apache.camel.spi.ManagementStrategy; @@ -163,26 +165,22 @@ public class ManagedEventNotifier extends NotificationBroadcasterSupport impleme @Override public MBeanNotificationInfo[] getNotificationInfo() { - // all the class names in the event package - String[] names = { - "CamelContextStartedEvent", "CamelContextStartingEvent", "CamelContextStartupFailureEvent", - "CamelContextStopFailureEvent", "CamelContextStoppedEvent", "CamelContextStoppingEvent", - "CamelContextSuspendingEvent", "CamelContextSuspendedEvent", "CamelContextResumingEvent", - "CamelContextResumedEvent", - "CamelContextResumeFailureEvent", "ExchangeCompletedEvent", "ExchangeCreatedEvent", "ExchangeFailedEvent", - "ExchangeFailureHandledEvent", "ExchangeRedeliveryEvents", "ExchangeSendingEvent", "ExchangeSentEvent", - "RouteStartedEvent", - "RouteStoppedEvent", "ServiceStartupFailureEvent", "ServiceStopFailureEvent", - "StepStartedEvent", "StepCompletedEvent", "StepFailedEvent" }; - + // JmxNotificationEventNotifier uses the simple class name of the event as the notification type List<MBeanNotificationInfo> infos = new ArrayList<>(); - for (String name : names) { - MBeanNotificationInfo info = new MBeanNotificationInfo( - new String[] { "org.apache.camel.management.event" }, - "org.apache.camel.management.event." + name, "The event " + name + " occurred"); - infos.add(info); + for (CamelEvent.Type type : CamelEvent.Type.values()) { + if (type == CamelEvent.Type.Custom) { + // custom events have their own class names + continue; + } + String name = type.name(); + if (name.startsWith("Routes")) { + // the routes events of the CamelContext + name = "CamelContext" + name; + } + name = name + "Event"; + infos.add(new MBeanNotificationInfo( + new String[] { name }, Notification.class.getName(), "The event " + name + " occurred")); } - return infos.toArray(new MBeanNotificationInfo[0]); } diff --git a/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMaximum.java b/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMaximum.java index 01aef94189f0..39e752c048d7 100644 --- a/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMaximum.java +++ b/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMaximum.java @@ -24,12 +24,14 @@ public class StatisticMaximum extends Statistic { @Override public void updateValue(long newValue) { - // its okay its not 100% thread safe (these jmx counters are not guaranteed to be accurate for min/max values) - // if we use the atomic operation updateAndGet then the JVM creates a new lambda per call which creates a new object - // in the JVM and causes higher memory footprint + // compare and set in a loop (and not updateAndGet, which creates a new lambda per call), so a concurrent + // update cannot be lost long current = value.get(); - if (current == -1 || current < newValue) { - value.set(newValue); + while (current == -1 || current < newValue) { + if (value.compareAndSet(current, newValue)) { + return; + } + current = value.get(); } } diff --git a/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMinimum.java b/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMinimum.java index 3e7872a07c9b..abacb4cbcd38 100644 --- a/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMinimum.java +++ b/core/camel-management/src/main/java/org/apache/camel/management/mbean/StatisticMinimum.java @@ -24,12 +24,14 @@ public class StatisticMinimum extends Statistic { @Override public void updateValue(long newValue) { - // its okay its not 100% thread safe (these jmx counters are not guaranteed to be accurate for min/max values) - // if we use the atomic operation updateAndGet then the JVM creates a new lambda per call which creates a new object - // in the JVM and causes higher memory footprint + // compare and set in a loop (and not updateAndGet, which creates a new lambda per call), so a concurrent + // update cannot be lost long current = value.get(); - if (current == -1 || current > newValue) { - value.set(newValue); + while (current == -1 || current > newValue) { + if (value.compareAndSet(current, newValue)) { + return; + } + current = value.get(); } } diff --git a/core/camel-management/src/test/java/org/apache/camel/management/JmxNotificationEventNotifierTest.java b/core/camel-management/src/test/java/org/apache/camel/management/JmxNotificationEventNotifierTest.java index 1785833dc61b..278a5ec7f86d 100644 --- a/core/camel-management/src/test/java/org/apache/camel/management/JmxNotificationEventNotifierTest.java +++ b/core/camel-management/src/test/java/org/apache/camel/management/JmxNotificationEventNotifierTest.java @@ -16,6 +16,11 @@ */ package org.apache.camel.management; +import java.util.Arrays; +import java.util.HashSet; +import java.util.Set; + +import javax.management.MBeanNotificationInfo; import javax.management.Notification; import javax.management.NotificationFilter; import javax.management.NotificationListener; @@ -30,7 +35,9 @@ import org.junit.jupiter.api.condition.OS; import static org.apache.camel.management.DefaultManagementObjectNameStrategy.TYPE_EVENT_NOTIFIER; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; @DisabledOnOs(OS.AIX) public class JmxNotificationEventNotifierTest extends ManagementTestSupport { @@ -77,6 +84,18 @@ public class JmxNotificationEventNotifierTest extends ManagementTestSupport { assertEquals(8, listener.getEventCounter(), "Get a wrong number of events"); + // the notification types that were sent are advertised by the MBean + Set<String> advertised = new HashSet<>(); + for (MBeanNotificationInfo info : context.getManagementStrategy().getManagementAgent().getMBeanServer() + .getMBeanInfo(on).getNotifications()) { + assertEquals(Notification.class.getName(), info.getName()); + advertised.addAll(Arrays.asList(info.getNotifTypes())); + } + assertFalse(listener.getTypes().isEmpty()); + for (String type : listener.getTypes()) { + assertTrue(advertised.contains(type), "Notification type " + type + " should be advertised"); + } + context.stop(); } @@ -117,11 +136,17 @@ public class JmxNotificationEventNotifierTest extends ManagementTestSupport { private class MyNotificationListener implements NotificationListener { private int eventCounter; + private final Set<String> types = new HashSet<>(); @Override public void handleNotification(Notification notification, Object handback) { log.debug("Get the notification : {}", notification); eventCounter++; + types.add(notification.getType()); + } + + public Set<String> getTypes() { + return types; } public int getEventCounter() { diff --git a/core/camel-management/src/test/java/org/apache/camel/management/ManagedEndpointServiceRegistryRuntimeTest.java b/core/camel-management/src/test/java/org/apache/camel/management/ManagedEndpointServiceRegistryRuntimeTest.java new file mode 100644 index 000000000000..0d4d545af397 --- /dev/null +++ b/core/camel-management/src/test/java/org/apache/camel/management/ManagedEndpointServiceRegistryRuntimeTest.java @@ -0,0 +1,37 @@ +/* + * 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.management; + +import org.apache.camel.CamelContext; +import org.apache.camel.impl.engine.DefaultRuntimeEndpointRegistry; +import org.junit.jupiter.api.condition.DisabledOnOs; +import org.junit.jupiter.api.condition.OS; + +/** + * The same as {@link ManagedEndpointServiceRegistryTest} with the runtime endpoint registry enabled, which gives the + * hits and route of each endpoint service. + */ +@DisabledOnOs(OS.AIX) +public class ManagedEndpointServiceRegistryRuntimeTest extends ManagedEndpointServiceRegistryTest { + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext context = super.createCamelContext(); + context.setRuntimeEndpointRegistry(new DefaultRuntimeEndpointRegistry()); + return context; + } +} diff --git a/core/camel-management/src/test/java/org/apache/camel/management/ManagedEndpointServiceRegistryTest.java b/core/camel-management/src/test/java/org/apache/camel/management/ManagedEndpointServiceRegistryTest.java new file mode 100644 index 000000000000..836644a4d50c --- /dev/null +++ b/core/camel-management/src/test/java/org/apache/camel/management/ManagedEndpointServiceRegistryTest.java @@ -0,0 +1,125 @@ +/* + * 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.management; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; + +import javax.management.MBeanServer; +import javax.management.ObjectName; +import javax.management.openmbean.CompositeData; +import javax.management.openmbean.TabularData; + +import org.apache.camel.CamelContext; +import org.apache.camel.Consumer; +import org.apache.camel.Endpoint; +import org.apache.camel.MultipleConsumersSupport; +import org.apache.camel.Processor; +import org.apache.camel.Producer; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.spi.EndpointServiceLocation; +import org.apache.camel.support.DefaultComponent; +import org.apache.camel.support.DefaultConsumer; +import org.apache.camel.support.DefaultEndpoint; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.condition.DisabledOnOs; +import org.junit.jupiter.api.condition.OS; + +import static org.apache.camel.management.DefaultManagementObjectNameStrategy.TYPE_SERVICE; +import static org.junit.jupiter.api.Assertions.assertEquals; + +@DisabledOnOs(OS.AIX) +public class ManagedEndpointServiceRegistryTest extends ManagementTestSupport { + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext context = super.createCamelContext(); + context.addComponent("service", new DefaultComponent() { + @Override + protected Endpoint createEndpoint(String uri, String remaining, Map<String, Object> parameters) { + return new ServiceEndpoint(uri, this); + } + }); + // the registry is created on first use, so create it before the context is started to have it managed + context.getCamelContextExtension().getEndpointServiceRegistry(); + return context; + } + + @Test + public void testTwoRoutesConsumeSameService() throws Exception { + MBeanServer mbeanServer = getMBeanServer(); + ObjectName on = getCamelObjectName(TYPE_SERVICE, "DefaultEndpointServiceRegistry"); + + // both routes consume the same service, which failed with KeyAlreadyExistsException + TabularData data = (TabularData) mbeanServer.invoke(on, "listEndpointServices", null, null); + + List<String> routeIds = new ArrayList<>(); + for (Object row : data.values()) { + CompositeData cd = (CompositeData) row; + assertEquals("localhost:8080", cd.get("serviceUrl")); + if ("in".equals(cd.get("dir"))) { + routeIds.add((String) cd.get("routeId")); + } + } + routeIds.sort(null); + assertEquals(List.of("a", "b"), routeIds); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("service:foo").routeId("a").to("mock:a"); + from("service:foo").routeId("b").to("mock:b"); + } + }; + } + + private static class ServiceEndpoint extends DefaultEndpoint implements EndpointServiceLocation, MultipleConsumersSupport { + + ServiceEndpoint(String uri, DefaultComponent component) { + super(uri, component); + } + + @Override + public String getServiceUrl() { + return "localhost:8080"; + } + + @Override + public String getServiceProtocol() { + return "tcp"; + } + + @Override + public boolean isMultipleConsumersSupported() { + return true; + } + + @Override + public Producer createProducer() { + throw new UnsupportedOperationException(); + } + + @Override + public Consumer createConsumer(Processor processor) { + return new DefaultConsumer(this, processor); + } + } +} diff --git a/core/camel-management/src/test/java/org/apache/camel/management/mbean/StatisticMinimumMaximumTest.java b/core/camel-management/src/test/java/org/apache/camel/management/mbean/StatisticMinimumMaximumTest.java new file mode 100644 index 000000000000..f9994aeecc4a --- /dev/null +++ b/core/camel-management/src/test/java/org/apache/camel/management/mbean/StatisticMinimumMaximumTest.java @@ -0,0 +1,95 @@ +/* + * 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.management.mbean; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.RepeatedTest; +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 StatisticMinimumMaximumTest { + + private static final int THREADS = 8; + private static final int UPDATES = 20000; + + @Test + public void testMinimumMaximum() { + StatisticMinimum min = new StatisticMinimum(); + StatisticMaximum max = new StatisticMaximum(); + assertFalse(min.isUpdated()); + assertFalse(max.isUpdated()); + assertEquals(0, min.getValue()); + assertEquals(0, max.getValue()); + + for (long v : new long[] { 5, 3, 8, 3, 6 }) { + min.updateValue(v); + max.updateValue(v); + } + assertTrue(min.isUpdated()); + assertEquals(3, min.getValue()); + assertEquals(8, max.getValue()); + + min.reset(); + max.reset(); + assertFalse(min.isUpdated()); + assertFalse(max.isUpdated()); + } + + @RepeatedTest(5) + public void testConcurrentUpdates() throws Exception { + StatisticMinimum min = new StatisticMinimum(); + StatisticMaximum max = new StatisticMaximum(); + + // each thread updates with its own range of values, in the order that makes each update a new min and max + ExecutorService executor = Executors.newFixedThreadPool(THREADS); + try { + CountDownLatch start = new CountDownLatch(1); + List<Future<?>> futures = new ArrayList<>(); + for (int t = 0; t < THREADS; t++) { + final int thread = t; + futures.add(executor.submit(() -> { + start.await(); + for (int i = 0; i < UPDATES; i++) { + max.updateValue((long) i * THREADS + thread + 1); + min.updateValue((long) (UPDATES - i) * THREADS - thread); + } + return null; + })); + } + start.countDown(); + for (Future<?> future : futures) { + future.get(30, TimeUnit.SECONDS); + } + } finally { + executor.shutdownNow(); + } + + // a lost update would leave a value from another thread that is not the real min or max + assertEquals((long) UPDATES * THREADS, max.getValue()); + assertEquals(1L, min.getValue()); + } +} 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 04d6027f2ded..85483aafbcbd 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 @@ -531,6 +531,13 @@ Some tabular data returned by the JMX MBeans had a key that was not unique, so t `KeyAlreadyExistsException`. The tabular data of the Choice EIP and doTry EIP `extendedInformation` and of `listTasks` of the task manager registry now have an `index` item as their key, the exchange factories of `listStatistics` are keyed by `url` and `routeId`, and endpoints that only differ in a secret (which is masked) are shown once. +The `listEndpointServices` operation of the endpoint service registry is keyed by `routeId` as well, so two routes +that consume from the same service are both listed. + +The notification types advertised by the event notifier MBean (`getNotificationInfo`) are now the types that +`JmxNotificationEventNotifier` sends, which is the simple class name of the event (such as `ExchangeCompletedEvent`), +and the notification class is `javax.management.Notification`. Before, the advertised types used a +`org.apache.camel.management.event.` prefix and several events were missing. === camel-exec
