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 629ebe9  SLING-13136 - Provide separate metric for largest (non) 
clearable queue (#185)
629ebe9 is described below

commit 629ebe96578ab6dba01f310d846f9a2d517cfadb
Author: Christian Schneider <[email protected]>
AuthorDate: Tue Mar 24 10:30:21 2026 +0100

    SLING-13136 - Provide separate metric for largest (non) clearable queue 
(#185)
    
    * SLING-13136 - Compute queue size by clearable not editable to make 
meaning more clear
    
    * SLING-13136 - Provide separate metric for largest (non) clearable queue
    
    * Fix Sonar code smells in queue refresh and publisher test
    
    - PubQueueProviderImpl: iterate cachedMaxQueueSizes.entrySet() instead of
      keySet() plus get (performance rule)
    - DistributionPublisherTest: remove redundant eq() in getMaxQueueSize stubs
    
    Made-with: Cursor
---
 README.md                                          |  2 +-
 docs/metrics_overview.md                           | 13 +++--
 .../impl/publisher/DistributionPublisher.java      |  5 +-
 .../journal/impl/publisher/PublishMetrics.java     | 15 +++++-
 .../journal/queue/PubQueueProvider.java            |  9 ++--
 .../journal/queue/impl/PubQueueProviderImpl.java   | 53 +++++++++++-------
 .../impl/publisher/DistributionPublisherTest.java  |  4 +-
 .../journal/queue/impl/PubQueueProviderTest.java   | 62 ++++++++--------------
 8 files changed, 90 insertions(+), 73 deletions(-)

diff --git a/README.md b/README.md
index 3e72ccc..8e6044f 100644
--- a/README.md
+++ b/README.md
@@ -17,7 +17,7 @@ Publisher metrics (prefixed with 
`sling_distribution_journal_publisher_`) track
 - **Package Export**: 
`sling_distribution_journal_publisher_exported_package_size` (histogram)
 - **Request Handling**: 
`sling_distribution_journal_publisher_accepted_requests`, 
`sling_distribution_journal_publisher_dropped_requests` (meters)
 - **Package Building**: 
`sling_distribution_journal_publisher_build_package_duration`, 
`sling_distribution_journal_publisher_enqueue_package_duration` (timers)
-- **Queue Operations**: `sling_distribution_journal_publisher_queue_size` 
(gauge), `sling_distribution_journal_publisher_queue_cache_fetch_count`, 
`sling_distribution_journal_publisher_queue_access_error_count` (counters)
+- **Queue Operations**: `sling_distribution_journal_publisher_queue_size` 
(gauge, tagged `pub_name` and `clearable`), 
`sling_distribution_journal_publisher_queue_cache_fetch_count`, 
`sling_distribution_journal_publisher_queue_access_error_count` (counters)
 - **Subscriber Discovery**: 
`sling_distribution_journal_publisher_subscriber_count` (gauge)
 
 ### Subscriber Metrics
diff --git a/docs/metrics_overview.md b/docs/metrics_overview.md
index c9f2f6c..ebed07f 100644
--- a/docs/metrics_overview.md
+++ b/docs/metrics_overview.md
@@ -4,9 +4,11 @@ This document provides a comprehensive overview of all metrics 
in the Apache Sli
 
 ## Publisher Metrics
 
-All publisher metrics are prefixed with 
`sling_distribution_journal_publisher_` and include the following tag:
+Most publisher metrics are prefixed with 
`sling_distribution_journal_publisher_` and include:
 - `pub_name`: Name of the publish agent
 
+The `queue_size` gauge additionally includes `clearable` (see below).
+
 ### Package Export Metrics
 
 #### `sling_distribution_journal_publisher_exported_package_size` (Histogram)
@@ -45,14 +47,15 @@ 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.
-- **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.
+- **Description**: Max backlog depth for a subscriber cohort, measured from 
the minimum `lastProcessedOffset` in that cohort (tail size in the journal 
cache). **Clearable** cohort: subscribers with a non-null clear callback on 
`QueueState`. **Non-clearable** cohort: subscribers with a `QueueState` and a 
null clear callback. Two time series per publisher (see tags). Values are 
cached and refreshed in the background every 30 seconds. Publisher throttling 
uses the **clearable** series only  [...]
+- **Tags**: `pub_name`, `clearable` (`true` or `false`)
+- **Migration**: Previously this metric used only `pub_name`. Series are now 
distinguished by `clearable`; update dashboards and alerts to include the 
`clearable` tag (e.g. `clearable=true` for the former single series).
+- **Staleness**: The value can be up to ~30 seconds stale under normal 
conditions. When queue size computation 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.
+- **Description**: Duration of computing **both** clearable and non-clearable 
max queue sizes for an agent during each 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/DistributionPublisher.java
 
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java
index 5829126..33ba863 100644
--- 
a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java
+++ 
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisher.java
@@ -140,7 +140,8 @@ public class DistributionPublisher implements 
DistributionAgent {
         requireNonNull(metricsService);
         this.publishMetrics = new PublishMetrics(metricsService, pubAgentName);
         this.pubQueueProvider = pubQueueProvider;
-        this.publishMetrics.queueSize(() -> 
pubQueueProvider.getMaxQueueSize(pubAgentName));
+        this.publishMetrics.queueSizeByClearable(() -> 
pubQueueProvider.getMaxQueueSize(pubAgentName, true), true);
+        this.publishMetrics.queueSizeByClearable(() -> 
pubQueueProvider.getMaxQueueSize(pubAgentName, false), false);
 
         distLog = new DefaultDistributionLog(pubAgentName, this.getClass(), 
DefaultDistributionLog.LogLevel.INFO);
         distributionLogEventListener = new 
DistributionLogEventListener(context, distLog, pubAgentName);
@@ -212,7 +213,7 @@ public class DistributionPublisher implements 
DistributionAgent {
             distLog.info(msg);
             return new 
SimpleDistributionResponse(DistributionRequestState.DROPPED, msg);
         }
-        int queueSize = pubQueueProvider.getMaxQueueSize(pubAgentName);
+        int queueSize = pubQueueProvider.getMaxQueueSize(pubAgentName, true);
         int sleepMs = getSleepTime(queueSize);
         sleep(sleepMs);
         final PackageMessage pkg = buildPackage(resourceResolver, request);
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 d6d1bec..628c631 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
@@ -33,6 +33,7 @@ import org.apache.sling.distribution.journal.metrics.Tag;
 
 public class PublishMetrics {
     private static final String TAG_AGENT_NAME = "pub_name";
+    private static final String TAG_CLEARABLE = "clearable";
 
     public static final String PUB_COMPONENT = 
"distribution.journal.publisher.";
     private static final String EXPORTED_PACKAGE_SIZE = PUB_COMPONENT + 
"exported_package_size";
@@ -47,10 +48,12 @@ public class PublishMetrics {
     /** 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 String pubAgentName;
     private final List<Tag> tags;
     private final MetricsService metricsService;
 
     public PublishMetrics(MetricsService metricsService, String pubAgentName) {
+        this.pubAgentName = pubAgentName;
         this.tags = Arrays.asList(Tag.of(TAG_AGENT_NAME, pubAgentName));
         this.metricsService = metricsService;
     }
@@ -122,8 +125,16 @@ public class PublishMetrics {
         metricsService.gauge(getMetricName(SUBSCRIBER_COUNT, tags), 
subscriberCountCallback);
     }
 
-    public void queueSize(Supplier<Integer> queueSizeCallback) {
-        metricsService.gauge(getMetricName(QUEUE_SIZE, tags), 
queueSizeCallback);
+    /**
+     * Gauge of max queue backlog for subscribers in the given cohort.
+     *
+     * @param clearable {@code true} for clearable subscribers, {@code false} 
for non-clearable
+     */
+    public void queueSizeByClearable(Supplier<Integer> queueSizeCallback, 
boolean clearable) {
+        List<Tag> queueSizeTags = Arrays.asList(
+                Tag.of(TAG_AGENT_NAME, pubAgentName),
+                Tag.of(TAG_CLEARABLE, Boolean.toString(clearable)));
+        metricsService.gauge(getMetricName(QUEUE_SIZE, queueSizeTags), 
queueSizeCallback);
     }
 
 }
diff --git 
a/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java
 
b/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java
index 3da54df..7cefb90 100644
--- 
a/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java
+++ 
b/src/main/java/org/apache/sling/distribution/journal/queue/PubQueueProvider.java
@@ -38,11 +38,14 @@ public interface PubQueueProvider extends Closeable {
     DistributionQueue getQueue(String pubAgentName, String queueName);
     
     /**
-     * Get maximum size of all queues for a pubAgentName
+     * Maximum backlog depth for a subscriber cohort (min {@code 
lastProcessedOffset} in cohort, then journal tail size).
+     *
      * @param pubAgentName name of the pub agent
-     * @return max size of all queues or 0 if there are none
+     * @param clearable {@code true} for clearable subscribers (non-null clear 
callback); {@code false} for non-clearable
+     *                  (queue state present, null clear callback). Publisher 
throttling uses {@code clearable == true}.
+     * @return max size for that cohort or 0 if there are none
      */
-    int getMaxQueueSize(String pubAgentName);
+    int getMaxQueueSize(String pubAgentName, boolean clearable);
 
     @Nonnull
     OffsetQueue<DistributionQueueItem> getOffsetQueue(String pubAgentName, 
long minOffset);
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..851f6b5 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;
@@ -87,7 +86,7 @@ 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 Map<String, CachedMaxQueueSizes> cachedMaxQueueSizes = new 
ConcurrentHashMap<>();
 
     private final ScheduledExecutorService queueSizeExecutor;
 
@@ -239,17 +238,17 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
     
 
     @Override
-    public int getMaxQueueSize(String pubAgentName) {
-        return cachedQueueSizes.computeIfAbsent(pubAgentName, key -> 0);
+    public int getMaxQueueSize(String pubAgentName, boolean clearable) {
+        CachedMaxQueueSizes entry = 
cachedMaxQueueSizes.computeIfAbsent(pubAgentName, k -> new 
CachedMaxQueueSizes());
+        return clearable ? entry.clearable : entry.nonClearable;
     }
 
-    private int computeMaxQueueSize(String pubAgentName) {
-        Optional<Long> minOffset = getMinEditableQueueOffset(pubAgentName);
+    private int computeMaxQueueSize(String pubAgentName, boolean 
clearableCohort) {
+        Optional<Long> minOffset = getMinQueueOffsetForCohort(pubAgentName, 
clearableCohort);
         if (minOffset.isPresent()) {
             return getOffsetQueue(pubAgentName, 
minOffset.get()).getMinOffsetQueue(minOffset.get()).getSize();
-        } else {
-            return 0;
         }
+        return 0;
     }
 
     /**
@@ -260,34 +259,45 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
     }
 
     private void refreshQueueSizes() {
-        for (String agentName : cachedQueueSizes.keySet()) {
+        for (Map.Entry<String, CachedMaxQueueSizes> e : 
cachedMaxQueueSizes.entrySet()) {
+            String agentName = e.getKey();
+            CachedMaxQueueSizes entry = e.getValue();
             try {
                 long startNanos = System.nanoTime();
-                int size = computeMaxQueueSize(agentName);
+                int clearableSize = computeMaxQueueSize(agentName, true);
+                int nonClearableSize = computeMaxQueueSize(agentName, false);
                 long elapsedMs = (System.nanoTime() - startNanos) / 1_000_000;
-                cachedQueueSizes.put(agentName, size);
+                entry.clearable = clearableSize;
+                entry.nonClearable = nonClearableSize;
 
                 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);
+            } catch (Exception ex) {
+                LOG.warn("Failed to refresh queue size for agent {}", 
agentName, ex);
             }
         }
     }
 
-    private Optional<Long> getMinEditableQueueOffset(String pubAgentName) {
+    private Optional<Long> getMinQueueOffsetForCohort(String pubAgentName, 
boolean clearableCohort) {
         return callback.getSubscribedAgentIds(pubAgentName).stream()
-            .filter(subAgentName -> isEditable(pubAgentName, subAgentName))
+            .filter(subAgentName -> inClearableCohort(pubAgentName, 
subAgentName, clearableCohort))
             .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();
+
+    /**
+     * @param clearableCohort {@code true} for clearable subscribers, {@code 
false} for non-clearable (state present, no clear callback).
+     */
+    private boolean inClearableCohort(String pubAgentName, String 
subAgentName, boolean clearableCohort) {
+        QueueState queueState = callback.getQueueState(pubAgentName, 
subAgentName);
+        if (queueState == null) {
+            return false;
+        }
+        boolean clearable = queueState.getClearCallback() != null;
+        return clearableCohort ? clearable : !clearable;
     }
 
     private long lastProcessedOffset(String pubAgentName, String subAgentName) 
{
@@ -310,4 +320,9 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
         return new PubQueueCache(queuedNotifier, callback);
     }
 
+    private static final class CachedMaxQueueSizes {
+        volatile int clearable;
+        volatile int nonClearable;
+    }
+
 }
diff --git 
a/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java
 
b/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java
index 200d564..0b97197 100644
--- 
a/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java
+++ 
b/src/test/java/org/apache/sling/distribution/journal/impl/publisher/DistributionPublisherTest.java
@@ -175,7 +175,7 @@ public class DistributionPublisherTest {
     @Test
     public void testQueueSizeLimitHalf() throws IOException, 
DistributionException {
         int queueSize = DEFAULT_QUEUE_SIZE_LIMIT + DEFAULT_QUEUE_SIZE_LIMIT / 
2;
-        
when(pubQueueProvider.getMaxQueueSize(PUB1AGENT1)).thenReturn(queueSize);
+        when(pubQueueProvider.getMaxQueueSize(PUB1AGENT1, 
true)).thenReturn(queueSize);
         DistributionRequest request = new 
SimpleDistributionRequest(DistributionRequestType.ADD, "/test");
         long time = distribute(request);
         assertThat(time, greaterThanOrEqualTo(500L));
@@ -185,7 +185,7 @@ public class DistributionPublisherTest {
     @Test
     public void testDoubleQueueSizeLimitReached() throws IOException, 
DistributionException {
         int queueSize = DEFAULT_QUEUE_SIZE_LIMIT * 2;
-        
when(pubQueueProvider.getMaxQueueSize(PUB1AGENT1)).thenReturn(queueSize);
+        when(pubQueueProvider.getMaxQueueSize(PUB1AGENT1, 
true)).thenReturn(queueSize);
         DistributionRequest request = new 
SimpleDistributionRequest(DistributionRequestType.ADD, "/test");
         long time = distribute(request);
         assertThat(time, greaterThanOrEqualTo(1000L));
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..05108b0 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,14 +177,12 @@ 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);
-        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+            .thenReturn(new QueueState(0, -1, 0, mock(ClearCallback.class)));
+        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true);
         queueProvider.triggerQueueSizeRefreshForTest();
-        int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+        int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true);
         assertThat(size, equalTo(2));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, false), 
equalTo(0));
     }
     
     @Test
@@ -194,7 +192,7 @@ public class PubQueueProviderTest {
         when(callback.getQueueState(Mockito.eq(PUB1_AGENT_NAME), 
Mockito.any()))
             .thenReturn(new QueueState(0, -1, 0, null));
 
-        int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME);
+        int size = queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true);
         assertThat(size, equalTo(0));
     }
 
@@ -205,18 +203,15 @@ 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.getMaxQueueSize(PUB1_AGENT_NAME, true);
         queueProvider.triggerQueueSizeRefreshForTest();
-        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), 
equalTo(2));
 
         handler.handle(info(3L), packageMessage("packageid4", 
PUB1_AGENT_NAME));
 
