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