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 fe67aa8989d8 CAMEL-24943: camel-core - Aggregate EIP: the recover task
should not re-send an exchange that is being completed
fe67aa8989d8 is described below
commit fe67aa8989d8196c7eb944d9cb61ba07a84fe43c
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 22:20:24 2026 +0530
CAMEL-24943: camel-core - Aggregate EIP: the recover task should not
re-send an exchange that is being completed
With a RecoverableAggregationRepository an aggregated exchange could be
sent twice without any failure. The completing thread removes the group
from the repository under the aggregation lock, which moves it into the
recovery store, but registers it in inProgressCompleteExchanges only later
in onSubmitCompletion, after the lock is released. A recover task scan in
between re-delivered the exchange too. Optimistic locking exposed timeout
completions as well. This is a remaining case of the race addressed by
CAMEL-6097 and CAMEL-8010.
Before remove on a recoverable repository the exchange is now marked as
being completed, with a per-id counter so that a thread losing an
optimistic compare-and-set cannot clear the winner's mark. The recover
task treats marked ids as in progress, and onSubmitCompletion clears the
mark. The mark is also cleared whenever the exchange is not submitted
(remove throws, the exchange is discarded or the group force-discarded),
so those exchanges are still recovered and delivery stays at-least-once.
Closes #26788
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.../processor/aggregate/AggregateProcessor.java | 48 ++-
.../AggregatePreCompleteLostGroupTest.java | 50 ++-
.../aggregator/AggregateRecoverInProgressTest.java | 397 +++++++++++++++++++++
3 files changed, 488 insertions(+), 7 deletions(-)
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
index 572eeb75d559..574c03376e9b 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
@@ -132,6 +132,9 @@ public class AggregateProcessor extends BaseProcessorSupport
private final WaitableInteger inProgressCount = new WaitableInteger();
private final Set<String> unconfirmedCompleteExchanges =
ConcurrentHashMap.newKeySet();
private final Set<String> inProgressCompleteExchangesForRecoveryTask =
ConcurrentHashMap.newKeySet();
+ // aggregated exchanges being completed (exchange id -> count): moved to
the completed store of a recoverable
+ // repository, but not yet registered in inProgressCompleteExchanges by
onSubmitCompletion
+ private final Map<String, Integer> completingExchanges = new
ConcurrentHashMap<>();
private final Map<String, RedeliveryData> redeliveryState = new
ConcurrentHashMap<>();
private final AggregateProcessorStatistics statistics = new Statistics();
@@ -858,6 +861,28 @@ public class AggregateProcessor extends
BaseProcessorSupport
protected Exchange onCompletion(
final String key, final Exchange original, final Exchange
aggregated, boolean fromTimeout,
boolean aggregateFailed) {
+ // a recoverable repository moves the exchange to its completed store
when it is removed, so mark it as being
+ // completed until onSubmitCompletion has registered it as in
progress, as otherwise the recover task could
+ // recover and send it as well
+ final boolean completing = original != null &&
isRecoverableRepository();
+ if (completing) {
+ markCompleting(aggregated.getExchangeId());
+ }
+ Exchange answer = null;
+ try {
+ answer = doOnCompletion(key, original, aggregated, fromTimeout,
aggregateFailed);
+ return answer;
+ } finally {
+ if (completing && answer == null) {
+ // not removed, or not to be sent
+ unmarkCompleting(aggregated.getExchangeId());
+ }
+ }
+ }
+
+ private Exchange doOnCompletion(
+ final String key, final Exchange original, final Exchange
aggregated, boolean fromTimeout,
+ boolean aggregateFailed) {
// store the correlation key as property before we remove so the
repository has that information
if (original != null) {
original.setProperty(ExchangePropertyKey.AGGREGATED_CORRELATION_KEY, key);
@@ -907,6 +932,14 @@ public class AggregateProcessor extends
BaseProcessorSupport
return answer;
}
+ private void markCompleting(String exchangeId) {
+ completingExchanges.merge(exchangeId, 1, Integer::sum);
+ }
+
+ private void unmarkCompleting(String exchangeId) {
+ completingExchanges.computeIfPresent(exchangeId, (id, count) -> count
> 1 ? count - 1 : null);
+ }
+
private void discard(String key, Exchange aggregated) {
// this exchange is discarded
discarded.incrementAndGet();
@@ -928,6 +961,8 @@ public class AggregateProcessor extends BaseProcessorSupport
if (recoveryInProgress.get()) {
inProgressCompleteExchangesForRecoveryTask.add(exchange.getExchangeId());
}
+ // registered as in progress, so no longer needs to be marked as being
completed
+ unmarkCompleting(exchange.getExchangeId());
// invoke the on completion callback
aggregationStrategy.onCompletion(exchange);
@@ -1542,8 +1577,10 @@ public class AggregateProcessor extends
BaseProcessorSupport
lock.lock();
try {
// consider in progress if it was in progress before
we did the scan, or currently after we did the scan
+ // (or is being completed and not yet registered as in
progress)
// its safer to consider it in progress than risk
duplicates due both in progress + recovered
- final boolean inProgress =
inProgressCompleteExchangesForRecoveryTask.contains(exchangeId);
+ final boolean inProgress =
inProgressCompleteExchangesForRecoveryTask.contains(exchangeId)
+ || completingExchanges.containsKey(exchangeId);
if (inProgress) {
LOG.trace("Aggregated exchange with id: {} is
already in progress.", exchangeId);
if
(unconfirmedCompleteExchanges.contains(exchangeId)) {
@@ -1867,6 +1904,7 @@ public class AggregateProcessor extends
BaseProcessorSupport
// cleanup when shutting down
inProgressCompleteExchanges.clear();
+ completingExchanges.clear();
inProgressCount.reset();
if (shutdownExecutorService) {
@@ -2048,7 +2086,13 @@ public class AggregateProcessor extends
BaseProcessorSupport
// onCompletion only discards on aggregation failure when
discardOnAggregationFailure is enabled,
// so discard here, as otherwise the group is removed without
being confirmed (and a recoverable
// repository would recover and send it later)
- discard(key, answer);
+ try {
+ discard(key, answer);
+ } finally {
+ // onCompletion returned the exchange, so it is still
marked as being completed, but it is
+ // discarded instead of passed to onSubmitCompletion, so
clear the mark here
+ unmarkCompleting(answer.getExchangeId());
+ }
}
return true;
} catch
(OptimisticLockingAggregationRepository.OptimisticLockingException e) {
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregatePreCompleteLostGroupTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregatePreCompleteLostGroupTest.java
index aa806fd1f3b1..4c6d72d5758c 100644
---
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregatePreCompleteLostGroupTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregatePreCompleteLostGroupTest.java
@@ -27,16 +27,19 @@ import org.apache.camel.AsyncProcessor;
import org.apache.camel.CamelContext;
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.processor.SendProcessor;
import org.apache.camel.processor.aggregate.AggregateProcessor;
import org.apache.camel.processor.aggregate.MemoryAggregationRepository;
import org.apache.camel.processor.aggregate.OptimisticLockRetryPolicy;
import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.KeyValueAggregationRepository;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import static org.awaitility.Awaitility.await;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -85,7 +88,7 @@ public class AggregatePreCompleteLostGroupTest extends
ContextTestSupport {
super.remove(camelContext, key, exchange);
if
(Thread.currentThread().getName().equals("producer-START-b") &&
pauseRemove.getAndSet(false)) {
removed.countDown();
- await(releaseRemove);
+ awaitLatch(releaseRemove);
}
}
};
@@ -103,7 +106,7 @@ public class AggregatePreCompleteLostGroupTest extends
ContextTestSupport {
Exchange b = createExchange("START-b");
Thread producer = new Thread(() -> process(ap, b), "producer-START-b");
producer.start();
- await(removed);
+ awaitLatch(removed);
// c starts a new group, so START-b fails to add its new group and is
retried
ap.process(createExchange("c"));
@@ -136,6 +139,40 @@ public class AggregatePreCompleteLostGroupTest extends
ContextTestSupport {
ap.stop();
}
+ @Test
+ public void
testAggregationFailureAfterPreCompletionWithRecoverableRepository() throws
Exception {
+ // the pre-completed group is sent at once, and as it is registered as
in progress before it can be seen as
+ // completed in the repository, the recover task must not send it a
second time
+ KeyValueAggregationRepository repository = new
KeyValueAggregationRepository();
+ repository.setUseRecovery(true);
+ repository.setRecoveryInterval(50);
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"),
createStrategy()).aggregationRepository(repository).completionTimeout(60000)
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("a1");
+
+ template.sendBodyAndHeader("direct:start", "a1", "id", 1);
+ Exchange bad = template.send("direct:start", e -> {
+ e.getIn().setBody("START-fail");
+ e.getIn().setHeader("id", 1);
+ });
+
+ assertNotNull(bad.getException());
+ assertMockEndpointsSatisfied();
+ // several recovery runs later the group is still delivered only once,
and it is confirmed
+ await().atMost(5, TimeUnit.SECONDS).until(() ->
repository.scan(context).isEmpty());
+ await().during(500, TimeUnit.MILLISECONDS).atMost(5, TimeUnit.SECONDS)
+ .until(() -> mock.getReceivedCounter() == 1);
+ }
+
@Test
public void testAggregationFailureDiscardedAfterPreCompletion() throws
Exception {
MockEndpoint mock = getMockEndpoint("mock:result");
@@ -157,8 +194,12 @@ public class AggregatePreCompleteLostGroupTest extends
ContextTestSupport {
private AggregateProcessor createProcessor() {
AsyncProcessor done = new
SendProcessor(context.getEndpoint("mock:result"));
+ return new AggregateProcessor(context, done, header("id"),
createStrategy(), executorService, true);
+ }
+
+ private static AggregationStrategy createStrategy() {
// pre-completes the current group when a body starting with START
arrives
- AggregationStrategy strategy = new AggregationStrategy() {
+ return new AggregationStrategy() {
@Override
public boolean canPreComplete() {
return true;
@@ -182,7 +223,6 @@ public class AggregatePreCompleteLostGroupTest extends
ContextTestSupport {
return oldExchange;
}
};
- return new AggregateProcessor(context, done, header("id"), strategy,
executorService, true);
}
private Exchange createExchange(String body) {
@@ -200,7 +240,7 @@ public class AggregatePreCompleteLostGroupTest extends
ContextTestSupport {
}
}
- private static void await(CountDownLatch latch) {
+ private static void awaitLatch(CountDownLatch latch) {
try {
if (!latch.await(10, TimeUnit.SECONDS)) {
fail("Timeout waiting for latch");
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateRecoverInProgressTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateRecoverInProgressTest.java
new file mode 100644
index 000000000000..80d33ed3bee1
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateRecoverInProgressTest.java
@@ -0,0 +1,397 @@
+/*
+ * 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.aggregator;
+
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.CamelContext;
+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.processor.aggregate.AggregateController;
+import org.apache.camel.processor.aggregate.DefaultAggregateController;
+import org.apache.camel.spi.OptimisticLockingAggregationRepository;
+import org.apache.camel.spi.RecoverableAggregationRepository;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.DefaultExchangeHolder;
+import org.apache.camel.support.service.ServiceSupport;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.fail;
+
+/**
+ * The recover task must not recover an aggregated exchange that has been
moved to the completed store of a recoverable
+ * repository but has not been submitted yet by the thread that completed it.
+ */
+public class AggregateRecoverInProgressTest extends ContextTestSupport {
+
+ private final CountDownLatch paused = new CountDownLatch(1);
+ private final CountDownLatch release = new CountDownLatch(1);
+ private final RecoverableRepository repository = new
RecoverableRepository();
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Override
+ @AfterEach
+ public void tearDown() throws Exception {
+ release.countDown();
+ super.tearDown();
+ }
+
+ @Test
+ public void testCompletedByIncomingExchange() throws Exception {
+ // the completed exchange B is in the completed store, and not
submitted while the exchange A+C is submitted
+ // (completion by batch consumer completes several groups at once)
+ AggregationStrategy strategy = new BodyStrategy() {
+ @Override
+ public void onCompletion(Exchange exchange) {
+ if ("A+C".equals(exchange.getIn().getBody(String.class))) {
+ pause();
+ }
+ }
+ };
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").aggregate(header("id"),
strategy).aggregationRepository(repository)
+ .completionFromBatchConsumer()
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceivedInAnyOrder("A+C", "B");
+
+ send("A", "1");
+ send("B", "2");
+ Thread producer = new Thread(() -> send("C", "1"), "producer-C");
+ producer.start();
+
+ assertNotRecoveredWhilePaused(producer);
+ assertMockEndpointsSatisfied();
+ assertDeliveredOnce(mock);
+ }
+
+ @Test
+ public void testOptimisticLockingCompletedBySize() throws Exception {
+ repository.pauseAfterRemove("producer-B");
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").aggregate(header("id"), new
BodyStrategy()).aggregationRepository(repository)
+ .optimisticLocking().completionSize(2)
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("A+B");
+
+ send("A", "1");
+ Thread producer = new Thread(() -> send("B", "1"), "producer-B");
+ producer.start();
+
+ assertNotRecoveredWhilePaused(producer);
+ assertMockEndpointsSatisfied();
+ assertDeliveredOnce(mock);
+ }
+
+ @Test
+ public void testOptimisticLockingCompletedByTimeout() throws Exception {
+ repository.pauseAfterRemove("AggregateTimeoutChecker");
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").aggregate(header("id"), new
BodyStrategy()).aggregationRepository(repository)
+
.optimisticLocking().completionTimeout(100).completionTimeoutCheckerInterval(10)
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("A+B");
+
+ send("A", "1");
+ send("B", "1");
+
+ assertNotRecoveredWhilePaused(null);
+ assertMockEndpointsSatisfied();
+ assertDeliveredOnce(mock);
+ }
+
+ @Test
+ public void testForceDiscardingNotConfirmedIsRecovered() throws Exception {
+ // a force discarded group is not sent, so it must not be left marked
as being completed: if it could not be
+ // confirmed, it is still in the completed store, and the recover task
must recover it
+ AggregateController controller = new DefaultAggregateController();
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").aggregate(header("id"), new
BodyStrategy()).aggregationRepository(repository)
+ .completionSize(10).aggregateController(controller)
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("A");
+ mock.expectedHeaderReceived(Exchange.REDELIVERED, true);
+
+ send("A", "1");
+ repository.failConfirm.set(true);
+ assertThrows(IllegalStateException.class, () ->
controller.forceDiscardingOfGroup("1"));
+ assertFalse(repository.completed.isEmpty(), "The discarded exchange
should still be in the completed store");
+
+ assertMockEndpointsSatisfied();
+ await().atMost(20,
TimeUnit.SECONDS).until(repository.completed::isEmpty);
+ assertEquals(1, repository.recovered.get());
+ }
+
+ private void assertNotRecoveredWhilePaused(Thread producer) throws
Exception {
+ awaitLatch(paused);
+ assertFalse(repository.completed.isEmpty(), "The completed exchange
should be in the completed store");
+ // wait until the recover task has run completely at least once after
the exchange was moved to the completed
+ // store (it runs with a fixed delay, so the run with the next scan is
complete when the scan after it starts)
+ int scans = repository.scans.get();
+ await().atMost(20, TimeUnit.SECONDS).until(() ->
repository.scans.get() >= scans + 2);
+ assertEquals(0, repository.recovered.get(), "The recover task should
not recover an exchange being completed");
+
+ release.countDown();
+ if (producer != null) {
+ producer.join(10000);
+ }
+ }
+
+ private void assertDeliveredOnce(MockEndpoint mock) {
+ // the completed exchanges are confirmed, and the recover task finds
nothing more to recover
+ await().atMost(20,
TimeUnit.SECONDS).until(repository.completed::isEmpty);
+ int scans = repository.scans.get();
+ await().atMost(20, TimeUnit.SECONDS).until(() ->
repository.scans.get() >= scans + 2);
+ assertEquals(0, repository.recovered.get());
+ assertEquals(mock.getExpectedCount(), mock.getReceivedCounter());
+ }
+
+ private void send(String body, String id) {
+ template.send("direct:start", exchange -> {
+ exchange.getIn().setBody(body);
+ exchange.getIn().setHeader("id", id);
+ // used by completionFromBatchConsumer
+ exchange.setProperty(Exchange.BATCH_SIZE, 3);
+ });
+ }
+
+ private void pause() {
+ paused.countDown();
+ try {
+ if (!release.await(20, TimeUnit.SECONDS)) {
+ fail("Not released");
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ private static void awaitLatch(CountDownLatch latch) throws
InterruptedException {
+ assertTrue(latch.await(20, TimeUnit.SECONDS), "Timeout waiting for
latch");
+ }
+
+ private static class BodyStrategy implements AggregationStrategy {
+ @Override
+ public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
+ if (oldExchange == null) {
+ return newExchange;
+ }
+ oldExchange.getIn().setBody(
+ oldExchange.getIn().getBody(String.class) + "+" +
newExchange.getIn().getBody(String.class));
+ return oldExchange;
+ }
+ }
+
+ /**
+ * An in-memory recoverable repository which moves a removed exchange to a
completed store (like the JDBC
+ * repository), with optimistic locking based on a version.
+ */
+ private final class RecoverableRepository extends ServiceSupport
+ implements RecoverableAggregationRepository,
OptimisticLockingAggregationRepository {
+
+ private static final String VERSION = "RecoverableRepositoryVersion";
+
+ private final Map<String, DefaultExchangeHolder> groups = new
HashMap<>();
+ private final Map<String, Long> versions = new HashMap<>();
+ private final Map<String, DefaultExchangeHolder> completed = new
ConcurrentHashMap<>();
+ private final AtomicLong version = new AtomicLong();
+ private final AtomicInteger scans = new AtomicInteger();
+ private final AtomicInteger recovered = new AtomicInteger();
+ private final AtomicBoolean pauseAfterRemove = new AtomicBoolean();
+ private final AtomicBoolean failConfirm = new AtomicBoolean();
+ private volatile String pauseThread;
+
+ void pauseAfterRemove(String threadName) {
+ pauseThread = threadName;
+ pauseAfterRemove.set(true);
+ }
+
+ @Override
+ public synchronized Exchange add(CamelContext camelContext, String
key, Exchange exchange) {
+ Exchange old = get(camelContext, key);
+ groups.put(key, DefaultExchangeHolder.marshal(exchange, true));
+ return old;
+ }
+
+ @Override
+ public synchronized Exchange add(CamelContext camelContext, String
key, Exchange oldExchange, Exchange newExchange) {
+ Long current = versions.get(key);
+ Long expected = oldExchange == null ? null :
oldExchange.getProperty(VERSION, Long.class);
+ if (current == null ? expected != null :
!current.equals(expected)) {
+ throw new OptimisticLockingException();
+ }
+ long next = version.incrementAndGet();
+ newExchange.setProperty(VERSION, next);
+ groups.put(key, DefaultExchangeHolder.marshal(newExchange, true));
+ versions.put(key, next);
+ return oldExchange;
+ }
+
+ @Override
+ public synchronized Exchange get(CamelContext camelContext, String
key) {
+ DefaultExchangeHolder holder = groups.get(key);
+ if (holder == null) {
+ return null;
+ }
+ Exchange answer = unmarshal(camelContext, holder);
+ if (versions.containsKey(key)) {
+ answer.setProperty(VERSION, versions.get(key));
+ }
+ return answer;
+ }
+
+ @Override
+ public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ synchronized (this) {
+ Long expected = exchange.getProperty(VERSION, Long.class);
+ if (expected != null && !expected.equals(versions.get(key))) {
+ throw new OptimisticLockingException();
+ }
+ groups.remove(key);
+ versions.remove(key);
+ completed.put(exchange.getExchangeId(),
DefaultExchangeHolder.marshal(exchange, true));
+ }
+ if (pauseThread != null &&
Thread.currentThread().getName().contains(pauseThread)
+ && pauseAfterRemove.compareAndSet(true, false)) {
+ pause();
+ }
+ }
+
+ @Override
+ public void confirm(CamelContext camelContext, String exchangeId) {
+ if (failConfirm.compareAndSet(true, false)) {
+ throw new IllegalStateException("Simulated failure to confirm
" + exchangeId);
+ }
+ completed.remove(exchangeId);
+ }
+
+ @Override
+ public synchronized Set<String> getKeys() {
+ return Set.copyOf(groups.keySet());
+ }
+
+ @Override
+ public Set<String> scan(CamelContext camelContext) {
+ scans.incrementAndGet();
+ return Set.copyOf(completed.keySet());
+ }
+
+ @Override
+ public Exchange recover(CamelContext camelContext, String exchangeId) {
+ DefaultExchangeHolder holder = completed.get(exchangeId);
+ if (holder == null) {
+ return null;
+ }
+ recovered.incrementAndGet();
+ return unmarshal(camelContext, holder);
+ }
+
+ private Exchange unmarshal(CamelContext camelContext,
DefaultExchangeHolder holder) {
+ Exchange answer = new DefaultExchange(camelContext);
+ DefaultExchangeHolder.unmarshal(answer, holder);
+ return answer;
+ }
+
+ @Override
+ public void setRecoveryInterval(long interval, TimeUnit timeUnit) {
+ }
+
+ @Override
+ public void setRecoveryInterval(long interval) {
+ }
+
+ @Override
+ public long getRecoveryInterval() {
+ return 100;
+ }
+
+ @Override
+ public void setUseRecovery(boolean useRecovery) {
+ }
+
+ @Override
+ public boolean isUseRecovery() {
+ return true;
+ }
+
+ @Override
+ public void setDeadLetterUri(String deadLetterUri) {
+ }
+
+ @Override
+ public String getDeadLetterUri() {
+ return null;
+ }
+
+ @Override
+ public void setMaximumRedeliveries(int maximumRedeliveries) {
+ }
+
+ @Override
+ public int getMaximumRedeliveries() {
+ return 0;
+ }
+ }
+}