This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-25058 in repository https://gitbox.apache.org/repos/asf/camel.git
commit 8e7efed79cf4707347dd9b36e7e974c960302a80 Author: Claus Ibsen <[email protected]> AuthorDate: Sun Sep 27 20:39:33 2026 +0200 CAMEL-25058: camel-saga - saga:complete and saga:compensate only consult Long-Running-Action for a saga service that uses it Aligns SagaProducer with SagaProcessor (CAMEL-24449): when the exchange is not bound to a saga, the Long-Running-Action header is only used when the configured CamelSagaService returns true from isLongRunningActionHeaderSupported(), as LRASagaService does. Depends on CAMEL-25011 so that exchange copies (seda, wire tap, ...) keep their saga binding. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Signed-off-by: Claus Ibsen <[email protected]> --- .../apache/camel/component/saga/SagaProducer.java | 5 +- .../SagaProducerLongRunningActionHeaderTest.java | 110 +++++++++++++++++++++ .../ROOT/pages/camel-4x-upgrade-guide-4_18.adoc | 13 +++ .../ROOT/pages/camel-4x-upgrade-guide-4_22.adoc | 13 +++ .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 13 +++ 5 files changed, 153 insertions(+), 1 deletion(-) diff --git a/components/camel-saga/src/main/java/org/apache/camel/component/saga/SagaProducer.java b/components/camel-saga/src/main/java/org/apache/camel/component/saga/SagaProducer.java index 4aa208c648f7..fa83306012d7 100644 --- a/components/camel-saga/src/main/java/org/apache/camel/component/saga/SagaProducer.java +++ b/components/camel-saga/src/main/java/org/apache/camel/component/saga/SagaProducer.java @@ -44,8 +44,11 @@ public class SagaProducer extends DefaultAsyncProducer { @Override public boolean process(Exchange exchange, AsyncCallback callback) { + // try internal state first (survives removeHeaders("*")) String sagaId = exchange.getExchangeExtension().getSagaLongRunningAction(); - if (sagaId == null) { + if (sagaId == null && camelSagaService.isLongRunningActionHeaderSupported()) { + // fall back to header only for a saga service that takes part in a protocol carrying the id that way + // (e.g. LRA), same as the Saga EIP. The header is outside the Camel namespace that consumers filter. sagaId = exchange.getIn().getHeader(SagaConstants.SAGA_LONG_RUNNING_ACTION, String.class); } if (sagaId == null) { diff --git a/core/camel-core/src/test/java/org/apache/camel/component/saga/SagaProducerLongRunningActionHeaderTest.java b/core/camel-core/src/test/java/org/apache/camel/component/saga/SagaProducerLongRunningActionHeaderTest.java new file mode 100644 index 000000000000..35a2a2cb5221 --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/component/saga/SagaProducerLongRunningActionHeaderTest.java @@ -0,0 +1,110 @@ +/* + * 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.component.saga; + +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; +import org.apache.camel.saga.InMemorySagaService; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * An exchange that is not bound to a saga must not pick a saga to complete or compensate via the Long-Running-Action + * header, unless the saga service takes part in a protocol that carries the id that way. + */ +public class SagaProducerLongRunningActionHeaderTest extends ContextTestSupport { + + private final LongRunningActionHeaderSagaService sagaService = new LongRunningActionHeaderSagaService(); + + @Test + public void testHeaderIgnoredWhenNotSupported() throws Exception { + String sagaId = startManualSaga(); + + MockEndpoint compensated = getMockEndpoint("mock:compensated"); + compensated.expectedMessageCount(0); + + CamelExecutionException e = assertThrows(CamelExecutionException.class, + () -> template.sendBodyAndHeader("direct:compensate", "cancel", Exchange.SAGA_LONG_RUNNING_ACTION, + sagaId)); + IllegalStateException cause = assertInstanceOf(IllegalStateException.class, e.getCause()); + assertTrue(cause.getMessage().contains("not bound to a saga context")); + + compensated.assertIsSatisfied(200, TimeUnit.MILLISECONDS); + } + + @Test + public void testHeaderUsedWhenSupported() throws Exception { + sagaService.setHeaderSupported(true); + String sagaId = startManualSaga(); + + MockEndpoint compensated = getMockEndpoint("mock:compensated"); + compensated.expectedMessageCount(1); + + template.sendBodyAndHeader("direct:compensate", "cancel", Exchange.SAGA_LONG_RUNNING_ACTION, sagaId); + + compensated.assertIsSatisfied(); + } + + private String startManualSaga() { + Exchange out = template.request("direct:start", e -> e.getIn().setBody("order")); + String sagaId = out.getMessage().getHeader(Exchange.SAGA_LONG_RUNNING_ACTION, String.class); + assertNotNull(sagaId); + return sagaId; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() throws Exception { + context.addService(sagaService); + + from("direct:start") + .saga().compensation("mock:compensated").completion("mock:completed") + .completionMode(SagaCompletionMode.MANUAL) + .to("mock:start"); + + from("direct:compensate") + .to("saga:compensate"); + } + }; + } + + private static final class LongRunningActionHeaderSagaService extends InMemorySagaService { + + private boolean headerSupported; + + void setHeaderSupported(boolean headerSupported) { + this.headerSupported = headerSupported; + } + + @Override + public boolean isLongRunningActionHeaderSupported() { + return headerSupported; + } + } +} diff --git a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc index a80230d15c7d..5bb6b53b7adb 100644 --- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc +++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc @@ -256,6 +256,19 @@ are unaffected. Routes that set the header by its literal string name, or that u `allowTemplateFromHeader=true` with the old header names, must switch to the new `Camel`-prefixed names. +=== camel-saga + +The `saga:complete` and `saga:compensate` endpoints now follow the same rule as the Saga EIP: when the +exchange is not bound to a saga, the `Long-Running-Action` message header is consulted only if the +configured `CamelSagaService` returns `true` from `isLongRunningActionHeaderSupported()`, as +`LRASagaService` does. With the default `InMemorySagaService` the header is ignored, and sending an +exchange that is not bound to a saga to these endpoints fails with `IllegalStateException`. + +Exchanges created inside a saga, including copies made by EIPs such as Wire Tap, Multicast or SEDA, +stay bound to it and are not affected. A route that completes or compensates an in-memory saga from an +unrelated exchange must bind it explicitly from trusted code, for example with +`exchange.getExchangeExtension().setSagaLongRunningAction(id)`. + == Upgrading from 4.18.3 to 4.18.4 === camel-core - Multicast EIP honors UseOriginalAggregationStrategy diff --git a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc index 025e8d4e3f68..53fdb509f1a3 100644 --- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc +++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc @@ -61,6 +61,19 @@ healthy while no longer receiving any change event. Deployments that use readine now see a Debezium route whose engine has died reported as `DOWN`, where it was previously reported as `UP`. The engine is still not restarted automatically. +=== camel-saga + +The `saga:complete` and `saga:compensate` endpoints now follow the same rule as the Saga EIP: when the +exchange is not bound to a saga, the `Long-Running-Action` message header is consulted only if the +configured `CamelSagaService` returns `true` from `isLongRunningActionHeaderSupported()`, as +`LRASagaService` does. With the default `InMemorySagaService` the header is ignored, and sending an +exchange that is not bound to a saga to these endpoints fails with `IllegalStateException`. + +Exchanges created inside a saga, including copies made by EIPs such as Wire Tap, Multicast or SEDA, +stay bound to it and are not affected. A route that completes or compensates an in-memory saga from an +unrelated exchange must bind it explicitly from trusted code, for example with +`exchange.getExchangeExtension().setSagaLongRunningAction(id)`. + == Upgrading from 4.22.0 to 4.22.1 === camel-docling 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 5f750d4fb38e..f2e2f970eedc 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 @@ -1969,6 +1969,19 @@ 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-saga + +The `saga:complete` and `saga:compensate` endpoints now follow the same rule as the Saga EIP: when the +exchange is not bound to a saga, the `Long-Running-Action` message header is consulted only if the +configured `CamelSagaService` returns `true` from `isLongRunningActionHeaderSupported()`, as +`LRASagaService` does. With the default `InMemorySagaService` the header is ignored, and sending an +exchange that is not bound to a saga to these endpoints fails with `IllegalStateException`. + +Exchanges created inside a saga, including copies made by EIPs such as Wire Tap, Multicast or SEDA, +stay bound to it and are not affected. A route that completes or compensates an in-memory saga from an +unrelated exchange must bind it explicitly from trusted code, for example with +`exchange.getExchangeExtension().setSagaLongRunningAction(id)`. + === camel-microprofile-health The `error.stacktrace` entry of a failed health check is now only included in the response when
