This is an automated email from the ASF dual-hosted git repository.
Croway 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 e44ef154a5d5 CAMEL-25267: camel-management - find the processors of a
route without reading every processor mbean (#27280)
e44ef154a5d5 is described below
commit e44ef154a5d51ca8c5a353fb838c52815d197ae1
Author: Federico Mariani <[email protected]>
AuthorDate: Fri Oct 2 15:28:36 2026 +0200
CAMEL-25267: camel-management - find the processors of a route without
reading every processor mbean (#27280)
* CAMEL-25267: camel-management - find the processors of a route without
reading every processor mbean
ManagedRoute found the processors of its route by reading the route id of
every processor mbean of the CamelContext (processorIds,
dumpRouteStatsAsXml/
JSon, dumpStepStatsAsXml, dumpRouteSourceLocationsAsXml and reset), and
ManagedCamelContext.getManagedProcessor/getManagedStep walked every route to
find the processor of an id. The consoles call them for every route and
every
processor, which was quadratic in the size of the integration.
The default management agent now indexes the processor and step mbeans it
registers by route and by id, and these lookups use it. A custom agent, or
an
id that is not indexed, still uses the previous lookup.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
* CAMEL-25267: use {@code null} in javadoc
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---------
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.../camel/management/DefaultManagementAgent.java | 90 ++++++++++-
.../camel/management/ManagedCamelContextImpl.java | 29 +++-
.../camel/management/mbean/ManagedRoute.java | 164 +++++++++------------
.../ManagedRouteProcessorsOfOtherRoutesTest.java | 135 +++++++++++++++++
4 files changed, 314 insertions(+), 104 deletions(-)
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/DefaultManagementAgent.java
b/core/camel-management/src/main/java/org/apache/camel/management/DefaultManagementAgent.java
index fb6c3ef83e08..2d184167f0f3 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/DefaultManagementAgent.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/DefaultManagementAgent.java
@@ -17,9 +17,11 @@
package org.apache.camel.management;
import java.lang.management.ManagementFactory;
+import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@@ -36,6 +38,8 @@ import org.apache.camel.CamelContextAware;
import org.apache.camel.ManagementMBeansLevel;
import org.apache.camel.ManagementStatisticsLevel;
import org.apache.camel.api.management.JmxSystemPropertyKeys;
+import org.apache.camel.api.management.mbean.ManagedProcessorMBean;
+import org.apache.camel.api.management.mbean.ManagedStepMBean;
import org.apache.camel.spi.ManagementAgent;
import org.apache.camel.spi.ManagementMBeanAssembler;
import org.apache.camel.support.management.DefaultManagementMBeanAssembler;
@@ -60,6 +64,13 @@ public class DefaultManagementAgent extends ServiceSupport
implements Management
// need a name -> actual name mapping as some servers changes the names
(such as WebSphere)
private final ConcurrentMap<ObjectName, ObjectName> mbeansRegistered = new
ConcurrentHashMap<>();
+ // the registered processor and step mbeans by the route they belong to
and by their id, so they can be found
+ // without reading every processor mbean of the CamelContext (or walking
every route) which is slow with many routes
+ private final ConcurrentMap<ObjectName, ProcessorMBean> processorMBeans =
new ConcurrentHashMap<>();
+ private final ConcurrentMap<String, Set<ObjectName>> routeProcessors = new
ConcurrentHashMap<>();
+ private final ConcurrentMap<String, Set<ObjectName>> routeSteps = new
ConcurrentHashMap<>();
+ private final ConcurrentMap<String, ManagedProcessorMBean> processorsById
= new ConcurrentHashMap<>();
+ private final ConcurrentMap<String, ManagedProcessorMBean> stepsById = new
ConcurrentHashMap<>();
private String mBeanServerDefaultDomain = DEFAULT_DOMAIN;
private String mBeanObjectDomainName = DEFAULT_DOMAIN;
@@ -355,21 +366,30 @@ public class DefaultManagementAgent extends
ServiceSupport implements Management
@Override
public void register(Object obj, ObjectName name, boolean
forceRegistration) throws JMException {
+ ObjectName registeredName;
try {
- registerMBeanWithServer(obj, name, forceRegistration);
+ registeredName = registerMBeanWithServer(obj, name,
forceRegistration);
} catch (NotCompliantMBeanException e) {
// If this is not a "normal" MBean, then try to deploy it using
JMX annotations
ObjectHelper.notNull(assembler, "ManagementMBeanAssembler",
camelContext);
Object mbean = assembler.assemble(server, obj, name);
+ registeredName = null;
if (mbean != null) {
// and register the mbean
- registerMBeanWithServer(mbean, name, forceRegistration);
+ registeredName = registerMBeanWithServer(mbean, name,
forceRegistration);
+ }
+ }
+ if (registeredName != null) {
+ removeProcessorMBean(name);
+ if (obj instanceof ManagedProcessorMBean mp) {
+ addProcessorMBean(name, registeredName, mp);
}
}
}
@Override
public void unregister(ObjectName name) throws JMException {
+ removeProcessorMBean(name);
if (isRegistered(name)) {
ObjectName on = mbeansRegistered.remove(name);
server.unregisterMBean(on);
@@ -379,6 +399,65 @@ public class DefaultManagementAgent extends ServiceSupport
implements Management
}
}
+ /**
+ * Gets the names of the registered processor mbeans that belong to the
given route.
+ * <p/>
+ * This avoids querying (and reading the route id of) every processor
mbean of the CamelContext, which is slow when
+ * there are many routes.
+ *
+ * @param routeId the route id
+ * @param steps whether to get the step mbeans, instead of the
processor mbeans
+ * @return the names the mbeans are registered with in the mbean
server, or an empty list
+ */
+ public List<ObjectName> getRouteProcessorMBeanNames(String routeId,
boolean steps) {
+ Set<ObjectName> names = (steps ? routeSteps :
routeProcessors).get(routeId);
+ return names != null ? new ArrayList<>(names) : new ArrayList<>();
+ }
+
+ /**
+ * Gets the managed object of the registered processor mbean with the
given id.
+ * <p/>
+ * This avoids walking every route to find the processor and its
definition, which is slow when there are many
+ * routes.
+ *
+ * @param id the processor id
+ * @param steps whether to get a step mbean, instead of a processor mbean
+ * @return the managed object, or {@code null} if no such mbean is
registered
+ */
+ public ManagedProcessorMBean getProcessorMBean(String id, boolean steps) {
+ return (steps ? stepsById : processorsById).get(id);
+ }
+
+ private void addProcessorMBean(ObjectName name, ObjectName registeredName,
ManagedProcessorMBean mbean) {
+ ProcessorMBean pm = new ProcessorMBean(
+ registeredName, mbean, mbean.getRouteId(),
mbean.getProcessorId(),
+ mbean instanceof ManagedStepMBean);
+ processorMBeans.put(name, pm);
+ if (pm.routeId() != null) {
+ (pm.step() ? routeSteps :
routeProcessors).computeIfAbsent(pm.routeId(), k ->
ConcurrentHashMap.newKeySet())
+ .add(registeredName);
+ }
+ if (pm.id() != null) {
+ // keep the first if a custom name strategy registers more mbeans
with the same id
+ (pm.step() ? stepsById : processorsById).putIfAbsent(pm.id(),
mbean);
+ }
+ }
+
+ private void removeProcessorMBean(ObjectName name) {
+ ProcessorMBean old = processorMBeans.remove(name);
+ if (old != null) {
+ if (old.routeId() != null) {
+ (old.step() ? routeSteps :
routeProcessors).computeIfPresent(old.routeId(), (k, names) -> {
+ names.remove(old.registeredName());
+ return names.isEmpty() ? null : names;
+ });
+ }
+ if (old.id() != null) {
+ (old.step() ? stepsById : processorsById).remove(old.id(),
old.mbean());
+ }
+ }
+ }
+
@Override
public boolean isRegistered(ObjectName name) {
if (server == null) {
@@ -450,7 +529,7 @@ public class DefaultManagementAgent extends ServiceSupport
implements Management
ServiceHelper.stopService(assembler);
}
- private void registerMBeanWithServer(Object obj, ObjectName name, boolean
forceRegistration)
+ private ObjectName registerMBeanWithServer(Object obj, ObjectName name,
boolean forceRegistration)
throws JMException {
// have we already registered the bean, there can be shared instances
in the camel routes
@@ -477,7 +556,9 @@ public class DefaultManagementAgent extends ServiceSupport
implements Management
ObjectName registeredName = instance.getObjectName();
LOG.debug("Registered MBean with ObjectName: {}", registeredName);
mbeansRegistered.put(name, registeredName);
+ return registeredName;
}
+ return null;
}
protected void createMBeanServer() {
@@ -506,4 +587,7 @@ public class DefaultManagementAgent extends ServiceSupport
implements Management
return MBeanServerFactory.createMBeanServer(mBeanServerDefaultDomain);
}
+ private record ProcessorMBean(
+ ObjectName registeredName, ManagedProcessorMBean mbean, String
routeId, String id, boolean step) {
+ }
}
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/ManagedCamelContextImpl.java
b/core/camel-management/src/main/java/org/apache/camel/management/ManagedCamelContextImpl.java
index 4a1e17ae2aa0..296f9c19956a 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/ManagedCamelContextImpl.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/ManagedCamelContextImpl.java
@@ -34,6 +34,7 @@ import
org.apache.camel.api.management.mbean.ManagedProcessorMBean;
import org.apache.camel.api.management.mbean.ManagedRouteGroupMBean;
import org.apache.camel.api.management.mbean.ManagedRouteMBean;
import org.apache.camel.api.management.mbean.ManagedStepMBean;
+import org.apache.camel.management.mbean.ManagedProcessor;
import org.apache.camel.model.Model;
import org.apache.camel.model.ProcessorDefinition;
import org.apache.camel.spi.ManagementStrategy;
@@ -60,9 +61,17 @@ public class ManagedCamelContextImpl implements
ManagedCamelContext {
return null;
}
- Processor processor = camelContext.getProcessor(id);
- ProcessorDefinition<?> def
- =
camelContext.getCamelContextExtension().getContextPlugin(Model.class).getProcessorDefinition(id);
+ Processor processor;
+ ProcessorDefinition<?> def;
+ if (getManagementStrategy().getManagementAgent() instanceof
DefaultManagementAgent agent
+ && agent.getProcessorMBean(id, false) instanceof
ManagedProcessor mp) {
+ // the processor is managed, so use its processor and definition
instead of walking every route
+ processor = mp.getProcessor();
+ def = mp.getDefinition();
+ } else {
+ processor = camelContext.getProcessor(id);
+ def =
camelContext.getCamelContextExtension().getContextPlugin(Model.class).getProcessorDefinition(id);
+ }
// processor may be null if its anonymous inner class or as lambda
if (def != null) {
@@ -85,9 +94,17 @@ public class ManagedCamelContextImpl implements
ManagedCamelContext {
return null;
}
- Processor processor = camelContext.getProcessor(id);
- ProcessorDefinition<?> def
- =
camelContext.getCamelContextExtension().getContextPlugin(Model.class).getProcessorDefinition(id);
+ Processor processor;
+ ProcessorDefinition<?> def;
+ if (getManagementStrategy().getManagementAgent() instanceof
DefaultManagementAgent agent
+ && agent.getProcessorMBean(id, true) instanceof
ManagedProcessor mp) {
+ // the step is managed, so use its processor and definition
instead of walking every route
+ processor = mp.getProcessor();
+ def = mp.getDefinition();
+ } else {
+ processor = camelContext.getProcessor(id);
+ def =
camelContext.getCamelContextExtension().getContextPlugin(Model.class).getProcessorDefinition(id);
+ }
// processor may be null if its anonymous inner class or as lambda
if (def != null) {
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
index 85efb01eb4cc..083ff5caa506 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedRoute.java
@@ -26,16 +26,11 @@ import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import java.util.Set;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeUnit;
-import javax.management.AttributeValueExp;
import javax.management.MBeanServer;
import javax.management.ObjectName;
-import javax.management.Query;
-import javax.management.QueryExp;
-import javax.management.StringValueExp;
import javax.management.openmbean.CompositeData;
import javax.management.openmbean.CompositeDataSupport;
import javax.management.openmbean.CompositeType;
@@ -55,11 +50,13 @@ import
org.apache.camel.api.management.mbean.ManagedProcessorMBean;
import org.apache.camel.api.management.mbean.ManagedRouteMBean;
import org.apache.camel.api.management.mbean.ManagedStepMBean;
import org.apache.camel.api.management.mbean.RouteError;
+import org.apache.camel.management.DefaultManagementAgent;
import org.apache.camel.model.Model;
import org.apache.camel.model.ModelCamelContext;
import org.apache.camel.model.RouteDefinition;
import org.apache.camel.model.RoutesDefinition;
import org.apache.camel.spi.InflightRepository;
+import org.apache.camel.spi.ManagementAgent;
import org.apache.camel.spi.ManagementStrategy;
import org.apache.camel.spi.RoutePolicy;
import org.apache.camel.support.ExchangeHelper;
@@ -530,21 +527,8 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
sb.append(" <processorStats>\n");
MBeanServer server =
getContext().getManagementStrategy().getManagementAgent().getMBeanServer();
if (server != null) {
- // get all the processor mbeans and sort them accordingly to
their index
- String prefix =
getContext().getManagementStrategy().getManagementAgent().getIncludeHostName()
? "*/" : "";
- ObjectName query = ObjectName.getInstance(
- jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=processors,*");
- Set<ObjectName> names = server.queryNames(query, null);
- List<ManagedProcessorMBean> mps = new ArrayList<>();
- for (ObjectName on : names) {
- ManagedProcessorMBean processor =
context.getManagementStrategy().getManagementAgent().newProxyClient(on,
- ManagedProcessorMBean.class);
-
- // the processor must belong to this route
- if (getRouteId().equals(processor.getRouteId())) {
- mps.add(processor);
- }
- }
+ // get all the processor mbeans of this route and sort them
accordingly to their index
+ List<ManagedProcessorMBean> mps = routeProcessorMBeans(false,
ManagedProcessorMBean.class);
mps.sort(new OrderProcessorMBeans());
// walk the processors in reverse order, and calculate the
accumulated total time
@@ -649,21 +633,8 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
arr = new JsonArray();
MBeanServer server =
getContext().getManagementStrategy().getManagementAgent().getMBeanServer();
if (server != null) {
- // get all the processor mbeans and sort them accordingly to
their index
- String prefix =
getContext().getManagementStrategy().getManagementAgent().getIncludeHostName()
? "*/" : "";
- ObjectName query = ObjectName.getInstance(
- jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=processors,*");
- Set<ObjectName> names = server.queryNames(query, null);
- List<ManagedProcessorMBean> mps = new ArrayList<>();
- for (ObjectName on : names) {
- ManagedProcessorMBean processor =
context.getManagementStrategy().getManagementAgent().newProxyClient(on,
- ManagedProcessorMBean.class);
-
- // the processor must belong to this route
- if (getRouteId().equals(processor.getRouteId())) {
- mps.add(processor);
- }
- }
+ // get all the processor mbeans of this route and sort them
accordingly to their index
+ List<ManagedProcessorMBean> mps = routeProcessorMBeans(false,
ManagedProcessorMBean.class);
mps.sort(new OrderProcessorMBeans());
// walk the processors in reverse order, and calculate the
accumulated total time
@@ -724,21 +695,8 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
sb.append(" <stepStats>\n");
MBeanServer server =
getContext().getManagementStrategy().getManagementAgent().getMBeanServer();
if (server != null) {
- // get all the processor mbeans and sort them accordingly to their
index
- String prefix =
getContext().getManagementStrategy().getManagementAgent().getIncludeHostName()
? "*/" : "";
- ObjectName query = ObjectName
- .getInstance(jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=steps,*");
- Set<ObjectName> names = server.queryNames(query, null);
- List<ManagedStepMBean> mps = new ArrayList<>();
- for (ObjectName on : names) {
- ManagedStepMBean step
- =
context.getManagementStrategy().getManagementAgent().newProxyClient(on,
ManagedStepMBean.class);
-
- // the step must belong to this route
- if (getRouteId().equals(step.getRouteId())) {
- mps.add(step);
- }
- }
+ // get all the step mbeans of this route and sort them accordingly
to their index
+ List<ManagedStepMBean> mps = routeProcessorMBeans(true,
ManagedStepMBean.class);
mps.sort(new OrderProcessorMBeans());
// and now add the sorted list of steps to the xml output
@@ -790,20 +748,8 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
MBeanServer server =
getContext().getManagementStrategy().getManagementAgent().getMBeanServer();
if (server != null) {
- String prefix =
getContext().getManagementStrategy().getManagementAgent().getIncludeHostName()
? "*/" : "";
- List<ManagedProcessorMBean> processors = new ArrayList<>();
- // gather all the processors for this CamelContext, which requires
JMX
- ObjectName query = ObjectName
- .getInstance(jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=processors,*");
- Set<ObjectName> names = server.queryNames(query, null);
- for (ObjectName on : names) {
- ManagedProcessorMBean processor
- =
context.getManagementStrategy().getManagementAgent().newProxyClient(on,
ManagedProcessorMBean.class);
- // the processor must belong to this route
- if (getRouteId().equals(processor.getRouteId())) {
- processors.add(processor);
- }
- }
+ // gather all the processors for this route, which requires JMX
+ List<ManagedProcessorMBean> processors =
routeProcessorMBeans(false, ManagedProcessorMBean.class);
processors.sort(new OrderProcessorMBeans());
// grab route consumer
@@ -818,16 +764,13 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
escapeXml(route.getRouteId()), escapeXml(id),
0, escapeXml(location), line));
}
for (ManagedProcessorMBean processor : processors) {
- // the step must belong to this route
- if (route.getRouteId().equals(processor.getRouteId())) {
- int line = processor.getSourceLineNumber() != null ?
processor.getSourceLineNumber() : -1;
- String location = processor.getSourceLocation() != null ?
processor.getSourceLocation() : "";
- sb.append("\n <routeLocation")
- .append(String.format(
- " routeId=\"%s\" id=\"%s\" index=\"%s\"
sourceLocation=\"%s\" sourceLineNumber=\"%s\"/>",
- escapeXml(route.getRouteId()),
escapeXml(processor.getProcessorId()), processor.getIndex(),
- escapeXml(location), line));
- }
+ int line = processor.getSourceLineNumber() != null ?
processor.getSourceLineNumber() : -1;
+ String location = processor.getSourceLocation() != null ?
processor.getSourceLocation() : "";
+ sb.append("\n <routeLocation")
+ .append(String.format(
+ " routeId=\"%s\" id=\"%s\" index=\"%s\"
sourceLocation=\"%s\" sourceLineNumber=\"%s\"/>",
+ escapeXml(route.getRouteId()),
escapeXml(processor.getProcessorId()), processor.getIndex(),
+ escapeXml(location), line));
}
}
sb.append("\n</routeLocations>");
@@ -848,17 +791,12 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
if (includeProcessors) {
MBeanServer server =
getContext().getManagementStrategy().getManagementAgent().getMBeanServer();
if (server != null) {
- // get all the processor mbeans and sort them accordingly to
their index
- String prefix =
getContext().getManagementStrategy().getManagementAgent().getIncludeHostName()
? "*/" : "";
- // the route id must be equal (match would treat * and ? in
the route id as wildcards)
- QueryExp queryExp = Query.eq(new AttributeValueExp("RouteId"),
new StringValueExp(getRouteId()));
// steps are registered as their own type
- for (String type : new String[] { "processors", "steps" }) {
- ObjectName query = ObjectName.getInstance(
- jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=" + type + ",*");
- Set<ObjectName> names = server.queryNames(query, queryExp);
- for (ObjectName name : names) {
- server.invoke(name, "reset", null, null);
+ for (boolean steps : new boolean[] { false, true }) {
+ for (ObjectName name : routeProcessorMBeanNames(steps)) {
+ if (server.isRegistered(name)) {
+ server.invoke(name, "reset", null, null);
+ }
}
}
}
@@ -1007,22 +945,58 @@ public class ManagedRoute extends
ManagedPerformanceCounter implements ManagedRo
MBeanServer server =
getContext().getManagementStrategy().getManagementAgent().getMBeanServer();
if (server != null) {
- String prefix =
getContext().getManagementStrategy().getManagementAgent().getIncludeHostName()
? "*/" : "";
- // gather all the processors for this CamelContext, which requires
JMX
+ // gather all the processors for this route, which requires JMX
+ for (ManagedProcessorMBean processor : routeProcessorMBeans(false,
ManagedProcessorMBean.class)) {
+ ids.add(processor.getProcessorId());
+ }
+ }
+
+ return ids;
+ }
+
+ /**
+ * Gets the processor (or step) mbeans of this route.
+ */
+ private <T extends ManagedProcessorMBean> List<T>
routeProcessorMBeans(boolean steps, Class<T> type) throws Exception {
+ ManagementAgent agent =
getContext().getManagementStrategy().getManagementAgent();
+ List<T> answer = new ArrayList<>();
+ for (ObjectName on : routeProcessorMBeanNames(steps)) {
+ T mp = agent.newProxyClient(on, type);
+ if (mp != null) {
+ answer.add(mp);
+ }
+ }
+ return answer;
+ }
+
+ /**
+ * Gets the names of the processor (or step) mbeans of this route.
+ * <p/>
+ * The default management agent knows the mbeans of each route, otherwise
all the processor mbeans of the
+ * CamelContext are queried, which is slow with many routes as the route
id of each mbean must be read.
+ */
+ private List<ObjectName> routeProcessorMBeanNames(boolean steps) throws
Exception {
+ ManagementAgent agent =
getContext().getManagementStrategy().getManagementAgent();
+ if (agent instanceof DefaultManagementAgent dma) {
+ return dma.getRouteProcessorMBeanNames(getRouteId(), steps);
+ }
+
+ List<ObjectName> answer = new ArrayList<>();
+ MBeanServer server = agent.getMBeanServer();
+ if (server != null) {
+ String prefix = agent.getIncludeHostName() ? "*/" : "";
+ String type = steps ? "steps" : "processors";
ObjectName query = ObjectName
- .getInstance(jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=processors,*");
- Set<ObjectName> names = server.queryNames(query, null);
- for (ObjectName on : names) {
- ManagedProcessorMBean processor
- =
context.getManagementStrategy().getManagementAgent().newProxyClient(on,
ManagedProcessorMBean.class);
+ .getInstance(jmxDomain + ":context=" + prefix +
getContext().getManagementName() + ",type=" + type + ",*");
+ for (ObjectName on : server.queryNames(query, null)) {
+ ManagedProcessorMBean processor = agent.newProxyClient(on,
ManagedProcessorMBean.class);
// the processor must belong to this route
- if (getRouteId().equals(processor.getRouteId())) {
- ids.add(processor.getProcessorId());
+ if (processor != null &&
getRouteId().equals(processor.getRouteId())) {
+ answer.add(on);
}
}
}
-
- return ids;
+ return answer;
}
private Integer getInflightExchanges() {
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedRouteProcessorsOfOtherRoutesTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRouteProcessorsOfOtherRoutesTest.java
new file mode 100644
index 000000000000..6b8cab55a078
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedRouteProcessorsOfOtherRoutesTest.java
@@ -0,0 +1,135 @@
+/*
+ * 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.management.ManagementFactory;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Proxy;
+import java.util.Set;
+import java.util.TreeSet;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.api.management.ManagedCamelContext;
+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.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+/**
+ * A route mbean must not read the processor mbeans of the other routes to
find its own processors, as calling it for
+ * every route is then quadratic in the size of the integration (CAMEL-25267).
+ */
+@DisabledOnOs(OS.AIX)
+public class ManagedRouteProcessorsOfOtherRoutesTest extends
ManagementTestSupport {
+
+ // the processor and step mbeans that have been read
+ private final Set<String> touched = new TreeSet<>();
+ private volatile boolean recording;
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext context = super.createCamelContext();
+ MBeanServer platform = ManagementFactory.getPlatformMBeanServer();
+ MBeanServer server = (MBeanServer)
Proxy.newProxyInstance(getClass().getClassLoader(),
+ new Class<?>[] { MBeanServer.class }, (proxy, method, args) ->
{
+ if (recording && args != null && args.length > 1 &&
args[0] instanceof ObjectName on) {
+ String name = method.getName();
+ if (name.equals("getAttribute") ||
name.equals("getAttributes") || name.equals("invoke")) {
+ record(on);
+ } else if (name.equals("queryNames") && args[1] !=
null) {
+ // the query expression is evaluated against every
mbean matching the pattern
+ platform.queryNames(on,
null).forEach(this::record);
+ }
+ }
+ try {
+ return method.invoke(platform, args);
+ } catch (InvocationTargetException e) {
+ throw e.getCause();
+ }
+ });
+
context.getManagementStrategy().getManagementAgent().setMBeanServer(server);
+ return context;
+ }
+
+ private void record(ObjectName on) {
+ String type = on.getKeyProperty("type");
+ if ("processors".equals(type) || "steps".equals(type)) {
+ touched.add(ObjectName.unquote(on.getKeyProperty("name")));
+ }
+ }
+
+ private ManagedCamelContext mcc() {
+ return
context.getCamelContextExtension().getContextPlugin(ManagedCamelContext.class);
+ }
+
+ @Test
+ public void testProcessorIds() throws Exception {
+ recording = true;
+ Set<String> ids = new
TreeSet<>(mcc().getManagedRoute("foo").processorIds());
+ recording = false;
+
+ assertEquals(Set.of("foo-log", "foo-to"), ids);
+ assertEquals(Set.of("foo-log", "foo-to"), touched);
+ }
+
+ @Test
+ public void testReset() throws Exception {
+ template.sendBody("direct:foo", "Hello");
+ template.sendBody("direct:bar", "Hello");
+
+ recording = true;
+ mcc().getManagedRoute("foo").reset(true);
+ recording = false;
+
+ assertEquals(Set.of("foo-log", "foo-step", "foo-to"), touched);
+ assertEquals(0,
mcc().getManagedProcessor("foo-to").getExchangesTotal());
+ assertEquals(0, mcc().getManagedStep("foo-step").getExchangesTotal());
+ assertEquals(1,
mcc().getManagedProcessor("bar-to").getExchangesTotal());
+ }
+
+ @Test
+ public void testRemoveRoute() throws Exception {
+ context.getRouteController().stopRoute("bar");
+ context.removeRoute("bar");
+
+ assertNull(mcc().getManagedProcessor("bar-to"));
+ assertNull(mcc().getManagedStep("bar-step"));
+ assertEquals(Set.of("foo-log", "foo-to"), new
TreeSet<>(mcc().getManagedRoute("foo").processorIds()));
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:foo").routeId("foo")
+ .step("foo-step").log("${body}").id("foo-log").end()
+ .to("mock:foo").id("foo-to");
+
+ from("direct:bar").routeId("bar")
+ .step("bar-step").log("${body}").id("bar-log").end()
+ .to("mock:bar").id("bar-to");
+ }
+ };
+ }
+}