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 b187159c1a14 CAMEL-25073: camel-management - Do not unregister an
endpoint MBean that belongs to another endpoint (#27461)
b187159c1a14 is described below
commit b187159c1a1400def85c950fe40cfee06e087948
Author: allthingssecurity <[email protected]>
AuthorDate: Wed Oct 7 13:40:38 2026 +0530
CAMEL-25073: camel-management - Do not unregister an endpoint MBean that
belongs to another endpoint (#27461)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../management/JmxManagementLifecycleStrategy.java | 60 +++++++++++-
.../ManagedAddEndpointMaskedSecretRaceTest.java | 106 +++++++++++++++++++++
.../ManagedRemoveEndpointMaskedSecretTest.java | 96 +++++++++++++++++++
.../ManagedRemoveInterceptedEndpointTest.java | 65 +++++++++++++
4 files changed, 325 insertions(+), 2 deletions(-)
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/JmxManagementLifecycleStrategy.java
b/core/camel-management/src/main/java/org/apache/camel/management/JmxManagementLifecycleStrategy.java
index e81ece038129..2b0ecc33282b 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/JmxManagementLifecycleStrategy.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/JmxManagementLifecycleStrategy.java
@@ -24,8 +24,11 @@ import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
import javax.management.JMException;
import javax.management.MBeanServer;
@@ -104,6 +107,7 @@ import org.apache.camel.spi.ErrorRegistry;
import org.apache.camel.spi.EventNotifier;
import org.apache.camel.spi.ExchangeFactoryManager;
import org.apache.camel.spi.InflightRepository;
+import org.apache.camel.spi.InterceptSendToEndpoint;
import org.apache.camel.spi.InternalProcessor;
import org.apache.camel.spi.LifecycleStrategy;
import org.apache.camel.spi.ManagementAgent;
@@ -157,6 +161,12 @@ public class JmxManagementLifecycleStrategy extends
ServiceSupport implements Li
private final Map<BacklogTracer, ManagedBacklogTracer>
managedBacklogTracers = new HashMap<>();
private final Map<DefaultBacklogDebugger, ManagedBacklogDebugger>
managedBacklogDebuggers = new HashMap<>();
private final Map<Object, Object> managedThreadPools = new HashMap<>();
+ // the endpoint whose MBean is registered with the given name, as
endpoints that only differ in a masked secret
+ // (such as a password) get the same name, and removing one of them must
not unregister the MBean of the other
+ private final Map<ObjectName, Endpoint> managedEndpoints = new
ConcurrentHashMap<>();
+ // held while an endpoint MBean is registered or unregistered together
with the change to managedEndpoints, so that
+ // endpoints with the same name that are added or removed at the same time
agree on the endpoint that owns the MBean
+ private final Lock managedEndpointsLock = new ReentrantLock();
// route group MBean is shared by all routes in the same group, so its
performance counters
// aggregate the statistics across all the member routes
private final Map<String, ManagedRouteGroup> managedRouteGroups = new
HashMap<>();
@@ -488,7 +498,17 @@ public class JmxManagementLifecycleStrategy extends
ServiceSupport implements Li
// endpoint should not be managed
return;
}
- manageObject(me);
+ ObjectName on =
getManagementStrategy().getManagementObjectNameStrategy().getObjectName(me);
+ managedEndpointsLock.lock();
+ try {
+ boolean exists = on != null &&
getManagementStrategy().isManagedName(on);
+ manageObject(me);
+ if (on != null && !exists) {
+ managedEndpoints.put(on, originalEndpoint(endpoint));
+ }
+ } finally {
+ managedEndpointsLock.unlock();
+ }
} catch (Exception e) {
LOG.warn("Could not register Endpoint MBean for endpoint: {}. This
exception will be ignored.", endpoint, e);
}
@@ -503,12 +523,47 @@ public class JmxManagementLifecycleStrategy extends
ServiceSupport implements Li
try {
Object me =
getManagementObjectStrategy().getManagedObjectForEndpoint(camelContext,
endpoint);
- unmanageObject(me);
+ ObjectName on = me != null ?
getManagementStrategy().getManagementObjectNameStrategy().getObjectName(me) :
null;
+ managedEndpointsLock.lock();
+ try {
+ if (on != null) {
+ Endpoint owner = managedEndpoints.get(on);
+ if (owner != null && owner != originalEndpoint(endpoint)
&& isEndpointRegistered(owner)) {
+ // the MBean belongs to another endpoint with the same
name (only differs in a masked secret)
+ LOG.debug("Not unregistering Endpoint MBean: {} as it
belongs to another endpoint", on);
+ return;
+ }
+ managedEndpoints.remove(on);
+ }
+ unmanageObject(me);
+ } finally {
+ managedEndpointsLock.unlock();
+ }
} catch (Exception e) {
LOG.warn("Could not unregister Endpoint MBean for endpoint: {}.
This exception will be ignored.", endpoint, e);
}
}
+ private boolean isEndpointRegistered(Endpoint endpoint) {
+ // the endpoint may have been replaced in the registry (addEndpoint)
without being removed
+ for (Endpoint registered : camelContext.getEndpoints()) {
+ if (originalEndpoint(registered) == endpoint) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static Endpoint originalEndpoint(Endpoint endpoint) {
+ // the endpoint registry may hold an intercepting endpoint
(interceptSendToEndpoint, mock endpoints) that wraps
+ // the endpoint whose MBean was registered
+ Endpoint answer = endpoint;
+ while (answer instanceof InterceptSendToEndpoint intercept) {
+ answer = intercept.getOriginalEndpoint();
+ }
+ return answer;
+ }
+
@Override
public void onServiceAdd(CamelContext context, Service service, Route
route) {
if (!initialized) {
@@ -1234,6 +1289,7 @@ public class JmxManagementLifecycleStrategy extends
ServiceSupport implements Li
managedBacklogTracers.clear();
managedBacklogDebuggers.clear();
managedThreadPools.clear();
+ managedEndpoints.clear();
managedRouteGroups.clear();
}
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedAddEndpointMaskedSecretRaceTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedAddEndpointMaskedSecretRaceTest.java
new file mode 100644
index 000000000000..efef6ce406d9
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedAddEndpointMaskedSecretRaceTest.java
@@ -0,0 +1,106 @@
+/*
+ * 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.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Proxy;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.Endpoint;
+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_ENDPOINT;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Two endpoints that only differ in a masked secret are added at the same
time: the endpoint recorded as the owner of
+ * the MBean must be the one whose MBean was registered.
+ */
+@DisabledOnOs(OS.AIX)
+class ManagedAddEndpointMaskedSecretRaceTest extends ManagementTestSupport {
+
+ @Test
+ void testAddEndpointsWithTheSameNameAtTheSameTime() throws Exception {
+
context.getManagementStrategy().getManagementAgent().setRegisterAlways(true);
+ ObjectName on = getCamelObjectName(TYPE_ENDPOINT,
"stub://race\\?password=xxxxxx");
+
+ // the first time the thread adding endpoint x asks whether the name
is registered, it waits for the test after
+ // the answer, so that endpoint y can be added meanwhile
+ MBeanServer server = getMBeanServer();
+ AtomicReference<Thread> adder = new AtomicReference<>();
+ CountDownLatch entered = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ MBeanServer gated = (MBeanServer)
Proxy.newProxyInstance(MBeanServer.class.getClassLoader(),
+ new Class<?>[] { MBeanServer.class }, (proxy, method, args) ->
{
+ Object answer;
+ try {
+ answer = method.invoke(server, args);
+ } catch (InvocationTargetException e) {
+ throw e.getCause();
+ }
+ if ("isRegistered".equals(method.getName()) &&
on.equals(args[0])
+ && Thread.currentThread() == adder.get() &&
entered.getCount() > 0) {
+ entered.countDown();
+ assertTrue(release.await(20, TimeUnit.SECONDS), "The
test should release the adder");
+ }
+ return answer;
+ });
+
context.getManagementStrategy().getManagementAgent().setMBeanServer(gated);
+ try {
+ AtomicReference<Endpoint> x = new AtomicReference<>();
+ Thread threadX = new Thread(() ->
x.set(context.getEndpoint("stub:race?password=x")), "adder-x");
+ adder.set(threadX);
+ threadX.start();
+ assertTrue(entered.await(20, TimeUnit.SECONDS), "Endpoint x should
be added");
+
+ AtomicReference<Endpoint> y = new AtomicReference<>();
+ Thread threadY = new Thread(() ->
y.set(context.getEndpoint("stub:race?password=y")), "adder-y");
+ threadY.start();
+ // endpoint y is added meanwhile, or waits for endpoint x
+ await().atMost(20, TimeUnit.SECONDS).until(() -> !threadY.isAlive()
+ || threadY.getState() == Thread.State.WAITING ||
threadY.getState() == Thread.State.BLOCKED);
+ release.countDown();
+ threadX.join(20000);
+ threadY.join(20000);
+ assertFalse(threadX.isAlive(), "Endpoint x should be added");
+ assertFalse(threadY.isAlive(), "Endpoint y should be added");
+ assertTrue(server.isRegistered(on), "Should be registered");
+
+ // endpoint x was the first to ask, so the MBean is the one of
endpoint x, and removing endpoint y keeps it
+ context.removeEndpoint(y.get());
+ assertTrue(server.isRegistered(on), "The MBean of endpoint x
should still be registered");
+ assertEquals("Started", server.getAttribute(on, "State"),
+ "The MBean should be the one of endpoint x, which is still
in use");
+
+ context.removeEndpoint(x.get());
+ assertFalse(server.isRegistered(on), "Should no longer be
registered");
+ } finally {
+ release.countDown();
+
context.getManagementStrategy().getManagementAgent().setMBeanServer(server);
+ }
+ }
+}
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedRemoveEndpointMaskedSecretTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRemoveEndpointMaskedSecretTest.java
new file mode 100644
index 000000000000..48fc8eb9ad2e
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRemoveEndpointMaskedSecretTest.java
@@ -0,0 +1,96 @@
+/*
+ * 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 javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.Endpoint;
+import org.apache.camel.builder.RouteBuilder;
+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_ENDPOINT;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Endpoints that only differ in a secret get the same (masked) MBean name,
and removing one of them should not
+ * unregister the MBean of the other.
+ */
+@DisabledOnOs(OS.AIX)
+class ManagedRemoveEndpointMaskedSecretTest extends ManagementTestSupport {
+
+ @Test
+ void testRemoveEndpointCreatedAtRuntime() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+ ObjectName on = getCamelObjectName(TYPE_ENDPOINT,
"stub://foo\\?password=xxxxxx");
+ assertTrue(mbeanServer.isRegistered(on), "Should be registered");
+
+ // not registered as it is created after CamelContext has been started
+ Endpoint other = context.getEndpoint("stub:foo?password=other");
+ context.removeEndpoint(other);
+ assertTrue(mbeanServer.isRegistered(on), "The MBean of the endpoint of
route a should still be registered");
+
+ context.getRouteController().stopRoute("a");
+ context.removeRoute("a");
+ context.removeEndpoint(context.hasEndpoint("stub:foo?password=a"));
+ assertFalse(mbeanServer.isRegistered(on), "Should no longer be
registered");
+ }
+
+ @Test
+ void testRemoveEndpointOfOtherRoute() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+ ObjectName on = getCamelObjectName(TYPE_ENDPOINT,
"stub://foo\\?password=xxxxxx");
+ assertTrue(mbeanServer.isRegistered(on), "Should be registered");
+
+ context.getRouteController().stopRoute("b");
+ context.removeRoute("b");
+ context.removeEndpoint(context.hasEndpoint("stub:foo?password=b"));
+ assertTrue(mbeanServer.isRegistered(on), "The MBean of the endpoint of
route a should still be registered");
+
+ context.getRouteController().stopRoute("a");
+ context.removeRoute("a");
+ context.removeEndpoint(context.hasEndpoint("stub:foo?password=a"));
+ assertFalse(mbeanServer.isRegistered(on), "Should no longer be
registered");
+ }
+
+ @Test
+ void testRemoveEndpointThatReplacedTheOwner() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+ ObjectName on = getCamelObjectName(TYPE_ENDPOINT,
"stub://foo\\?password=xxxxxx");
+ assertTrue(mbeanServer.isRegistered(on), "Should be registered");
+
+ // replaces the endpoint of route a in the registry (the replaced
endpoint is not removed)
+ Endpoint replacement =
context.getComponent("stub").createEndpoint("stub://foo?password=a");
+ context.addEndpoint("stub:foo?password=a", replacement);
+ context.removeEndpoint(replacement);
+ assertFalse(mbeanServer.isRegistered(on), "Should no longer be
registered");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("stub:foo?password=a").routeId("a").to("mock:a");
+ from("stub:foo?password=b").routeId("b").to("mock:b");
+ }
+ };
+ }
+}
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedRemoveInterceptedEndpointTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRemoveInterceptedEndpointTest.java
new file mode 100644
index 000000000000..8afb14ce4ba0
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRemoveInterceptedEndpointTest.java
@@ -0,0 +1,65 @@
+/*
+ * 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 javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.builder.RouteBuilder;
+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_ENDPOINT;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The endpoint registry holds an intercepting endpoint that wraps the
endpoint whose MBean was registered, and removing
+ * another intercepted endpoint with the same masked name should not
unregister that MBean.
+ */
+@DisabledOnOs(OS.AIX)
+class ManagedRemoveInterceptedEndpointTest extends ManagementTestSupport {
+
+ @Test
+ void testRemoveInterceptedEndpointWithMaskedSecret() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+ ObjectName on = getCamelObjectName(TYPE_ENDPOINT,
"stub://bar\\?password=xxxxxx");
+ assertTrue(mbeanServer.isRegistered(on), "Should be registered");
+
+ context.getRouteController().stopRoute("c");
+ context.removeRoute("c");
+ assertTrue(mbeanServer.isRegistered(on), "The MBean of the endpoint of
route b should still be registered");
+
+ context.getRouteController().stopRoute("b");
+ context.removeRoute("b");
+ assertFalse(mbeanServer.isRegistered(on), "Should no longer be
registered");
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ interceptSendToEndpoint("stub:*").to("mock:intercepted");
+
+ from("direct:b").routeId("b").to("stub:bar?password=b");
+ from("direct:c").routeId("c").to("stub:bar?password=c");
+ }
+ };
+ }
+}