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

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

commit ebff68380997327cbc775d8fdb7be07259d38e03
Author: Christian Schneider <[email protected]>
AuthorDate: Mon Mar 23 17:13:01 2026 +0100

    SLING-13136 - Compute queue size by clearable not editable to make meaning 
more clear
---
 docs/metrics_overview.md                           |  2 +-
 .../journal/queue/impl/PubQueueProviderImpl.java   | 15 +++++-----
 .../journal/queue/impl/PubQueueProviderTest.java   | 32 +++++-----------------
 3 files changed, 15 insertions(+), 34 deletions(-)

diff --git a/docs/metrics_overview.md b/docs/metrics_overview.md
index c9f2f6c..2935746 100644
--- a/docs/metrics_overview.md
+++ b/docs/metrics_overview.md
@@ -45,7 +45,7 @@ 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). Values are 
cached and refreshed in the background every 30 seconds.
+- **Description**: Current size of the queue (maximum queue size). The backlog 
is measured from the minimum `lastProcessedOffset` among **clearable** 
subscribers (those with a non-null clear callback on `QueueState`). 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.
 
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 181efc8..2b47f57 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
@@ -43,7 +43,6 @@ 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;
@@ -244,7 +243,7 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
     }
 
     private int computeMaxQueueSize(String pubAgentName) {
-        Optional<Long> minOffset = getMinEditableQueueOffset(pubAgentName);
+        Optional<Long> minOffset = getMinClearableQueueOffset(pubAgentName);
         if (minOffset.isPresent()) {
             return getOffsetQueue(pubAgentName, 
minOffset.get()).getMinOffsetQueue(minOffset.get()).getSize();
         } else {
@@ -278,16 +277,16 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
         }
     }
 
-    private Optional<Long> getMinEditableQueueOffset(String pubAgentName) {
+    private Optional<Long> getMinClearableQueueOffset(String pubAgentName) {
         return callback.getSubscribedAgentIds(pubAgentName).stream()
-            .filter(subAgentName -> isEditable(pubAgentName, subAgentName))
+            .filter(subAgentName -> isClearable(pubAgentName, subAgentName))
             .map(subAgentName -> lastProcessedOffset(pubAgentName, 
subAgentName))
             .min(Long::compare);
     }
-    
-    private boolean isEditable(String pubAgentName, String subAgentName) {
-        State state = callback.getState(pubAgentName, subAgentName);
-        return state == null ? false : state.isEditable();
+
+    private boolean isClearable(String pubAgentName, String subAgentName) {
+        QueueState queueState = callback.getQueueState(pubAgentName, 
subAgentName);
+        return queueState != null && queueState.getClearCallback() != null;
     }
 
     private long lastProcessedOffset(String pubAgentName, String 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 dec721c..2f25602 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
@@ -53,12 +53,12 @@ import org.apache.sling.distribution.journal.MessageInfo;
 import org.apache.sling.distribution.journal.MessageSender;
 import org.apache.sling.distribution.journal.MessagingProvider;
 import org.apache.sling.distribution.journal.Reset;
-import org.apache.sling.distribution.journal.impl.discovery.State;
 import org.apache.sling.distribution.journal.messages.PackageMessage;
 import org.apache.sling.distribution.journal.messages.PackageMessage.ReqType;
 import org.apache.sling.distribution.journal.messages.PackageStatusMessage;
 import 
org.apache.sling.distribution.journal.messages.PackageStatusMessage.Status;
 import org.apache.sling.distribution.journal.queue.CacheCallback;
+import org.apache.sling.distribution.journal.queue.ClearCallback;
 import org.apache.sling.distribution.journal.queue.QueueState;
 import org.apache.sling.distribution.journal.shared.Topics;
 import org.apache.sling.distribution.queue.DistributionQueueEntry;
@@ -177,10 +177,7 @@ public class PubQueueProviderTest {
         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);
+            .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class)));
         queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
         queueProvider.triggerQueueSizeRefreshForTest();
         int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
@@ -205,10 +202,7 @@ public class PubQueueProviderTest {
         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);
+            .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class)));
 
         queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
         queueProvider.triggerQueueSizeRefreshForTest();
@@ -276,10 +270,7 @@ public class PubQueueProviderTest {
         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);
+            .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class)));
 
         assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0));
         queueProvider.triggerQueueSizeRefreshForTest();
@@ -302,10 +293,7 @@ public class PubQueueProviderTest {
         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);
+            .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class)));
 
         providerWithMetrics.getMaxQueueSize(PUB1_AGENT_NAME);
         providerWithMetrics.triggerQueueSizeRefreshForTest();
@@ -323,10 +311,7 @@ public class PubQueueProviderTest {
         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);
+            .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class)));
 
         queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
         queueProvider.getMaxQueueSize(PUB2_AGENT_NAME);
@@ -335,16 +320,13 @@ public class PubQueueProviderTest {
     }
 
     @Test
-    public void testQueueSizeWithNoEditableSubscribers() {
+    public void testQueueSizeWithNoClearableSubscribers() {
         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();

Reply via email to