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 aabd31342e31 CAMEL-25067, CAMEL-25071, CAMEL-25072: camel-core - small
follow-ups from the deep review (#27071)
aabd31342e31 is described below
commit aabd31342e31a904f77cc316a386de61a8258383
Author: Claus Ibsen <[email protected]>
AuthorDate: Tue Sep 29 13:59:58 2026 +0200
CAMEL-25067, CAMEL-25071, CAMEL-25072: camel-core - small follow-ups from
the deep review (#27071)
* CAMEL-25067, CAMEL-25071, CAMEL-25072: camel-core - small follow-ups from
the deep review
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
---
.../org/apache/camel/model/ThreadsDefinition.java | 3 +-
.../org/apache/camel/reifier/ThreadsReifier.java | 28 ++++++-
.../camel/processor/ThreadsKeepAliveTimeTest.java | 63 +++++++++++++++
.../management/mbean/ManagedCamelContext.java | 1 +
.../camel/management/mbean/ManagedComponent.java | 2 +-
.../ManagedCamelContextDumpStatsAsJSonTest.java | 93 ++++++++++++++++++++++
.../camel/management/ManagedComponentTest.java | 27 +++++++
7 files changed, 212 insertions(+), 5 deletions(-)
diff --git
a/core/camel-core-model/src/main/java/org/apache/camel/model/ThreadsDefinition.java
b/core/camel-core-model/src/main/java/org/apache/camel/model/ThreadsDefinition.java
index 7e0686c345cd..e10a3aaa835e 100644
---
a/core/camel-core-model/src/main/java/org/apache/camel/model/ThreadsDefinition.java
+++
b/core/camel-core-model/src/main/java/org/apache/camel/model/ThreadsDefinition.java
@@ -198,7 +198,8 @@ public class ThreadsDefinition extends
NoOutputDefinition<ThreadsDefinition>
}
/**
- * Sets the keep alive time for idle threads
+ * Sets the keep alive time for idle threads. A plain number is in the
time unit (seconds by default), and a
+ * duration such as 30s or 1m30s is converted to the time unit.
*
* @param keepAliveTime keep alive time
* @return the builder
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ThreadsReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ThreadsReifier.java
index 67db283c8f55..05ac19f08b25 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ThreadsReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/ThreadsReifier.java
@@ -61,9 +61,7 @@ public class ThreadsReifier extends
ProcessorReifier<ThreadsDefinition> {
ThreadPoolProfile profile = new ThreadPoolProfile(name);
profile.setPoolSize(definition.getPoolSize() != null ?
parseInt(definition.getPoolSize()) : null);
profile.setMaxPoolSize(definition.getMaxPoolSize() != null ?
parseInt(definition.getMaxPoolSize()) : null);
- profile.setKeepAliveTime(
- definition.getKeepAliveTime() != null ?
parseDuration(definition.getKeepAliveTime()) : null);
- profile.setTimeUnit(definition.getTimeUnit() != null ?
parse(TimeUnit.class, definition.getTimeUnit()) : null);
+ configureKeepAliveTime(profile);
profile.setMaxQueueSize(definition.getMaxQueueSize() != null ?
parseInt(definition.getMaxQueueSize()) : null);
profile.setRejectedPolicy(policy);
profile.setAllowCoreThreadTimeOut(definition.getAllowCoreThreadTimeOut() != null
@@ -105,6 +103,30 @@ public class ThreadsReifier extends
ProcessorReifier<ThreadsDefinition> {
return answer;
}
+ private void configureKeepAliveTime(ThreadPoolProfile profile) {
+ TimeUnit unit = definition.getTimeUnit() != null ?
parse(TimeUnit.class, definition.getTimeUnit()) : null;
+ String text = parseString(definition.getKeepAliveTime());
+ Long keepAliveTime = null;
+ if (text != null) {
+ text = text.trim();
+ if (!text.isEmpty() && text.chars().allMatch(Character::isDigit)) {
+ // a plain number is in the time unit (seconds by default)
+ keepAliveTime = Long.parseLong(text);
+ } else {
+ // a duration such as 30s or 1m5s is in milliseconds, so
convert it to the time unit
+ long millis = parseDuration(text);
+ if (unit != null) {
+ keepAliveTime = unit.convert(millis,
TimeUnit.MILLISECONDS);
+ } else {
+ keepAliveTime = millis;
+ unit = TimeUnit.MILLISECONDS;
+ }
+ }
+ }
+ profile.setKeepAliveTime(keepAliveTime);
+ profile.setTimeUnit(unit);
+ }
+
protected ThreadPoolRejectedPolicy resolveRejectedPolicy() {
String ref = parseString(definition.getExecutorService());
if (ref != null && definition.getRejectedPolicy() == null) {
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/ThreadsKeepAliveTimeTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/ThreadsKeepAliveTimeTest.java
new file mode 100644
index 000000000000..0d89c01c7d35
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/ThreadsKeepAliveTimeTest.java
@@ -0,0 +1,63 @@
+/*
+ * 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.processor;
+
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+
+public class ThreadsKeepAliveTimeTest extends ContextTestSupport {
+
+ @Test
+ public void testKeepAliveTime() {
+ // a plain number is in the time unit, which is seconds by default
+ assertKeepAliveTime("plain", 10);
+ assertKeepAliveTime("plainMinutes", 120);
+ // a duration is converted to the time unit
+ assertKeepAliveTime("duration", 90);
+ assertKeepAliveTime("durationSeconds", 90);
+ assertKeepAliveTime("durationMillis", 1);
+ }
+
+ private void assertKeepAliveTime(String id, long expectedSeconds) {
+ ThreadsProcessor threads = context.getProcessor(id,
ThreadsProcessor.class);
+ ThreadPoolExecutor pool = assertInstanceOf(ThreadPoolExecutor.class,
threads.getExecutorService());
+ assertEquals(expectedSeconds, pool.getKeepAliveTime(TimeUnit.SECONDS),
id);
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:plain").threads(1,
2).keepAliveTime(10).id("plain").to("mock:result");
+ from("direct:plainMinutes").threads(1,
2).keepAliveTime(2).timeUnit(TimeUnit.MINUTES).id("plainMinutes")
+ .to("mock:result");
+ from("direct:duration").threads(1,
2).keepAliveTime("1m30s").id("duration").to("mock:result");
+ from("direct:durationSeconds").threads(1,
2).keepAliveTime("1m30s").timeUnit(TimeUnit.SECONDS)
+ .id("durationSeconds").to("mock:result");
+ from("direct:durationMillis").threads(1,
2).keepAliveTime("1500ms").id("durationMillis").to("mock:result");
+ }
+ };
+ }
+}
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
index c6ca85c46fc1..d85a2e391bcb 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedCamelContext.java
@@ -802,6 +802,7 @@ public class ManagedCamelContext extends
ManagedPerformanceCounter implements Ma
// use substring as we only want the attributes
route.statsAsJSon(jo, fullStats);
+ jo.put("exchangesInflight", route.getExchangesInflight());
// add processor details if needed
if (includeProcessors) {
diff --git
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedComponent.java
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedComponent.java
index b601bade15d6..e9988da3d577 100644
---
a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedComponent.java
+++
b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedComponent.java
@@ -215,7 +215,7 @@ public class ManagedComponent implements ManagedInstance,
ManagedComponentMBean
return StandardCode.ILLEGAL_PARAMETER_VALUE;
} else if (code
==
org.apache.camel.component.extension.ComponentVerifierExtension.VerificationError.StandardCode.INCOMPLETE_PARAMETER_GROUP)
{
- return StandardCode.ILLEGAL_PARAMETER_GROUP_COMBINATION;
+ return StandardCode.INCOMPLETE_PARAMETER_GROUP;
} else if (code
==
org.apache.camel.component.extension.ComponentVerifierExtension.VerificationError.StandardCode.UNSUPPORTED)
{
return StandardCode.UNSUPPORTED;
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextDumpStatsAsJSonTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextDumpStatsAsJSonTest.java
new file mode 100644
index 000000000000..442e7e03eae6
--- /dev/null
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedCamelContextDumpStatsAsJSonTest.java
@@ -0,0 +1,93 @@
+/*
+ * 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.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import javax.management.MBeanServer;
+import javax.management.ObjectName;
+
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.util.json.JsonArray;
+import org.apache.camel.util.json.JsonObject;
+import org.apache.camel.util.json.Jsoner;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.condition.DisabledOnOs;
+import org.junit.jupiter.api.condition.OS;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+@DisabledOnOs(OS.AIX)
+public class ManagedCamelContextDumpStatsAsJSonTest extends
ManagementTestSupport {
+
+ private final CountDownLatch latch = new CountDownLatch(1);
+
+ @Test
+ public void testRouteExchangesInflight() throws Exception {
+ MBeanServer mbeanServer = getMBeanServer();
+ ObjectName on = getContextObjectName();
+
+ getMockEndpoint("mock:foo").expectedMessageCount(1);
+ getMockEndpoint("mock:bar").expectedMessageCount(1);
+
+ // the exchange in route foo waits on the latch, so it is inflight
while the stats are dumped
+ template.asyncSendBody("direct:start", "Hello World");
+ template.sendBody("direct:bar", "Bye World");
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
context.getInflightRepository().size("foo") == 1);
+
+ try {
+ String json = (String) mbeanServer.invoke(on,
"dumpRouteStatsAsJSon", new Object[] { false, true },
+ new String[] { "boolean", "boolean" });
+ log.info(json);
+
+ JsonObject root = (JsonObject) Jsoner.deserialize(json);
+ assertNotNull(root);
+ assertEquals(1, root.getInteger("exchangesInflight"));
+
+ JsonArray routes = (JsonArray) root.getCollection("routes");
+ assertEquals(2, routes.size());
+ for (Object o : routes) {
+ JsonObject route = (JsonObject) o;
+ int expected = "foo".equals(route.getString("id")) ? 1 : 0;
+ assertEquals(expected, route.getInteger("exchangesInflight"),
"route " + route.getString("id"));
+ }
+ } finally {
+ latch.countDown();
+ }
+
+ assertMockEndpointsSatisfied();
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").routeId("foo")
+ .process(e -> latch.await(20,
TimeUnit.SECONDS)).id("a")
+ .to("mock:foo").id("b");
+
+ from("direct:bar").routeId("bar")
+ .to("mock:bar").id("c");
+ }
+ };
+ }
+
+}
diff --git
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedComponentTest.java
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedComponentTest.java
index 1d65d4cef172..3cca38df5920 100644
---
a/core/camel-management/src/test/java/org/apache/camel/management/ManagedComponentTest.java
+++
b/core/camel-management/src/test/java/org/apache/camel/management/ManagedComponentTest.java
@@ -28,9 +28,12 @@ import org.apache.camel.Endpoint;
import org.apache.camel.api.management.mbean.ComponentVerifierExtension;
import org.apache.camel.api.management.mbean.ComponentVerifierExtension.Result;
import org.apache.camel.api.management.mbean.ComponentVerifierExtension.Scope;
+import
org.apache.camel.api.management.mbean.ComponentVerifierExtension.VerificationError;
import org.apache.camel.component.direct.DirectComponent;
+import
org.apache.camel.component.extension.ComponentVerifierExtension.VerificationError.StandardCode;
import
org.apache.camel.component.extension.verifier.DefaultComponentVerifierExtension;
import org.apache.camel.component.extension.verifier.ResultBuilder;
+import org.apache.camel.component.extension.verifier.ResultErrorBuilder;
import org.apache.camel.support.DefaultComponent;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.DisabledOnOs;
@@ -98,6 +101,24 @@ public class ManagedComponentTest extends
ManagementTestSupport {
assertEquals(Scope.PARAMETERS, res.getScope());
}
+ @Test
+ public void testVerifyErrorCode() throws Exception {
+ MBeanServerConnection mbeanServer = getMBeanServer();
+
+ ObjectName on = getCamelObjectName(TYPE_COMPONENT,
"my-verifiable-component");
+
+ // each standard code is returned as the same code
+ for (StandardCode code : new StandardCode[] {
+ StandardCode.INCOMPLETE_PARAMETER_GROUP,
StandardCode.ILLEGAL_PARAMETER_GROUP_COMBINATION }) {
+ ComponentVerifierExtension.Result res = invoke(mbeanServer, on,
"verify",
+ new Object[] { "parameters", Map.of("errorCode", code) },
VERIFY_SIGNATURE);
+ assertEquals(Result.Status.ERROR, res.getStatus());
+ assertEquals(1, res.getErrors().size());
+ VerificationError.Code actual = res.getErrors().get(0).getCode();
+ assertEquals(code.getName(), actual.getName());
+ }
+ }
+
// ***********************************
//
// ***********************************
@@ -112,6 +133,12 @@ public class ManagedComponentTest extends
ManagementTestSupport {
@Override
protected Result verifyParameters(Map<String, Object>
parameters) {
+ Object code = parameters.get("errorCode");
+ if (code != null) {
+ return
ResultBuilder.withStatusAndScope(Result.Status.ERROR, Scope.PARAMETERS)
+
.error(ResultErrorBuilder.withCode((StandardCode) code).build())
+ .build();
+ }
return ResultBuilder.withStatusAndScope(Result.Status.OK,
Scope.PARAMETERS).build();
}
});