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 de7277581262 CAMEL-25215: camel-ehcache, camel-jcache - atomic 
putIfAbsent, replace and delete in the key-value repositories (#27164)
de7277581262 is described below

commit de7277581262bce38ed82d2ed3a324ce7bf85500
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 21:10:17 2026 +0530

    CAMEL-25215: camel-ehcache, camel-jcache - atomic putIfAbsent, replace and 
delete in the key-value repositories (#27164)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../processor/EhcacheKeyValueRepository.java       |  67 +++++++++-
 .../EhcacheKeyValueRepositoryByValueTest.java      | 103 +++++++++++++++
 ...eyValueRepositoryConcurrentPutIfAbsentTest.java | 130 +++++++++++++++++++
 .../processor/EhcacheKeyValueRepositoryTest.java   |  26 ++++
 .../jcache/processor/JCacheKeyValueRepository.java |  67 +++++++++-
 .../JCacheKeyValueRepositoryByValueTest.java       | 142 +++++++++++++++++++++
 ...eyValueRepositoryConcurrentPutIfAbsentTest.java | 104 +++++++++++++++
 .../processor/JCacheKeyValueRepositoryTest.java    |  26 ++++
 .../org/apache/camel/support/KeyValueTtlValue.java |  26 ++++
 9 files changed, 681 insertions(+), 10 deletions(-)

diff --git 
a/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepository.java
 
b/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepository.java
index 30fc8ae237ef..760372bc9739 100644
--- 
a/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepository.java
+++ 
b/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepository.java
@@ -19,6 +19,7 @@ package org.apache.camel.component.ehcache.processor;
 import java.time.Duration;
 import java.util.HashSet;
 import java.util.Iterator;
+import java.util.Objects;
 import java.util.Set;
 
 import org.apache.camel.api.management.ManagedAttribute;
@@ -108,7 +109,8 @@ public class EhcacheKeyValueRepository extends 
ServiceSupport implements KeyValu
             return null;
         }
         if (entry.isExpired()) {
-            cache.remove(key);
+            // remove only the expired entry, not a new entry stored for the 
key in the meantime
+            cache.remove(key, entry);
             return null;
         }
         return entry.value();
@@ -125,9 +127,9 @@ public class EhcacheKeyValueRepository extends 
ServiceSupport implements KeyValu
     @Override
     @ManagedOperation(description = "Put a key-value pair with optional TTL")
     public Object put(String key, Object value, Duration ttl) {
-        long expiresAt = hasPositiveTtl(ttl) ? System.currentTimeMillis() + 
ttl.toMillis() : Long.MAX_VALUE;
+        KeyValueTtlValue entry = newEntry(value, ttl);
         KeyValueTtlValue previous = cache.get(key);
-        cache.put(key, new KeyValueTtlValue(value, expiresAt));
+        cache.put(key, entry);
         if (previous == null || previous.isExpired()) {
             return null;
         }
@@ -153,7 +155,8 @@ public class EhcacheKeyValueRepository extends 
ServiceSupport implements KeyValu
             return false;
         }
         if (entry.isExpired()) {
-            cache.remove(key);
+            // remove only the expired entry, not a new entry stored for the 
key in the meantime
+            cache.remove(key, entry);
             return false;
         }
         return true;
@@ -168,12 +171,61 @@ public class EhcacheKeyValueRepository extends 
ServiceSupport implements KeyValu
             if (!entry.getValue().isExpired()) {
                 keys.add(entry.getKey());
             } else {
-                cache.remove(entry.getKey());
+                cache.remove(entry.getKey(), entry.getValue());
             }
         }
         return Set.copyOf(keys);
     }
 
