Merge DOMExtensibleService into DOMService
[mdsal.git] / dom / mdsal-dom-broker / src / main / java / org / opendaylight / mdsal / dom / broker / DOMNotificationRouter.java
index 5e739888307dd06d025c7b09fba27daa1fd47821..32f218c62462444a33f417121fe3c170478c06de 100644 (file)
@@ -10,7 +10,6 @@ package org.opendaylight.mdsal.dom.broker;
 import com.google.common.annotations.VisibleForTesting;
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMultimap;
-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;
@@ -24,7 +23,6 @@ import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledFuture;
 import java.util.concurrent.ScheduledThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
 import javax.annotation.PreDestroy;
@@ -42,8 +40,8 @@ import org.opendaylight.yangtools.concepts.ListenerRegistration;
 import org.opendaylight.yangtools.concepts.Registration;
 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.common.Empty;
 import org.opendaylight.yangtools.yang.model.api.stmt.SchemaNodeIdentifier.Absolute;
 import org.osgi.service.component.annotations.Activate;
 import org.osgi.service.component.annotations.Component;
@@ -62,14 +60,12 @@ import org.slf4j.LoggerFactory;
  * Internal implementation one by using a {@link QueuedNotificationManager}.
  *</p>
  */
-@Component(immediate = true, configurationPid = "org.opendaylight.mdsal.dom.notification", service = {
-    DOMNotificationService.class, DOMNotificationPublishService.class,
-    DOMNotificationSubscriptionListenerRegistry.class
+@Component(configurationPid = "org.opendaylight.mdsal.dom.notification", service = {
+    DOMNotificationRouter.class, DOMNotificationSubscriptionListenerRegistry.class
 })
 @Designate(ocd = DOMNotificationRouter.Config.class)
 // Non-final for testing
-public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPublishService,
-        DOMNotificationService, DOMNotificationSubscriptionListenerRegistry {
+public class DOMNotificationRouter implements DOMNotificationSubscriptionListenerRegistry, AutoCloseable {
     @ObjectClassDefinition()
     public @interface Config {
         @AttributeDefinition(name = "notification-queue-depth")
@@ -105,17 +101,107 @@ public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPubl
         }
     }
 
+    private final class PublishFacade implements DOMNotificationPublishService {
+        @Override
+        public ListenableFuture<? extends Object> putNotification(final DOMNotification notification)
+                throws InterruptedException {
+            return putNotificationImpl(notification);
+        }
+
+        @Override
+        public ListenableFuture<? extends Object> offerNotification(final DOMNotification notification) {
+            final var subscribers = listeners.get(notification.getType());
+            return subscribers.isEmpty() ? NO_LISTENERS : publish(notification, subscribers);
+        }
+
+        @Override
+        public ListenableFuture<? extends Object> offerNotification(final DOMNotification notification,
+                final long timeout, final TimeUnit unit) throws InterruptedException {
+            final var subscribers = listeners.get(notification.getType());
+            if (subscribers.isEmpty()) {
+                return NO_LISTENERS;
+            }
+            // Attempt to perform a non-blocking publish first
+            final var noBlock = publish(notification, subscribers);
+            if (!DOMNotificationPublishService.REJECTED.equals(noBlock)) {
+                return noBlock;
+            }
+
+            try {
+                final var publishThread = Thread.currentThread();
+                final var timerTask = observer.schedule(publishThread::interrupt, timeout, unit);
+                final var withBlock = putNotificationImpl(notification);
+                timerTask.cancel(true);
+                if (observer.getQueue().size() > 50) {
+                    observer.purge();
+                }
+                return withBlock;
+            } catch (InterruptedException e) {
+                return DOMNotificationPublishService.REJECTED;
+            }
+        }
+    }
+
+    private final class SubscribeFacade implements DOMNotificationService {
+        @Override
+        public <T extends DOMNotificationListener> ListenerRegistration<T> registerNotificationListener(
+                final T listener, final Collection<Absolute> types) {
+            synchronized (DOMNotificationRouter.this) {
+                final var reg = new SingleReg<>(listener);
+
+                if (!types.isEmpty()) {
+                    final var b = ImmutableMultimap.<Absolute, Reg<?>>builder();
+                    b.putAll(listeners);
+
+                    for (var t : types) {
+                        b.put(t, reg);
+                    }
+
+                    replaceListeners(b.build());
+                }
+
+                return reg;
+            }
+        }
+
+        @Override
+        public synchronized Registration registerNotificationListeners(
+                final Map<Absolute, DOMNotificationListener> typeToListener) {
+            synchronized (DOMNotificationRouter.this) {
+                final var b = ImmutableMultimap.<Absolute, Reg<?>>builder();
+                b.putAll(listeners);
+
+                final var tmp = new HashMap<DOMNotificationListener, ComponentReg>();
+                for (var e : typeToListener.entrySet()) {
+                    b.put(e.getKey(), tmp.computeIfAbsent(e.getValue(), ComponentReg::new));
+                }
+                replaceListeners(b.build());
+
+                final var regs = List.copyOf(tmp.values());
+                return new AbstractRegistration() {
+                    @Override
+                    protected void removeRegistration() {
+                        regs.forEach(ComponentReg::close);
+                        removeRegistrations(regs);
+                    }
+                };
+            }
+        }
+    }
+
     private static final Logger LOG = LoggerFactory.getLogger(DOMNotificationRouter.class);
-    private static final ListenableFuture<Void> NO_LISTENERS = FluentFutures.immediateNullFluentFuture();
+    private static final @NonNull ListenableFuture<?> NO_LISTENERS = Futures.immediateFuture(Empty.value());
 
     private final ListenerRegistry<DOMNotificationSubscriptionListener> subscriptionListeners =
             ListenerRegistry.create();
     private final EqualityQueuedNotificationManager<AbstractListenerRegistration<? extends DOMNotificationListener>,
                 DOMNotificationRouterEvent> queueNotificationManager;
+    private final @NonNull DOMNotificationPublishService notificationPublishService = new PublishFacade();
+    private final @NonNull DOMNotificationService notificationService = new SubscribeFacade();
     private final ScheduledThreadPoolExecutor observer;
     private final ExecutorService executor;
 
-    private volatile Multimap<Absolute, Reg<?>> listeners = ImmutableMultimap.of();
+    private volatile ImmutableMultimap<Absolute, Reg<?>> listeners = ImmutableMultimap.of();
 
     @Inject
     public DOMNotificationRouter(final int maxQueueCapacity) {
@@ -129,57 +215,20 @@ public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPubl
             .build());
         queueNotificationManager = new EqualityQueuedNotificationManager<>("DOMNotificationRouter", executor,
                 maxQueueCapacity, DOMNotificationRouter::deliverEvents);
+        LOG.info("DOM Notification Router started");
     }
 
     @Activate
     public DOMNotificationRouter(final Config config) {
         this(config.queueDepth());
-        LOG.info("DOM Notification Router started");
     }
 
-    @Deprecated(forRemoval = true)
-    public static DOMNotificationRouter create(final int maxQueueCapacity) {
-        return new DOMNotificationRouter(maxQueueCapacity);
-    }
-
-    @Override
-    public synchronized <T extends DOMNotificationListener> ListenerRegistration<T> registerNotificationListener(
-            final T listener, final Collection<Absolute> types) {
-        final var reg = new SingleReg<>(listener);
-
-        if (!types.isEmpty()) {
-            final var b = ImmutableMultimap.<Absolute, Reg<?>>builder();
-            b.putAll(listeners);
-
-            for (var t : types) {
-                b.put(t, reg);
-            }
-
-            replaceListeners(b.build());
-        }
-
-        return reg;
+    public @NonNull DOMNotificationService notificationService() {
+        return notificationService;
     }
 
-    @Override
-    public synchronized Registration registerNotificationListeners(
-            final Map<Absolute, DOMNotificationListener> typeToListener) {
-        final var b = ImmutableMultimap.<Absolute, Reg<?>>builder();
-        b.putAll(listeners);
-
-        final var tmp = new HashMap<DOMNotificationListener, ComponentReg>();
-        for (var e : typeToListener.entrySet()) {
-            b.put(e.getKey(), tmp.computeIfAbsent(e.getValue(), ComponentReg::new));
-        }
-
-        final var regs = List.copyOf(tmp.values());
-        return new AbstractRegistration() {
-            @Override
-            protected void removeRegistration() {
-                regs.forEach(ComponentReg::close);
-                removeRegistrations(regs);
-            }
-        };
+    public @NonNull DOMNotificationPublishService notificationPublishService() {
+        return notificationPublishService;
     }
 
     private synchronized void removeRegistration(final SingleReg<?> reg) {
@@ -195,17 +244,16 @@ public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPubl
      *
      * @param newListeners is used to notify listenerTypes changed
      */
-    private void replaceListeners(final Multimap<Absolute, Reg<?>> newListeners) {
+    private void replaceListeners(final ImmutableMultimap<Absolute, Reg<?>> newListeners) {
         listeners = newListeners;
         notifyListenerTypesChanged(newListeners.keySet());
     }
 
     @SuppressWarnings("checkstyle:IllegalCatch")
     private void notifyListenerTypesChanged(final Set<Absolute> typesAfter) {
-        final List<? extends DOMNotificationSubscriptionListener> listenersAfter =
-                subscriptionListeners.streamListeners().collect(ImmutableList.toImmutableList());
+        final var listenersAfter = subscriptionListeners.streamListeners().collect(ImmutableList.toImmutableList());
         executor.execute(() -> {
-            for (final DOMNotificationSubscriptionListener subListener : listenersAfter) {
+            for (var subListener : listenersAfter) {
                 try {
                     subListener.onSubscriptionChanged(typesAfter);
                 } catch (final Exception e) {
@@ -218,72 +266,30 @@ public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPubl
     @Override
     public <L extends DOMNotificationSubscriptionListener> ListenerRegistration<L> registerSubscriptionListener(
             final L listener) {
-        final Set<Absolute> initialTypes = listeners.keySet();
+        final var initialTypes = listeners.keySet();
         executor.execute(() -> listener.onSubscriptionChanged(initialTypes));
         return subscriptionListeners.register(listener);
     }
 
     @VisibleForTesting
-    ListenableFuture<? extends Object> publish(final DOMNotification notification,
-            final Collection<Reg<?>> subscribers) {
-        final List<ListenableFuture<Void>> futures = new ArrayList<>(subscribers.size());
+    @NonNull ListenableFuture<? extends Object> putNotificationImpl(final DOMNotification notification)
+            throws InterruptedException {
+        final var subscribers = listeners.get(notification.getType());
+        return subscribers.isEmpty() ? NO_LISTENERS : publish(notification, subscribers);
+    }
+
+    @VisibleForTesting
+    @NonNull ListenableFuture<?> publish(final DOMNotification notification, final Collection<Reg<?>> subscribers) {
+        final var futures = new ArrayList<ListenableFuture<?>>(subscribers.size());
         subscribers.forEach(subscriber -> {
-            final DOMNotificationRouterEvent event = new DOMNotificationRouterEvent(notification);
+            final var event = new DOMNotificationRouterEvent(notification);
             futures.add(event.future());
             queueNotificationManager.submitNotification(subscriber, event);
         });
-        return Futures.transform(Futures.successfulAsList(futures), ignored -> (Void)null,
+        return Futures.transform(Futures.successfulAsList(futures), ignored -> Empty.value(),
             MoreExecutors.directExecutor());
     }
 
-    @Override
-    public ListenableFuture<? extends Object> putNotification(final DOMNotification notification)
-            throws InterruptedException {
-        final var subscribers = listeners.get(notification.getType());
-        if (subscribers.isEmpty()) {
-            return NO_LISTENERS;
-        }
-
-        return publish(notification, subscribers);
-    }
-
-    @Override
-    public ListenableFuture<? extends Object> offerNotification(final DOMNotification notification) {
-        final var subscribers = listeners.get(notification.getType());
-        if (subscribers.isEmpty()) {
-            return NO_LISTENERS;
-        }
-
-        return publish(notification, subscribers);
-    }
-
-    @Override
-    public ListenableFuture<? extends Object> offerNotification(final DOMNotification notification, final long timeout,
-            final TimeUnit unit) throws InterruptedException {
-        final var subscribers = listeners.get(notification.getType());
-        if (subscribers.isEmpty()) {
-            return NO_LISTENERS;
-        }
-        // Attempt to perform a non-blocking publish first
-        final ListenableFuture<?> noBlock = publish(notification, subscribers);
-        if (!DOMNotificationPublishService.REJECTED.equals(noBlock)) {
-            return noBlock;
-        }
-
-        try {
-            final Thread publishThread = Thread.currentThread();
-            ScheduledFuture<?> timerTask = observer.schedule(publishThread::interrupt, timeout, unit);
-            final ListenableFuture<?> withBlock = putNotification(notification);
-            timerTask.cancel(true);
-            if (observer.getQueue().size() > 50) {
-                observer.purge();
-            }
-            return withBlock;
-        } catch (InterruptedException e) {
-            return DOMNotificationPublishService.REJECTED;
-        }
-    }
-
     @PreDestroy
     @Deactivate
     @Override
@@ -304,7 +310,7 @@ public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPubl
     }
 
     @VisibleForTesting
-    Multimap<Absolute, ?> listeners() {
+    ImmutableMultimap<Absolute, ?> listeners() {
         return listeners;
     }
 
@@ -316,8 +322,8 @@ public class DOMNotificationRouter implements AutoCloseable, DOMNotificationPubl
     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) {
+            final var listener = reg.getInstance();
+            for (var event : events) {
                 event.deliverTo(listener);
             }
         } else {