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

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


The following commit(s) were added to refs/heads/master by this push:
     new 27105b7  SLING-13135 - Compute queue size in background (#184)
27105b7 is described below

commit 27105b7d5810f766dd09db04d0d255bc6b62500d
Author: Christian Schneider <[email protected]>
AuthorDate: Thu Mar 12 14:18:20 2026 +0100

    SLING-13135 - Compute queue size in background (#184)
    
    * SLING-13135 - Compute queue size in background
    
    * SLING-13135 - Fixes from review
    
    * SLING-13135 - Add metric for queue size computation
    
    * SLING-13135 - Fixes from review
    
    * SLING-13135 - Fixes from review
    
    * SLING-13135 - Fixes from review
    
    * SLING-13135 - Fixes from review
---
 docs/metrics_overview.md                           |   9 +-
 .../journal/impl/publisher/PublishMetrics.java     |   2 +
 .../queue/impl/PubQueueProviderFactoryImpl.java    |  28 ++++--
 .../journal/queue/impl/PubQueueProviderImpl.java   |  61 ++++++++++++
 .../journal/queue/impl/PubQueueProviderTest.java   | 110 +++++++++++++++++++++
 5 files changed, 199 insertions(+), 11 deletions(-)

diff --git a/docs/metrics_overview.md b/docs/metrics_overview.md
index 395e7f2..c9f2f6c 100644
--- a/docs/metrics_overview.md
+++ b/docs/metrics_overview.md
@@ -45,7 +45,14 @@ All publisher metrics are prefixed with 
`sling_distribution_journal_publisher_`
 
 #### `sling_distribution_journal_publisher_queue_size` (Gauge)
 - **Type**: Gauge
-- **Description**: Current size of the queue (maximum queue size)
+- **Description**: Current size of the queue (maximum queue size). Values are 
cached and refreshed in the background every 30 seconds.
+- **Tags**: `pub_name`
+- **Staleness**: The value can be up to ~30 seconds stale under normal 
conditions. When `computeMaxQueueSize()` scales linearly (O(n)) with queue size 
and takes longer than the refresh interval, staleness increases: refresh cycles 
can back up, and the displayed value may lag further behind the true queue 
size. See 
`sling_distribution_journal_publisher_queue_size_computation_duration` for 
monitoring computation cost.
+
+#### `sling_distribution_journal_publisher_queue_size_computation_duration` 
(Timer)
+- **Type**: Timer
+- **Unit**: Milliseconds
+- **Description**: Duration of computing the max queue size per agent during 
background refresh. The computation is O(n) with queue size; slow values (>30s) 
indicate increased staleness.
 - **Tags**: `pub_name`
 
 #### `sling_distribution_journal_publisher_queue_cache_fetch_count` (Counter)
diff --git 
a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java
 
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java
index 5692450..d6d1bec 100644
--- 
a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java
+++ 
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublishMetrics.java
@@ -44,6 +44,8 @@ public class PublishMetrics {
     private static final String QUEUE_ACCESS_ERROR_COUNT = PUB_COMPONENT + 
"queue_access_error_count";
     private static final String SUBSCRIBER_COUNT = PUB_COMPONENT + 
"subscriber_count";
     private static final String QUEUE_SIZE = PUB_COMPONENT + "queue_size";
+    /** Metric name for queue size computation duration (use with 
Tag.of("pub_name", agentName)). */
+    public static final String QUEUE_SIZE_COMPUTATION_DURATION = PUB_COMPONENT 
+ "queue_size_computation_duration";
 
     private final List<Tag> tags;
     private final MetricsService metricsService;
diff --git 
a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
 
b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
index 671cf77..2c4f001 100644
--- 
a/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
+++ 
b/src/main/java/org/apache/sling/distribution/journal/queue/impl/PubQueueProviderFactoryImpl.java
@@ -18,32 +18,40 @@
  */
 package org.apache.sling.distribution.journal.queue.impl;
 
+import org.apache.sling.commons.metrics.MetricsService;
 import org.apache.sling.distribution.journal.queue.CacheCallback;
 import org.apache.sling.distribution.journal.queue.PubQueueProvider;
 import org.apache.sling.distribution.journal.queue.PubQueueProviderFactory;
 import org.osgi.framework.BundleContext;
+import org.osgi.service.component.annotations.Activate;
 import org.osgi.service.component.annotations.Component;
 import org.osgi.service.component.annotations.Reference;
+import org.osgi.service.component.annotations.ReferenceCardinality;
 import org.osgi.service.event.EventAdmin;
 
 @Component
 public class PubQueueProviderFactoryImpl implements PubQueueProviderFactory {
-    
-    @Reference
-    private EventAdmin eventAdmin;
 
-    @Reference
-    private QueueErrors queueErrors;
-    
-    private BundleContext context;
+    private final EventAdmin eventAdmin;
+    private final QueueErrors queueErrors;
+    private final MetricsService metricsService;
+    private final BundleContext context;
 
-    public void activate(BundleContext context) {
+    @Activate
+    public PubQueueProviderFactoryImpl(
+            @Reference EventAdmin eventAdmin,
+            @Reference QueueErrors queueErrors,
+            @Reference(cardinality = ReferenceCardinality.OPTIONAL) 
MetricsService metricsService,
+            BundleContext context) {
+        this.eventAdmin = eventAdmin;
+        this.queueErrors = queueErrors;
+        this.metricsService = metricsService;
         this.context = context;
     }
 
     @Override
     public PubQueueProvider create(CacheCallback callback) {
-        return new PubQueueProviderImpl(eventAdmin, queueErrors, callback, 
context);
+        return new PubQueueProviderImpl(eventAdmin, queueErrors, callback, 
context, metricsService);
     }
-    
+
 }
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..181efc8 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
@@ -20,6 +20,7 @@ package org.apache.sling.distribution.journal.queue.impl;
 
 import static 
org.apache.sling.commons.scheduler.Scheduler.PROPERTY_SCHEDULER_CONCURRENT;
 import static 
org.apache.sling.commons.scheduler.Scheduler.PROPERTY_SCHEDULER_PERIOD;
+import static 
org.apache.sling.distribution.journal.metrics.TaggedMetrics.getMetricName;
 
 import java.util.Dictionary;
 import java.util.HashSet;
@@ -28,6 +29,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;
@@ -35,12 +39,16 @@ import javax.annotation.ParametersAreNonnullByDefault;
 
 import org.apache.commons.io.IOUtils;
 import org.apache.commons.lang3.StringUtils;
+import org.apache.sling.commons.metrics.MetricsService;
+import org.apache.sling.commons.metrics.Timer;
 import org.apache.sling.commons.scheduler.Scheduler;
 import org.apache.sling.distribution.journal.MessageInfo;
 import org.apache.sling.distribution.journal.impl.discovery.State;
 import 
org.apache.sling.distribution.journal.impl.publisher.PackageQueuedNotifier;
 import org.apache.sling.distribution.journal.messages.PackageStatusMessage;
 import 
org.apache.sling.distribution.journal.messages.PackageStatusMessage.Status;
+import org.apache.sling.distribution.journal.impl.publisher.PublishMetrics;
+import org.apache.sling.distribution.journal.metrics.Tag;
 import org.apache.sling.distribution.journal.queue.CacheCallback;
 import org.apache.sling.distribution.journal.queue.OffsetQueue;
 import org.apache.sling.distribution.journal.queue.PubQueueProvider;
@@ -62,6 +70,8 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
      */
     private static final int CLEANUP_THRESHOLD = 10_000;
 
+    static final long QUEUE_SIZE_REFRESH_INTERVAL_SECONDS = 30;
+
     private static final Logger LOG = 
LoggerFactory.getLogger(PubQueueProviderImpl.class);
     
     private final PackageQueuedNotifier queuedNotifier;
@@ -77,13 +87,33 @@ 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;
+
+    @Nullable
+    private final MetricsService metricsService;
+
     private ServiceRegistration<?> reg;
 
     public PubQueueProviderImpl(EventAdmin eventAdmin, QueueErrors 
queueErrors, CacheCallback callback, BundleContext context) {
+        this(eventAdmin, queueErrors, callback, context, null);
+    }
+
+    public PubQueueProviderImpl(EventAdmin eventAdmin, QueueErrors 
queueErrors, CacheCallback callback, BundleContext context,
+            @Nullable MetricsService metricsService) {
         queuedNotifier = new PackageQueuedNotifier(eventAdmin);
         this.queueErrors = queueErrors;
         this.callback = callback;
+        this.metricsService = metricsService;
         cache = newCache();
+        queueSizeExecutor = Executors.newSingleThreadScheduledExecutor(r -> {
+            Thread t = new Thread(r, "queue-size-refresh");
+            t.setDaemon(true);
+            return t;
+        });
+        queueSizeExecutor.scheduleAtFixedRate(this::refreshQueueSizes,
+                0, QUEUE_SIZE_REFRESH_INTERVAL_SECONDS, TimeUnit.SECONDS);
         startCleanupTask(context);
         LOG.info("Started Publisher queue provider service");
     }
@@ -99,6 +129,7 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
 
     @Override
     public void close() {
+        queueSizeExecutor.shutdownNow();
         PubQueueCache queueCache = this.cache;
         if (queueCache != null) {
             queueCache.close();
@@ -209,6 +240,10 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
 
     @Override
     public int getMaxQueueSize(String pubAgentName) {
+        return cachedQueueSizes.computeIfAbsent(pubAgentName, key -> 0);
+    }
+
+    private int computeMaxQueueSize(String pubAgentName) {
         Optional<Long> minOffset = getMinEditableQueueOffset(pubAgentName);
         if (minOffset.isPresent()) {
             return getOffsetQueue(pubAgentName, 
minOffset.get()).getMinOffsetQueue(minOffset.get()).getSize();
@@ -217,6 +252,32 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
         }
     }
 
+    /**
+     * Package-private for tests. Triggers an immediate synchronous refresh of 
cached queue sizes.
+     */
+    void triggerQueueSizeRefreshForTest() {
+        refreshQueueSizes();
+    }
+
+    private void refreshQueueSizes() {
+        for (String agentName : cachedQueueSizes.keySet()) {
+            try {
+                long startNanos = System.nanoTime();
+                int size = computeMaxQueueSize(agentName);
+                long elapsedMs = (System.nanoTime() - startNanos) / 1_000_000;
+                cachedQueueSizes.put(agentName, size);
+
+                if (metricsService != null) {
+                    Timer timer = metricsService.timer(getMetricName(
+                            PublishMetrics.QUEUE_SIZE_COMPUTATION_DURATION, 
Tag.of("pub_name", agentName)));
+                    timer.update(elapsedMs, TimeUnit.MILLISECONDS);
+                }
+            } 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..dec721c 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
@@ -20,6 +20,8 @@ package org.apache.sling.distribution.journal.queue.impl;
 
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.equalTo;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.atLeast;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
@@ -29,6 +31,7 @@ import java.io.Closeable;
 import java.io.IOException;
 import java.lang.management.ManagementFactory;
 import java.util.Collections;
+import java.util.concurrent.TimeUnit;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Set;
@@ -42,6 +45,8 @@ import javax.management.ObjectInstance;
 import javax.management.ObjectName;
 import javax.management.ReflectionException;
 
+import org.apache.sling.commons.metrics.MetricsService;
+import org.apache.sling.commons.metrics.Timer;
 import org.apache.sling.distribution.journal.HandlerAdapter;
 import org.apache.sling.distribution.journal.MessageHandler;
 import org.apache.sling.distribution.journal.MessageInfo;
@@ -176,6 +181,8 @@ public class PubQueueProviderTest {
         State state = Mockito.mock(State.class);
         when(state.isEditable()).thenReturn(true);
         when(callback.getState(Mockito.eq(PUB1_AGENT_NAME), 
Mockito.anyString())).thenReturn(state);
+        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+        queueProvider.triggerQueueSizeRefreshForTest();
         int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
         assertThat(size, equalTo(2));
     }
@@ -190,6 +197,27 @@ public class PubQueueProviderTest {
         int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
         assertThat(size, equalTo(0));
     }
+
+    @Test
+    public void testQueueSizeIsCached() {
+        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);
+
+        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+        queueProvider.triggerQueueSizeRefreshForTest();
+        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
@@ -241,6 +269,88 @@ public class PubQueueProviderTest {
         assertThat(queueSize(), equalTo(1));
     }
 
+    @Test
+    public void testQueueSizeNewAgentReturnsZeroUntilRefresh() {
+        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(0));
+        queueProvider.triggerQueueSizeRefreshForTest();
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+    }
+
+    @Test
+    public void testQueueSizeWithMetricsServiceRecordsTimer() {
+        MetricsService mockMetrics = mock(MetricsService.class);
+        Timer mockTimer = mock(Timer.class);
+        when(mockMetrics.timer(Mockito.anyString())).thenReturn(mockTimer);
+
+        QueueErrors queueErrors = mock(QueueErrors.class);
+        PubQueueProviderImpl providerWithMetrics = new PubQueueProviderImpl(
+            eventAdmin, queueErrors, callback, context, mockMetrics);
+        MessageHandler<PackageMessage> handler2 = 
handlerCaptor.getAllValues().get(1);
+
+        handler2.handle(info(0L), packageMessage("packageid1", 
PUB1_AGENT_NAME));
+        handler2.handle(info(1L), packageMessage("packageid2", 
PUB2_AGENT_NAME));
+        handler2.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);
+
+        providerWithMetrics.getMaxQueueSize(PUB1_AGENT_NAME);
+        providerWithMetrics.triggerQueueSizeRefreshForTest();
+        verify(mockTimer).update(anyLong(), eq(TimeUnit.MILLISECONDS));
+        providerWithMetrics.close();
+    }
+
+    @Test
+    public void testRefreshQueueSizesHandlesExceptionForOneAgent() {
+        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))
+            .thenThrow(new RuntimeException("test"));
+        when(callback.getSubscribedAgentIds(PUB2_AGENT_NAME))
+            .thenReturn(Collections.singleton("sub1"));
+        when(callback.getQueueState(Mockito.eq(PUB2_AGENT_NAME), 
Mockito.eq("sub1")))
+            .thenReturn(new QueueState(0, -1, 0, null));
+        State state = Mockito.mock(State.class);
+        when(state.isEditable()).thenReturn(true);
+        when(callback.getState(Mockito.eq(PUB2_AGENT_NAME), 
Mockito.eq(SUB_AGENT_NAME))).thenReturn(state);
+
+        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+        queueProvider.getMaxQueueSize(PUB2_AGENT_NAME);
+        queueProvider.triggerQueueSizeRefreshForTest();
+        assertThat(queueProvider.getMaxQueueSize(PUB2_AGENT_NAME), equalTo(1));
+    }
+
+    @Test
+    public void testQueueSizeWithNoEditableSubscribers() {
+        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(false);
+        when(callback.getState(Mockito.eq(PUB1_AGENT_NAME), 
Mockito.anyString())).thenReturn(state);
+
+        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+        queueProvider.triggerQueueSizeRefreshForTest();
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0));
+    }
+
     private int queueSize() {
         return queueProvider.getOffsetQueue(PUB1_AGENT_NAME, 0).getSize();
     }

Reply via email to