+    @Override
+    public Object putIfAbsent(String key, Object value, Duration ttl) {
+        KeyValueTtlValue entry = newEntry(value, ttl);
+        while (true) {
+            KeyValueTtlValue existing = cache.putIfAbsent(key, entry);
+            if (existing == null) {
+                return null;
+            }
+            if (!existing.isExpired()) {
+                return existing.value();
+            }
+            // an expired entry counts as absent: replace it, unless it was 
changed in the meantime. KeyValueTtlValue
+            // equality is per write, so the swap matches the stored entry 
also when the cache stores copies, and it
+            // fails only if another thread wrote or removed the key: the loop 
ends when no other thread does
+            if (cache.replace(key, existing, entry)) {
+                return null;
+            }
+        }
+    }
+
+    @Override
+    public boolean replace(String key, Object expectedOldValue, Object 
newValue, Duration ttl) {
+        KeyValueTtlValue entry = newEntry(newValue, ttl);
+        while (true) {
+            KeyValueTtlValue current = cache.get(key);
+            if (current == null || current.isExpired() || 
!Objects.equals(current.value(), expectedOldValue)) {
+                return false;
+            }
+            // compare-and-swap, so a concurrent update of the key is not 
overwritten
+            if (cache.replace(key, current, entry)) {
+                return true;
+            }
+        }
+    }
+
+    @Override
+    public boolean delete(String key, Object expectedValue) {
+        while (true) {
+            KeyValueTtlValue current = cache.get(key);
+            if (current == null || current.isExpired() || 
!Objects.equals(current.value(), expectedValue)) {
+                return false;
+            }
+            // compare-and-remove, so a concurrent update of the key is not 
removed
+            if (cache.remove(key, current)) {
+                return true;
+            }
+        }
+    }
+
     @Override
     @ManagedOperation(description = "Clear all entries")
     public void clear() {
@@ -186,6 +238,11 @@ public class EhcacheKeyValueRepository extends 
ServiceSupport implements KeyValu
         return keys().size();
     }
 
+    private static KeyValueTtlValue newEntry(Object value, Duration ttl) {
+        long expiresAt = hasPositiveTtl(ttl) ? System.currentTimeMillis() + 
ttl.toMillis() : Long.MAX_VALUE;
+        return new KeyValueTtlValue(value, expiresAt);
+    }
+
     private static boolean hasPositiveTtl(Duration ttl) {
         return ttl != null && !ttl.isZero() && !ttl.isNegative();
     }
