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();
