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 054611d9be95 CAMEL-24150, CAMEL-24151: Saga EIP fixes — shared service
lifecycle, code review findings
054611d9be95 is described below
commit 054611d9be95efa31c59ab374ab49999b4582110
Author: Claus Ibsen <[email protected]>
AuthorDate: Fri Jul 17 22:43:33 2026 +0200
CAMEL-24150, CAMEL-24151: Saga EIP fixes — shared service lifecycle, code
review findings
Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
.../org/apache/camel/service/lra/LRAClient.java | 21 ++++--
.../camel/service/lra/LRASagaCoordinator.java | 4 +-
.../apache/camel/service/lra/LRASagaService.java | 24 +++---
.../src/main/docs/modules/eips/pages/saga-eip.adoc | 4 +-
.../processor/saga/MandatorySagaProcessor.java | 5 +-
.../camel/processor/saga/NeverSagaProcessor.java | 5 +-
.../processor/saga/NotSupportedSagaProcessor.java | 5 +-
.../processor/saga/RequiredSagaProcessor.java | 5 +-
.../processor/saga/RequiresNewSagaProcessor.java | 5 +-
.../apache/camel/processor/saga/SagaProcessor.java | 24 ++----
.../camel/processor/saga/SagaProcessorBuilder.java | 20 ++---
.../processor/saga/SupportsSagaProcessor.java | 5 +-
.../java/org/apache/camel/reifier/SagaReifier.java | 12 ++-
.../processor/SagaSharedServiceRouteStopTest.java | 86 ++++++++++++++++++++++
.../apache/camel/saga/InMemorySagaCoordinator.java | 34 ++++++---
.../java/org/apache/camel/main/KameletMain.java | 2 +-
.../apache/camel/main/download/SagaDownloader.java | 14 +++-
17 files changed, 191 insertions(+), 84 deletions(-)
diff --git
a/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRAClient.java
b/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRAClient.java
index 93b3b627d955..3528a0ef62ac 100644
---
a/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRAClient.java
+++
b/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRAClient.java
@@ -125,7 +125,8 @@ public class LRAClient implements Closeable {
String lraEndpoint = lra.toString();
if (step.getTimeoutInMilliseconds().isPresent()) {
- lraEndpoint = lraEndpoint + "?" + HEADER_TIME_LIMIT + "=" +
step.getTimeoutInMilliseconds().get();
+ String separator = lraEndpoint.contains("?") ? "&" : "?";
+ lraEndpoint = lraEndpoint + separator + HEADER_TIME_LIMIT +
"=" + step.getTimeoutInMilliseconds().get();
}
HttpRequest request = prepareRequest(URI.create(lraEndpoint),
exchange)
.setHeader(HEADER_LINK, link.toString())
@@ -138,8 +139,10 @@ public class LRAClient implements Closeable {
}, sagaService.getExecutorService())
.thenCompose(Function.identity())
.thenApply(response -> {
- if (response.statusCode() != HttpURLConnection.HTTP_OK) {
- throw new RuntimeCamelException("Cannot join LRA");
+ int status = response.statusCode();
+ if (status >= HttpURLConnection.HTTP_BAD_REQUEST) {
+ throw new RuntimeCamelException(
+ "Cannot join LRA " + lra + " (HTTP " + status
+ "): " + response.body());
}
return null;
@@ -155,8 +158,10 @@ public class LRAClient implements Closeable {
CompletableFuture<HttpResponse<String>> future =
client.sendAsync(request, HttpResponse.BodyHandlers.ofString());
return future.thenApply(response -> {
- if (response.statusCode() != HttpURLConnection.HTTP_OK) {
- throw new RuntimeCamelException("Cannot complete LRA");
+ int status = response.statusCode();
+ if (status >= HttpURLConnection.HTTP_BAD_REQUEST) {
+ throw new RuntimeCamelException(
+ "Cannot complete LRA " + lra + " (HTTP " + status +
"): " + response.body());
}
return null;
@@ -172,8 +177,10 @@ public class LRAClient implements Closeable {
CompletableFuture<HttpResponse<String>> future =
client.sendAsync(request, HttpResponse.BodyHandlers.ofString());
return future.thenApply(response -> {
- if (response.statusCode() != HttpURLConnection.HTTP_OK) {
- throw new RuntimeCamelException("Cannot compensate LRA");
+ int status = response.statusCode();
+ if (status >= HttpURLConnection.HTTP_BAD_REQUEST) {
+ throw new RuntimeCamelException(
+ "Cannot compensate LRA " + lra + " (HTTP " + status +
"): " + response.body());
}
return null;
diff --git
a/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaCoordinator.java
b/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaCoordinator.java
index 67de593cd285..fd89af76c8b1 100644
---
a/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaCoordinator.java
+++
b/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaCoordinator.java
@@ -41,9 +41,7 @@ public class LRASagaCoordinator implements
CamelSagaCoordinator {
try {
sagaStep = LRASagaStep.fromCamelSagaStep(step, exchange);
} catch (RuntimeException ex) {
- return CompletableFuture.supplyAsync(() -> {
- throw ex;
- });
+ return CompletableFuture.failedFuture(ex);
}
return sagaService.getClient().join(this.lraURL, sagaStep, exchange);
}
diff --git
a/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaService.java
b/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaService.java
index e10eed2b07c9..d7ef791feb75 100644
---
a/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaService.java
+++
b/components/camel-lra/src/main/java/org/apache/camel/service/lra/LRASagaService.java
@@ -79,13 +79,23 @@ public class LRASagaService extends ServiceSupport
implements StaticService, Cam
@Override
public void registerStep(CamelSagaStep step) {
- // Register which uris should be exposed
-
step.getCompensation().map(Endpoint::getEndpointUri).map(this.sagaURIs::add);
-
step.getCompletion().map(Endpoint::getEndpointUri).map(this.sagaURIs::add);
+
step.getCompensation().map(Endpoint::getEndpointUri).ifPresent(this.sagaURIs::add);
+
step.getCompletion().map(Endpoint::getEndpointUri).ifPresent(this.sagaURIs::add);
}
@Override
protected void doStart() throws Exception {
+ if (coordinatorUrl == null) {
+ throw new IllegalStateException("coordinatorUrl must be configured
on the LRA saga service");
+ }
+ if (localParticipantUrl == null) {
+ throw new IllegalStateException("localParticipantUrl must be
configured on the LRA saga service");
+ }
+
+ if (this.routes == null) {
+ this.routes = new LRASagaRoutes(this);
+ camelContext.addRoutes(this.routes);
+ }
if (this.executorService == null) {
this.executorService = camelContext.getExecutorServiceManager()
.newDefaultScheduledThreadPool(this, "saga-lra");
@@ -120,14 +130,6 @@ public class LRASagaService extends ServiceSupport
implements StaticService, Cam
@Override
public void setCamelContext(CamelContext camelContext) {
this.camelContext = camelContext;
- if (this.routes == null) {
- this.routes = new LRASagaRoutes(this);
- try {
- this.camelContext.addRoutes(this.routes);
- } catch (Exception ex) {
- throw RuntimeCamelException.wrapRuntimeException(ex);
- }
- }
}
@Override
diff --git
a/core/camel-core-engine/src/main/docs/modules/eips/pages/saga-eip.adoc
b/core/camel-core-engine/src/main/docs/modules/eips/pages/saga-eip.adoc
index 95a62fe373e5..124128079450 100644
--- a/core/camel-core-engine/src/main/docs/modules/eips/pages/saga-eip.adoc
+++ b/core/camel-core-engine/src/main/docs/modules/eips/pages/saga-eip.adoc
@@ -711,7 +711,9 @@ Sagas are long-running actions, but this does not mean that
they should not have
[NOTE]
====
-The Saga EIP implementation may have a default timeout set on all Sagas that
don't specify it explicitly
+There is no default timeout on Sagas.
+If no timeout is specified, a Saga may remain open indefinitely in the case of
failure.
+Always set an explicit timeout, especially when using `MANUAL` completion mode.
====
When the timeout expires, the Saga EIP will decide to *cancel the Saga* (and
compensate all participants), unless a different decision has been taken before.
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/MandatorySagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/MandatorySagaProcessor.java
index 8be0a37d016c..c568f946f0d3 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/MandatorySagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/MandatorySagaProcessor.java
@@ -17,7 +17,6 @@
package org.apache.camel.processor.saga;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.CamelExchangeException;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
@@ -29,9 +28,9 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class MandatorySagaProcessor extends SagaProcessor {
- public MandatorySagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ public MandatorySagaProcessor(Processor childProcessor, CamelSagaService
sagaService,
SagaCompletionMode completionMode,
CamelSagaStep step) {
- super(camelContext, childProcessor, sagaService, completionMode, step);
+ super(childProcessor, sagaService, completionMode, step);
}
@Override
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NeverSagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NeverSagaProcessor.java
index 1eb346666582..61587ec02116 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NeverSagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NeverSagaProcessor.java
@@ -17,7 +17,6 @@
package org.apache.camel.processor.saga;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.CamelExchangeException;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
@@ -29,9 +28,9 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class NeverSagaProcessor extends SagaProcessor {
- public NeverSagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ public NeverSagaProcessor(Processor childProcessor, CamelSagaService
sagaService,
SagaCompletionMode completionMode, CamelSagaStep
step) {
- super(camelContext, childProcessor, sagaService, completionMode, step);
+ super(childProcessor, sagaService, completionMode, step);
if (!step.isEmpty()) {
throw new IllegalArgumentException("Saga configuration is not
allowed when propagation is set to NEVER");
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NotSupportedSagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NotSupportedSagaProcessor.java
index 53e2c8f2f5a3..1dd096b7f023 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NotSupportedSagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/NotSupportedSagaProcessor.java
@@ -17,7 +17,6 @@
package org.apache.camel.processor.saga;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.saga.CamelSagaService;
@@ -28,9 +27,9 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class NotSupportedSagaProcessor extends SagaProcessor {
- public NotSupportedSagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ public NotSupportedSagaProcessor(Processor childProcessor,
CamelSagaService sagaService,
SagaCompletionMode completionMode,
CamelSagaStep step) {
- super(camelContext, childProcessor, sagaService, completionMode, step);
+ super(childProcessor, sagaService, completionMode, step);
if (!step.isEmpty()) {
throw new IllegalArgumentException("Saga configuration is not
allowed when propagation is set to NOT_SUPPORTED");
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiredSagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiredSagaProcessor.java
index 091cc690b81c..f3f14cd7e7ae 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiredSagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiredSagaProcessor.java
@@ -19,7 +19,6 @@ package org.apache.camel.processor.saga;
import java.util.concurrent.CompletableFuture;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.saga.CamelSagaCoordinator;
@@ -31,9 +30,9 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class RequiredSagaProcessor extends SagaProcessor {
- public RequiredSagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ public RequiredSagaProcessor(Processor childProcessor, CamelSagaService
sagaService,
SagaCompletionMode completionMode,
CamelSagaStep step) {
- super(camelContext, childProcessor, sagaService, completionMode, step);
+ super(childProcessor, sagaService, completionMode, step);
}
@Override
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiresNewSagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiresNewSagaProcessor.java
index 54f268df7ad9..53f3c4f8bb89 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiresNewSagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/RequiresNewSagaProcessor.java
@@ -17,7 +17,6 @@
package org.apache.camel.processor.saga;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.saga.CamelSagaService;
@@ -28,9 +27,9 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class RequiresNewSagaProcessor extends SagaProcessor {
- public RequiresNewSagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ public RequiresNewSagaProcessor(Processor childProcessor, CamelSagaService
sagaService,
SagaCompletionMode completionMode,
CamelSagaStep step) {
- super(camelContext, childProcessor, sagaService, completionMode, step);
+ super(childProcessor, sagaService, completionMode, step);
}
@Override
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessor.java
index d88c5965dc40..7e022195f459 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessor.java
@@ -19,7 +19,6 @@ package org.apache.camel.processor.saga;
import java.util.concurrent.CompletableFuture;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.Traceable;
@@ -30,7 +29,6 @@ import org.apache.camel.saga.CamelSagaStep;
import org.apache.camel.spi.IdAware;
import org.apache.camel.spi.RouteIdAware;
import org.apache.camel.spi.StepIdAware;
-import org.apache.camel.support.service.ServiceHelper;
import org.apache.camel.util.ObjectHelper;
/**
@@ -46,7 +44,7 @@ public abstract class SagaProcessor extends
BaseDelegateProcessorSupport
private String routeId;
private String stepId;
- protected SagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ protected SagaProcessor(Processor childProcessor, CamelSagaService
sagaService,
SagaCompletionMode completionMode, CamelSagaStep
step) {
super(ObjectHelper.notNull(childProcessor, "childProcessor"));
this.sagaService = ObjectHelper.notNull(sagaService, "sagaService");
@@ -144,7 +142,7 @@ public abstract class SagaProcessor extends
BaseDelegateProcessorSupport
@Override
public String toString() {
- return "id";
+ return id;
}
@Override
@@ -167,19 +165,13 @@ public abstract class SagaProcessor extends
BaseDelegateProcessorSupport
callback.done(false);
}
} else {
- code.run();
+ try {
+ code.run();
+ } catch (Exception e) {
+ exchange.setException(e);
+ callback.done(false);
+ }
}
}
- @Override
- protected void doStart() throws Exception {
- super.doStart();
- ServiceHelper.startService(sagaService);
- }
-
- @Override
- protected void doStop() throws Exception {
- super.doStop();
- ServiceHelper.stopService(sagaService);
- }
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessorBuilder.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessorBuilder.java
index 43493edd50d6..d93c17c615b8 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessorBuilder.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SagaProcessorBuilder.java
@@ -16,7 +16,6 @@
*/
package org.apache.camel.processor.saga;
-import org.apache.camel.CamelContext;
import org.apache.camel.Processor;
import org.apache.camel.saga.CamelSagaService;
import org.apache.camel.saga.CamelSagaStep;
@@ -26,8 +25,6 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class SagaProcessorBuilder {
- private CamelContext camelContext;
-
private Processor childProcessor;
private CamelSagaService sagaService;
@@ -41,11 +38,6 @@ public class SagaProcessorBuilder {
public SagaProcessorBuilder() {
}
- public SagaProcessorBuilder camelContext(CamelContext camelContext) {
- this.camelContext = camelContext;
- return this;
- }
-
public SagaProcessorBuilder childProcessor(Processor childProcessor) {
this.childProcessor = childProcessor;
return this;
@@ -78,17 +70,17 @@ public class SagaProcessorBuilder {
switch (propagation) {
case REQUIRED:
- return new RequiredSagaProcessor(camelContext, childProcessor,
sagaService, completionMode, step);
+ return new RequiredSagaProcessor(childProcessor, sagaService,
completionMode, step);
case REQUIRES_NEW:
- return new RequiresNewSagaProcessor(camelContext,
childProcessor, sagaService, completionMode, step);
+ return new RequiresNewSagaProcessor(childProcessor,
sagaService, completionMode, step);
case SUPPORTS:
- return new SupportsSagaProcessor(camelContext, childProcessor,
sagaService, completionMode, step);
+ return new SupportsSagaProcessor(childProcessor, sagaService,
completionMode, step);
case NOT_SUPPORTED:
- return new NotSupportedSagaProcessor(camelContext,
childProcessor, sagaService, completionMode, step);
+ return new NotSupportedSagaProcessor(childProcessor,
sagaService, completionMode, step);
case NEVER:
- return new NeverSagaProcessor(camelContext, childProcessor,
sagaService, completionMode, step);
+ return new NeverSagaProcessor(childProcessor, sagaService,
completionMode, step);
case MANDATORY:
- return new MandatorySagaProcessor(camelContext,
childProcessor, sagaService, completionMode, step);
+ return new MandatorySagaProcessor(childProcessor, sagaService,
completionMode, step);
default:
throw new IllegalStateException("Unsupported propagation mode:
" + propagation);
}
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SupportsSagaProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SupportsSagaProcessor.java
index 29641e7ae660..70b8304bdc12 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SupportsSagaProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/saga/SupportsSagaProcessor.java
@@ -17,7 +17,6 @@
package org.apache.camel.processor.saga;
import org.apache.camel.AsyncCallback;
-import org.apache.camel.CamelContext;
import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.apache.camel.saga.CamelSagaService;
@@ -28,9 +27,9 @@ import org.apache.camel.saga.CamelSagaStep;
*/
public class SupportsSagaProcessor extends SagaProcessor {
- public SupportsSagaProcessor(CamelContext camelContext, Processor
childProcessor, CamelSagaService sagaService,
+ public SupportsSagaProcessor(Processor childProcessor, CamelSagaService
sagaService,
SagaCompletionMode completionMode,
CamelSagaStep step) {
- super(camelContext, childProcessor, sagaService, completionMode, step);
+ super(childProcessor, sagaService, completionMode, step);
if (completionMode != null && completionMode !=
SagaCompletionMode.defaultCompletionMode()) {
throw new IllegalArgumentException("CompletionMode cannot be
specified when propagation is SUPPORTS");
}
diff --git
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SagaReifier.java
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SagaReifier.java
index 75737ae210d5..2a3568b1ce9c 100644
---
a/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SagaReifier.java
+++
b/core/camel-core-reifier/src/main/java/org/apache/camel/reifier/SagaReifier.java
@@ -36,9 +36,13 @@ import org.apache.camel.processor.saga.SagaPropagation;
import org.apache.camel.saga.CamelSagaService;
import org.apache.camel.saga.CamelSagaStep;
import org.apache.camel.support.EndpointHelper;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
public class SagaReifier extends ProcessorReifier<SagaDefinition> {
+ private static final Logger LOG =
LoggerFactory.getLogger(SagaReifier.class);
+
public SagaReifier(Route route, ProcessorDefinition<?> definition) {
super(route, (SagaDefinition) definition);
}
@@ -104,13 +108,19 @@ public class SagaReifier extends
ProcessorReifier<SagaDefinition> {
completionMode = SagaCompletionMode.defaultCompletionMode();
}
+ if (completionMode == SagaCompletionMode.MANUAL && timeout == null) {
+ LOG.warn("Saga in route '{}' uses MANUAL completion without a
timeout."
+ + " The saga will remain open indefinitely if never
completed or compensated manually.",
+ route.getRouteId());
+ }
+
Processor childProcessor = this.createChildProcessor(true);
CamelSagaService camelSagaService = resolveSagaService();
CamelContextAware.trySetCamelContext(camelSagaService,
getCamelContext());
camelSagaService.registerStep(step);
- SagaProcessor answer = new
SagaProcessorBuilder().camelContext(camelContext).childProcessor(childProcessor)
+ SagaProcessor answer = new
SagaProcessorBuilder().childProcessor(childProcessor)
.sagaService(camelSagaService).step(step)
.propagation(propagation).completionMode(completionMode).build();
answer.setDisabled(isDisabled(camelContext, definition));
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/SagaSharedServiceRouteStopTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/SagaSharedServiceRouteStopTest.java
new file mode 100644
index 000000000000..0eb76a7fa2e7
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/SagaSharedServiceRouteStopTest.java
@@ -0,0 +1,86 @@
+/*
+ * 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.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.saga.InMemorySagaService;
+import org.apache.camel.support.service.ServiceHelper;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+public class SagaSharedServiceRouteStopTest extends ContextTestSupport {
+
+ @Test
+ public void testStoppingOneSagaRouteKeepsSharedServiceRunning() throws
Exception {
+ InMemorySagaService sagaService =
context.hasService(InMemorySagaService.class);
+
+ assertTrue(ServiceHelper.isStarted(sagaService), "Saga service should
be running");
+
+ context.getRouteController().stopRoute("route-a");
+
+ assertTrue(ServiceHelper.isStarted(sagaService),
+ "Saga service should still be running after stopping one
route");
+ }
+
+ @Test
+ public void
testCompensationOnOtherRouteStillWorksAfterStoppingOneSagaRoute() throws
Exception {
+ MockEndpoint compensated = getMockEndpoint("mock:compensate-b");
+ compensated.expectedMessageCount(1);
+
+ context.getRouteController().stopRoute("route-a");
+
+ try {
+ template.sendBody("direct:route-b", "trigger-fail");
+ } catch (Exception e) {
+ // expected — the saga processor throws after compensation is
triggered
+ }
+
+ MockEndpoint.assertIsSatisfied(context, 10, TimeUnit.SECONDS);
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() throws Exception {
+ return new RouteBuilder() {
+ @Override
+ public void configure() throws Exception {
+ InMemorySagaService sagaService = new InMemorySagaService();
+ context.addService(sagaService);
+
+ from("direct:route-a").routeId("route-a")
+ .saga().compensation("direct:compensate-a")
+ .log("Route A");
+
+ from("direct:compensate-a")
+ .log("Compensating A");
+
+ from("direct:route-b").routeId("route-b")
+ .saga().compensation("direct:compensate-b")
+ .process(e -> {
+ throw new RuntimeException("forced failure");
+ });
+
+ from("direct:compensate-b")
+ .to("mock:compensate-b");
+ }
+ };
+ }
+}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
b/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
index 85f900ca6dc2..7019d7fa5b20 100644
---
a/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
+++
b/core/camel-support/src/main/java/org/apache/camel/saga/InMemorySagaCoordinator.java
@@ -24,6 +24,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
@@ -57,6 +58,7 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
private final InMemorySagaService sagaService;
private final String sagaId;
private final List<StepEnlistment> enlistments;
+ private final List<ScheduledFuture<?>> timeoutFutures;
private final AtomicReference<Status> currentStatus;
public InMemorySagaCoordinator(CamelContext camelContext,
InMemorySagaService sagaService, String sagaId) {
@@ -64,6 +66,7 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
this.sagaService = ObjectHelper.notNull(sagaService, "sagaService");
this.sagaId = ObjectHelper.notNull(sagaId, "sagaId");
this.enlistments = new CopyOnWriteArrayList<>();
+ this.timeoutFutures = new CopyOnWriteArrayList<>();
this.currentStatus = new AtomicReference<>(Status.RUNNING);
}
@@ -92,10 +95,9 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
values.put(entry.getKey(), value);
}
} catch (Exception ex) {
- return CompletableFuture.supplyAsync(() -> {
- throw new RuntimeCamelException(
- "Cannot evaluate saga option '" +
entry.getKey() + "'", ex);
- });
+ return CompletableFuture.failedFuture(
+ new RuntimeCamelException(
+ "Cannot evaluate saga option '" +
entry.getKey() + "'", ex));
}
}
}
@@ -103,12 +105,13 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
this.enlistments.add(new StepEnlistment(step, values));
if (step.getTimeoutInMilliseconds().isPresent()) {
- sagaService.getExecutorService().schedule(() -> {
+ ScheduledFuture<?> timeoutFuture =
sagaService.getExecutorService().schedule(() -> {
boolean doAction = currentStatus.compareAndSet(Status.RUNNING,
Status.COMPENSATING);
if (doAction) {
doCompensate(exchange);
}
}, step.getTimeoutInMilliseconds().get(), TimeUnit.MILLISECONDS);
+ timeoutFutures.add(timeoutFuture);
}
return CompletableFuture.completedFuture(null);
@@ -119,6 +122,7 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
boolean doAction = currentStatus.compareAndSet(Status.RUNNING,
Status.COMPENSATING);
if (doAction) {
+ cancelTimeouts();
return doCompensate(exchange).thenApply(res -> {
if (!res) {
throw new RuntimeCamelException(
@@ -143,6 +147,7 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
boolean doAction = currentStatus.compareAndSet(Status.RUNNING,
Status.COMPLETING);
if (doAction) {
+ cancelTimeouts();
return doComplete(exchange).thenApply(res -> {
if (!res) {
throw new RuntimeCamelException(
@@ -164,19 +169,23 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
public CompletableFuture<Boolean> doCompensate(final Exchange exchange) {
return doFinalize(exchange, CamelSagaStep::getCompensation,
"compensation")
- .thenApply(res -> {
+ .whenComplete((res, ex) -> {
+ if (ex != null || !Boolean.TRUE.equals(res)) {
+ LOG.warn("Saga {} compensation did not fully succeed —
manual intervention may be needed", sagaId);
+ }
currentStatus.set(Status.COMPENSATED);
sagaService.removeSaga(sagaId);
- return res;
});
}
public CompletableFuture<Boolean> doComplete(final Exchange exchange) {
return doFinalize(exchange, CamelSagaStep::getCompletion, "completion")
- .thenApply(res -> {
+ .whenComplete((res, ex) -> {
+ if (ex != null || !Boolean.TRUE.equals(res)) {
+ LOG.warn("Saga {} completion did not fully succeed —
manual intervention may be needed", sagaId);
+ }
currentStatus.set(Status.COMPLETED);
sagaService.removeSaga(sagaId);
- return res;
});
}
@@ -256,6 +265,13 @@ public class InMemorySagaCoordinator implements
CamelSagaCoordinator {
return answer;
}
+ private void cancelTimeouts() {
+ for (ScheduledFuture<?> future : timeoutFutures) {
+ future.cancel(false);
+ }
+ timeoutFutures.clear();
+ }
+
private <T> List<T> reversed(List<T> list) {
List<T> reversed = new ArrayList<>(list);
Collections.reverse(reversed);
diff --git
a/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/KameletMain.java
b/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/KameletMain.java
index c3ca1353fdd3..4c666e8756b1 100644
---
a/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/KameletMain.java
+++
b/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/KameletMain.java
@@ -493,7 +493,7 @@ public class KameletMain extends MainCommandLineSupport {
TransactedDownloader.registerDownloadReifiers(this);
// in case we use saga
- SagaDownloader.registerDownloadReifiers(this);
+ SagaDownloader.registerDownloadReifiers();
// if transforming DSL then disable processors as we just want to work
on the model (not runtime processors)
if (transform) {
diff --git
a/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/download/SagaDownloader.java
b/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/download/SagaDownloader.java
index 3fea84822389..d23f368a303e 100644
---
a/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/download/SagaDownloader.java
+++
b/dsl/camel-kamelet-main/src/main/java/org/apache/camel/main/download/SagaDownloader.java
@@ -16,10 +16,11 @@
*/
package org.apache.camel.main.download;
-import org.apache.camel.main.KameletMain;
+import org.apache.camel.CamelContext;
import org.apache.camel.model.SagaDefinition;
import org.apache.camel.reifier.ProcessReifier;
import org.apache.camel.reifier.ProcessorReifier;
+import org.apache.camel.saga.CamelSagaService;
import org.apache.camel.saga.InMemorySagaService;
/**
@@ -30,7 +31,7 @@ public class SagaDownloader {
private SagaDownloader() {
}
- public static void registerDownloadReifiers(KameletMain main) {
+ public static void registerDownloadReifiers() {
ProcessorReifier.registerReifier(SagaDefinition.class,
(route, processorDefinition) -> {
if (processorDefinition instanceof SagaDefinition) {
@@ -42,7 +43,14 @@ public class SagaDownloader {
route.getCamelContext().getVersion());
}
}
- main.bind("inMemorySagaService", new
InMemorySagaService());
+ CamelContext ctx = route.getCamelContext();
+ if (ctx.hasService(CamelSagaService.class) == null) {
+ try {
+ ctx.addService(new InMemorySagaService());
+ } catch (Exception e) {
+ throw new RuntimeException("Failed to add
InMemorySagaService", e);
+ }
+ }
return ProcessReifier.coreReifier(route,
processorDefinition);
});
}