diff --git 
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryByValueTest.java
 
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryByValueTest.java
new file mode 100644
index 000000000000..0a4729ddff5c
--- /dev/null
+++ 
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryByValueTest.java
@@ -0,0 +1,103 @@
+/*
+ * 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.ehcache.processor;
+
+import java.time.Duration;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.support.KeyValueTtlValue;
+import org.ehcache.Cache;
+import org.ehcache.CacheManager;
+import org.ehcache.config.builders.CacheConfigurationBuilder;
+import org.ehcache.config.builders.CacheManagerBuilder;
+import org.ehcache.config.builders.ResourcePoolsBuilder;
+import org.ehcache.impl.copy.SerializingCopier;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+
+/**
+ * A cache that stores copies of its values (here a serializing value copier) 
returns and compares copies of the
+ * entries, and a byte array value has no value equality: the compare-and-swap 
operations must still match the stored
+ * entry.
+ */
+class EhcacheKeyValueRepositoryByValueTest {
+
+    private static final String CACHE_NAME = "test-kvrepo-by-value";
+
+    private CacheManager cacheManager;
+    private Cache<String, KeyValueTtlValue> cache;
+    private EhcacheKeyValueRepository repository;
+
+    @BeforeEach
+    void setUp() throws Exception {
+        cacheManager = CacheManagerBuilder.newCacheManagerBuilder()
+                .withCache(CACHE_NAME,
+                        CacheConfigurationBuilder.newCacheConfigurationBuilder(
+                                String.class,
+                                KeyValueTtlValue.class,
+                                ResourcePoolsBuilder.heap(100))
+                                
.withValueCopier(SerializingCopier.<KeyValueTtlValue> asCopierClass()))
+                .build(true);
+        cache = cacheManager.getCache(CACHE_NAME, String.class, 
KeyValueTtlValue.class);
+
+        repository = new EhcacheKeyValueRepository(cacheManager, CACHE_NAME);
+        repository.start();
+    }
+
+    @AfterEach
+    void tearDown() throws Exception {
+        repository.stop();
+        cacheManager.close();
+    }
+
+    @Test
+    void testPutIfAbsentReplacesExpiredByteArrayValue() {
+        repository.put("key1", new byte[] { 1 }, Duration.ofMillis(100));
+        // wait on the cache itself, as reading through the repository would 
remove the expired entry
+        await().atMost(2, TimeUnit.SECONDS).until(() -> 
cache.get("key1").isExpired());
+
+        Object previous = assertTimeoutPreemptively(Duration.ofSeconds(5),
+                () -> repository.putIfAbsent("key1", new byte[] { 2 }, null));
+
+        assertThat(previous).isNull();
+        assertThat((byte[]) repository.get("key1")).containsExactly(2);
+    }
+
+    @Test
+    void testGetRemovesExpiredByteArrayValue() {
+        repository.put("key1", new byte[] { 1 }, Duration.ofMillis(100));
+        await().atMost(2, TimeUnit.SECONDS).until(() -> 
cache.get("key1").isExpired());
+
+        assertThat(repository.get("key1")).isNull();
+        assertThat(cache.get("key1")).isNull();
+    }
+
+    @Test
+    void testReplaceAndDeleteMatchTheStoredCopy() {
+        repository.put("key1", "value1", null);
+
+        assertThat(repository.replace("key1", "value1", "value2", 
null)).isTrue();
+        assertThat(repository.get("key1")).isEqualTo("value2");
+        assertThat(repository.delete("key1", "value2")).isTrue();
+        assertThat(cache.get("key1")).isNull();
+    }
+}
diff --git 
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryConcurrentPutIfAbsentTest.java
 
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryConcurrentPutIfAbsentTest.java
new file mode 100644
index 000000000000..e37a15d3478d
--- /dev/null
+++ 
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryConcurrentPutIfAbsentTest.java
@@ -0,0 +1,130 @@
+/*
+ * 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.ehcache.processor;
+
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.util.List;
+import java.util.concurrent.BrokenBarrierException;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import org.apache.camel.support.KeyValueIdempotentRepository;
+import org.apache.camel.support.KeyValueTtlValue;
+import org.ehcache.Cache;
+import org.ehcache.CacheManager;
+import org.ehcache.config.builders.CacheConfigurationBuilder;
+import org.ehcache.config.builders.CacheManagerBuilder;
+import org.ehcache.config.builders.ResourcePoolsBuilder;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * The Idempotent Consumer calls {@code add} without a lock, so {@link 
KeyValueIdempotentRepository} relies on an atomic
+ * {@code putIfAbsent} of the key-value repository.
+ */
+class EhcacheKeyValueRepositoryConcurrentPutIfAbsentTest {
+
+    private static final String CACHE_NAME = "test-kvrepo-concurrent";
+
+    private final CyclicBarrier bothRead = new CyclicBarrier(2);
+
+    private CacheManager cacheManager;
+    private EhcacheKeyValueRepository repository;
+    private ExecutorService executor;
+
+    @BeforeEach
+    void setUp() throws Exception {
+        cacheManager = CacheManagerBuilder.newCacheManagerBuilder()
+                .withCache(CACHE_NAME,
+                        CacheConfigurationBuilder.newCacheConfigurationBuilder(
+                                String.class,
+                                KeyValueTtlValue.class,
+                                ResourcePoolsBuilder.heap(100)))
+                .build(true);
+
+        repository = new 
EhcacheKeyValueRepository(readersMeetCacheManager(cacheManager), CACHE_NAME);
+        repository.start();
+        executor = Executors.newFixedThreadPool(2);
+    }
+
+    @AfterEach
+    void tearDown() throws Exception {
+        executor.shutdownNow();
+        repository.stop();
+        cacheManager.close();
+    }
+
+    @Test
+    void testConcurrentAddOfTheSameKeyAddsItOnce() throws Exception {
+        KeyValueIdempotentRepository idempotentRepository = new 
KeyValueIdempotentRepository(repository);
+        idempotentRepository.start();
+
+        Future<Boolean> first = executor.submit(() -> 
idempotentRepository.add("message-1"));
+        Future<Boolean> second = executor.submit(() -> 
idempotentRepository.add("message-1"));
+
+        assertThat(List.of(first.get(10, TimeUnit.SECONDS), second.get(10, 
TimeUnit.SECONDS)))
+                .containsExactlyInAnyOrder(true, false);
+    }
+
+    /**
+     * A cache manager whose caches let a read of a key return only when a 
second thread has read it as well (or after a
+     * timeout), so two check-then-put sequences both check before either puts.
+     */
+    @SuppressWarnings("unchecked")
+    private CacheManager readersMeetCacheManager(CacheManager delegate) {
+        return (CacheManager) 
Proxy.newProxyInstance(getClass().getClassLoader(), new Class<?>[] { 
CacheManager.class },
+                (proxy, method, args) -> {
+                    Object answer = invoke(delegate, method, args);
+                    if ("getCache".equals(method.getName()) && answer != null) 
{
+                        return readersMeetCache((Cache<String, 
KeyValueTtlValue>) answer);
+                    }
+                    return answer;
+                });
+    }
+
+    private Cache<?, ?> readersMeetCache(Cache<String, KeyValueTtlValue> 
delegate) {
+        return (Cache<?, ?>) 
Proxy.newProxyInstance(getClass().getClassLoader(), new Class<?>[] { 
Cache.class },
+                (proxy, method, args) -> {
+                    Object answer = invoke(delegate, method, args);
+                    if ("get".equals(method.getName())) {
+                        try {
+                            bothRead.await(2, TimeUnit.SECONDS);
+                        } catch (TimeoutException | BrokenBarrierException e) {
+                            // only one reader, nothing to wait for
+                        }
+                    }
+                    return answer;
+                });
+    }
+
+    private static Object invoke(Object target, Method method, Object[] args) 
throws Throwable {
+        try {
+            return method.invoke(target, args);
+        } catch (InvocationTargetException e) {
+            throw e.getCause();
+        }
+    }
+}
diff --git 
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryTest.java
 
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryTest.java
index 9dbde03ca1ae..f892dd324549 100644
--- 
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryTest.java
+++ 
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/EhcacheKeyValueRepositoryTest.java
@@ -214,6 +214,32 @@ class EhcacheKeyValueRepositoryTest {
         assertThat(repository.get("boolean")).isEqualTo(Boolean.TRUE);
     }
 
