Replace MessagePublisher with new class TopicEventPublisher

* TopicEventPublisher use a internal class called TopicEvent to be
able to keep track of the topic related to the event and a flag that
keep track if the event has been published.
* Publisher::getEventPublisher is removed because it only caused
extra indirection and getName and getProperties are removed because
they are not needed.

Change-Id: I060f289481e5852997b601756a4a01819b3a61a6
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Manager.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Manager.java
index 8290d05..8535d5f 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Manager.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Manager.java
@@ -23,8 +23,8 @@
 import com.googlesource.gerrit.plugins.rabbitmq.config.Properties;
 import com.googlesource.gerrit.plugins.rabbitmq.config.PropertiesFactory;
 import com.googlesource.gerrit.plugins.rabbitmq.config.section.Gerrit;
+import com.googlesource.gerrit.plugins.rabbitmq.message.GerritEventPublisherFactory;
 import com.googlesource.gerrit.plugins.rabbitmq.message.Publisher;
-import com.googlesource.gerrit.plugins.rabbitmq.message.PublisherFactory;
 import com.googlesource.gerrit.plugins.rabbitmq.worker.DefaultEventWorker;
 import com.googlesource.gerrit.plugins.rabbitmq.worker.EventWorker;
 import com.googlesource.gerrit.plugins.rabbitmq.worker.EventWorkerFactory;
@@ -48,7 +48,7 @@
   private final Path pluginDataDir;
   private final EventWorker defaultEventWorker;
   private final EventWorker userEventWorker;
-  private final PublisherFactory publisherFactory;
+  private final GerritEventPublisherFactory publisherFactory;
   private final PropertiesFactory propFactory;
   private final List<Publisher> publisherList = new ArrayList<>();
 
@@ -58,7 +58,7 @@
       @PluginData final File pluginData,
       final DefaultEventWorker defaultEventWorker,
       final EventWorkerFactory eventWorkerFactory,
