*/
package org.opendaylight.mdsal.dom.broker;
-import static com.google.common.base.Preconditions.checkArgument;
-
import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMultimap;
import com.google.common.collect.ImmutableMultimap.Builder;
import com.google.common.collect.Multimap;
import com.google.common.collect.Multimaps;
+import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
+import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
-import com.lmax.disruptor.EventHandler;
-import com.lmax.disruptor.InsufficientCapacityException;
-import com.lmax.disruptor.PhasedBackoffWaitStrategy;
-import com.lmax.disruptor.WaitStrategy;
-import com.lmax.disruptor.dsl.Disruptor;
-import com.lmax.disruptor.dsl.ProducerType;
+import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Set;
import org.opendaylight.yangtools.concepts.AbstractListenerRegistration;
import org.opendaylight.yangtools.concepts.ListenerRegistration;
import org.opendaylight.yangtools.util.ListenerRegistry;
+import org.opendaylight.yangtools.util.concurrent.EqualityQueuedNotificationManager;
import org.opendaylight.yangtools.util.concurrent.FluentFutures;
+import org.opendaylight.yangtools.util.concurrent.QueuedNotificationManager;
import org.opendaylight.yangtools.yang.model.api.stmt.SchemaNodeIdentifier.Absolute;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
* routing of notifications from publishers to subscribers.
*
*<p>
- * Internal implementation works by allocating a two-handler Disruptor. The first handler delivers notifications
- * to subscribed listeners and the second one notifies whoever may be listening on the returned future. Registration
- * state tracking is performed by a simple immutable multimap -- when a registration or unregistration occurs we
- * re-generate the entire map from scratch and set it atomically. While registrations/unregistrations synchronize
- * on this instance, notifications do not take any locks here.
- *
- *<p>
- * The fully-blocking {@link #publish(long, DOMNotification, Collection)}
- * and non-blocking {@link #offerNotification(DOMNotification)}
- * are realized using the Disruptor's native operations. The bounded-blocking {@link
- * #offerNotification(DOMNotification, long, TimeUnit)}
- * is realized by arming a background wakeup interrupt.
+ * Internal implementation one by using a {@link QueuedNotificationManager}.
+ *</p>
*/
public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPublishService,
DOMNotificationService, DOMNotificationSubscriptionListenerRegistry {
private static final Logger LOG = LoggerFactory.getLogger(DOMNotificationRouter.class);
private static final ListenableFuture<Void> NO_LISTENERS = FluentFutures.immediateNullFluentFuture();
- private static final WaitStrategy DEFAULT_STRATEGY = PhasedBackoffWaitStrategy.withLock(
- 1L, 30L, TimeUnit.MILLISECONDS);
- private static final EventHandler<DOMNotificationRouterEvent> DISPATCH_NOTIFICATIONS =
- (event, sequence, endOfBatch) -> event.deliverNotification();
- private static final EventHandler<DOMNotificationRouterEvent> NOTIFY_FUTURE =
- (event, sequence, endOfBatch) -> event.setFuture();
private final ListenerRegistry<DOMNotificationSubscriptionListener> subscriptionListeners =
ListenerRegistry.create();
- private final Disruptor<DOMNotificationRouterEvent> disruptor;
+ private final EqualityQueuedNotificationManager<AbstractListenerRegistration<? extends DOMNotificationListener>,
+ DOMNotificationRouterEvent> queueNotificationManager;
private final ScheduledThreadPoolExecutor observer;
private final ExecutorService executor;
ImmutableMultimap.of();
@VisibleForTesting
- DOMNotificationRouter(final int queueDepth, final WaitStrategy strategy) {
+ DOMNotificationRouter(int maxQueueCapacity) {
observer = new ScheduledThreadPoolExecutor(1,
new ThreadFactoryBuilder().setDaemon(true).setNameFormat("DOMNotificationRouter-observer-%d").build());
executor = Executors.newCachedThreadPool(
new ThreadFactoryBuilder().setDaemon(true).setNameFormat("DOMNotificationRouter-listeners-%d").build());
- disruptor = new Disruptor<>(DOMNotificationRouterEvent.FACTORY, queueDepth,
- new ThreadFactoryBuilder().setNameFormat("DOMNotificationRouter-disruptor-%d").build(),
- ProducerType.MULTI, strategy);
- disruptor.handleEventsWith(DISPATCH_NOTIFICATIONS);
- disruptor.after(DISPATCH_NOTIFICATIONS).handleEventsWith(NOTIFY_FUTURE);
- disruptor.start();
+ queueNotificationManager = new EqualityQueuedNotificationManager<>("DOMNotificationRouter", executor,
+ maxQueueCapacity, DOMNotificationRouter::deliverEvents);
}
- public static DOMNotificationRouter create(final int queueDepth) {
- return new DOMNotificationRouter(queueDepth, DEFAULT_STRATEGY);
- }
-
- public static DOMNotificationRouter create(final int queueDepth, final long spinTime, final long parkTime,
- final TimeUnit unit) {
- checkArgument(Long.lowestOneBit(queueDepth) == Long.highestOneBit(queueDepth),
- "Queue depth %s is not power-of-two", queueDepth);
- return new DOMNotificationRouter(queueDepth, PhasedBackoffWaitStrategy.withLock(spinTime, parkTime, unit));
+ public static DOMNotificationRouter create(int maxQueueCapacity) {
+ return new DOMNotificationRouter(maxQueueCapacity);
}
@Override
return subscriptionListeners.register(listener);
}
- private ListenableFuture<Void> publish(final long seq, final DOMNotification notification,
+
+ @VisibleForTesting
+ ListenableFuture<? extends Object> publish(DOMNotification notification,
final Collection<AbstractListenerRegistration<? extends DOMNotificationListener>> subscribers) {
- final DOMNotificationRouterEvent event = disruptor.get(seq);
- final ListenableFuture<Void> future = event.initialize(notification, subscribers);
- disruptor.getRingBuffer().publish(seq);
- return future;
+ final List<ListenableFuture<Void>> futures = new ArrayList<>(subscribers.size());
+ subscribers.forEach(subscriber -> {
+ final DOMNotificationRouterEvent event = new DOMNotificationRouterEvent(notification);
+ futures.add(event.future());
+ queueNotificationManager.submitNotification(subscriber, event);
+ });
+ return Futures.transform(Futures.successfulAsList(futures), ignored -> (Void)null,
+ MoreExecutors.directExecutor());
}
@Override
return NO_LISTENERS;
}
- final long seq = disruptor.getRingBuffer().next();
- return publish(seq, notification, subscribers);
- }
-
- @SuppressWarnings("checkstyle:IllegalCatch")
- @VisibleForTesting
- ListenableFuture<? extends Object> tryPublish(final DOMNotification notification,
- final Collection<AbstractListenerRegistration<? extends DOMNotificationListener>> subscribers) {
- final long seq;
- try {
- seq = disruptor.getRingBuffer().tryNext();
- } catch (final InsufficientCapacityException e) {
- return DOMNotificationPublishService.REJECTED;
- }
-
- return publish(seq, notification, subscribers);
+ return publish(notification, subscribers);
}
@Override
return NO_LISTENERS;
}
- return tryPublish(notification, subscribers);
+ return publish(notification, subscribers);
}
@Override
return NO_LISTENERS;
}
// Attempt to perform a non-blocking publish first
- final ListenableFuture<?> noBlock = tryPublish(notification, subscribers);
+ final ListenableFuture<?> noBlock = publish(notification, subscribers);
if (!DOMNotificationPublishService.REJECTED.equals(noBlock)) {
return noBlock;
}
@Override
public void close() {
- disruptor.shutdown();
observer.shutdown();
executor.shutdown();
}
ListenerRegistry<DOMNotificationSubscriptionListener> subscriptionListeners() {
return subscriptionListeners;
}
+
+ private static void deliverEvents(final AbstractListenerRegistration<? extends DOMNotificationListener> reg,
+ final ImmutableList<DOMNotificationRouterEvent> events) {
+ if (reg.notClosed()) {
+ final DOMNotificationListener listener = reg.getInstance();
+ for (DOMNotificationRouterEvent event : events) {
+ event.deliverTo(listener);
+ }
+ } else {
+ events.forEach(DOMNotificationRouterEvent::clear);
+ }
+ }
}
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
-import com.lmax.disruptor.EventFactory;
-import java.util.Collection;
+import org.eclipse.jdt.annotation.NonNull;
import org.opendaylight.mdsal.dom.api.DOMNotification;
import org.opendaylight.mdsal.dom.api.DOMNotificationListener;
-import org.opendaylight.yangtools.concepts.AbstractListenerRegistration;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
- * A single notification event in the disruptor ringbuffer. These objects are reused, so they do have mutable state.
+ * A single notification event in the notification router.
*/
final class DOMNotificationRouterEvent {
private static final Logger LOG = LoggerFactory.getLogger(DOMNotificationRouterEvent.class);
- static final EventFactory<DOMNotificationRouterEvent> FACTORY = DOMNotificationRouterEvent::new;
+ private final SettableFuture<Void> future = SettableFuture.create();
+ private final @NonNull DOMNotification notification;
- private Collection<AbstractListenerRegistration<? extends DOMNotificationListener>> subscribers;
- private DOMNotification notification;
- private SettableFuture<Void> future;
-
- private DOMNotificationRouterEvent() {
- // Hidden on purpose, initialized in initialize()
+ DOMNotificationRouterEvent(final DOMNotification notification) {
+ this.notification = requireNonNull(notification);
}
- @SuppressWarnings("checkstyle:hiddenField")
- ListenableFuture<Void> initialize(final DOMNotification notification,
- final Collection<AbstractListenerRegistration<? extends DOMNotificationListener>> subscribers) {
- this.notification = requireNonNull(notification);
- this.subscribers = requireNonNull(subscribers);
- this.future = SettableFuture.create();
- return this.future;
+ ListenableFuture<Void> future() {
+ return future;
}
@SuppressWarnings("checkstyle:illegalCatch")
- void deliverNotification() {
- for (AbstractListenerRegistration<? extends DOMNotificationListener> reg : subscribers) {
- if (reg.notClosed()) {
- final DOMNotificationListener listener = reg.getInstance();
- try {
- listener.onNotification(notification);
- } catch (Exception e) {
- LOG.warn("Listener {} failed during notification delivery", listener, e);
- }
- }
+ void deliverTo(DOMNotificationListener listener) {
+ try {
+ listener.onNotification(notification);
+ } catch (Exception e) {
+ LOG.warn("Listener {} failed during notification delivery", listener, e);
+ } finally {
+ clear();
}
}
- void setFuture() {
+ void clear() {
future.set(null);
}
}
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;
import com.google.common.collect.Multimap;
import com.google.common.util.concurrent.ListenableFuture;
-import com.lmax.disruptor.PhasedBackoffWaitStrategy;
-import com.lmax.disruptor.WaitStrategy;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
public class DOMNotificationRouterTest extends TestUtils {
- private static final WaitStrategy DEFAULT_STRATEGY = PhasedBackoffWaitStrategy.withLock(
- 1L, 30L, TimeUnit.MILLISECONDS);
-
@Test
public void create() throws Exception {
- assertNotNull(DOMNotificationRouter.create(1,1,1,TimeUnit.SECONDS));
- assertNotNull(DOMNotificationRouter.create(1));
+ assertNotNull(DOMNotificationRouter.create(1024));
}
@SuppressWarnings("checkstyle:IllegalCatch")
public void complexTest() throws Exception {
final DOMNotificationSubscriptionListener domNotificationSubscriptionListener =
mock(DOMNotificationSubscriptionListener.class);
+ doNothing().when(domNotificationSubscriptionListener).onSubscriptionChanged(any());
+
final CountDownLatch latch = new CountDownLatch(1);
final DOMNotificationListener domNotificationListener = new TestListener(latch);
- final DOMNotificationRouter domNotificationRouter = DOMNotificationRouter.create(1);
+ final DOMNotificationRouter domNotificationRouter = DOMNotificationRouter.create(1024);
Multimap<Absolute, ?> listeners = domNotificationRouter.listeners();
@Test
public void offerNotification() throws Exception {
- final DOMNotificationRouter domNotificationRouter = DOMNotificationRouter.create(1);
+ final DOMNotificationRouter domNotificationRouter = DOMNotificationRouter.create(1024);
final DOMNotification domNotification = mock(DOMNotification.class);
doReturn(Absolute.of(TestModel.TEST_QNAME)).when(domNotification).getType();
doReturn(TEST_CHILD).when(domNotification).getBody();
assertEquals("Received notifications", 1, testListener.getReceivedNotifications().size());
assertEquals(DOMNotificationPublishService.REJECTED,
- testRouter.offerNotification(domNotification, 1, TimeUnit.SECONDS));
+ testRouter.offerNotification(domNotification, 1, TimeUnit.SECONDS));
assertEquals("Received notifications", 1, testListener.getReceivedNotifications().size());
}
}
@Test
public void close() throws Exception {
- final DOMNotificationRouter domNotificationRouter = DOMNotificationRouter.create(1);
+ final DOMNotificationRouter domNotificationRouter = DOMNotificationRouter.create(1024);
final ExecutorService executor = domNotificationRouter.executor();
final ExecutorService observer = domNotificationRouter.observer();
}
private static class TestRouter extends DOMNotificationRouter {
+
+ private boolean triggerRejected = false;
+
TestRouter(final int queueDepth) {
- super(queueDepth, DEFAULT_STRATEGY);
+ super(queueDepth);
}
@Override
- protected ListenableFuture<? extends Object> tryPublish(final DOMNotification notification,
- final Collection<AbstractListenerRegistration<? extends DOMNotificationListener>> subscribers) {
- return DOMNotificationPublishService.REJECTED;
+ ListenableFuture<? extends Object> publish(DOMNotification notification,
+ Collection<AbstractListenerRegistration<? extends DOMNotificationListener>> subscribers) {
+ if (triggerRejected) {
+ return REJECTED;
+ }
+
+ triggerRejected = true;
+ return super.publish(notification, subscribers);
}
@Override