+    @Test
+    void testPutIfAbsentAddsMissingKey() {
+        assertThat(repository.putIfAbsent("key1", "value1", null)).isNull();
+
+        assertThat(repository.get("key1")).isEqualTo("value1");
+    }
+
+    @Test
+    void testPutIfAbsentKeepsExistingValue() {
+        repository.put("key1", "value1", null);
+
+        assertThat(repository.putIfAbsent("key1", "value2", 
null)).isEqualTo("value1");
+        assertThat(repository.get("key1")).isEqualTo("value1");
+    }
+
+    @Test
+    void testPutIfAbsentReplacesExpiredValue() {
+        repository.put("key1", "value1", Duration.ofMillis(100));
+        // wait on the cache itself, as reading through the repository would 
remove the expired entry
+        await().atMost(2, TimeUnit.SECONDS)
+                .until(() -> cacheManager.getCache(CACHE_NAME, String.class, 
KeyValueTtlValue.class).get("key1").isExpired());
+
+        assertThat(repository.putIfAbsent("key1", "value2", null)).isNull();
+        assertThat(repository.get("key1")).isEqualTo("value2");
+    }
+
     @Test
     void testReplaceMatchingValue() {
         repository.put("key1", "value1", null);
diff --git 
a/components/camel-jcache/src/main/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepository.java
 
b/components/camel-jcache/src/main/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepository.java
index 5212d6fcb47c..cae9ac7b9a2d 100644
--- 
a/components/camel-jcache/src/main/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepository.java
+++ 
b/components/camel-jcache/src/main/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepository.java
@@ -19,6 +19,7 @@ package org.apache.camel.component.jcache.processor;
 import java.time.Duration;
 import java.util.HashSet;
 import java.util.Iterator;
+import java.util.Objects;
 import java.util.Set;
 
 import javax.cache.Cache;
@@ -114,7 +115,8 @@ public class JCacheKeyValueRepository extends 
ServiceSupport implements CamelCon
             return null;
         }
         if (entry.isExpired()) {
-            cache.remove(key);
+            // remove only the expired entry, not a new entry stored for the 
key in the meantime
+            cache.remove(key, entry);
             return null;
         }
         return entry.value();
