This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-25067-25071-25072-small in repository https://gitbox.apache.org/repos/asf/camel.git
commit e4db0cc56fe80e5a9bc5ad4a8d61bc2ec1294bf5 Author: Claus Ibsen <[email protected]> AuthorDate: Tue Sep 29 11:46:39 2026 +0200 CAMEL-25067, CAMEL-25071, CAMEL-25072: camel-core - small follow-ups from the deep review - threads: a keepAliveTime given as a duration (such as 1m30s) is converted to the time unit, instead of the milliseconds being used as seconds. A plain number is still in the time unit. - The JSON route stats dump of the CamelContext includes exchangesInflight for each route, as the XML dump does. - The component verifier MBean returns INCOMPLETE_PARAMETER_GROUP for that error code, instead of ILLEGAL_PARAMETER_GROUP_COMBINATION. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Signed-off-by: Claus Ibsen <[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(); } });
