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