@@ -131,9 +133,9 @@ public class JCacheKeyValueRepository extends 
ServiceSupport implements CamelCon
     @Override
     @ManagedOperation(description = "Put a key-value pair with optional TTL")
     public Object put(String key, Object value, Duration ttl) {
-        long expiresAt = hasPositiveTtl(ttl) ? System.currentTimeMillis() + 
ttl.toMillis() : Long.MAX_VALUE;
+        KeyValueTtlValue entry = newEntry(value, ttl);
         KeyValueTtlValue previous = cache.get(key);
-        cache.put(key, new KeyValueTtlValue(value, expiresAt));
+        cache.put(key, entry);
         if (previous == null || previous.isExpired()) {
             return null;
         }
@@ -159,7 +161,8 @@ public class JCacheKeyValueRepository extends 
ServiceSupport implements CamelCon
             return false;
         }
         if (entry.isExpired()) {
-            cache.remove(key);
+            // remove only the expired entry, not a new entry stored for the 
key in the meantime
+            cache.remove(key, entry);
             return false;
         }
         return true;
@@ -174,12 +177,61 @@ public class JCacheKeyValueRepository extends 
ServiceSupport implements CamelCon
             if (!entry.getValue().isExpired()) {
                 keys.add(entry.getKey());
             } else {
-                cache.remove(entry.getKey());
+                cache.remove(entry.getKey(), entry.getValue());
             }
         }
         return Set.copyOf(keys);
     }
 
+    @Override
+    public Object putIfAbsent(String key, Object value, Duration ttl) {
+        KeyValueTtlValue entry = newEntry(value, ttl);
+        while (true) {
+            if (cache.putIfAbsent(key, entry)) {
+                return null;
+            }
+            KeyValueTtlValue existing = cache.get(key);
+            if (existing != null && !existing.isExpired()) {
+                return existing.value();
+            }
+            // an expired entry counts as absent: replace it, unless it was 
changed in the meantime. KeyValueTtlValue
+            // equality is per write, so the swap matches the stored entry 
also when the cache stores copies, and it
+            // fails only if another thread wrote or removed the key: the loop 
ends when no other thread does
+            if (existing != null && cache.replace(key, existing, entry)) {
+                return null;
+            }
+        }
+    }
+
+    @Override
+    public boolean replace(String key, Object expectedOldValue, Object 
newValue, Duration ttl) {
+        KeyValueTtlValue entry = newEntry(newValue, ttl);
+        while (true) {
+            KeyValueTtlValue current = cache.get(key);
+            if (current == null || current.isExpired() || 
!Objects.equals(current.value(), expectedOldValue)) {
+                return false;
+            }
+            // compare-and-swap, so a concurrent update of the key is not 
overwritten
+            if (cache.replace(key, current, entry)) {
+                return true;
+            }
+        }
+    }
+
+    @Override
+    public boolean delete(String key, Object expectedValue) {
+        while (true) {
+            KeyValueTtlValue current = cache.get(key);
+            if (current == null || current.isExpired() || 
!Objects.equals(current.value(), expectedValue)) {
+                return false;
+            }
+            // compare-and-remove, so a concurrent update of the key is not 
removed
+            if (cache.remove(key, current)) {
+                return true;
+            }
+        }
+    }
+
     @Override
     @ManagedOperation(description = "Clear all entries")
     public void clear() {
@@ -192,6 +244,11 @@ public class JCacheKeyValueRepository extends 
ServiceSupport implements CamelCon
         return keys().size();
     }
 
+    private static KeyValueTtlValue newEntry(Object value, Duration ttl) {
+        long expiresAt = hasPositiveTtl(ttl) ? System.currentTimeMillis() + 
ttl.toMillis() : Long.MAX_VALUE;
+        return new KeyValueTtlValue(value, expiresAt);
+    }
+
     private static boolean hasPositiveTtl(Duration ttl) {
         return ttl != null && !ttl.isZero() && !ttl.isNegative();
     }
diff --git 
a/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryByValueTest.java
 