-        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), 
equalTo(2));
     }
     
     @SuppressWarnings("null")
@@ -276,14 +271,11 @@ 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));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), 
equalTo(0));
         queueProvider.triggerQueueSizeRefreshForTest();
-        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(2));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), 
equalTo(2));
     }
 
     @Test
@@ -302,12 +294,9 @@ 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.getMaxQueueSize(PUB1_AGENT_NAME, true);
         providerWithMetrics.triggerQueueSizeRefreshForTest();
         verify(mockTimer).update(anyLong(), eq(TimeUnit.MILLISECONDS));
         providerWithMetrics.close();
@@ -323,32 +312,27 @@ 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);
+        queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true);
+        queueProvider.getMaxQueueSize(PUB2_AGENT_NAME, true);
         queueProvider.triggerQueueSizeRefreshForTest();
-        assertThat(queueProvider.getMaxQueueSize(PUB2_AGENT_NAME), equalTo(1));
+        assertThat(queueProvider.getMaxQueueSize(PUB2_AGENT_NAME, true), 
equalTo(1));
     }
 
     @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.getMaxQueueSize(PUB1_AGENT_NAME, true);
         queueProvider.triggerQueueSizeRefreshForTest();
-        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME), equalTo(0));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, true), 
equalTo(0));
+        assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, false), 
equalTo(2));
     }
 
     private int queueSize() {

Reply via email to