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 f0092776a765 CAMEL-25006: camel-core - Saga EIP: fail a REQUIRED or
SUPPORTS step of a saga that has already ended (#26866)
f0092776a765 is described below
commit f0092776a76566611363e5fd277d78f8f9cd86f4
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 17:51:44 2026 +0530
CAMEL-25006: camel-core - Saga EIP: fail a REQUIRED or SUPPORTS step of a
saga that has already ended (#26866)
Since CAMEL-24144 (4.22.0) InMemorySagaService removes a saga once it is
completed or compensated, and getSaga(id) then returns null. The saga
processors treat null as "the exchange is not in a saga". So after a
saga timed out and was compensated, a REQUIRED step on an exchange of
that saga started a new saga of its own and completed it, and a
SUPPORTS step ran outside of any saga. For example, a payment was taken
and confirmed although the order it belonged to was compensated.
Before CAMEL-24144, and still while the saga is being compensated, such
a step failed with "Cannot begin: status is ...".
RequiredSagaProcessor and SupportsSagaProcessor now fail with
IllegalStateException when the exchange carries the id of a saga that
the saga service no longer knows. The saga is still removed from the
service, so this keeps the memory fix of CAMEL-24144. MANDATORY already
fails, and REQUIRES_NEW, NOT_SUPPORTED and NEVER are unchanged.
Regression from CAMEL-24144.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../processor/saga/RequiredSagaProcessor.java | 1 +
.../apache/camel/processor/saga/SagaProcessor.java | 33 ++++++++++--
.../processor/saga/SupportsSagaProcessor.java | 1 +
.../apache/camel/processor/SagaTimeoutTest.java | 61 +++++++++++++++++++++-
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 15 ++++++
5 files changed, 105 insertions(+), 6 deletions(-)
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 f3f14cd7e7ae..eb931fd3d256 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
@@ -39,6 +39,7 @@ public class RequiredSagaProcessor extends SagaProcessor {
public boolean process(Exchange exchange, AsyncCallback callback) {
getCurrentSagaCoordinator(exchange)
.whenComplete((existingCoordinator, ex) -> ifNotException(ex,
exchange, callback, () -> {
+ checkSagaIsActive(exchange, existingCoordinator);
CompletableFuture<CamelSagaCoordinator> coordinatorFuture;
final boolean inheritedCoordinator;
if (existingCoordinator != null) {
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 cd14c5e97e22..efcb288144e9 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
@@ -53,6 +53,33 @@ public abstract class SagaProcessor extends
BaseDelegateProcessorSupport
}
protected CompletableFuture<CamelSagaCoordinator>
getCurrentSagaCoordinator(Exchange exchange) {
+ String currentSaga = getCurrentSagaId(exchange);
+ if (currentSaga != null) {
+ return sagaService.getSaga(currentSaga);
+ }
+
+ return CompletableFuture.completedFuture(null);
+ }
+
+ /**
+ * Checks that the saga of the exchange, if any, is still known by the
saga service. The in-memory saga service
+ * removes a saga once it is completed or compensated (for example after a
timeout), and then a step of that saga
+ * must fail, instead of starting a new saga or running outside of a saga.
+ *
+ * @param exchange the exchange
+ * @param coordinator the coordinator of the saga found for the
exchange, or <tt>null</tt> if none
+ * @throws IllegalStateException if the exchange belongs to a saga that is
no longer active
+ */
+ protected void checkSagaIsActive(Exchange exchange, CamelSagaCoordinator
coordinator) {
+ if (coordinator == null) {
+ String currentSaga = getCurrentSagaId(exchange);
+ if (currentSaga != null) {
+ throw new IllegalStateException("Cannot begin: saga " +
currentSaga + " is not active or not known");
+ }
+ }
+ }
+
+ private String getCurrentSagaId(Exchange exchange) {
// try internal state first (survives removeHeaders("*"))
String currentSaga =
exchange.getExchangeExtension().getSagaLongRunningAction();
if (currentSaga == null &&
sagaService.isLongRunningActionHeaderSupported()) {
@@ -62,11 +89,7 @@ public abstract class SagaProcessor extends
BaseDelegateProcessorSupport
// message pick which saga its exchange joins.
currentSaga =
exchange.getIn().getHeader(Exchange.SAGA_LONG_RUNNING_ACTION, String.class);
}
- if (currentSaga != null) {
- return sagaService.getSaga(currentSaga);
- }
-
- return CompletableFuture.completedFuture(null);
+ return currentSaga;
}
protected void setCurrentSagaCoordinator(Exchange exchange,
CamelSagaCoordinator coordinator) {
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 70b8304bdc12..8e7916fb982d 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
@@ -38,6 +38,7 @@ public class SupportsSagaProcessor extends SagaProcessor {
@Override
public boolean process(Exchange exchange, AsyncCallback callback) {
getCurrentSagaCoordinator(exchange).whenComplete((coordinator, ex) ->
ifNotException(ex, exchange, callback, () -> {
+ checkSagaIsActive(exchange, coordinator);
if (coordinator != null) {
coordinator.beginStep(exchange, step)
.whenComplete((done, ex2) -> ifNotException(ex2,
exchange, callback, () -> {
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/SagaTimeoutTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/SagaTimeoutTest.java
index 87eebbef1cd0..757f88dcb392 100644
---
a/core/camel-core/src/test/java/org/apache/camel/processor/SagaTimeoutTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/SagaTimeoutTest.java
@@ -21,6 +21,7 @@ import java.util.concurrent.TimeUnit;
import org.apache.camel.CamelExecutionException;
import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.component.mock.MockEndpoint;
import org.apache.camel.model.SagaCompletionMode;
@@ -28,11 +29,15 @@ import org.apache.camel.model.SagaPropagation;
import org.apache.camel.saga.InMemorySagaService;
import org.junit.jupiter.api.Test;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class SagaTimeoutTest extends ContextTestSupport {
+ private InMemorySagaService sagaService;
+
@Test
public void testTimeoutCalledCorrectly() throws Exception {
MockEndpoint compensate = getMockEndpoint("mock:compensate");
@@ -94,12 +99,46 @@ public class SagaTimeoutTest extends ContextTestSupport {
compensate.assertIsSatisfied();
}
+ @Test
+ void testRequiredStepAfterTimeoutDoesNotStartNewSaga() throws Exception {
+ assertStepAfterTimeoutFails("direct:saga-timeout-required");
+ }
+
+ @Test
+ void testSupportsStepAfterTimeoutDoesNotRunOutsideSaga() throws Exception {
+ assertStepAfterTimeoutFails("direct:saga-timeout-supports");
+ }
+
+ private void assertStepAfterTimeoutFails(String uri) throws Exception {
+ MockEndpoint compensate = getMockEndpoint("mock:compensate");
+ compensate.expectedMessageCount(1);
+
+ MockEndpoint payment = getMockEndpoint("mock:payment");
+ payment.expectedMessageCount(0);
+
+ MockEndpoint complete = getMockEndpoint("mock:complete");
+ complete.expectedMessageCount(0);
+
+ CamelExecutionException ex =
assertThrows(CamelExecutionException.class, () -> template.sendBody(uri,
"Hello"));
+
+ MockEndpoint.assertIsSatisfied(context);
+ assertInstanceOf(IllegalStateException.class, ex.getCause());
+ assertTrue(ex.getCause().getMessage().endsWith("is not active or not
known"), ex.getCause().getMessage());
+ }
+
+ private void awaitSagaEnded(Exchange exchange) {
+ // a slow call outlasts the saga timeout: the saga is compensated and
removed from the saga service
+ String sagaId =
exchange.getExchangeExtension().getSagaLongRunningAction();
+ await().atMost(5, TimeUnit.SECONDS).until(() ->
sagaService.getSaga(sagaId).get() == null);
+ }
+
@Override
protected RouteBuilder createRouteBuilder() {
return new RouteBuilder() {
@Override
public void configure() throws Exception {
- context.addService(new InMemorySagaService());
+ sagaService = new InMemorySagaService();
+ context.addService(sagaService);
from("direct:saga").saga().timeout(100,
TimeUnit.MILLISECONDS).option("id", constant("myid"))
.completionMode(SagaCompletionMode.MANUAL)
@@ -130,6 +169,26 @@ public class SagaTimeoutTest extends ContextTestSupport {
.propagation(SagaPropagation.MANDATORY).timeout(500,
TimeUnit.MILLISECONDS)
.compensation("mock:compensate").completion("mock:complete")
.to("mock:end");
+
+ from("direct:saga-timeout-required")
+ .saga().timeout(100,
TimeUnit.MILLISECONDS).compensation("mock:compensate")
+ .process(SagaTimeoutTest.this::awaitSagaEnded)
+ .to("direct:payment-required");
+
+ from("direct:payment-required")
+ .saga().propagation(SagaPropagation.REQUIRED)
+
.compensation("mock:compensate-payment").completion("mock:complete")
+ .to("mock:payment");
+
+ from("direct:saga-timeout-supports")
+ .saga().timeout(100,
TimeUnit.MILLISECONDS).compensation("mock:compensate")
+ .process(SagaTimeoutTest.this::awaitSagaEnded)
+ .to("direct:payment-supports");
+
+ from("direct:payment-supports")
+ .saga().propagation(SagaPropagation.SUPPORTS)
+ .compensation("mock:compensate-payment")
+ .to("mock:payment");
}
};
}
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 6000e7a460a9..ea0f34c2c218 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
@@ -2095,6 +2095,21 @@ A custom `CamelSagaService` that relies on the header to
join sagas started by a
override the new method. Everything else is unaffected: the header is still
set on the exchange, and
routes reading it continue to work.
+=== camel-core - a saga step of a saga that has already ended fails
+
+Since 4.22, the `InMemorySagaService` removes a saga once it is completed or
compensated, for example after
+its timeout, or after `saga:complete` for a `MANUAL` saga. A saga step with
propagation `REQUIRED` on an exchange
+of such a saga then started a new saga of its own and completed it, and a step
with propagation `SUPPORTS` ran
+outside of any saga, although the saga of the exchange had already ended.
+
+Such a step now fails with `IllegalStateException` (`Cannot begin: saga <id>
is not active or not known`), as it
+did before 4.22 and as it does while the saga is still being completed or
compensated. A step with propagation
+`MANDATORY` fails as before.
+
+The check applies to any `CamelSagaService` that returns no coordinator for
the saga id of the exchange. For
+example, a custom saga service that reads the `Long-Running-Action` header,
and receives the id of a saga it does
+not know, now fails such a step instead of starting a new saga.
+
=== camel-saga
The `saga:complete` and `saga:compensate` endpoints now follow the same rule
as the Saga EIP: when the