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