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");
+            }
+        };
+    }
+}

Reply via email to