-      final PublisherFactory publisherFactory,
+      final GerritEventPublisherFactory publisherFactory,
       final PropertiesFactory propFactory) {
     this.pluginName = pluginName;
     this.pluginDataDir = pluginData.toPath();
@@ -88,13 +88,9 @@
   public void stop() {
     for (Publisher publisher : publisherList) {
       publisher.stop();
-      String listenAs = publisher.getProperties().getSection(Gerrit.class).listenAs;
-      if (!listenAs.isEmpty()) {
-        userEventWorker.removePublisher(publisher);
-      } else {
-        defaultEventWorker.removePublisher(publisher);
-      }
     }
+    defaultEventWorker.clear();
+    userEventWorker.clear();
     publisherList.clear();
   }
 
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Module.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Module.java
index 126a036..1f07de4 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Module.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/Module.java
@@ -31,10 +31,10 @@
 import com.googlesource.gerrit.plugins.rabbitmq.config.section.Message;
 import com.googlesource.gerrit.plugins.rabbitmq.config.section.Monitor;
 import com.googlesource.gerrit.plugins.rabbitmq.config.section.Section;
+import com.googlesource.gerrit.plugins.rabbitmq.message.GerritEventPublisher;
+import com.googlesource.gerrit.plugins.rabbitmq.message.GerritEventPublisherFactory;
 import com.googlesource.gerrit.plugins.rabbitmq.message.GsonProvider;
-import com.googlesource.gerrit.plugins.rabbitmq.message.MessagePublisher;
 import com.googlesource.gerrit.plugins.rabbitmq.message.Publisher;
-import com.googlesource.gerrit.plugins.rabbitmq.message.PublisherFactory;
 import com.googlesource.gerrit.plugins.rabbitmq.session.SessionFactory;
 import com.googlesource.gerrit.plugins.rabbitmq.session.SessionFactoryProvider;
 import com.googlesource.gerrit.plugins.rabbitmq.worker.DefaultEventWorker;
@@ -57,8 +57,8 @@
 
     install(
         new FactoryModuleBuilder()
-            .implement(Publisher.class, MessagePublisher.class)
-            .build(PublisherFactory.class));
+            .implement(Publisher.class, GerritEventPublisher.class)
+            .build(GerritEventPublisherFactory.class));
     install(
         new FactoryModuleBuilder()
             .implement(Properties.class, PluginProperties.class)
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/GerritEventPublisher.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/GerritEventPublisher.java
new file mode 100644
index 0000000..ff18eac
--- /dev/null
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/GerritEventPublisher.java
@@ -0,0 +1,45 @@
+// Copyright (C) 2015 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package com.googlesource.gerrit.plugins.rabbitmq.message;
+
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.gerrit.server.events.Event;
+import com.google.gerrit.server.events.EventGson;
+import com.google.gson.Gson;
+import com.google.inject.Inject;
+import com.google.inject.assistedinject.Assisted;
+import com.googlesource.gerrit.plugins.rabbitmq.config.Properties;
+import com.googlesource.gerrit.plugins.rabbitmq.config.section.Message;
+import com.googlesource.gerrit.plugins.rabbitmq.session.SessionFactoryProvider;
+import java.util.Optional;
+
+public class GerritEventPublisher extends MessagePublisher {
+  private final Optional<String> defaultTopic;
+
+  @Inject
+  public GerritEventPublisher(
+      @Assisted final Properties properties,
+      SessionFactoryProvider sessionFactoryProvider,
+      @EventGson Gson gson) {
+    super(properties, sessionFactoryProvider, gson);
+    String routingKey = properties.getSection(Message.class).routingKey;
+    this.defaultTopic =
+        routingKey != null && !routingKey.isEmpty() ? Optional.of(routingKey) : Optional.empty();
+  }
+
+  public ListenableFuture<Boolean> publish(String topic, Event event) {
+    return super.publish(defaultTopic.orElse(topic), event);
+  }
+}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/PublisherFactory.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/GerritEventPublisherFactory.java
similarity index 94%
rename from src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/PublisherFactory.java
rename to src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/GerritEventPublisherFactory.java
index 225f201..639101e 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/PublisherFactory.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/GerritEventPublisherFactory.java
@@ -16,6 +16,6 @@
 
 import com.googlesource.gerrit.plugins.rabbitmq.config.Properties;
 
-public interface PublisherFactory {
+public interface GerritEventPublisherFactory {
   Publisher create(Properties properties);
 }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/MessagePublisher.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/MessagePublisher.java
index 8096556..43beb1a 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/MessagePublisher.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/MessagePublisher.java
@@ -15,9 +15,11 @@
 package com.googlesource.gerrit.plugins.rabbitmq.message;
 
 import com.google.common.flogger.FluentLogger;
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.common.util.concurrent.SettableFuture;
 import com.google.gerrit.extensions.events.LifecycleListener;
 import com.google.gerrit.server.events.Event;
-import com.google.gerrit.server.events.EventListener;
+import com.google.gerrit.server.events.EventGson;
 import com.google.gson.Gson;
 import com.google.inject.Inject;
 import com.google.inject.assistedinject.Assisted;
@@ -40,58 +42,37 @@
   private static final String END_OF_STREAM = "END-OF-STREAM_$F7;XTSUQ(Dv#N6]g+gd,,uzRp%G-P";
   private static final Event EOS = new Event(END_OF_STREAM) {};
 
-  private final Session session;
   private final Properties properties;
+  private final Session session;
   private final Gson gson;
   private final Timer monitorTimer = new Timer();
-  private final LinkedBlockingQueue<Event> queue = new LinkedBlockingQueue<>(MAX_EVENTS);
+  private final LinkedBlockingQueue<TopicEvent> queue = new LinkedBlockingQueue<>(MAX_EVENTS);
   private final Object sessionMon = new Object();
-  private EventListener eventListener;
   private GracefullyCancelableRunnable publisher;
   private Thread publisherThread;
+  private final Object lostEventCountLock = new Object();
+  private int lostEventCount = 0;
 
   @Inject
   public MessagePublisher(
       @Assisted final Properties properties,
       SessionFactoryProvider sessionFactoryProvider,
-      Gson gson) {
+      @EventGson Gson gson) {
     this.session = sessionFactoryProvider.get().create(properties);
     this.properties = properties;
     this.gson = gson;
-    this.eventListener =
-        new EventListener() {
-          private int lostEventCount = 0;
-
-          @Override
-          public void onEvent(Event event) {
-            if (!publisherThread.isAlive()) {
-              ensurePublisherThreadStarted();
-            }
-
-            if (queue.offer(event)) {
-              if (lostEventCount > 0) {
-                logger.atWarning().log(
-                    "Event queue is no longer full, %d events were lost", lostEventCount);
-                lostEventCount = 0;
-              }
-            } else {
-              if (lostEventCount++ % 10 == 0) {
-                logger.atSevere().log("Event queue is full, lost %d event(s)", lostEventCount);
-              }
-            }
-          }
-        };
     this.publisher =
         new GracefullyCancelableRunnable() {
 
-          volatile boolean canceled = false;
+          volatile boolean canceled;
 
           @Override
           public void run() {
+            canceled = false;
             while (!canceled) {
               try {
-                Event event = queue.take();
-                if (event.getType().equals(END_OF_STREAM)) {
+                TopicEvent topicEvent = queue.take();
+                if (topicEvent.event.getType().equals(END_OF_STREAM)) {
                   continue;
                 }
                 while (!isConnected() && !canceled) {
@@ -99,8 +80,8 @@
                     sessionMon.wait(1000);
                   }
                 }
-                if (!publishEvent(event) && !queue.offer(event)) {
-                  logger.atSevere().log("Event lost: %s", gson.toJson(event));
+                if (!publishEvent(topicEvent) && !queue.offer(topicEvent)) {
+                  logger.atSevere().log("Event lost: %s", gson.toJson(topicEvent.event));
                 }
               } catch (InterruptedException e) {
                 logger.atWarning().withCause(e).log(
@@ -113,7 +94,7 @@
           public void cancel() {
             canceled = true;
             if (queue.isEmpty()) {
-              queue.offer(EOS);
+              queue.offer(new TopicEvent(null, EOS, null));
             }
           }
 
@@ -162,26 +143,41 @@
   }
 
   @Override
-  public Properties getProperties() {
-    return properties;
+  public ListenableFuture<Boolean> publish(String topic, Event event) {
+    SettableFuture<Boolean> future = SettableFuture.create();
+    publish(new TopicEvent(topic, event, future));
+    return future;
   }
 
-  @Override
-  public String getName() {
-    return properties.getName();
-  }
-
-  @Override
-  public EventListener getEventListener() {
-    return this.eventListener;
+  private void publish(TopicEvent topicEvent) {
+    if (!publisherThread.isAlive()) {
+      ensurePublisherThreadStarted();
+    }
+    logger.atFine().log(
+        "Adding event %s for topic %s to publisher queue", topicEvent.event, topicEvent.topic);
+    synchronized (lostEventCountLock) {
+      if (queue.offer(topicEvent)) {
+        if (lostEventCount > 0) {
+          logger.atWarning().log(
+              "Event queue is no longer full, %d events were lost", lostEventCount);
+          lostEventCount = 0;
+        }
+      } else {
+        if (lostEventCount++ % 10 == 0) {
+          logger.atSevere().log("Event queue is full, lost %d event(s)", lostEventCount);
+        }
+      }
+    }
   }
 
   private boolean isConnected() {
     return session != null && session.isOpen();
   }
 
-  private boolean publishEvent(Event event) {
-    return session.publish(gson.toJson(event), event.type);
+  private boolean publishEvent(TopicEvent topicEvent) {
+    boolean published = session.publish(gson.toJson(topicEvent.event), topicEvent.topic);
+    topicEvent.published.set(published);
+    return published;
   }
 
   private void connect() {
@@ -205,4 +201,16 @@
     /** Gracefully cancels the Runnable after completing ongoing task. */
     public void cancel();
   }
+
+  private class TopicEvent {
+    String topic;
+    Event event;
+    SettableFuture<Boolean> published;
+
+    TopicEvent(String topic, Event event, SettableFuture<Boolean> published) {
+      this.topic = topic;
+      this.event = event;
+      this.published = published;
+    }
+  }
 }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/Publisher.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/Publisher.java
index 8ef1b3e..95e74e9 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/Publisher.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/message/Publisher.java
@@ -1,16 +1,26 @@
+// Copyright (C) 2023 The Android Open Source Project
+//
+// Licensed under the Apache License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
 package com.googlesource.gerrit.plugins.rabbitmq.message;
 
-import com.google.gerrit.server.events.EventListener;
-import com.googlesource.gerrit.plugins.rabbitmq.config.Properties;
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.gerrit.server.events.Event;
 
 public interface Publisher {
   void start();
 
   void stop();
 
-  Properties getProperties();
-
-  String getName();
-
-  EventListener getEventListener();
+  ListenableFuture<Boolean> publish(String topic, Event message);
 }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/DefaultEventWorker.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/DefaultEventWorker.java
index 53ed103..bd6ea22 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/DefaultEventWorker.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/DefaultEventWorker.java
@@ -52,7 +52,7 @@
   @Override
   public void onEvent(Event event) {
     for (Publisher publisher : publishers) {
-      publisher.getEventListener().onEvent(event);
+      publisher.publish(event.type, event);
     }
   }
 }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/UserEventWorker.java b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/UserEventWorker.java
index 01b403c..82bd00e 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/UserEventWorker.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/rabbitmq/worker/UserEventWorker.java
@@ -86,13 +86,13 @@
                         });
                 try {
                   final IdentifiedUser user = accountResolver.resolve(userName).asUniqueUser();
-                  RegistrationHandle registration =
+                  RegistrationHandle handle =
                       eventListeners.add(
                           pluginName,
                           new UserScopedEventListener() {
                             @Override
                             public void onEvent(Event event) {
-                              publisher.getEventListener().onEvent(event);
+                              publisher.publish(event.type, event);
                             }
 
                             @Override
@@ -100,7 +100,7 @@
                               return user;
                             }
                           });
-                  eventListenerRegistrations.put(publisher, registration);
+                  eventListenerRegistrations.put(publisher, handle);
                   logger.atInfo().log("Listen events as : %s", userName);
                 } catch (UnresolvableAccountException uae) {
                   logger.atSevere().withCause(uae).log(
@@ -116,14 +116,20 @@
 
   @Override
   public void removePublisher(final Publisher publisher) {
-    RegistrationHandle registration = eventListenerRegistrations.remove(publisher);
-    if (registration != null) {
-      registration.remove();
+    RegistrationHandle handle = eventListenerRegistrations.remove(publisher);
+    if (handle != null) {
+      handle.remove();
     }
   }
 
   @Override
   public void clear() {
-    // no op.
+    for (Map.Entry<Publisher, RegistrationHandle> entry : eventListenerRegistrations.entrySet()) {
+      RegistrationHandle handle = entry.getValue();
+      if (handle != null) {
+        handle.remove();
+      }
+    }
+    eventListenerRegistrations.clear();
   }
 }