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

cschneider pushed a commit to branch SLING-13135
in repository 
https://gitbox.apache.org/repos/asf/sling-org-apache-sling-distribution-journal.git

commit 5b644c02bb4824857746c87376df70da513a6313
Author: Christian Schneider <[email protected]>
AuthorDate: Wed Mar 11 11:38:15 2026 +0100

    SLING-13135 - Compute queue size in background
---
 .../journal/queue/impl/PubQueueProviderImpl.java   | 31 ++++++++++++++++++++++
 .../journal/queue/impl/PubQueueProviderTest.java   | 19 +++++++++++++
 2 files changed, 50 insertions(+)

diff --git 
a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java
 
b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java
index 3b2d74c..b726df5 100644
--- 
a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java
+++ 
b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderImpl.java
@@ -28,6 +28,9 @@ import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
 
 import javax.annotation.Nonnull;
 import javax.annotation.Nullable;
@@ -62,6 +65,8 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
      */
     private static final int CLEANUP_THRESHOLD = 10_000;
 
+    static final long QUEUE_SIZE_REFRESH_SECONDS = 30;
+
     private static final Logger LOG = 
LoggerFactory.getLogger(PubQueueProviderImpl.class);
     
     private final PackageQueuedNotifier queuedNotifier;
@@ -77,6 +82,10 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
      */
     private final Map<String, OffsetQueue<Long>> errorQueues = new 
ConcurrentHashMap<>();
 
+    private final Map<String, Integer> cachedQueueSizes = new 
ConcurrentHashMap<>();
+
+    private final ScheduledExecutorService queueSizeExecutor;
+
     private ServiceRegistration<?> reg;
 
     public PubQueueProviderImpl(EventAdmin eventAdmin, QueueErrors 
queueErrors, CacheCallback callback, BundleContext context) {
@@ -84,6 +93,13 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
         this.queueErrors = queueErrors;
         this.callback = callback;
         cache = newCache();
+        queueSizeExecutor = Executors.newSingleThreadScheduledExecutor(r -> {
+            Thread t = new Thread(r, "queue-size-refresh");
+            t.setDaemon(true);
+            return t;
+        });
+        queueSizeExecutor.scheduleAtFixedRate(this::refreshQueueSizes,
+                QUEUE_SIZE_REFRESH_SECONDS, QUEUE_SIZE_REFRESH_SECONDS, 
TimeUnit.SECONDS);
         startCleanupTask(context);
         LOG.info("Started Publisher queue provider service");
     }
@@ -99,6 +115,7 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
 
     @Override
     public void close() {
+        queueSizeExecutor.shutdownNow();
         PubQueueCache queueCache = this.cache;
         if (queueCache != null) {
             queueCache.close();
@@ -209,6 +226,10 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
 
     @Override
     public int getMaxQueueSize(String pubAgentName) {
+        return cachedQueueSizes.computeIfAbsent(pubAgentName, 
this::computeMaxQueueSize);
+    }
+
+    private int computeMaxQueueSize(String pubAgentName) {
         Optional<Long> minOffset = getMinEditableQueueOffset(pubAgentName);
         if (minOffset.isPresent()) {
             return getOffsetQueue(pubAgentName, 
minOffset.get()).getMinOffsetQueue(minOffset.get()).getSize();
@@ -217,6 +238,16 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
         }
     }
 
+    private void refreshQueueSizes() {
+        for (String agentName : cachedQueueSizes.keySet()) {
+            try {
+                cachedQueueSizes.put(agentName, 
computeMaxQueueSize(agentName));
+            } catch (Exception e) {
+                LOG.warn("Failed to refresh queue size for agent {}", 
agentName, e);
+            }
+        }
+    }
+
     private Optional<Long> getMinEditableQueueOffset(String pubAgentName) {
         return callback.getSubscribedAgentIds(pubAgentName).stream()
             .filter(subAgentName -> isEditable(pubAgentName, subAgentName))
diff --git 
a/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java
 
b/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java
index 008610e..3f6801d 100644
--- 
a/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java
+++ 
b/src/test/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderTest.java
@@ -190,6 +190,25 @@ public class PubQueueProviderTest {
         int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
         assertThat(size, equalTo(0));
     }
+
+    @Test
+    public void testQueueSizeIsCached() throws Exception {
+        handler.handle(info(0L), packageMessage("packageid1", 
PUB1_AGENT_NAME));
+        handler.handle(info(1L), packageMessage("packageid2", 
PUB2_AGENT_NAME));
+        handler.handle(info(2L), packageMessage("packageid3", 
PUB1_AGENT_NAME));
+        
when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(Collections.singleton("sub1"));
+        when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), 
Mockito.any()))
+            .thenReturn(new QueueState(0, -1, 0, null));
+        State state = Mockito.mock(State.class);
+        when(state.isEditable()).thenReturn(true);
+        when(callback.getState(Mockito.eq(PUB1_AGENT_NAME), 
Mockito.anyString())).thenReturn(state);
+
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+
+        handler.handle(info(3L), packageMessage("packageid4", 
PUB1_AGENT_NAME));
+
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+    }
     
     @SuppressWarnings("null")
     @Test

Reply via email to