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;
+        }
+    }
+}

Reply via email to