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

joerghoh pushed a commit to branch SLING-13271
in repository 
https://gitbox.apache.org/repos/asf/sling-org-apache-sling-event.git

commit 3171054ee626c1fab51a3646d6f1d58c75c532d6
Author: Joerg Hoh <[email protected]>
AuthorDate: Sun Jul 26 13:09:18 2026 +0200

    SLING-13271 fix concurrency issues
---
 .../impl/jobs/config/TopologyCapabilities.java     | 19 ++---
 .../impl/jobs/config/TopologyCapabilitiesTest.java | 82 ++++++++++++++++++++++
 2 files changed, 92 insertions(+), 9 deletions(-)

diff --git 
a/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilities.java
 
b/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilities.java
index f338716..693ff32 100644
--- 
a/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilities.java
+++ 
b/src/main/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilities.java
@@ -25,6 +25,7 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.TreeMap;
+import java.util.concurrent.ConcurrentHashMap;
 
 import org.apache.sling.discovery.InstanceDescription;
 import org.apache.sling.discovery.TopologyView;
@@ -48,7 +49,7 @@ public class TopologyCapabilities {
     private final Map<String, List<InstanceDescription>> instanceCapabilities;
 
     /** Round robin map. */
-    private final Map<String, Integer> roundRobinMap = new HashMap<String, 
Integer>();
+    private final Map<String, Integer> roundRobinMap = new 
ConcurrentHashMap<String, Integer>();
 
     /** Instance map. */
     private final Map<String, InstanceDescription> instanceMap = new 
HashMap<String, InstanceDescription>();
@@ -302,14 +303,14 @@ public class TopologyCapabilities {
             }
             // TODO - this is a simple round robin which is not based on the 
actual load
             //        of the instances
-            Integer index = this.roundRobinMap.get(jobTopic);
-            if (index == null) {
-                index = 0;
-            }
-            if (index >= potentialTargets.size()) {
-                index = 0;
-            }
-            this.roundRobinMap.put(jobTopic, index + 1);
+            final int size = potentialTargets.size();
+            final int nextCounter = this.roundRobinMap.merge(jobTopic, 1, 
(previous, increment) -> {
+                if (previous >= size) {
+                    return 1;
+                }
+                return previous + 1;
+            });
+            final int index = nextCounter - 1;
             final String result = potentialTargets.get(index).getSlingId();
             logger.debug("Target for {} : {}", jobTopic, result);
             return result;
diff --git 
a/src/test/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilitiesTest.java
 
b/src/test/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilitiesTest.java
index d68fac0..570e66c 100644
--- 
a/src/test/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilitiesTest.java
+++ 
b/src/test/java/org/apache/sling/event/impl/jobs/config/TopologyCapabilitiesTest.java
@@ -18,11 +18,24 @@
  */
 package org.apache.sling.event.impl.jobs.config;
 
+import java.util.ArrayList;
 import java.util.Collections;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
 
 import org.apache.sling.discovery.ClusterView;
 import org.apache.sling.discovery.InstanceDescription;
 import org.apache.sling.discovery.TopologyView;
+import org.apache.sling.event.jobs.QueueConfiguration;
 import org.junit.Before;
 import org.junit.Test;
 import org.mockito.Mockito;
@@ -70,4 +83,73 @@ public class TopologyCapabilitiesTest {
         assertEquals(0, caps.getPotentialTargets("x/y").size());
         assertEquals(1, caps.getPotentialTargets("d/1/2").size());
     }
+
+    @Test
+    public void testConcurrentDetectTargetIsRaceFree() throws Exception {
+        final int instanceCount = 50;
+        final int threads = 10;
+        final int callsPerThread = 500;
+        final int expectedPerTarget = (threads * callsPerThread) / 
instanceCount;
+
+        final ClusterView cv = Mockito.mock(ClusterView.class);
+        Mockito.when(cv.getId()).thenReturn("cluster");
+
+        final Set<InstanceDescription> instances = new LinkedHashSet<>();
+        for (int i = 0; i < instanceCount; i++) {
+            final InstanceDescription desc = 
Mockito.mock(InstanceDescription.class);
+            Mockito.when(desc.getSlingId()).thenReturn("instance-" + i);
+            
Mockito.when(desc.getProperty(TopologyCapabilities.PROPERTY_TOPICS)).thenReturn("foo");
+            Mockito.when(desc.getClusterView()).thenReturn(cv);
+            instances.add(desc);
+        }
+        final InstanceDescription local = instances.iterator().next();
+        Mockito.when(local.isLeader()).thenReturn(true);
+
+        final TopologyView tv = Mockito.mock(TopologyView.class);
+        Mockito.when(tv.getInstances()).thenReturn(instances);
+        Mockito.when(tv.getLocalInstance()).thenReturn(local);
+
+        final JobManagerConfiguration config = 
Mockito.mock(JobManagerConfiguration.class);
+        final TopologyCapabilities localCaps = new TopologyCapabilities(tv, 
config);
+
+        final InternalQueueConfiguration queueConfig = 
Mockito.mock(InternalQueueConfiguration.class);
+        
Mockito.when(queueConfig.getType()).thenReturn(QueueConfiguration.Type.UNORDERED);
+        
Mockito.when(queueConfig.isPreferRunOnCreationInstance()).thenReturn(false);
+        final QueueConfigurationManager.QueueInfo queueInfo = new 
QueueConfigurationManager.QueueInfo();
+        queueInfo.queueConfiguration = queueConfig;
+        queueInfo.queueName = "foo";
+
+        // count how often each instance is selected as a target
+        final Map<String, AtomicInteger> hits = new ConcurrentHashMap<>();
+        final ExecutorService executor = Executors.newFixedThreadPool(threads);
+        try {
+            // release all threads at once (via the latch) to maximize 
contention on detectTarget
+            final CountDownLatch start = new CountDownLatch(1);
+            final List<Future<?>> futures = new ArrayList<>();
+            for (int t = 0; t < threads; t++) {
+                futures.add(executor.submit(() -> {
+                    start.await();
+                    for (int i = 0; i < callsPerThread; i++) {
+                        final String target = localCaps.detectTarget("foo", 
null, queueInfo);
+                        hits.computeIfAbsent(target, k -> new 
AtomicInteger()).incrementAndGet();
+                    }
+                    return null;
+                }));
+            }
+            start.countDown();
+            // wait for all callers to finish (and surface any exception 
thrown in a worker)
+            for (final Future<?> f : futures) {
+                f.get(30, TimeUnit.SECONDS);
+            }
+        } finally {
+            executor.shutdownNow();
+        }
+
+        // total calls are a whole multiple of the target count, so a correct 
atomic round robin
+        // hits every target exactly the same number of times; a lost update 
would skew the counts
+        assertEquals(instanceCount, hits.size());
+        for (final AtomicInteger hit : hits.values()) {
+            assertEquals(expectedPerTarget, hit.get());
+        }
+    }
 }

Reply via email to