This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch feature/CAMEL-25502-route-stopping-exception in repository https://gitbox.apache.org/repos/asf/camel.git
commit 6fbd6392e2524d742bbe90a47b9d5c89d04d9872 Author: Claus Ibsen <[email protected]> AuthorDate: Fri Oct 9 23:39:39 2026 +0200 CAMEL-25502: camel-core - RouteStoppingException for an exchange cut off by a route or CamelContext stop An exchange cut off because its route or the CamelContext is being stopped (a route stop, a dev-mode reload, the graceful shutdown timing out) failed with a plain RejectedExecutionException, which could not be told apart from a real rejection (a full thread pool or queue): CAMEL-25484 recognised it by its message, and the error registry recorded it as an error. - org.apache.camel.RouteStoppingException extends RejectedExecutionException, so onException and catch blocks for RejectedExecutionException still match - set where an in-flight exchange is cut off by a stop: the error handler, a direct producer interrupted while waiting for a consumer, Delay, Throttle, Threads, Failover, Resequencer, and the internal processor on a forced shutdown; real rejections stay RejectedExecutionException - the error handler logs any such cut-off as one WARN line by its type (also a CamelContext stop) - the error registry leaves it out, unless camel.errorRegistry.includeRouteStopping=true; the option is shown in JMX and the errors dev console Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_01STT6whBgK1AqsSsUKrnE8m --- .../apache/camel/catalog/dev-consoles-openapi.json | 7 ++- .../apache/camel/catalog/dev-consoles/errors.json | 7 ++- .../org/apache/camel/catalog/docs/main.adoc | 3 +- .../main/camel-main-configuration-metadata.json | 1 + .../camel/component/direct/DirectProducer.java | 12 ++--- .../org/apache/camel/RouteStoppingException.java | 39 +++++++++++++++ .../java/org/apache/camel/spi/ErrorRegistry.java | 14 ++++++ .../camel/impl/engine/CamelInternalProcessor.java | 4 +- .../camel/impl/engine/DefaultErrorRegistry.java | 17 +++++++ .../impl/engine/SharedCamelInternalProcessor.java | 4 +- .../org/apache/camel/dev-console/errors.json | 7 ++- .../camel/impl/console/ErrorRegistryConsole.java | 4 +- .../impl/console/ErrorRegistryConsoleTest.java | 2 + .../apache/camel/processor/AbstractThrottler.java | 4 +- .../processor/ConcurrentRequestsThrottler.java | 3 +- .../camel/processor/DelayProcessorSupport.java | 7 +-- .../apache/camel/processor/ThreadsProcessor.java | 5 +- .../camel/processor/TotalRequestsThrottler.java | 3 +- .../errorhandler/RedeliveryErrorHandler.java | 19 ++++---- .../loadbalancer/FailOverLoadBalancer.java | 4 +- .../processor/resequencer/ResequencerEngine.java | 7 +-- .../camel/impl/engine/ShutdownWaitingAtTest.java | 6 +-- .../errorhandler/RouteStopCutOffLogTest.java | 39 ++++++++++++++- ...rRegistryConfigurationPropertiesConfigurer.java | 7 +++ .../camel-main-configuration-metadata.json | 1 + core/camel-main/src/main/docs/main.adoc | 3 +- .../org/apache/camel/main/BaseMainSupport.java | 1 + .../main/ErrorRegistryConfigurationProperties.java | 25 ++++++++++ .../ErrorRegistryConfigurationPropertiesTest.java | 55 ++++++++++++++++++++++ .../mbean/ManagedErrorRegistryMBean.java | 6 +++ .../management/mbean/ManagedErrorRegistry.java | 10 ++++ .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 22 +++++++-- .../modules/ROOT/pages/error-registry.adoc | 1 + 33 files changed, 304 insertions(+), 45 deletions(-) diff --git a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles-openapi.json b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles-openapi.json index c413cba69c2b..bc39ea57c610 100644 --- a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles-openapi.json +++ b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles-openapi.json @@ -1735,6 +1735,10 @@ "type": "string", "description": "The time to live for entries" }, + "includeRouteStopping": { + "type": "boolean", + "description": "Whether exchanges cut off by a route stop, a route reload or a CamelContext stop are also captured" + }, "errors": { "type": "array", "items": { @@ -1747,7 +1751,8 @@ "required": [ "enabled", "size", - "maximumEntries" + "maximumEntries", + "includeRouteStopping" ] } } diff --git a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles/errors.json b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles/errors.json index 8b44a9600807..de288992e4ec 100644 --- a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles/errors.json +++ b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/dev-consoles/errors.json @@ -118,6 +118,10 @@ "type": "string", "description": "The time to live for entries" }, + "includeRouteStopping": { + "type": "boolean", + "description": "Whether exchanges cut off by a route stop, a route reload or a CamelContext stop are also captured" + }, "errors": { "type": "array", "items": { @@ -130,7 +134,8 @@ "required": [ "enabled", "size", - "maximumEntries" + "maximumEntries", + "includeRouteStopping" ] } } diff --git a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/main.adoc b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/main.adoc index 5e3ffea4d8df..4d300fd6eb9c 100644 --- a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/main.adoc +++ b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/main.adoc @@ -748,7 +748,7 @@ The camel.lra supports 5 options, which are listed below. === Camel Error Registry configurations -The camel.errorRegistry supports 9 options, which are listed below. +The camel.errorRegistry supports 10 options, which are listed below. [width="100%",cols="2,5,^1,2",options="header"] |=== @@ -759,6 +759,7 @@ The camel.errorRegistry supports 9 options, which are listed below. | `camel.errorRegistry.enabled` | Whether the error registry is enabled to capture errors during message routing. | false | boolean | `camel.errorRegistry.includeExchangeProperties` | Whether to include the exchange properties in the captured error data. | true | boolean | `camel.errorRegistry.includeExchangeVariables` | Whether to include the exchange variables in the captured error data. | true | boolean +| `camel.errorRegistry.includeRouteStopping` | Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop (RouteStoppingException). Nothing in the route failed, but a consumer that does not roll back may lose the message, so enable it to see how often stops cut off work. | false | boolean | `camel.errorRegistry.maximumEntries` | The maximum number of error entries to keep in the registry. When the limit is exceeded, the oldest entries are evicted. | 100 | int | `camel.errorRegistry.maximumEntriesPerKind` | The maximum number of error entries of the same kind (same route, node and exception type) to keep, so a storm of one failure does not evict all the other errors. The counter of that kind keeps rising even when its older entries are evicted. | 3 | int | `camel.errorRegistry.timeToLiveSeconds` | The time-to-live in seconds for error entries. Entries older than this are evicted. The default value is 0 (disabled). | 0 | int diff --git a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/main/camel-main-configuration-metadata.json b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/main/camel-main-configuration-metadata.json index 8d59e0dc9e73..a1e0eca1dfa0 100644 --- a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/main/camel-main-configuration-metadata.json +++ b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/main/camel-main-configuration-metadata.json @@ -266,6 +266,7 @@ { "name": "camel.errorRegistry.enabled", "required": false, "description": "Whether the error registry is enabled to capture errors during message routing.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue": false, "secret": false }, { "name": "camel.errorRegistry.includeExchangeProperties", "required": false, "description": "Whether to include the exchange properties in the captured error data.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue": true, "secret": false }, { "name": "camel.errorRegistry.includeExchangeVariables", "required": false, "description": "Whether to include the exchange variables in the captured error data.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue": true, "secret": false }, + { "name": "camel.errorRegistry.includeRouteStopping", "required": false, "description": "Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop (RouteStoppingException). Nothing in the route failed, but a consumer that does not roll back may lose the message, so enable it to see how often stops cut off work.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue" [...] { "name": "camel.errorRegistry.maximumEntries", "required": false, "description": "The maximum number of error entries to keep in the registry. When the limit is exceeded, the oldest entries are evicted.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "integer", "javaType": "int", "defaultValue": 100, "secret": false }, { "name": "camel.errorRegistry.maximumEntriesPerKind", "required": false, "description": "The maximum number of error entries of the same kind (same route, node and exception type) to keep, so a storm of one failure does not evict all the other errors. The counter of that kind keeps rising even when its older entries are evicted.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "integer", "javaType": "int", "defaultValue": 3, "secret": false }, { "name": "camel.errorRegistry.timeToLiveSeconds", "required": false, "description": "The time-to-live in seconds for error entries. Entries older than this are evicted. The default value is 0 (disabled).", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "integer", "javaType": "int", "defaultValue": 0, "secret": false }, diff --git a/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java b/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java index 4f7adf8e5f1c..12be998da2f0 100644 --- a/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java +++ b/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java @@ -19,6 +19,7 @@ package org.apache.camel.component.direct; import org.apache.camel.AsyncCallback; import org.apache.camel.Exchange; import org.apache.camel.ExchangePropertyKey; +import org.apache.camel.RouteStoppingException; import org.apache.camel.support.DefaultAsyncProducer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -104,13 +105,12 @@ public class DirectProducer extends DefaultAsyncProducer { LOG.info("Interrupted while waiting for a consumer on {}: the route is being stopped or reloaded", endpoint.getEndpointUri()); Thread.currentThread().interrupt(); - DirectConsumerNotAvailableException cause = new DirectConsumerNotAvailableException( + // a cut-off by a route stop, not a failure (CAMEL-25502); the interruption is kept as the cause, so + // onException(InterruptedException.class) still matches + exchange.setException(new RouteStoppingException( "No consumers available on endpoint: " + endpoint - + " (interrupted while waiting for one, as the route is being stopped or reloaded)", - exchange); - // keep the interruption as the cause, so onException(InterruptedException.class) still matches - cause.initCause(e); - exchange.setException(cause); + + " (interrupted while waiting for one, as the route is being stopped or reloaded)", + e)); // stay marked as interrupted, as setException(InterruptedException) did, so the error handler stops // routing instead of handling a failure (onException, redelivery, dead letter channel, the ERROR log) exchange.getExchangeExtension().setInterrupted(true); diff --git a/core/camel-api/src/main/java/org/apache/camel/RouteStoppingException.java b/core/camel-api/src/main/java/org/apache/camel/RouteStoppingException.java new file mode 100644 index 000000000000..19b4507c72ca --- /dev/null +++ b/core/camel-api/src/main/java/org/apache/camel/RouteStoppingException.java @@ -0,0 +1,39 @@ +/* + * 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; + +import java.util.concurrent.RejectedExecutionException; + +/** + * An exchange cut off because its route, or the {@link CamelContext}, is being stopped while the exchange is in flight: + * a route stop, a route reload in dev mode, or the graceful shutdown timing out. Nothing in the route failed; a + * consumer that rolls back, such as file, delivers the message again. + * <p/> + * It extends {@link RejectedExecutionException}, which Camel used for this before, so an + * {@code onException(RejectedExecutionException.class)} or a catch of it still matches. A real rejection (a thread pool + * or a queue that is full) stays a plain {@link RejectedExecutionException}. + */ +public class RouteStoppingException extends RejectedExecutionException { + + public RouteStoppingException(String message) { + super(message); + } + + public RouteStoppingException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/core/camel-api/src/main/java/org/apache/camel/spi/ErrorRegistry.java b/core/camel-api/src/main/java/org/apache/camel/spi/ErrorRegistry.java index 17b505be71b4..ca3dd80f6801 100644 --- a/core/camel-api/src/main/java/org/apache/camel/spi/ErrorRegistry.java +++ b/core/camel-api/src/main/java/org/apache/camel/spi/ErrorRegistry.java @@ -165,4 +165,18 @@ public interface ErrorRegistry extends ErrorRegistryView, StaticService { * This is by default enabled. */ void setIncludeExchangeVariables(boolean includeExchangeVariables); + + /** + * Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop (a + * {@link org.apache.camel.RouteStoppingException}). Nothing in the route failed, but a consumer that does not roll + * back may lose the message, so enable it to see how often stops cut off work. + */ + boolean isIncludeRouteStopping(); + + /** + * Sets whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop. + * <p/> + * This is by default disabled. + */ + void setIncludeRouteStopping(boolean includeRouteStopping); } diff --git a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/CamelInternalProcessor.java b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/CamelInternalProcessor.java index bbaa2efb053b..1a22992d7ac1 100644 --- a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/CamelInternalProcessor.java +++ b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/CamelInternalProcessor.java @@ -23,7 +23,6 @@ import java.util.Objects; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CopyOnWriteArrayList; -import java.util.concurrent.RejectedExecutionException; import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; @@ -39,6 +38,7 @@ import org.apache.camel.NonManagedService; import org.apache.camel.Ordered; import org.apache.camel.Processor; import org.apache.camel.Route; +import org.apache.camel.RouteStoppingException; import org.apache.camel.StatefulService; import org.apache.camel.impl.debugger.BacklogTracer; import org.apache.camel.impl.debugger.DefaultBacklogTracerEventMessage; @@ -370,7 +370,7 @@ public class CamelInternalProcessor extends DelegateAsyncProcessor implements In + exchange; LOG.debug(msg); if (exchange.getException() == null) { - exchange.setException(new RejectedExecutionException(msg)); + exchange.setException(new RouteStoppingException(msg)); } // force shutdown so we should not continue originalCallback.done(true); diff --git a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java index c47a53cdfcbe..f21546f64038 100644 --- a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java +++ b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultErrorRegistry.java @@ -31,6 +31,7 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.camel.Exchange; import org.apache.camel.ExchangePropertyKey; import org.apache.camel.MessageHistory; +import org.apache.camel.RouteStoppingException; import org.apache.camel.spi.BacklogErrorEventMessage; import org.apache.camel.spi.CamelEvent; import org.apache.camel.spi.ErrorRegistry; @@ -38,6 +39,7 @@ import org.apache.camel.spi.ErrorRegistryView; import org.apache.camel.support.EventNotifierSupport; import org.apache.camel.support.LoggerHelper; import org.apache.camel.support.MessageHelper; +import org.apache.camel.util.ObjectHelper; import org.apache.camel.util.json.JsonObject; import org.apache.camel.util.json.Jsonable; import org.apache.camel.util.json.Jsoner; @@ -61,6 +63,7 @@ public class DefaultErrorRegistry extends EventNotifierSupport implements ErrorR private volatile boolean bodyIncludeFiles = true; private volatile boolean includeExchangeProperties = true; private volatile boolean includeExchangeVariables = true; + private volatile boolean includeRouteStopping; public DefaultErrorRegistry() { setIgnoreCamelContextEvents(true); @@ -120,6 +123,10 @@ public class DefaultErrorRegistry extends EventNotifierSupport implements ErrorR if (exception == null) { return; } + if (!includeRouteStopping && ObjectHelper.getException(RouteStoppingException.class, exception) != null) { + // cut off by a route stop, a reload or a CamelContext stop: nothing in the route failed (CAMEL-25502) + return; + } // for correlated copy exchanges (e.g., created by circuit breaker, multicast, splitter) // use the original exchange ID so the error is tracked under the parent exchange @@ -517,6 +524,16 @@ public class DefaultErrorRegistry extends EventNotifierSupport implements ErrorR this.includeExchangeVariables = includeExchangeVariables; } + @Override + public boolean isIncludeRouteStopping() { + return includeRouteStopping; + } + + @Override + public void setIncludeRouteStopping(boolean includeRouteStopping) { + this.includeRouteStopping = includeRouteStopping; + } + /** * A filtered view over entries for a specific route. */ diff --git a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/SharedCamelInternalProcessor.java b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/SharedCamelInternalProcessor.java index 5ab5b8ddf66d..bdf7a06db2bf 100644 --- a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/SharedCamelInternalProcessor.java +++ b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/SharedCamelInternalProcessor.java @@ -18,13 +18,13 @@ package org.apache.camel.impl.engine; import java.util.Objects; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.RejectedExecutionException; import org.apache.camel.AsyncCallback; import org.apache.camel.AsyncProcessor; import org.apache.camel.CamelContext; import org.apache.camel.Exchange; import org.apache.camel.Processor; +import org.apache.camel.RouteStoppingException; import org.apache.camel.spi.AsyncProcessorAwaitManager; import org.apache.camel.spi.CamelInternalProcessorAdvice; import org.apache.camel.spi.ReactiveExecutor; @@ -246,7 +246,7 @@ public class SharedCamelInternalProcessor implements SharedInternalProcessor { + exchange; LOG.debug(msg); if (exchange.getException() == null) { - exchange.setException(new RejectedExecutionException(msg)); + exchange.setException(new RouteStoppingException(msg)); } } return false; diff --git a/core/camel-console/src/generated/resources/META-INF/org/apache/camel/dev-console/errors.json b/core/camel-console/src/generated/resources/META-INF/org/apache/camel/dev-console/errors.json index 8b44a9600807..de288992e4ec 100644 --- a/core/camel-console/src/generated/resources/META-INF/org/apache/camel/dev-console/errors.json +++ b/core/camel-console/src/generated/resources/META-INF/org/apache/camel/dev-console/errors.json @@ -118,6 +118,10 @@ "type": "string", "description": "The time to live for entries" }, + "includeRouteStopping": { + "type": "boolean", + "description": "Whether exchanges cut off by a route stop, a route reload or a CamelContext stop are also captured" + }, "errors": { "type": "array", "items": { @@ -130,7 +134,8 @@ "required": [ "enabled", "size", - "maximumEntries" + "maximumEntries", + "includeRouteStopping" ] } } diff --git a/core/camel-console/src/main/java/org/apache/camel/impl/console/ErrorRegistryConsole.java b/core/camel-console/src/main/java/org/apache/camel/impl/console/ErrorRegistryConsole.java index 583853dcc37c..8e45d2606fe4 100644 --- a/core/camel-console/src/main/java/org/apache/camel/impl/console/ErrorRegistryConsole.java +++ b/core/camel-console/src/main/java/org/apache/camel/impl/console/ErrorRegistryConsole.java @@ -38,6 +38,7 @@ public class ErrorRegistryConsole extends AbstractDevConsole { @Metadata(description = "Number of captured errors") int size, @Metadata(description = "The maximum number of entries retained") int maximumEntries, @Metadata(description = "The time to live for entries") String timeToLive, + @Metadata(description = "Whether exchanges cut off by a route stop, a route reload or a CamelContext stop are also captured") boolean includeRouteStopping, @Metadata(description = "The captured errors; shape depends on the underlying message implementation") List<Map<String, Object>> errors) { } @@ -76,6 +77,7 @@ public class ErrorRegistryConsole extends AbstractDevConsole { ErrorRegistry registry = getCamelContext().getErrorRegistry(); sb.append(String.format("%n Enabled: %s", registry.isEnabled())); sb.append(String.format("%n Size: %s", registry.size())); + sb.append(String.format("%n Include Route Stopping: %s", registry.isIncludeRouteStopping())); List<BacklogErrorEventMessage> entries = fetchAndFilter(registry, options); @@ -133,7 +135,7 @@ public class ErrorRegistryConsole extends AbstractDevConsole { Response response = new Response( registry.isEnabled(), registry.size(), registry.getMaximumEntries(), registry.getTimeToLive().toString(), - errors); + registry.isIncludeRouteStopping(), errors); return JsonRecordSupport.toJsonObject(response); } diff --git a/core/camel-console/src/test/java/org/apache/camel/impl/console/ErrorRegistryConsoleTest.java b/core/camel-console/src/test/java/org/apache/camel/impl/console/ErrorRegistryConsoleTest.java index 73ae8a778742..c1811b0a81f5 100644 --- a/core/camel-console/src/test/java/org/apache/camel/impl/console/ErrorRegistryConsoleTest.java +++ b/core/camel-console/src/test/java/org/apache/camel/impl/console/ErrorRegistryConsoleTest.java @@ -21,6 +21,7 @@ import org.apache.camel.util.json.JsonArray; import org.apache.camel.util.json.JsonObject; import org.junit.jupiter.api.Test; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; public class ErrorRegistryConsoleTest extends AbstractDevConsoleTest { @@ -33,6 +34,7 @@ public class ErrorRegistryConsoleTest extends AbstractDevConsoleTest { assertNotNull(out.getBoolean("enabled")); assertNotNull(out.getInteger("size")); assertNotNull(out.getInteger("maximumEntries")); + assertEquals(Boolean.FALSE, out.getBoolean("includeRouteStopping")); assertNotNull(out.getString("timeToLive")); JsonArray errors = out.getJsonArray("errors"); diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/AbstractThrottler.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/AbstractThrottler.java index 725e51122aa8..3a7bb0b973cd 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/AbstractThrottler.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/AbstractThrottler.java @@ -16,13 +16,13 @@ */ package org.apache.camel.processor; -import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; import org.apache.camel.Exchange; import org.apache.camel.Expression; +import org.apache.camel.RouteStoppingException; import org.apache.camel.Traceable; import org.apache.camel.spi.IdAware; import org.apache.camel.spi.RouteIdAware; @@ -70,7 +70,7 @@ public abstract class AbstractThrottler extends BaseProcessorSupport String msg = "Run not allowed as ShutdownStrategy is forcing shutting down, will reject executing exchange: " + exchange; LOG.debug(msg); - exchange.setException(new RejectedExecutionException(msg, e)); + exchange.setException(new RouteStoppingException(msg, e)); } else { exchange.setException(e); } diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/ConcurrentRequestsThrottler.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/ConcurrentRequestsThrottler.java index afd90ce238e0..f89805341521 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/ConcurrentRequestsThrottler.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/ConcurrentRequestsThrottler.java @@ -32,6 +32,7 @@ import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; import org.apache.camel.Exchange; import org.apache.camel.Expression; +import org.apache.camel.RouteStoppingException; import org.apache.camel.RuntimeExchangeException; import org.apache.camel.spi.Synchronization; import org.apache.camel.util.ObjectHelper; @@ -89,7 +90,7 @@ public class ConcurrentRequestsThrottler extends AbstractThrottler { try { if (!isRunAllowed()) { - throw new RejectedExecutionException("Run is not allowed"); + throw new RouteStoppingException("Run is not allowed"); } return doProcess(exchange, callback, state, queuedStart, doneSync); diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/DelayProcessorSupport.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/DelayProcessorSupport.java index 0854e4c05d53..7193c9ab1a72 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/DelayProcessorSupport.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/DelayProcessorSupport.java @@ -25,6 +25,7 @@ import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; import org.apache.camel.Exchange; import org.apache.camel.Processor; +import org.apache.camel.RouteStoppingException; import org.apache.camel.util.ObjectHelper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -63,7 +64,7 @@ public abstract class DelayProcessorSupport extends BaseDelegateProcessorSupport LOG.trace("Delayed task woke up and continues routing for exchangeId: {}", exchange.getExchangeId()); } if (!isRunAllowed()) { - exchange.setException(new RejectedExecutionException("Run is not allowed")); + exchange.setException(new RouteStoppingException("Run is not allowed")); } // process the exchange now that we woke up @@ -131,7 +132,7 @@ public abstract class DelayProcessorSupport extends BaseDelegateProcessorSupport delayedCount.decrementAndGet(); if (isCallerRunsWhenRejected()) { if (!isRunAllowed()) { - exchange.setException(new RejectedExecutionException()); + exchange.setException(new RouteStoppingException("Run is not allowed")); } else { if (LOG.isDebugEnabled()) { LOG.debug( @@ -161,7 +162,7 @@ public abstract class DelayProcessorSupport extends BaseDelegateProcessorSupport @Override public boolean process(Exchange exchange, AsyncCallback callback) { if (!isRunAllowed()) { - exchange.setException(new RejectedExecutionException("Run is not allowed")); + exchange.setException(new RouteStoppingException("Run is not allowed")); callback.done(true); return true; } diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/ThreadsProcessor.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/ThreadsProcessor.java index a5a7e820475b..1629d4859a46 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/ThreadsProcessor.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/ThreadsProcessor.java @@ -24,6 +24,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; import org.apache.camel.Exchange; +import org.apache.camel.RouteStoppingException; import org.apache.camel.spi.IdAware; import org.apache.camel.spi.RouteIdAware; import org.apache.camel.spi.StepIdAware; @@ -78,7 +79,7 @@ public class ThreadsProcessor extends BaseProcessorSupport implements IdAware, R public void run() { LOG.trace("Continue routing exchange {}", exchange); if (shutdown.get()) { - exchange.setException(new RejectedExecutionException("ThreadsProcessor is not running.")); + exchange.setException(new RouteStoppingException("ThreadsProcessor is not running.")); } callback.done(done); } @@ -89,7 +90,7 @@ public class ThreadsProcessor extends BaseProcessorSupport implements IdAware, R exchange.setException(new RejectedExecutionException()); LOG.trace("Rejected routing exchange {}", exchange); if (shutdown.get()) { - exchange.setException(new RejectedExecutionException("ThreadsProcessor is not running.")); + exchange.setException(new RouteStoppingException("ThreadsProcessor is not running.")); } callback.done(done); } diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/TotalRequestsThrottler.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/TotalRequestsThrottler.java index a7f4708bad72..06363fa377cc 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/TotalRequestsThrottler.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/TotalRequestsThrottler.java @@ -32,6 +32,7 @@ import org.apache.camel.AsyncCallback; import org.apache.camel.CamelContext; import org.apache.camel.Exchange; import org.apache.camel.Expression; +import org.apache.camel.RouteStoppingException; import org.apache.camel.RuntimeExchangeException; import org.apache.camel.util.ObjectHelper; import org.slf4j.Logger; @@ -88,7 +89,7 @@ public class TotalRequestsThrottler extends AbstractThrottler { try { if (!isRunAllowed()) { - throw new RejectedExecutionException("Run is not allowed"); + throw new RouteStoppingException("Run is not allowed"); } String key = DEFAULT_KEY; diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java index e7494304853c..f0281e9857b6 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java @@ -35,6 +35,7 @@ import org.apache.camel.Navigate; import org.apache.camel.Predicate; import org.apache.camel.Processor; import org.apache.camel.Route; +import org.apache.camel.RouteStoppingException; import org.apache.camel.RuntimeCamelException; import org.apache.camel.processor.PooledExchangeTask; import org.apache.camel.processor.PooledExchangeTaskFactory; @@ -858,7 +859,7 @@ public abstract class RedeliveryErrorHandler extends ErrorHandlerSupport private void runNotAllowed() { LOG.trace("Run not allowed, will reject executing exchange: {}", exchange); if (exchange.getException() == null) { - exchange.setException(new RejectedExecutionException(notAllowedReason())); + exchange.setException(new RouteStoppingException(notAllowedReason())); } AsyncCallback cb = callback; taskFactory.release(this); @@ -1121,7 +1122,7 @@ public abstract class RedeliveryErrorHandler extends ErrorHandlerSupport if (!isRunAllowed()) { LOG.trace("Run not allowed, will reject executing exchange: {}", exchange); if (exchange.getException() == null) { - exchange.setException(new RejectedExecutionException(notAllowedReason())); + exchange.setException(new RouteStoppingException(notAllowedReason())); } AsyncCallback cb = callback; taskFactory.release(this); @@ -2231,16 +2232,16 @@ public abstract class RedeliveryErrorHandler extends ErrorHandlerSupport } /** - * Whether the failure is an exchange cut off by a route stop or a dev mode reload (not a forced CamelContext stop): - * nothing in the route failed, and a consumer that rolls back, such as file, delivers it again (CAMEL-25484). + * Whether the failure is an exchange cut off by a route stop, a dev mode reload or a CamelContext stop: nothing in + * the route failed, and a consumer that rolls back, such as file, delivers it again (CAMEL-25484, CAMEL-25502). */ static boolean isCutOffByRouteStop(Throwable e) { - return e instanceof RejectedExecutionException && ROUTE_STOPPING_REASON.equals(e.getMessage()); + return e instanceof RouteStoppingException; } /** - * Logs an exchange that a route stop or a dev mode reload cut off as one line at most WARN, without the message - * history and stack trace of a failure, as nothing in the route failed (CAMEL-25484). + * Logs an exchange that a route stop, a dev mode reload or a CamelContext stop cut off as one line at most WARN, + * without the message history and stack trace of a failure, as nothing in the route failed (CAMEL-25484). * * @return whether the exchange was cut off, and so is logged */ @@ -2248,8 +2249,10 @@ public abstract class RedeliveryErrorHandler extends ErrorHandlerSupport if (!isCutOffByRouteStop(e) || exchange.isRollbackOnly() || exchange.isRollbackOnlyLast()) { return false; } + String reason = shutdownStrategy.isForceShutdown() + ? "the CamelContext is being stopped" : "its route is being stopped or reloaded"; logger.log("Exchange cut off " + ExchangeHelper.logIds(exchange) + failureOrigin(exchange) - + ": its route is being stopped or reloaded, so it is not continued." + + ": " + reason + ", so it is not continued." + " A consumer that rolls back, such as file, delivers it again", level == LoggingLevel.ERROR ? LoggingLevel.WARN : level); return true; diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/loadbalancer/FailOverLoadBalancer.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/loadbalancer/FailOverLoadBalancer.java index d6be8d9a4ea4..1e565805a8da 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/loadbalancer/FailOverLoadBalancer.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/loadbalancer/FailOverLoadBalancer.java @@ -17,7 +17,6 @@ package org.apache.camel.processor.loadbalancer; import java.util.List; -import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.atomic.AtomicInteger; import org.apache.camel.AsyncCallback; @@ -25,6 +24,7 @@ import org.apache.camel.AsyncProcessor; import org.apache.camel.CamelContext; import org.apache.camel.CamelContextAware; import org.apache.camel.Exchange; +import org.apache.camel.RouteStoppingException; import org.apache.camel.Traceable; import org.apache.camel.support.ExchangeHelper; import org.apache.camel.util.ObjectHelper; @@ -220,7 +220,7 @@ public class FailOverLoadBalancer extends LoadBalancerSupport implements Traceab if (!isRunAllowed()) { LOG.trace("Run not allowed, will reject executing exchange: {}", exchange); if (exchange.getException() == null) { - exchange.setException(new RejectedExecutionException()); + exchange.setException(new RouteStoppingException("Run is not allowed")); } // we cannot process so invoke callback callback.done(false); diff --git a/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java b/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java index 5b5f910e7f18..780cfd200051 100644 --- a/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java +++ b/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java @@ -25,6 +25,7 @@ import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.Predicate; +import org.apache.camel.RouteStoppingException; import org.apache.camel.util.concurrent.ThreadHelper; /** @@ -193,7 +194,7 @@ public class ResequencerEngine<E> { lock.lock(); try { if (stopped) { - throw new RejectedExecutionException("Resequencer is stopped"); + throw new RouteStoppingException("Resequencer is stopped"); } if (pred.test(sequence)) { return; @@ -207,7 +208,7 @@ public class ResequencerEngine<E> { latch.await(); // compare with the stop count rather than read the flag, as a quick restart may have reset it already if (stopped || stopCount != stopCountAtWait) { - throw new RejectedExecutionException("Resequencer is stopped"); + throw new RouteStoppingException("Resequencer is stopped"); } } @@ -304,7 +305,7 @@ public class ResequencerEngine<E> { try { // a stopped resequencer has cancelled its timer, so the element could not be scheduled for timing out if (stopped) { - throw new RejectedExecutionException("Resequencer is stopped"); + throw new RouteStoppingException("Resequencer is stopped"); } // wrap object into internal element diff --git a/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java b/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java index 3e7e95683c1e..f3aa7a6d8372 100644 --- a/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java +++ b/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java @@ -25,8 +25,8 @@ import java.util.concurrent.atomic.AtomicReference; import org.apache.camel.CamelContext; import org.apache.camel.ContextTestSupport; import org.apache.camel.Exchange; +import org.apache.camel.RouteStoppingException; import org.apache.camel.builder.RouteBuilder; -import org.apache.camel.component.direct.DirectConsumerNotAvailableException; import org.apache.camel.component.log.ConsumingAppender; import org.apache.camel.spi.CamelEvent; import org.apache.camel.spi.RouteStartupOrder; @@ -104,8 +104,8 @@ public class ShutdownWaitingAtTest extends ContextTestSupport { await().atMost(5, TimeUnit.SECONDS).until(() -> failure.get() != null); Exchange exchange = failure.get(); - DirectConsumerNotAvailableException e - = assertInstanceOf(DirectConsumerNotAvailableException.class, exchange.getException()); + // a cut-off by the route stop, not a failure (CAMEL-25502) + RouteStoppingException e = assertInstanceOf(RouteStoppingException.class, exchange.getException()); assertTrue(e.getMessage().contains("direct://shipment"), e.getMessage()); assertTrue(e.getMessage().contains("interrupted while waiting for one, as the route is being stopped or reloaded"), e.getMessage()); diff --git a/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RouteStopCutOffLogTest.java b/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RouteStopCutOffLogTest.java index 70cc9a6946b8..4d33c34bee03 100644 --- a/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RouteStopCutOffLogTest.java +++ b/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RouteStopCutOffLogTest.java @@ -22,6 +22,7 @@ import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; import org.apache.camel.ContextTestSupport; +import org.apache.camel.RouteStoppingException; import org.apache.camel.builder.RouteBuilder; import org.apache.camel.component.log.ConsumingAppender; import org.apache.logging.log4j.Level; @@ -37,7 +38,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue; /** * CAMEL-25484: an exchange that a route stop or a dev mode reload cuts off is logged as one WARN line, not as an - * exhausted failure with message history and stack trace. + * exhausted failure with message history and stack trace; CAMEL-25502: and is not recorded in the error registry. */ public class RouteStopCutOffLogTest extends ContextTestSupport { @@ -84,6 +85,42 @@ public class RouteStopCutOffLogTest extends ContextTestSupport { assertFalse(message.endsWith("THROWN"), message); } + @Test + public void theCutOffIsNotInTheErrorRegistry() throws Exception { + context.getErrorRegistry().setEnabled(true); + context.getShutdownStrategy().setTimeout(1); + context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS); + context.getInflightRepository().setInflightBrowseEnabled(true); + + template.sendBody("seda:start", List.of("A", "B")); + await().atMost(5, TimeUnit.SECONDS).until(() -> context.getInflightRepository().browse().stream() + .anyMatch(e -> "shipment".equals(e.getNodeId()))); + context.getRouteController().stopRoute("slow"); + + // the cut-off is a RouteStoppingException (the error handler's and the direct producer's): nothing failed + await().atMost(5, TimeUnit.SECONDS).until(() -> context.getRouteController().getRouteStatus("slow").isStopped()); + assertEquals(0, context.getErrorRegistry().size(), () -> context.getErrorRegistry().browse().toString()); + } + + @Test + public void theCutOffIsInTheErrorRegistryWhenIncluded() throws Exception { + context.getErrorRegistry().setEnabled(true); + context.getErrorRegistry().setIncludeRouteStopping(true); + context.getShutdownStrategy().setTimeout(1); + context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS); + context.getInflightRepository().setInflightBrowseEnabled(true); + + template.sendBody("seda:start", List.of("A", "B")); + await().atMost(5, TimeUnit.SECONDS).until(() -> context.getInflightRepository().browse().stream() + .anyMatch(e -> "shipment".equals(e.getNodeId()))); + context.getRouteController().stopRoute("slow"); + + await().atMost(5, TimeUnit.SECONDS).until(() -> context.getErrorRegistry().size() > 0); + assertTrue(context.getErrorRegistry().browse().stream() + .allMatch(e -> RouteStoppingException.class.getName().equals(e.getExceptionType())), + () -> context.getErrorRegistry().browse().toString()); + } + @Override protected RouteBuilder createRouteBuilder() { return new RouteBuilder() { diff --git a/core/camel-main/src/generated/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesConfigurer.java b/core/camel-main/src/generated/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesConfigurer.java index 71765daf0be9..025a8ddcc363 100644 --- a/core/camel-main/src/generated/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesConfigurer.java +++ b/core/camel-main/src/generated/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesConfigurer.java @@ -28,6 +28,7 @@ public class ErrorRegistryConfigurationPropertiesConfigurer extends org.apache.c map.put("Enabled", boolean.class); map.put("IncludeExchangeProperties", boolean.class); map.put("IncludeExchangeVariables", boolean.class); + map.put("IncludeRouteStopping", boolean.class); map.put("MaximumEntries", int.class); map.put("MaximumEntriesPerKind", int.class); map.put("TimeToLiveSeconds", int.class); @@ -49,6 +50,8 @@ public class ErrorRegistryConfigurationPropertiesConfigurer extends org.apache.c case "includeExchangeProperties": target.setIncludeExchangeProperties(property(camelContext, boolean.class, value)); return true; case "includeexchangevariables": case "includeExchangeVariables": target.setIncludeExchangeVariables(property(camelContext, boolean.class, value)); return true; + case "includeroutestopping": + case "includeRouteStopping": target.setIncludeRouteStopping(property(camelContext, boolean.class, value)); return true; case "maximumentries": case "maximumEntries": target.setMaximumEntries(property(camelContext, int.class, value)); return true; case "maximumentriesperkind": @@ -78,6 +81,8 @@ public class ErrorRegistryConfigurationPropertiesConfigurer extends org.apache.c case "includeExchangeProperties": return boolean.class; case "includeexchangevariables": case "includeExchangeVariables": return boolean.class; + case "includeroutestopping": + case "includeRouteStopping": return boolean.class; case "maximumentries": case "maximumEntries": return int.class; case "maximumentriesperkind": @@ -103,6 +108,8 @@ public class ErrorRegistryConfigurationPropertiesConfigurer extends org.apache.c case "includeExchangeProperties": return target.isIncludeExchangeProperties(); case "includeexchangevariables": case "includeExchangeVariables": return target.isIncludeExchangeVariables(); + case "includeroutestopping": + case "includeRouteStopping": return target.isIncludeRouteStopping(); case "maximumentries": case "maximumEntries": return target.getMaximumEntries(); case "maximumentriesperkind": diff --git a/core/camel-main/src/generated/resources/META-INF/camel-main-configuration-metadata.json b/core/camel-main/src/generated/resources/META-INF/camel-main-configuration-metadata.json index 8d59e0dc9e73..a1e0eca1dfa0 100644 --- a/core/camel-main/src/generated/resources/META-INF/camel-main-configuration-metadata.json +++ b/core/camel-main/src/generated/resources/META-INF/camel-main-configuration-metadata.json @@ -266,6 +266,7 @@ { "name": "camel.errorRegistry.enabled", "required": false, "description": "Whether the error registry is enabled to capture errors during message routing.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue": false, "secret": false }, { "name": "camel.errorRegistry.includeExchangeProperties", "required": false, "description": "Whether to include the exchange properties in the captured error data.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue": true, "secret": false }, { "name": "camel.errorRegistry.includeExchangeVariables", "required": false, "description": "Whether to include the exchange variables in the captured error data.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue": true, "secret": false }, + { "name": "camel.errorRegistry.includeRouteStopping", "required": false, "description": "Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop (RouteStoppingException). Nothing in the route failed, but a consumer that does not roll back may lose the message, so enable it to see how often stops cut off work.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "boolean", "javaType": "boolean", "defaultValue" [...] { "name": "camel.errorRegistry.maximumEntries", "required": false, "description": "The maximum number of error entries to keep in the registry. When the limit is exceeded, the oldest entries are evicted.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "integer", "javaType": "int", "defaultValue": 100, "secret": false }, { "name": "camel.errorRegistry.maximumEntriesPerKind", "required": false, "description": "The maximum number of error entries of the same kind (same route, node and exception type) to keep, so a storm of one failure does not evict all the other errors. The counter of that kind keeps rising even when its older entries are evicted.", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "integer", "javaType": "int", "defaultValue": 3, "secret": false }, { "name": "camel.errorRegistry.timeToLiveSeconds", "required": false, "description": "The time-to-live in seconds for error entries. Entries older than this are evicted. The default value is 0 (disabled).", "sourceType": "org.apache.camel.main.ErrorRegistryConfigurationProperties", "type": "integer", "javaType": "int", "defaultValue": 0, "secret": false }, diff --git a/core/camel-main/src/main/docs/main.adoc b/core/camel-main/src/main/docs/main.adoc index 5e3ffea4d8df..4d300fd6eb9c 100644 --- a/core/camel-main/src/main/docs/main.adoc +++ b/core/camel-main/src/main/docs/main.adoc @@ -748,7 +748,7 @@ The camel.lra supports 5 options, which are listed below. === Camel Error Registry configurations -The camel.errorRegistry supports 9 options, which are listed below. +The camel.errorRegistry supports 10 options, which are listed below. [width="100%",cols="2,5,^1,2",options="header"] |=== @@ -759,6 +759,7 @@ The camel.errorRegistry supports 9 options, which are listed below. | `camel.errorRegistry.enabled` | Whether the error registry is enabled to capture errors during message routing. | false | boolean | `camel.errorRegistry.includeExchangeProperties` | Whether to include the exchange properties in the captured error data. | true | boolean | `camel.errorRegistry.includeExchangeVariables` | Whether to include the exchange variables in the captured error data. | true | boolean +| `camel.errorRegistry.includeRouteStopping` | Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop (RouteStoppingException). Nothing in the route failed, but a consumer that does not roll back may lose the message, so enable it to see how often stops cut off work. | false | boolean | `camel.errorRegistry.maximumEntries` | The maximum number of error entries to keep in the registry. When the limit is exceeded, the oldest entries are evicted. | 100 | int | `camel.errorRegistry.maximumEntriesPerKind` | The maximum number of error entries of the same kind (same route, node and exception type) to keep, so a storm of one failure does not evict all the other errors. The counter of that kind keeps rising even when its older entries are evicted. | 3 | int | `camel.errorRegistry.timeToLiveSeconds` | The time-to-live in seconds for error entries. Entries older than this are evicted. The default value is 0 (disabled). | 0 | int diff --git a/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java b/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java index 8799470ba586..5e470b81c067 100644 --- a/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java +++ b/core/camel-main/src/main/java/org/apache/camel/main/BaseMainSupport.java @@ -2730,6 +2730,7 @@ public abstract class BaseMainSupport extends BaseService { registry.setBodyIncludeFiles(config.isBodyIncludeFiles()); registry.setIncludeExchangeProperties(config.isIncludeExchangeProperties()); registry.setIncludeExchangeVariables(config.isIncludeExchangeVariables()); + registry.setIncludeRouteStopping(config.isIncludeRouteStopping()); } private void setAiObservabilityProperties( diff --git a/core/camel-main/src/main/java/org/apache/camel/main/ErrorRegistryConfigurationProperties.java b/core/camel-main/src/main/java/org/apache/camel/main/ErrorRegistryConfigurationProperties.java index 8404baf80097..fb445c0ebf80 100644 --- a/core/camel-main/src/main/java/org/apache/camel/main/ErrorRegistryConfigurationProperties.java +++ b/core/camel-main/src/main/java/org/apache/camel/main/ErrorRegistryConfigurationProperties.java @@ -46,6 +46,8 @@ public class ErrorRegistryConfigurationProperties implements BootstrapCloseable private boolean includeExchangeProperties = true; @Metadata(defaultValue = "true") private boolean includeExchangeVariables = true; + @Metadata + private boolean includeRouteStopping; public ErrorRegistryConfigurationProperties(MainConfigurationProperties parent) { this.parent = parent; @@ -166,6 +168,19 @@ public class ErrorRegistryConfigurationProperties implements BootstrapCloseable this.includeExchangeVariables = includeExchangeVariables; } + public boolean isIncludeRouteStopping() { + return includeRouteStopping; + } + + /** + * Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop + * (RouteStoppingException). Nothing in the route failed, but a consumer that does not roll back may lose the + * message, so enable it to see how often stops cut off work. + */ + public void setIncludeRouteStopping(boolean includeRouteStopping) { + this.includeRouteStopping = includeRouteStopping; + } + // -- fluent builder methods -- /** @@ -244,4 +259,14 @@ public class ErrorRegistryConfigurationProperties implements BootstrapCloseable this.includeExchangeVariables = includeExchangeVariables; return this; } + + /** + * Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop + * (RouteStoppingException). Nothing in the route failed, but a consumer that does not roll back may lose the + * message, so enable it to see how often stops cut off work. + */ + public ErrorRegistryConfigurationProperties withIncludeRouteStopping(boolean includeRouteStopping) { + this.includeRouteStopping = includeRouteStopping; + return this; + } } diff --git a/core/camel-main/src/test/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesTest.java b/core/camel-main/src/test/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesTest.java new file mode 100644 index 000000000000..26b049e44b2b --- /dev/null +++ b/core/camel-main/src/test/java/org/apache/camel/main/ErrorRegistryConfigurationPropertiesTest.java @@ -0,0 +1,55 @@ +/* + * 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.main; + +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * CAMEL-25502: the error registry captures an exchange cut off by a route stop only when configured to. + */ +class ErrorRegistryConfigurationPropertiesTest { + + @Test + void routeStoppingIsLeftOutByDefault() throws Exception { + Main main = new Main(); + main.addInitialProperty("camel.errorRegistry.enabled", "true"); + + main.start(); + try { + assertThat(main.getCamelContext().getErrorRegistry().isEnabled()).isTrue(); + assertThat(main.getCamelContext().getErrorRegistry().isIncludeRouteStopping()).isFalse(); + } finally { + main.stop(); + } + } + + @Test + void routeStoppingIsIncludedWhenConfigured() throws Exception { + Main main = new Main(); + main.addInitialProperty("camel.errorRegistry.enabled", "true"); + main.addInitialProperty("camel.errorRegistry.includeRouteStopping", "true"); + + main.start(); + try { + assertThat(main.getCamelContext().getErrorRegistry().isIncludeRouteStopping()).isTrue(); + } finally { + main.stop(); + } + } +} diff --git a/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/ManagedErrorRegistryMBean.java b/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/ManagedErrorRegistryMBean.java index 0f0dc2624c9a..04707ccd973c 100644 --- a/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/ManagedErrorRegistryMBean.java +++ b/core/camel-management-api/src/main/java/org/apache/camel/api/management/mbean/ManagedErrorRegistryMBean.java @@ -80,6 +80,12 @@ public interface ManagedErrorRegistryMBean extends ManagedServiceMBean { @ManagedAttribute(description = "Whether to include exchange variables") void setIncludeExchangeVariables(boolean includeExchangeVariables); + @ManagedAttribute(description = "Whether to also capture exchanges cut off by a route stop, reload or CamelContext stop") + boolean isIncludeRouteStopping(); + + @ManagedAttribute(description = "Whether to also capture exchanges cut off by a route stop, reload or CamelContext stop") + void setIncludeRouteStopping(boolean includeRouteStopping); + @ManagedOperation(description = "Browse all error entries") TabularData browse(); diff --git a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java index 759d45c46f37..d80968ff89c4 100644 --- a/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java +++ b/core/camel-management/src/main/java/org/apache/camel/management/mbean/ManagedErrorRegistry.java @@ -142,6 +142,16 @@ public class ManagedErrorRegistry extends ManagedService implements ManagedError errorRegistry.setIncludeExchangeVariables(includeExchangeVariables); } + @Override + public boolean isIncludeRouteStopping() { + return errorRegistry.isIncludeRouteStopping(); + } + + @Override + public void setIncludeRouteStopping(boolean includeRouteStopping) { + errorRegistry.setIncludeRouteStopping(includeRouteStopping); + } + @Override public TabularData browse() { return browseEntries(errorRegistry.browse()); diff --git a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc index 8195e780d257..f022c9640d33 100644 --- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc +++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc @@ -1546,9 +1546,25 @@ When a route is stopped or reloaded (for example a save in `camel run --dev`) wh the graceful shutdown times out, the exchange is cut off. The error handler logged this as a failed delivery: an ERROR with "Exhausted after delivery attempt", the message history and the stack trace of the `RejectedExecutionException`. It is now one line at WARN, "Exchange cut off ...: its route is being stopped or -reloaded", as nothing in the route failed, and a consumer that rolls back, such as file, delivers the message again. -A level lower than WARN configured for exhausted logging is kept. A forced stop of the CamelContext is logged as -before. The exception set on the exchange is unchanged. +reloaded" (or "the CamelContext is being stopped"), as nothing in the route failed, and a consumer that rolls back, +such as file, delivers the message again. A level lower than WARN configured for exhausted logging is kept. + +=== camel-core - RouteStoppingException for an exchange cut off by a stop + +An exchange cut off because its route, or the CamelContext, is being stopped (a route stop, a route reload in dev +mode, or the graceful shutdown timing out) now fails with the new `org.apache.camel.RouteStoppingException` instead +of a plain `RejectedExecutionException`. It extends `RejectedExecutionException`, so an +`onException(RejectedExecutionException.class)` or a catch of it still matches; a real rejection, a thread pool or a +queue that is full, stays a plain `RejectedExecutionException`. It is set by the error handler, the Delay, Throttle, +Threads, Failover load balancer and Resequencer EIPs while they stop, and by a forced shutdown. + +A `direct` producer that was waiting for a consumer when a forced stop interrupted it now sets a +`RouteStoppingException` (with the `InterruptedException` as its cause) instead of a +`DirectConsumerNotAvailableException`. + +Such a cut-off is not recorded in the error registry, as nothing in the route failed; set +`camel.errorRegistry.includeRouteStopping=true` to record it as well, for example to see how often stops cut off +work that a consumer which does not roll back may lose. === camel-core - stopping or removing all routes is one graceful shutdown diff --git a/docs/user-manual/modules/ROOT/pages/error-registry.adoc b/docs/user-manual/modules/ROOT/pages/error-registry.adoc index 354e7d00507b..7a50f4fbde54 100644 --- a/docs/user-manual/modules/ROOT/pages/error-registry.adoc +++ b/docs/user-manual/modules/ROOT/pages/error-registry.adoc @@ -36,6 +36,7 @@ camel.errorRegistry.enabled = true | `camel.errorRegistry.bodyIncludeFiles` | `true` | boolean | Whether to include file-based message bodies. | `camel.errorRegistry.includeExchangeProperties` | `true` | boolean | Whether to include exchange properties in the snapshot. | `camel.errorRegistry.includeExchangeVariables` | `true` | boolean | Whether to include exchange variables in the snapshot. +| `camel.errorRegistry.includeRouteStopping` | `false` | boolean | Whether to also capture an exchange cut off by a route stop, a route reload or a CamelContext stop (`RouteStoppingException`). Nothing in the route failed, but a consumer that does not roll back may lose the message, so enable it to see how often stops cut off work. |=== == Captured Data