b/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryByValueTest.java
new file mode 100644
index 000000000000..072685e15411
--- /dev/null
+++ 
b/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryByValueTest.java
@@ -0,0 +1,142 @@
+/*
+ * 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.jcache.processor;
+
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.time.Duration;
+import java.util.concurrent.TimeUnit;
+
+import javax.cache.Cache;
+
+import org.apache.camel.component.jcache.JCacheConfiguration;
+import org.apache.camel.component.jcache.JCacheHelper;
+import org.apache.camel.component.jcache.JCacheManager;
+import org.apache.camel.component.jcache.support.HazelcastTest;
+import org.apache.camel.support.KeyValueTtlValue;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+
+/**
+ * A JCache that stores values by value (the default of {@link 
JCacheConfiguration}) returns copies of the entries. A
+ * provider may compare them with {@code equals()} in {@code replace(key, old, 
new)} and {@code remove(key, old)}, as
+ * the Ehcache provider does, and a byte array value has no value equality: 
the compare-and-swap operations must still
+ * match the stored entry. The Hazelcast provider of the tests compares the 
serialized form, so the cache is wrapped to
+ * compare with {@code equals()}.
+ */
+@HazelcastTest
+class JCacheKeyValueRepositoryByValueTest extends CamelTestSupport {
+
+    private JCacheManager<String, KeyValueTtlValue> cacheManager;
+    private Cache<String, KeyValueTtlValue> cache;
+    private JCacheKeyValueRepository repository;
+
+    @Override
+    public void doPostSetup() throws Exception {
+        cacheManager = JCacheHelper.createManager(context, new 
JCacheConfiguration("kvrepo-by-value"));
+        cache = cacheManager.getCache();
+
+        repository = new JCacheKeyValueRepository();
+        repository.setCamelContext(context);
+        repository.setCache(comparingWithEquals(cache));
+        repository.start();
+    }
+
+    @Override
+    public void doPostTearDown() throws Exception {
+        if (repository != null) {
+            repository.stop();
+        }
+        if (cacheManager != null) {
+            cacheManager.close();
+        }
+    }
+
+    @Test
+    void testPutIfAbsentReplacesExpiredByteArrayValue() {
+        repository.put("key1", new byte[] { 1 }, Duration.ofMillis(100));
+        // wait on the cache itself, as reading through the repository would 
remove the expired entry
+        await().atMost(2, TimeUnit.SECONDS).until(() -> 
cache.get("key1").isExpired());
+
+        Object previous = assertTimeoutPreemptively(Duration.ofSeconds(5),
+                () -> repository.putIfAbsent("key1", new byte[] { 2 }, null));
+
+        assertThat(previous).isNull();
+        assertThat((byte[]) repository.get("key1")).containsExactly(2);
+    }
+
+    @Test
+    void testGetRemovesExpiredByteArrayValue() {
+        repository.put("key1", new byte[] { 1 }, Duration.ofMillis(100));
+        await().atMost(2, TimeUnit.SECONDS).until(() -> 
cache.get("key1").isExpired());
+
+        assertThat(repository.get("key1")).isNull();
+        assertThat(cache.get("key1")).isNull();
+    }
+
+    @Test
+    void testReplaceAndDeleteMatchTheStoredCopy() {
+        repository.put("key1", "value1", null);
+
+        assertThat(repository.replace("key1", "value1", "value2", 
null)).isTrue();
+        assertThat(repository.get("key1")).isEqualTo("value2");
+        assertThat(repository.delete("key1", "value2")).isTrue();
+        assertThat(cache.get("key1")).isNull();
+    }
+
+    /**
+     * Wraps the cache so that {@code replace(key, old, new)} and {@code 
remove(key, old)} compare the stored copy with
+     * {@code equals()}. The tests use a single thread, so the check and the 
write need no lock.
+     */
+    @SuppressWarnings("unchecked")
+    private static Cache<String, KeyValueTtlValue> 
comparingWithEquals(Cache<String, KeyValueTtlValue> delegate) {
+        return (Cache<String, KeyValueTtlValue>) Proxy.newProxyInstance(
+                JCacheKeyValueRepositoryByValueTest.class.getClassLoader(), 
new Class<?>[] { Cache.class },
+                (proxy, method, args) -> {
+                    if ("replace".equals(method.getName()) && args.length == 
3) {
+                        String key = (String) args[0];
+                        if (!args[1].equals(delegate.get(key))) {
+                            return false;
+                        }
+                        delegate.put(key, (KeyValueTtlValue) args[2]);
+                        return true;
+                    }
+                    if ("remove".equals(method.getName()) && args != null && 
args.length == 2) {
+                        String key = (String) args[0];
+                        if (!args[1].equals(delegate.get(key))) {
+                            return false;
+                        }
+                        delegate.remove(key);
+                        return true;
+                    }
+                    return invoke(delegate, method, args);
+                });
+    }
+
+    private static Object invoke(Object target, Method method, Object[] args) 
throws Throwable {
+        try {
+            return method.invoke(target, args);
+        } catch (InvocationTargetException e) {
+            throw e.getCause();
+        }
+    }
+}
diff --git 
a/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryConcurrentPutIfAbsentTest.java
 
