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