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