b/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryConcurrentPutIfAbsentTest.java
new file mode 100644
index 000000000000..0dbbcf8f93df
--- /dev/null
+++ 
b/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryConcurrentPutIfAbsentTest.java
@@ -0,0 +1,104 @@
+/*
+ * 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.jcache.processor;
+
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.util.List;
+import java.util.concurrent.BrokenBarrierException;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import javax.cache.Cache;
+
+import org.apache.camel.component.jcache.JCacheConfiguration;
+import org.apache.camel.component.jcache.JCacheHelper;
+import org.apache.camel.component.jcache.JCacheManager;
+import org.apache.camel.component.jcache.support.HazelcastTest;
+import org.apache.camel.support.KeyValueIdempotentRepository;
+import org.apache.camel.support.KeyValueTtlValue;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * The Idempotent Consumer calls {@code add} without a lock, so {@link 
KeyValueIdempotentRepository} relies on an atomic
+ * {@code putIfAbsent} of the key-value repository.
+ */
+@HazelcastTest
+class JCacheKeyValueRepositoryConcurrentPutIfAbsentTest extends 
CamelTestSupport {
+
+    private final CyclicBarrier bothRead = new CyclicBarrier(2);
+
+    @Test
+    void testConcurrentAddOfTheSameKeyAddsItOnce() throws Exception {
+        JCacheManager<String, KeyValueTtlValue> cacheManager
+                = JCacheHelper.createManager(context, new 
JCacheConfiguration("kvrepo-concurrent"));
+        ExecutorService executor = Executors.newFixedThreadPool(2);
+        JCacheKeyValueRepository repository = new JCacheKeyValueRepository();
+        repository.setCamelContext(context);
+        repository.setCache(readersMeetCache(cacheManager.getCache()));
+        KeyValueIdempotentRepository idempotentRepository = new 
KeyValueIdempotentRepository(repository);
+        idempotentRepository.start();
+        try {
+            Future<Boolean> first = executor.submit(() -> 
idempotentRepository.add("message-1"));
+            Future<Boolean> second = executor.submit(() -> 
idempotentRepository.add("message-1"));
+
+            assertThat(List.of(first.get(10, TimeUnit.SECONDS), second.get(10, 
TimeUnit.SECONDS)))
+                    .containsExactlyInAnyOrder(true, false);
+        } finally {
+            executor.shutdownNow();
+            idempotentRepository.stop();
+            cacheManager.close();
+        }
+    }
+
+    /**
+     * A cache that lets a read of a key return only when a second thread has 
read it as well (or after a timeout), so
+     * two check-then-put sequences both check before either puts.
+     */
+    @SuppressWarnings("unchecked")
+    private Cache<String, KeyValueTtlValue> readersMeetCache(Cache<String, 
KeyValueTtlValue> delegate) {
+        return (Cache<String, KeyValueTtlValue>) 
Proxy.newProxyInstance(getClass().getClassLoader(),
+                new Class<?>[] { Cache.class },
+                (proxy, method, args) -> {
+                    Object answer = invoke(delegate, method, args);
+                    if ("get".equals(method.getName())) {
+                        try {
+                            bothRead.await(2, TimeUnit.SECONDS);
+                        } catch (TimeoutException | BrokenBarrierException e) {
+                            // only one reader, nothing to wait for
+                        }
+                    }
+                    return answer;
+                });
+    }
+
+    private static Object invoke(Object target, Method method, Object[] args) 
throws Throwable {
+        try {
+            return method.invoke(target, args);
+        } catch (InvocationTargetException e) {
+            throw e.getCause();
+        }
+    }
+}
diff --git 
a/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryTest.java
 
