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 02060c3  SLING-13147 - Support aggregated queues (#186)
02060c3 is described below

commit 02060c31a9ecdf69aeb0575ffafd877435d42a18
Author: Christian Schneider <[email protected]>
AuthorDate: Mon Mar 30 18:18:11 2026 +0200

    SLING-13147 - Support aggregated queues (#186)
    
    * SLING-13147 - Support aggregated queues
    
    * SLING-13147 - Fix sonar issues
---
 .../impl/publisher/DistributionPublisher.java      | 14 +++-
 .../impl/publisher/PublisherConfiguration.java     |  6 ++
 .../journal/queue/PubQueueProvider.java            | 28 +++++++
 .../journal/queue/impl/PubQueueProviderImpl.java   | 90 ++++++++++++++++++++++
 .../impl/publisher/DistributionPublisherTest.java  | 37 +++++++++
 .../journal/queue/impl/PubQueueProviderTest.java   | 72 +++++++++++++++++
 6 files changed, 244 insertions(+), 3 deletions(-)

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 33ba863..2653c64 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
@@ -107,6 +107,8 @@ public class DistributionPublisher implements 
DistributionAgent {
 
     private final int maxQueueSizeDelay;
 
+    private final boolean aggregateSubscriberQueues;
+
     private final Consumer<PackageMessage> sender;
 
     private final DistributionLogEventListener distributionLogEventListener;
@@ -150,13 +152,14 @@ public class DistributionPublisher implements 
DistributionAgent {
         queuedTimeout = config.queuedTimeout();
         queueSizeLimit = config.queueSizeLimit();
         maxQueueSizeDelay = config.maxQueueSizeDelay();
+        aggregateSubscriberQueues = config.aggregateSubscriberQueues();
         pkgType = packageBuilder.getType();
 
         this.sender = messagingProvider.createSender(Topics.PACKAGE_TOPIC);
         publishMetrics.subscriberCount(() -> 
discoveryService.getSubscriberCount(pubAgentName));
         
-        distLog.info("Started Publisher agent={} with packageBuilder={}, 
limitEnabled={}, queuedTimeout={}, queueSizeLimit={}, maxQueueSizeDelay={}",
-                pubAgentName, pkgType, limitEnabled, queuedTimeout, 
queueSizeLimit, maxQueueSizeDelay);
+        distLog.info("Started Publisher agent={} with packageBuilder={}, 
limitEnabled={}, queuedTimeout={}, queueSizeLimit={}, maxQueueSizeDelay={}, 
aggregateSubscriberQueues={}",
+                pubAgentName, pkgType, limitEnabled, queuedTimeout, 
queueSizeLimit, maxQueueSizeDelay, aggregateSubscriberQueues);
     }
 
     @Deactivate
@@ -174,13 +177,18 @@ public class DistributionPublisher implements 
DistributionAgent {
     @Nonnull
     @Override
     public Iterable<String> getQueueNames() {
+        if (aggregateSubscriberQueues) {
+            return 
Collections.unmodifiableCollection(pubQueueProvider.getAggregatedQueueNames(pubAgentName));
+        }
         return 
Collections.unmodifiableCollection(pubQueueProvider.getQueueNames(pubAgentName));
     }
 
     @Override
     public DistributionQueue getQueue(String queueName) {
         try {
-            DistributionQueue queue = pubQueueProvider.getQueue(pubAgentName, 
queueName);
+            DistributionQueue queue = aggregateSubscriberQueues
+                    ? pubQueueProvider.getAggregatedQueue(pubAgentName, 
queueName)
+                    : pubQueueProvider.getQueue(pubAgentName, queueName);
             if (queue == null) {
                 publishMetrics.getQueueAccessErrorCount().increment();
             }
diff --git 
a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublisherConfiguration.java
 
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublisherConfiguration.java
index 2573a9f..b2ab84d 100644
--- 
a/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublisherConfiguration.java
+++ 
b/src/main/java/org/apache/sling/distribution/journal/impl/publisher/PublisherConfiguration.java
@@ -44,4 +44,10 @@ public @interface PublisherConfiguration {
     int maxQueueSizeDelay() default 20000;
     
     int queueSizeLimit() default DEFAULT_QUEUE_SIZE_LIMIT;
+
+    @AttributeDefinition(
+            name = "Aggregate subscriber queues",
+            description = "If true, list only virtual queues \"persisted\" 
(clearable cohort backlog) and \"public\" (all subscribers; read-only) "
+                    + "instead of per-subscriber queues; no error queues. 
Names \"persisted\" and \"public\" are reserved.")
+    boolean aggregateSubscriberQueues() default false;
 }
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 7cefb90..d1d1d81 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
@@ -34,6 +34,18 @@ import 
org.apache.sling.distribution.queue.spi.DistributionQueue;
 @ParametersAreNonnullByDefault
 public interface PubQueueProvider extends Closeable {
 
+    /**
+     * Virtual queue name for aggregate mode: backlog from the minimum {@code 
lastProcessedOffset}
+     * among clearable subscribers; clear applies to every clearable 
subscriber.
+     */
+    String AGGREGATED_QUEUE_PERSISTED = "persisted";
+
+    /**
+     * Virtual queue name for aggregate mode: backlog from the minimum {@code 
lastProcessedOffset}
+     * among all subscribers with queue state; not clearable.
+     */
+    String AGGREGATED_QUEUE_PUBLIC = "public";
+
     @Nullable
     DistributionQueue getQueue(String pubAgentName, String queueName);
     
@@ -57,6 +69,22 @@ public interface PubQueueProvider extends Closeable {
      */
     Set<String> getQueueNames(String pubAgentName);
 
+    /**
+     * Queue names when {@code aggregateSubscriberQueues} is enabled on the 
publisher:
+     * {@link #AGGREGATED_QUEUE_PERSISTED} (omitted if there is no clearable 
subscriber)
+     * and {@link #AGGREGATED_QUEUE_PUBLIC} (omitted if no subscriber has 
queue state).
+     */
+    @Nonnull
+    Set<String> getAggregatedQueueNames(String pubAgentName);
+
+    /**
+     * Resolve a virtual queue from {@link #getAggregatedQueueNames(String)}.
+     *
+     * @return {@code null} if {@code queueName} is not an aggregated queue or 
the cohort is empty
+     */
+    @Nullable
+    DistributionQueue getAggregatedQueue(String pubAgentName, String 
queueName);
+
     @Nonnull
     PackageQueuedNotifier getQueuedNotifier();
 
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 7505a1f..b8882ef 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
@@ -26,6 +26,7 @@ import java.util.Dictionary;
 import java.util.HashSet;
 import java.util.Hashtable;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
@@ -49,6 +50,7 @@ import 
org.apache.sling.distribution.journal.messages.PackageStatusMessage.Statu
 import org.apache.sling.distribution.journal.impl.publisher.PublishMetrics;
 import org.apache.sling.distribution.journal.metrics.Tag;
 import org.apache.sling.distribution.journal.queue.CacheCallback;
+import org.apache.sling.distribution.journal.queue.ClearCallback;
 import org.apache.sling.distribution.journal.queue.OffsetQueue;
 import org.apache.sling.distribution.journal.queue.PubQueueProvider;
 import org.apache.sling.distribution.journal.queue.QueueState;
@@ -178,6 +180,94 @@ public class PubQueueProviderImpl implements 
PubQueueProvider, Runnable {
         return queueNames;
     }
 
+    @Nonnull
+    @Override
+    public Set<String> getAggregatedQueueNames(String pubAgentName) {
+        Set<String> queueNames = new HashSet<>();
+        if (getMinQueueOffsetForCohort(pubAgentName, true).isPresent()) {
+            queueNames.add(AGGREGATED_QUEUE_PERSISTED);
+        }
+        if (getMinLastProcessedOffsetAll(pubAgentName).isPresent()) {
+            queueNames.add(AGGREGATED_QUEUE_PUBLIC);
+        }
+        return queueNames;
+    }
+
+    @Nullable
+    @Override
+    public DistributionQueue getAggregatedQueue(String pubAgentName, String 
queueName) {
+        if (AGGREGATED_QUEUE_PERSISTED.equals(queueName)) {
+            return buildAggregatedQueue(pubAgentName, 
AGGREGATED_QUEUE_PERSISTED, true);
+        }
+        if (AGGREGATED_QUEUE_PUBLIC.equals(queueName)) {
+            return buildAggregatedQueue(pubAgentName, AGGREGATED_QUEUE_PUBLIC, 
false);
+        }
+        return null;
+    }
+
+    @Nullable
+    private DistributionQueue buildAggregatedQueue(String pubAgentName, String 
virtualQueueName, boolean persisted) {
+        Optional<Long> minLastProcessed = persisted
+                ? getMinQueueOffsetForCohort(pubAgentName, true)
+                : getMinLastProcessedOffsetAll(pubAgentName);
+        if (!minLastProcessed.isPresent()) {
+            return null;
+        }
+        Optional<String> stragglerId = findStragglerSubAgentId(pubAgentName, 
persisted);
+        if (!stragglerId.isPresent()) {
+            return null;
+        }
+        QueueState state = callback.getQueueState(pubAgentName, 
stragglerId.get());
+        if (state == null) {
+            return null;
+        }
+        long minOffset = minLastProcessed.get() + 1;
+        OffsetQueue<DistributionQueueItem> agentQueue = 
getOffsetQueue(pubAgentName, minOffset);
+        Throwable error = queueErrors.getError(pubAgentName, 
stragglerId.get());
+        ClearCallback clearCallback = persisted ? offset -> 
clearAllClearableSubscribers(pubAgentName, offset) : null;
+        return new PubQueue(virtualQueueName, 
agentQueue.getMinOffsetQueue(minOffset), state.getHeadRetries(), error,
+                clearCallback);
+    }
+
+    private void clearAllClearableSubscribers(String pubAgentName, long 
offset) {
+        for (String subAgentId : callback.getSubscribedAgentIds(pubAgentName)) 
{
+            QueueState qs = callback.getQueueState(pubAgentName, subAgentId);
+            if (qs != null && qs.getClearCallback() != null) {
+                qs.getClearCallback().clear(offset);
+            }
+        }
+    }
+
+    private Optional<String> findStragglerSubAgentId(String pubAgentName, 
boolean persistedCohort) {
+        Optional<Long> cohortMin = persistedCohort
+                ? getMinQueueOffsetForCohort(pubAgentName, true)
+                : getMinLastProcessedOffsetAll(pubAgentName);
+        if (!cohortMin.isPresent()) {
+            return Optional.empty();
+        }
+        long minLp = cohortMin.get();
+        return callback.getSubscribedAgentIds(pubAgentName).stream()
+                .filter(subId -> {
+                    QueueState qs = callback.getQueueState(pubAgentName, 
subId);
+                    if (qs == null || qs.getLastProcessedOffset() != minLp) {
+                        return false;
+                    }
+                    if (persistedCohort) {
+                        return qs.getClearCallback() != null;
+                    }
+                    return true;
+                })
+                .min(String::compareTo);
+    }
+
+    private Optional<Long> getMinLastProcessedOffsetAll(String pubAgentName) {
+        return callback.getSubscribedAgentIds(pubAgentName).stream()
+                .map(subId -> callback.getQueueState(pubAgentName, subId))
+                .filter(Objects::nonNull)
+                .map(QueueState::getLastProcessedOffset)
+                .min(Long::compare);
+    }
+
     @Nonnull
     @Override
     public PackageQueuedNotifier getQueuedNotifier() {
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 0b97197..3d2c101 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
@@ -22,6 +22,7 @@ import static 
org.apache.sling.distribution.journal.impl.publisher.PublisherConf
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.contains;
 import static org.hamcrest.Matchers.containsInAnyOrder;
+import static org.hamcrest.Matchers.sameInstance;
 import static org.hamcrest.Matchers.containsString;
 import static org.hamcrest.Matchers.equalTo;
 import static org.hamcrest.Matchers.greaterThanOrEqualTo;
@@ -33,6 +34,7 @@ import static org.junit.Assert.assertNull;
 import static org.junit.Assert.fail;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
@@ -42,6 +44,7 @@ import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
 
@@ -261,6 +264,40 @@ public class DistributionPublisherTest {
         assertEquals("Wrong getQueue error counter",1, count);
     }
 
+    @Test
+    public void testAggregatedQueueNamesDelegatesToAggregatedProvider() {
+        Map<String, Object> props = Map.of(
+                "name", PUB1AGENT1,
+                "maxQueueSizeDelay", "1000",
+                "aggregateSubscriberQueues", Boolean.TRUE);
+        PublisherConfiguration config = 
Converters.standardConverter().convert(props).to(PublisherConfiguration.class);
+        DistributionPublisher aggPublisher = new 
DistributionPublisher(messagingProvider, packageBuilder, discoveryService, 
factory,
+                eventAdmin, metricsService, pubQueueProvider, 
Condition.INSTANCE, config, context.bundleContext());
+        
when(pubQueueProvider.getAggregatedQueueNames(PUB1AGENT1)).thenReturn(Set.of(PubQueueProvider.AGGREGATED_QUEUE_PERSISTED));
+        Iterable<String> names = aggPublisher.getQueueNames();
+        assertThat(names, 
contains(PubQueueProvider.AGGREGATED_QUEUE_PERSISTED));
+        verify(pubQueueProvider).getAggregatedQueueNames(PUB1AGENT1);
+        verify(pubQueueProvider, never()).getQueueNames(anyString());
+        aggPublisher.deactivate();
+    }
+
+    @Test
+    public void testAggregatedGetQueueDelegatesToAggregatedProvider() {
+        Map<String, Object> props = Map.of(
+                "name", PUB1AGENT1,
+                "maxQueueSizeDelay", "1000",
+                "aggregateSubscriberQueues", Boolean.TRUE);
+        PublisherConfiguration config = 
Converters.standardConverter().convert(props).to(PublisherConfiguration.class);
+        DistributionPublisher aggPublisher = new 
DistributionPublisher(messagingProvider, packageBuilder, discoveryService, 
factory,
+                eventAdmin, metricsService, pubQueueProvider, 
Condition.INSTANCE, config, context.bundleContext());
+        PubQueue q = new PubQueue(PubQueueProvider.AGGREGATED_QUEUE_PERSISTED, 
new OffsetQueueImpl<>(), 0, null, null);
+        when(pubQueueProvider.getAggregatedQueue(PUB1AGENT1, 
PubQueueProvider.AGGREGATED_QUEUE_PERSISTED)).thenReturn(q);
+        
assertThat(aggPublisher.getQueue(PubQueueProvider.AGGREGATED_QUEUE_PERSISTED), 
sameInstance(q));
+        verify(pubQueueProvider).getAggregatedQueue(PUB1AGENT1, 
PubQueueProvider.AGGREGATED_QUEUE_PERSISTED);
+        verify(pubQueueProvider, never()).getQueue(anyString(), anyString());
+        aggPublisher.deactivate();
+    }
+
     private long getQueueAccessErrorCount() {
         return new PublishMetrics(metricsService, 
PUB1AGENT1).getQueueAccessErrorCount().getCount();
     }
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 e2a6368..9997705 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
@@ -19,7 +19,10 @@
 package org.apache.sling.distribution.journal.queue.impl;
 
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.containsInAnyOrder;
 import static org.hamcrest.Matchers.equalTo;
+import static org.hamcrest.Matchers.notNullValue;
+import static 
org.apache.sling.distribution.queue.DistributionQueueCapabilities.CLEARABLE;
 import static org.mockito.ArgumentMatchers.anyLong;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.atLeast;
@@ -30,7 +33,9 @@ import static org.mockito.Mockito.when;
 import java.io.Closeable;
 import java.io.IOException;
 import java.lang.management.ManagementFactory;
+import java.util.Arrays;
 import java.util.Collections;
+import java.util.HashSet;
 import java.util.concurrent.TimeUnit;
 import java.util.Iterator;
 import java.util.List;
@@ -59,6 +64,7 @@ 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.PubQueueProvider;
 import org.apache.sling.distribution.journal.queue.QueueState;
 import org.apache.sling.distribution.journal.shared.Topics;
 import org.apache.sling.distribution.queue.DistributionQueueEntry;
@@ -337,6 +343,72 @@ public class PubQueueProviderTest {
         assertThat(queueProvider.getMaxQueueSize(PUB1_AGENT_NAME, false), 
equalTo(2));
     }
 
+    @Test
+    public void testAggregatedQueueNamesOmitsPersistedWhenNoClearable() {
+        handler.handle(info(1L), packageMessage("p1", PUB1_AGENT_NAME));
+        when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(new 
HashSet<>(Arrays.asList("sub1", "sub2")));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub1")).thenReturn(new 
QueueState(0, -1, 0, null));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub2")).thenReturn(new 
QueueState(0, -1, 0, null));
+        Set<String> names = 
queueProvider.getAggregatedQueueNames(PUB1_AGENT_NAME);
+        assertThat(names, 
equalTo(Collections.singleton(PubQueueProvider.AGGREGATED_QUEUE_PUBLIC)));
+    }
+
+    @Test
+    public void testAggregatedQueueNamesIncludesBothWhenMixed() {
+        handler.handle(info(1L), packageMessage("p1", PUB1_AGENT_NAME));
+        ClearCallback cb = mock(ClearCallback.class);
+        when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(new 
HashSet<>(Arrays.asList("sub1", "sub2")));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub1")).thenReturn(new 
QueueState(0, -1, 0, cb));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub2")).thenReturn(new 
QueueState(5, -1, 0, null));
+        Set<String> names = 
queueProvider.getAggregatedQueueNames(PUB1_AGENT_NAME);
+        assertThat(names, 
containsInAnyOrder(PubQueueProvider.AGGREGATED_QUEUE_PERSISTED, 
PubQueueProvider.AGGREGATED_QUEUE_PUBLIC));
+    }
+
+    @Test
+    public void testAggregatedPublicQueueNotClearable() {
+        handler.handle(info(1L), packageMessage("p1", PUB1_AGENT_NAME));
+        ClearCallback cb = mock(ClearCallback.class);
+        when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(new 
HashSet<>(Arrays.asList("sub1", "sub2")));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub1")).thenReturn(new 
QueueState(0, -1, 0, cb));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub2")).thenReturn(new 
QueueState(5, -1, 0, null));
+        DistributionQueue q = 
queueProvider.getAggregatedQueue(PUB1_AGENT_NAME, 
PubQueueProvider.AGGREGATED_QUEUE_PUBLIC);
+        assertThat(q, notNullValue());
+        assertThat(q.hasCapability(CLEARABLE), equalTo(false));
+    }
+
+    @Test
+    public void testAggregatedPersistedQueueIsClearable() {
+        handler.handle(info(1L), packageMessage("p1", PUB1_AGENT_NAME));
+        ClearCallback cb = mock(ClearCallback.class);
+        when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(new 
HashSet<>(Arrays.asList("sub1", "sub2")));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub1")).thenReturn(new 
QueueState(0, -1, 0, cb));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "sub2")).thenReturn(new 
QueueState(5, -1, 0, null));
+        DistributionQueue q = 
queueProvider.getAggregatedQueue(PUB1_AGENT_NAME, 
PubQueueProvider.AGGREGATED_QUEUE_PERSISTED);
+        assertThat(q, notNullValue());
+        assertThat(q.hasCapability(CLEARABLE), equalTo(true));
+    }
+
+    @Test
+    public void testAggregatedPersistedClearFansOutToAllClearable() {
+        handler.handle(info(1L), packageMessage("p1", PUB1_AGENT_NAME));
+        ClearCallback cb1 = mock(ClearCallback.class);
+        ClearCallback cb2 = mock(ClearCallback.class);
+        when(callback.getSubscribedAgentIds(PUB1_AGENT_NAME)).thenReturn(new 
HashSet<>(Arrays.asList("a", "b")));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "a")).thenReturn(new 
QueueState(0, -1, 0, cb1));
+        when(callback.getQueueState(PUB1_AGENT_NAME, "b")).thenReturn(new 
QueueState(0, -1, 0, cb2));
+        DistributionQueue q = 
queueProvider.getAggregatedQueue(PUB1_AGENT_NAME, 
PubQueueProvider.AGGREGATED_QUEUE_PERSISTED);
+        DistributionQueueEntry head = q.getHead();
+        assertThat(head, notNullValue());
+        q.remove(head.getId());
+        verify(cb1).clear(anyLong());
+        verify(cb2).clear(anyLong());
+    }
+
+    @Test
+    public void testAggregatedQueueUnknownNameReturnsNull() {
+        assertThat(queueProvider.getAggregatedQueue(PUB1_AGENT_NAME, 
"unknown"), equalTo(null));
+    }
+
     private int queueSize() {
         return queueProvider.getOffsetQueue(PUB1_AGENT_NAME, 0).getSize();
     }

Reply via email to