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

Reply via email to