b/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryTest.java
index bdc009bc3213..f6faa5579c4e 100644
--- 
a/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryTest.java
+++ 
b/components/camel-jcache/src/test/java/org/apache/camel/component/jcache/processor/JCacheKeyValueRepositoryTest.java
@@ -214,6 +214,32 @@ class JCacheKeyValueRepositoryTest extends 
CamelTestSupport {
         assertThat(repository.get("boolean")).isEqualTo(Boolean.TRUE);
     }
 
+    @Test
+    void testPutIfAbsentAddsMissingKey() {
+        assertThat(repository.putIfAbsent("key1", "value1", null)).isNull();
+
+        assertThat(repository.get("key1")).isEqualTo("value1");
+    }
+
+    @Test
+    void testPutIfAbsentKeepsExistingValue() {
+        repository.put("key1", "value1", null);
+
+        assertThat(repository.putIfAbsent("key1", "value2", 
null)).isEqualTo("value1");
+        assertThat(repository.get("key1")).isEqualTo("value1");
+    }
+
+    @Test
+    void testPutIfAbsentReplacesExpiredValue() {
+        repository.put("key1", "value1", Duration.ofMillis(100));
+        // wait on the cache itself, as reading through the repository would 
remove the expired entry
+        await().atMost(2, TimeUnit.SECONDS)
+                .until(() -> cache.get("key1").isExpired());
+
+        assertThat(repository.putIfAbsent("key1", "value2", null)).isNull();
+        assertThat(repository.get("key1")).isEqualTo("value2");
+    }
+
     @Test
     void testReplaceMatchingValue() {
         repository.put("key1", "value1", null);
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/support/KeyValueTtlValue.java
 
b/core/camel-support/src/main/java/org/apache/camel/support/KeyValueTtlValue.java
index 59a46b2de3c0..1f8e7139d5ab 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/support/KeyValueTtlValue.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/support/KeyValueTtlValue.java
@@ -18,6 +18,7 @@ package org.apache.camel.support;
 
 import java.io.Serial;
 import java.io.Serializable;
+import java.util.concurrent.ThreadLocalRandom;
 
 /**
  * A value wrapper that holds the actual value and an expiration timestamp. 
Used by KeyValueRepository implementations
@@ -32,10 +33,13 @@ public final class KeyValueTtlValue implements Serializable 
{
 
     private final Object value;
     private final long expiresAt;
+    // identifies this write, so a copy made by a cache that stores values by 
value equals the original
+    private final long token;
 
     public KeyValueTtlValue(Object value, long expiresAt) {
         this.value = value;
         this.expiresAt = expiresAt;
+        this.token = ThreadLocalRandom.current().nextLong();
     }
 
     public Object value() {
@@ -49,4 +53,26 @@ public final class KeyValueTtlValue implements Serializable {
     public boolean isExpired() {
         return System.currentTimeMillis() >= expiresAt;
     }
+
+    /**
+     * Two instances are equal when they are the same write: the original and 
the copies a cache that stores values by
+     * value makes of it (for example a serializing copier). The wrapped value 
is not compared, as it may have no value
+     * equality (a byte array, a POJO without equals), and a compare-and-swap 
of the cache must match the stored copy of
+     * an entry whatever its value.
+     */
+    @Override
+    public boolean equals(Object o) {
+        if (this == o) {
+            return true;
+        }
+        if (!(o instanceof KeyValueTtlValue other)) {
+            return false;
+        }
+        return token == other.token && expiresAt == other.expiresAt;
+    }
+
+    @Override
+    public int hashCode() {
+        return Long.hashCode(token);
+    }
 }

Reply via email to