Merge branch 'stable-2.13' * stable-2.13: Use queue to hold Events during connection glitches Remove obsolete manifest entries Tidy up dependencies Build with plugin API 2.13.2 Change-Id: I32aa58836e0d4f1dc01e30ee7cbedf93aeef65ee
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 45ef707..48465bc 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
@@ -14,6 +14,7 @@ package com.googlesource.gerrit.plugins.rabbitmq.message; +import com.google.gerrit.common.EventListener; import com.google.gerrit.extensions.events.LifecycleListener; import com.google.gerrit.server.events.Event; import com.google.gson.Gson; @@ -21,15 +22,20 @@ import com.google.inject.assistedinject.Assisted; import com.googlesource.gerrit.plugins.rabbitmq.config.Properties; +import com.googlesource.gerrit.plugins.rabbitmq.config.section.AMQP; +import com.googlesource.gerrit.plugins.rabbitmq.config.section.Gerrit; import com.googlesource.gerrit.plugins.rabbitmq.config.section.Monitor; import com.googlesource.gerrit.plugins.rabbitmq.session.Session; import com.googlesource.gerrit.plugins.rabbitmq.session.SessionFactoryProvider; +import com.google.gerrit.server.git.WorkQueue.CancelableRunnable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.Timer; import java.util.TimerTask; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; public class MessagePublisher implements Publisher, LifecycleListener { @@ -37,24 +43,84 @@ private static final int MONITOR_FIRSTTIME_DELAY = 15000; + private static final int MAX_EVENTS = 16384; private final Session session; private final Properties properties; private final Gson gson; private final Timer monitorTimer = new Timer(); private boolean available = true; + private EventListener eventListener; + private final LinkedBlockingQueue<Event> queue = + new LinkedBlockingQueue<>(MAX_EVENTS); + private CancelableRunnable publisher; + private Thread publisherThread; @Inject public MessagePublisher( - @Assisted Properties properties, + @Assisted final Properties properties, SessionFactoryProvider sessionFactoryProvider, Gson gson) { this.session = sessionFactoryProvider.get().create(properties); this.properties = properties; this.gson = gson; + this.eventListener = new EventListener() { + @Override + public void onEvent(Event event) { + try { + if (!publisherThread.isAlive()) { + publisherThread.start(); + } + queue.put(event); + } catch (InterruptedException e) { + LOGGER.warn("Failed to queue event", e); + } + } + }; + this.publisher = new CancelableRunnable() { + + boolean canceled = false; + + @Override + public void run() { + while (!canceled) { + try { + if (isEnable() && session.isOpen()) { + Event event = queue.poll(200, TimeUnit.MILLISECONDS); + if (event != null) { + if (isEnable() && session.isOpen()) { + publishEvent(event); + } else { + queue.put(event); + } + } + } else { + Thread.sleep(1000); + } + } catch (InterruptedException e) { + LOGGER.warn("Interupted while taking event", e); + } + } + } + + @Override + public void cancel() { + this.canceled = true; + } + + @Override + public String toString() { + return "Rabbitmq publisher: " + + properties.getSection(Gerrit.class).listenAs + + "-" + + properties.getSection(AMQP.class).uri; + } + }; } @Override public void start() { + publisherThread = new Thread(publisher); + publisherThread.start(); if (!session.isOpen()) { session.connect(); monitorTimer.schedule(new TimerTask() { @@ -73,18 +139,19 @@ @Override public void stop() { monitorTimer.cancel(); + publisher.cancel(); + if (publisherThread != null) { + try { + publisherThread.join(); + } catch (InterruptedException e) { + // Do nothing + } + } session.disconnect(); available = false; } @Override - public void onEvent(Event event) { - if (available && session.isOpen()) { - session.publish(gson.toJson(event)); - } - } - - @Override public void enable() { available = true; } @@ -113,4 +180,13 @@ public String getName() { return properties.getName(); } + + @Override + public EventListener getEventListener() { + return this.eventListener; + } + + private void publishEvent(Event event) { + session.publish(gson.toJson(event)); + } }
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 28a1d0a..b729dd3 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
@@ -5,7 +5,7 @@ import com.googlesource.gerrit.plugins.rabbitmq.config.Properties; import com.googlesource.gerrit.plugins.rabbitmq.session.Session; -public interface Publisher extends EventListener { +public interface Publisher { void start(); void stop(); void enable(); @@ -14,4 +14,5 @@ Session getSession(); Properties getProperties(); String getName(); + EventListener getEventListener(); }
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 45e0e97..5241ce3 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
@@ -56,7 +56,7 @@ @Override public void onEvent(Event event) { for (Publisher publisher : publishers) { - publisher.onEvent(event); + publisher.getEventListener().onEvent(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 2debaac..a919dbd 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
@@ -122,7 +122,7 @@ eventListeners.add(new UserScopedEventListener() { @Override public void onEvent(Event event) { - publisher.onEvent(event); + publisher.getEventListener().onEvent(event); } @Override public CurrentUser getUser() {