This is an automated email from the ASF dual-hosted git repository.

gnodet 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 2c0810d2c58f CAMEL-20199: Replace synchronized with ReentrantLock in 
AI components (#25245)
2c0810d2c58f is described below

commit 2c0810d2c58f2a19f2be7b8d4cb6546e7da1be4c
Author: Guillaume Nodet <[email protected]>
AuthorDate: Thu Jul 30 17:49:34 2026 +0200

    CAMEL-20199: Replace synchronized with ReentrantLock in AI components 
(#25245)
    
    CAMEL-20199: Replace synchronized with ReentrantLock in AI components
    
    Convert synchronized blocks to ReentrantLock in camel-a2a (A2AProgress,
    GuardedTaskStore, AbstractTaskPredictor) with proper store lock cleanup
    on endpoint shutdown.
    
    Co-Authored-By: Claude Opus 4.6 <[email protected]>
---
 .../apache/camel/component/a2a/A2AEndpoint.java    |  5 +++
 .../apache/camel/component/a2a/A2AProgress.java    | 37 ++++++++++++++++++++--
 .../component/a2a/state/GuardedTaskStore.java      | 29 ++++++++++++++---
 .../huggingface/tasks/AbstractTaskPredictor.java   |  7 +++-
 4 files changed, 71 insertions(+), 7 deletions(-)

diff --git 
a/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AEndpoint.java
 
b/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AEndpoint.java
index a0f8cca59ddc..c9868f484217 100644
--- 
a/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AEndpoint.java
+++ 
b/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AEndpoint.java
@@ -274,6 +274,11 @@ public class A2AEndpoint extends DefaultEndpoint {
 
     private void cleanupEndpointResources() {
         if (taskStoreOwned && taskStore != null) {
+            // Clean up the cached lock before stopping the store. Only owned 
stores
+            // (InMemoryTaskStore) can appear in STORE_LOCKS — 
registry-discovered stores
+            // are wrapped in GuardedTaskStore, which provides its own lock 
via getLock()
+            // and never enters the static map.
+            A2AProgress.removeStoreLock(taskStore);
             try {
                 ServiceHelper.stopService(taskStore);
             } catch (Exception e) {
diff --git 
a/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AProgress.java
 
b/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AProgress.java
index d67c4a9ece4c..e40eade78875 100644
--- 
a/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AProgress.java
+++ 
b/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/A2AProgress.java
@@ -18,6 +18,8 @@ package org.apache.camel.component.a2a;
 
 import java.util.ArrayList;
 import java.util.List;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
 
 import org.apache.camel.Endpoint;
 import org.apache.camel.Exchange;
@@ -30,6 +32,7 @@ import org.apache.camel.component.a2a.model.TaskState;
 import org.apache.camel.component.a2a.model.TaskStatus;
 import org.apache.camel.component.a2a.model.TextPart;
 import org.apache.camel.component.a2a.state.A2ATaskStore;
+import org.apache.camel.component.a2a.state.GuardedTaskStore;
 
 /**
  * Static utility for emitting A2A progress updates from any exchange. 
Resolves the task store via the exchange's
@@ -37,9 +40,31 @@ import org.apache.camel.component.a2a.state.A2ATaskStore;
  */
 public final class A2AProgress {
 
+    private static final ConcurrentHashMap<A2ATaskStore, ReentrantLock> 
STORE_LOCKS = new ConcurrentHashMap<>();
+
     private A2AProgress() {
     }
 
+    /**
+     * Returns the lock associated with the given store. If the store is a 
{@link GuardedTaskStore}, its own lock is
+     * returned so that callers coordinate with the store's lifecycle 
operations. Otherwise, a per-store fallback lock
+     * is created and cached.
+     */
+    private static ReentrantLock lockFor(A2ATaskStore store) {
+        if (store instanceof GuardedTaskStore guarded) {
+            return guarded.getLock();
+        }
+        return STORE_LOCKS.computeIfAbsent(store, k -> new ReentrantLock());
+    }
+
+    /**
+     * Removes the cached lock for the given store. Called during endpoint 
shutdown to prevent the static
+     * {@code STORE_LOCKS} map from leaking entries after a store is no longer 
in use.
+     */
+    static void removeStoreLock(A2ATaskStore store) {
+        STORE_LOCKS.remove(store);
+    }
+
     /**
      * Emit a status update with {@link TaskState#WORKING} state.
      *
@@ -92,7 +117,9 @@ public final class A2AProgress {
             return;
         }
 
-        synchronized (taskContext.store) {
+        ReentrantLock storeLock = lockFor(taskContext.store);
+        storeLock.lock();
+        try {
             Task task = taskContext.store.get(taskContext.taskId);
             if (task != null) {
                 List<Artifact> artifacts = task.artifacts() != null ? new 
ArrayList<>(task.artifacts()) : new ArrayList<>();
@@ -105,6 +132,8 @@ public final class A2AProgress {
                         
.artifact(artifact).append(append).lastChunk(lastChunk).build();
                 taskContext.store.notifySubscribers(taskContext.taskId, 
StreamResponse.ofArtifactUpdate(event));
             }
+        } finally {
+            storeLock.unlock();
         }
     }
 
@@ -120,7 +149,9 @@ public final class A2AProgress {
             return;
         }
 
-        synchronized (taskContext.store) {
+        ReentrantLock storeLock = lockFor(taskContext.store);
+        storeLock.lock();
+        try {
             Task task = taskContext.store.get(taskContext.taskId);
             if (task != null) {
                 List<Message> history = task.history() != null ? new 
ArrayList<>(task.history()) : new ArrayList<>();
@@ -129,6 +160,8 @@ public final class A2AProgress {
                 taskContext.store.put(taskContext.taskId, updated);
                 taskContext.store.notifySubscribers(taskContext.taskId, 
StreamResponse.ofMessage(message));
             }
+        } finally {
+            storeLock.unlock();
         }
     }
 
diff --git 
a/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/state/GuardedTaskStore.java
 
b/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/state/GuardedTaskStore.java
index 810b946bcbc1..565ad60ed804 100644
--- 
a/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/state/GuardedTaskStore.java
+++ 
b/components/camel-ai/camel-a2a/src/main/java/org/apache/camel/component/a2a/state/GuardedTaskStore.java
@@ -21,6 +21,7 @@ import java.util.Set;
 import java.util.UUID;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.locks.ReentrantLock;
 
 import org.apache.camel.component.a2a.model.StreamResponse;
 import org.apache.camel.component.a2a.model.Task;
@@ -49,7 +50,7 @@ public final class GuardedTaskStore implements A2ATaskStore {
 
     private final A2ATaskStore delegate;
     private final boolean allowLocalWebhookUrls;
-    private final Object lifecycleLock = new Object();
+    private final ReentrantLock lock = new ReentrantLock();
     private final ConcurrentHashMap<String, 
CopyOnWriteArrayList<A2ATaskSubscriber>> subscribers
             = new ConcurrentHashMap<>();
 
@@ -58,15 +59,29 @@ public final class GuardedTaskStore implements A2ATaskStore 
{
         this.allowLocalWebhookUrls = allowLocalWebhookUrls;
     }
 
+    /**
+     * Returns the lock used by this store for serializing lifecycle 
operations. External callers (such as
+     * {@code A2AProgress}) that need to perform atomic read-modify-write 
sequences on this store should acquire this
+     * lock to coordinate with the store's own write operations.
+     *
+     * @return the reentrant lock guarding lifecycle state transitions
+     */
+    public ReentrantLock getLock() {
+        return lock;
+    }
+
     @Override
     public void put(String taskId, Task task) {
-        synchronized (lifecycleLock) {
+        lock.lock();
+        try {
             Task existing = delegate.get(taskId);
             if (isTerminal(existing)) {
                 LOG.debug("Ignoring put for task {} - already in terminal 
state {}", taskId, existing.status().state());
                 return;
             }
             delegate.put(taskId, task);
+        } finally {
+            lock.unlock();
         }
     }
 
@@ -179,7 +194,8 @@ public final class GuardedTaskStore implements A2ATaskStore 
{
     @Override
     public void updateStatusAndNotify(String taskId, TaskStatus status) {
         StreamResponse event;
-        synchronized (lifecycleLock) {
+        lock.lock();
+        try {
             Task task = delegate.get(taskId);
             if (task == null) {
                 return;
@@ -196,13 +212,16 @@ public final class GuardedTaskStore implements 
A2ATaskStore {
                     .contextId(task.contextId())
                     .status(status)
                     .build());
+        } finally {
+            lock.unlock();
         }
         notifySubscribers(taskId, event);
     }
 
     @Override
     public Task cancelIfNotTerminal(String taskId) {
-        synchronized (lifecycleLock) {
+        lock.lock();
+        try {
             Task task = delegate.get(taskId);
             if (task == null) {
                 return null;
@@ -214,6 +233,8 @@ public final class GuardedTaskStore implements A2ATaskStore 
{
             Task canceled = Task.builder(task).status(new 
TaskStatus(TaskState.CANCELED)).build();
             delegate.put(taskId, canceled);
             return canceled;
+        } finally {
+            lock.unlock();
         }
     }
 
diff --git 
a/components/camel-ai/camel-huggingface/src/main/java/org/apache/camel/component/huggingface/tasks/AbstractTaskPredictor.java
 
b/components/camel-ai/camel-huggingface/src/main/java/org/apache/camel/component/huggingface/tasks/AbstractTaskPredictor.java
index 7a5bafed3d47..46dc72bdc9af 100644
--- 
a/components/camel-ai/camel-huggingface/src/main/java/org/apache/camel/component/huggingface/tasks/AbstractTaskPredictor.java
+++ 
b/components/camel-ai/camel-huggingface/src/main/java/org/apache/camel/component/huggingface/tasks/AbstractTaskPredictor.java
@@ -27,6 +27,7 @@ import java.nio.file.attribute.BasicFileAttributes;
 import java.nio.file.attribute.PosixFilePermission;
 import java.nio.file.attribute.PosixFilePermissions;
 import java.util.EnumSet;
+import java.util.concurrent.locks.ReentrantLock;
 
 import ai.djl.inference.Predictor;
 import ai.djl.modality.Input;
@@ -45,6 +46,7 @@ import org.slf4j.LoggerFactory;
 public abstract class AbstractTaskPredictor implements TaskPredictor {
 
     protected static final Logger LOG = 
LoggerFactory.getLogger(AbstractTaskPredictor.class);
+    private final ReentrantLock lock = new ReentrantLock();
     protected HuggingFaceEndpoint endpoint;
     protected HuggingFaceConfiguration config;
     protected ZooModel<Input, Output> model;
@@ -131,12 +133,15 @@ public abstract class AbstractTaskPredictor implements 
TaskPredictor {
     public void predict(Exchange exchange) throws Exception {
         Input input = prepareInput(exchange);
         if (config.isPooling()) {
-            synchronized (this) {
+            lock.lock();
+            try {
                 if (predictor == null) {
                     predictor = model.newPredictor();
                 }
                 Output output = predictor.predict(input);
                 processOutput(exchange, output);
+            } finally {
+                lock.unlock();
             }
         } else {
             try (Predictor<Input, Output> djlPredictor = model.newPredictor()) 
{

Reply via email to