|
|
@@ -1,241 +1,236 @@
|
|
|
-/*
|
|
|
- * @author Ranides Atterwim <ranides@gmail.com>
|
|
|
- * @copyright Ranides Atterwim
|
|
|
- * @license WTFPL
|
|
|
- * @url http://ranides.net/projects/assira
|
|
|
- */
|
|
|
-package net.ranides.assira.events;
|
|
|
-
|
|
|
-import java.util.concurrent.ArrayBlockingQueue;
|
|
|
-import java.util.concurrent.BlockingQueue;
|
|
|
-import java.util.concurrent.LinkedBlockingQueue;
|
|
|
-import java.util.concurrent.TimeUnit;
|
|
|
-import net.ranides.assira.trace.LoggerUtils;
|
|
|
-import net.ranides.assira.trace.ThreadUtils;
|
|
|
-import org.slf4j.Logger;
|
|
|
-
|
|
|
-/**
|
|
|
- * EventRouter kolejkujący zdarzenia. Wywołuje procedury obsługi w bliżej
|
|
|
- * nieokreślonym czasie po zasygnalizowaniu zdarzenia.
|
|
|
- * <p>
|
|
|
- * Procedury obsługi zdarzenia zarejestrowanych obserwatorów są wywoływane
|
|
|
- * w oddzielnym wątku, zarządzanym przez {@code EventRouter}. Należy mieć to na
|
|
|
- * uwadze, ponieważ jest to zupełnie inny wątek, niż ten, który zarejestrował obserwatora,
|
|
|
- * albo który zasygnalizował event.
|
|
|
- * </p>
|
|
|
- *
|
|
|
- * <div class="message-note">thread-safe method</div>
|
|
|
- * @author ranides
|
|
|
- */
|
|
|
-public class EventReactor extends EventDispatcher {
|
|
|
-
|
|
|
- private static final Logger LOGGER = LoggerUtils.getLogger();
|
|
|
-
|
|
|
- private final BlockingQueue<Event> events;
|
|
|
- private final long maxtime;
|
|
|
- private final String name;
|
|
|
- private final Thread distributor;
|
|
|
- private final Events.Stop exit;
|
|
|
-
|
|
|
- /**
|
|
|
- * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
- * {@link #dispose() }. Przechowuje nieograniczoną ilość komunikatów oraz nie
|
|
|
- * narzuca limitów czasowych obsługi zdarzanie.
|
|
|
- * @param name nazwa wątku, który będzie przetwarzał zdarzenia
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static EventReactor newInstance(String name) {
|
|
|
- return newInstance(name, null);
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
- * {@link #dispose() }. Przechowuje nieograniczoną ilość komunikatów oraz nie
|
|
|
- * narzuca limitów czasowych obsługi zdarzanie.
|
|
|
- * <p>
|
|
|
- * Utworzony obiekt {@code EventRouter} obsługuje mechanizm złączania zdarzeń
|
|
|
- * w oparciu o strategię realizowaną przez {@link EventJoiner} przekazany jako
|
|
|
- * argument.
|
|
|
- * </p>
|
|
|
- * @param name nazwa wątku, który będzie przetwarzał zdarzenia
|
|
|
- * @param joiner procedura implementująca mechanizm łączenia wielu zdarzeń w jedno
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static EventReactor newInstance(String name, EventJoiner joiner) {
|
|
|
- return newInstance(name, 0, 0, joiner);
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
- * {@link #dispose() }.
|
|
|
- * @param name nazwa wątku, w którym obsługiwane są zdarzenia
|
|
|
- * @param size maksymalny rozmiar wewnętrznej kolejki komunikatów
|
|
|
- * @param maxtime maksymalny czas oczekiwania na miejsce w kolejce komunikatów
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static EventReactor newInstance(String name, int size, long maxtime) {
|
|
|
- return newInstance(name, size, maxtime, null);
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
- * {@link #dispose() }.
|
|
|
- * <p>
|
|
|
- * Utworzony obiekt {@code EventRouter} obsługuje mechanizm złączania zdarzeń
|
|
|
- * w oparciu o strategię realizowaną przez {@link EventJoiner} przekazany jako
|
|
|
- * argument.
|
|
|
- * </p>
|
|
|
- * @param name nazwa wątku, który będzie przetwarzał zdarzenia
|
|
|
- * @param size maksymalny rozmiar wewnętrznej kolejki komunikatów
|
|
|
- * @param maxtime maksymalny czas oczekiwania na miejsce w kolejce komunikatów
|
|
|
- * @param joiner procedura implementująca mechanizm łączenia wielu zdarzeń w jedno
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static EventReactor newInstance(String name, int size, long maxtime, EventJoiner joiner) {
|
|
|
- return new EventReactor(name, size, maxtime, joiner).start();
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Konstruktor chroniony, żeby niechcący nie utworzyć routera, który nic nie robi.
|
|
|
- * {@code EventRouter} zaczyna działać dopiero po wywołaniu metody {@link #start()}
|
|
|
- * @param name nazwa wątku
|
|
|
- * @param maxtime
|
|
|
- * @param size
|
|
|
- * @param joiner
|
|
|
- */
|
|
|
- protected EventReactor(final String name, int size, long maxtime, final EventJoiner joiner) {
|
|
|
- this.maxtime = maxtime;
|
|
|
- this.events = createQueue(size);
|
|
|
- this.name = name;
|
|
|
- this.exit = Events.stop(this);
|
|
|
- this.distributor = new Thread(name) {
|
|
|
-
|
|
|
- @Override
|
|
|
- public void run() {
|
|
|
- boolean interrupted = true;
|
|
|
- try {
|
|
|
- while (true) {
|
|
|
- Event event = events.take();
|
|
|
- while(null != joiner && joiner.isJoinable(event, events.peek()) ) {
|
|
|
- event = joiner.join(event, events.poll());
|
|
|
- }
|
|
|
- dispatchEvent(event);
|
|
|
- if( dispatchExit(event) ) { interrupted=false; break; }
|
|
|
- }
|
|
|
- } catch (InterruptedException _e) {
|
|
|
- dispatchEvent(Events.interrupt(EventReactor.this));
|
|
|
- interrupted = false;
|
|
|
- }
|
|
|
- if(!interrupted) {
|
|
|
- dispatchEvent(Events.shutdown(EventReactor.this));
|
|
|
- // dispatchEvent(CoreEvent.dispose(EventReactor.this));
|
|
|
- }
|
|
|
- }
|
|
|
- };
|
|
|
- }
|
|
|
-
|
|
|
- private static BlockingQueue<Event> createQueue(int size) {
|
|
|
- if( 0==size || Integer.MAX_VALUE==size) {
|
|
|
- return new LinkedBlockingQueue<>();
|
|
|
- } else {
|
|
|
- return new ArrayBlockingQueue<>(size, true);
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- private EventReactor start() {
|
|
|
- distributor.start();
|
|
|
- return this;
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Wyłącza {@code EventRouter}. Metoda blokująca - czeka, aż EventReactor
|
|
|
- * poprawnie zwolni wszystkie zasoby.
|
|
|
- * <div class="message-note">thread-safe method</div>
|
|
|
- * @return
|
|
|
- */
|
|
|
- public EventReactor stop() {
|
|
|
- signalEvent(exit);
|
|
|
- join();
|
|
|
- super.dispose();
|
|
|
- return this;
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda blokująca - czeka, aż event router zakończy swoje działanie.
|
|
|
- * <p>
|
|
|
- * Uwaga! Żeby metoda kiedykolwiek się zakończyła, w aplikacji musi istnieć
|
|
|
- * co najmniej jeden dodatkowy wątek, który wyśle do routera {@link ExitEvent}
|
|
|
- * lub wywoła metodę {@link #stop}.
|
|
|
- * </p>
|
|
|
- * <div class="message-note">thread-safe method</div>
|
|
|
- * @return
|
|
|
- */
|
|
|
- public EventReactor join() {
|
|
|
- try { distributor.join(); }
|
|
|
- catch (InterruptedException _e) { /* do nothing */ }
|
|
|
- return this;
|
|
|
- }
|
|
|
-
|
|
|
- public String name() {
|
|
|
- return name;
|
|
|
- }
|
|
|
-
|
|
|
- @SuppressWarnings("PMD.CompareObjectsWithEquals")
|
|
|
- private boolean dispatchExit(Event event) {
|
|
|
- if(exit == event) { return true; }
|
|
|
- return (event instanceof Events.Dispose) && this == ((Events.Dispose)event).router();
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * {@inheritDoc}.
|
|
|
- * <p>
|
|
|
- * Wstawia jedynie zdarzenie do kolejki komunikatów.
|
|
|
- * Jeśli kolejka jest w całości wypełniona, tzn osiągnęłą swoją maksymalną pojemność,
|
|
|
- * to czeka na zwolnienie w niej miejsca przez maksymalny dopuszczalny czas.
|
|
|
- * Obie wartości (rozmiar kolejki, długość czekania) są definiowane podczas
|
|
|
- * konstrukcji routera. Jeśli w kolejce nadal nie ma miejsca - to kończy działanie
|
|
|
- * zwracając {@code false}.
|
|
|
- * </p><p>
|
|
|
- * Metoda nie czeka na obsłużenie komunikatu przez zarejestrowanych obserwatorów.
|
|
|
- * Obsługa komunikatu jest uruchamiana w oddzielnym wątku zarządzanym przez router.
|
|
|
- * </p>
|
|
|
- * <div class="message-note">thread-safe method</div>
|
|
|
- * @param event
|
|
|
- * @return true - jeśli zdarzenie zostało zlecone do przekazania.
|
|
|
- */
|
|
|
- @Override
|
|
|
- public boolean signalEvent(Event event) {
|
|
|
- ThreadUtils.dump(LOGGER.isTraceEnabled(), LOGGER::trace);
|
|
|
- if(0 == maxtime) {
|
|
|
- return events.offer(event);
|
|
|
- }
|
|
|
- try {
|
|
|
- return events.offer(event, maxtime, TimeUnit.MILLISECONDS);
|
|
|
- } catch (InterruptedException _e) {
|
|
|
- return false;
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Zleca wyłączenie i zwalnienie zasobów routerowi. Nie czeka na zakończenie
|
|
|
- * jego pracy.
|
|
|
- * <div class="message-note">thread-safe method</div>
|
|
|
- */
|
|
|
- @Override
|
|
|
- public void dispose() {
|
|
|
- signalEvent(exit);
|
|
|
- new Thread(){
|
|
|
- @Override
|
|
|
- public void run() {
|
|
|
- EventReactor.this.join();
|
|
|
- EventReactor.super.dispose();
|
|
|
- }
|
|
|
- }.start();
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public String toString() {
|
|
|
- return "EventReactor<" + name + ">";
|
|
|
- }
|
|
|
-
|
|
|
-}
|
|
|
+/*
|
|
|
+ * @author Ranides Atterwim <ranides@gmail.com>
|
|
|
+ * @copyright Ranides Atterwim
|
|
|
+ * @license WTFPL
|
|
|
+ * @url http://ranides.net/projects/assira
|
|
|
+ */
|
|
|
+package net.ranides.assira.events;
|
|
|
+
|
|
|
+import java.util.concurrent.ArrayBlockingQueue;
|
|
|
+import java.util.concurrent.BlockingQueue;
|
|
|
+import java.util.concurrent.LinkedBlockingQueue;
|
|
|
+import java.util.concurrent.TimeUnit;
|
|
|
+import net.ranides.assira.trace.LoggerUtils;
|
|
|
+import net.ranides.assira.trace.ThreadUtils;
|
|
|
+import org.slf4j.Logger;
|
|
|
+
|
|
|
+/**
|
|
|
+ * EventRouter kolejkujący zdarzenia. Wywołuje procedury obsługi w bliżej
|
|
|
+ * nieokreślonym czasie po zasygnalizowaniu zdarzenia.
|
|
|
+ * <p>
|
|
|
+ * Procedury obsługi zdarzenia zarejestrowanych obserwatorów są wywoływane
|
|
|
+ * w oddzielnym wątku, zarządzanym przez {@code EventRouter}. Należy mieć to na
|
|
|
+ * uwadze, ponieważ jest to zupełnie inny wątek, niż ten, który zarejestrował obserwatora,
|
|
|
+ * albo który zasygnalizował event.
|
|
|
+ * </p>
|
|
|
+ *
|
|
|
+ * <div class="message-note">thread-safe method</div>
|
|
|
+ * @author ranides
|
|
|
+ */
|
|
|
+public class EventReactor extends EventDispatcher {
|
|
|
+
|
|
|
+ private static final Logger LOGGER = LoggerUtils.getLogger();
|
|
|
+
|
|
|
+ private final BlockingQueue<Event> events;
|
|
|
+ private final long maxtime;
|
|
|
+ private final Thread distributor;
|
|
|
+ private final Events.Stop exit;
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
+ * {@link #dispose() }. Przechowuje nieograniczoną ilość komunikatów oraz nie
|
|
|
+ * narzuca limitów czasowych obsługi zdarzanie.
|
|
|
+ * @param name nazwa wątku, który będzie przetwarzał zdarzenia
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static EventReactor newInstance(String name) {
|
|
|
+ return newInstance(name, null);
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
+ * {@link #dispose() }. Przechowuje nieograniczoną ilość komunikatów oraz nie
|
|
|
+ * narzuca limitów czasowych obsługi zdarzanie.
|
|
|
+ * <p>
|
|
|
+ * Utworzony obiekt {@code EventRouter} obsługuje mechanizm złączania zdarzeń
|
|
|
+ * w oparciu o strategię realizowaną przez {@link EventJoiner} przekazany jako
|
|
|
+ * argument.
|
|
|
+ * </p>
|
|
|
+ * @param name nazwa wątku, który będzie przetwarzał zdarzenia
|
|
|
+ * @param joiner procedura implementująca mechanizm łączenia wielu zdarzeń w jedno
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static EventReactor newInstance(String name, EventJoiner joiner) {
|
|
|
+ return newInstance(name, 0, 0, joiner);
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
+ * {@link #dispose() }.
|
|
|
+ * @param name nazwa wątku, w którym obsługiwane są zdarzenia
|
|
|
+ * @param size maksymalny rozmiar wewnętrznej kolejki komunikatów
|
|
|
+ * @param maxtime maksymalny czas oczekiwania na miejsce w kolejce komunikatów
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static EventReactor newInstance(String name, int size, long maxtime) {
|
|
|
+ return newInstance(name, size, maxtime, null);
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Tworzy nowy obiekt {@code EventRouter}. Utworzony obiekt należy zniszczyć po użyciu za pomocą
|
|
|
+ * {@link #dispose() }.
|
|
|
+ * <p>
|
|
|
+ * Utworzony obiekt {@code EventRouter} obsługuje mechanizm złączania zdarzeń
|
|
|
+ * w oparciu o strategię realizowaną przez {@link EventJoiner} przekazany jako
|
|
|
+ * argument.
|
|
|
+ * </p>
|
|
|
+ * @param name nazwa wątku, który będzie przetwarzał zdarzenia
|
|
|
+ * @param size maksymalny rozmiar wewnętrznej kolejki komunikatów
|
|
|
+ * @param maxtime maksymalny czas oczekiwania na miejsce w kolejce komunikatów
|
|
|
+ * @param joiner procedura implementująca mechanizm łączenia wielu zdarzeń w jedno
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static EventReactor newInstance(String name, int size, long maxtime, EventJoiner joiner) {
|
|
|
+ return new EventReactor(name, size, maxtime, joiner).start();
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Konstruktor chroniony, żeby niechcący nie utworzyć routera, który nic nie robi.
|
|
|
+ * {@code EventRouter} zaczyna działać dopiero po wywołaniu metody {@link #start()}
|
|
|
+ * @param name nazwa wątku
|
|
|
+ * @param maxtime
|
|
|
+ * @param size
|
|
|
+ * @param joiner
|
|
|
+ */
|
|
|
+ protected EventReactor(final String name, int size, long maxtime, final EventJoiner joiner) {
|
|
|
+ super(name);
|
|
|
+ this.maxtime = maxtime;
|
|
|
+ this.events = createQueue(size);
|
|
|
+ this.exit = Events.stop(this);
|
|
|
+ this.distributor = new Thread(name) {
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void run() {
|
|
|
+ boolean interrupted = true;
|
|
|
+ try {
|
|
|
+ while (true) {
|
|
|
+ Event event = events.take();
|
|
|
+ while(null != joiner && joiner.isJoinable(event, events.peek()) ) {
|
|
|
+ event = joiner.join(event, events.poll());
|
|
|
+ }
|
|
|
+ dispatchEvent(event);
|
|
|
+ if( dispatchExit(event) ) { interrupted=false; break; }
|
|
|
+ }
|
|
|
+ } catch (InterruptedException _e) {
|
|
|
+ dispatchEvent(Events.interrupt(EventReactor.this));
|
|
|
+ interrupted = false;
|
|
|
+ }
|
|
|
+ if(!interrupted) {
|
|
|
+ dispatchEvent(Events.shutdown(EventReactor.this));
|
|
|
+ // dispatchEvent(CoreEvent.dispose(EventReactor.this));
|
|
|
+ }
|
|
|
+ }
|
|
|
+ };
|
|
|
+ }
|
|
|
+
|
|
|
+ private static BlockingQueue<Event> createQueue(int size) {
|
|
|
+ if( 0==size || Integer.MAX_VALUE==size) {
|
|
|
+ return new LinkedBlockingQueue<>();
|
|
|
+ } else {
|
|
|
+ return new ArrayBlockingQueue<>(size, true);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ private EventReactor start() {
|
|
|
+ distributor.start();
|
|
|
+ return this;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Wyłącza {@code EventRouter}. Metoda blokująca - czeka, aż EventReactor
|
|
|
+ * poprawnie zwolni wszystkie zasoby.
|
|
|
+ * <div class="message-note">thread-safe method</div>
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public EventReactor stop() {
|
|
|
+ signalEvent(exit);
|
|
|
+ join();
|
|
|
+ super.dispose();
|
|
|
+ return this;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda blokująca - czeka, aż event router zakończy swoje działanie.
|
|
|
+ * <p>
|
|
|
+ * Uwaga! Żeby metoda kiedykolwiek się zakończyła, w aplikacji musi istnieć
|
|
|
+ * co najmniej jeden dodatkowy wątek, który wyśle do routera {@link ExitEvent}
|
|
|
+ * lub wywoła metodę {@link #stop}.
|
|
|
+ * </p>
|
|
|
+ * <div class="message-note">thread-safe method</div>
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public EventReactor join() {
|
|
|
+ try { distributor.join(); }
|
|
|
+ catch (InterruptedException _e) { /* do nothing */ }
|
|
|
+ return this;
|
|
|
+ }
|
|
|
+
|
|
|
+ @SuppressWarnings("PMD.CompareObjectsWithEquals")
|
|
|
+ private boolean dispatchExit(Event event) {
|
|
|
+ if(exit == event) { return true; }
|
|
|
+ return (event instanceof Events.Dispose) && this == ((Events.Dispose)event).router();
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * {@inheritDoc}.
|
|
|
+ * <p>
|
|
|
+ * Wstawia jedynie zdarzenie do kolejki komunikatów.
|
|
|
+ * Jeśli kolejka jest w całości wypełniona, tzn osiągnęłą swoją maksymalną pojemność,
|
|
|
+ * to czeka na zwolnienie w niej miejsca przez maksymalny dopuszczalny czas.
|
|
|
+ * Obie wartości (rozmiar kolejki, długość czekania) są definiowane podczas
|
|
|
+ * konstrukcji routera. Jeśli w kolejce nadal nie ma miejsca - to kończy działanie
|
|
|
+ * zwracając {@code false}.
|
|
|
+ * </p><p>
|
|
|
+ * Metoda nie czeka na obsłużenie komunikatu przez zarejestrowanych obserwatorów.
|
|
|
+ * Obsługa komunikatu jest uruchamiana w oddzielnym wątku zarządzanym przez router.
|
|
|
+ * </p>
|
|
|
+ * <div class="message-note">thread-safe method</div>
|
|
|
+ * @param event
|
|
|
+ * @return true - jeśli zdarzenie zostało zlecone do przekazania.
|
|
|
+ */
|
|
|
+ @Override
|
|
|
+ public boolean signalEvent(Event event) {
|
|
|
+ ThreadUtils.dump(LOGGER.isTraceEnabled(), LOGGER::trace);
|
|
|
+ if(0 == maxtime) {
|
|
|
+ return events.offer(event);
|
|
|
+ }
|
|
|
+ try {
|
|
|
+ return events.offer(event, maxtime, TimeUnit.MILLISECONDS);
|
|
|
+ } catch (InterruptedException _e) {
|
|
|
+ return false;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Zleca wyłączenie i zwalnienie zasobów routerowi. Nie czeka na zakończenie
|
|
|
+ * jego pracy.
|
|
|
+ * <div class="message-note">thread-safe method</div>
|
|
|
+ */
|
|
|
+ @Override
|
|
|
+ public void dispose() {
|
|
|
+ signalEvent(exit);
|
|
|
+ new Thread(){
|
|
|
+ @Override
|
|
|
+ public void run() {
|
|
|
+ EventReactor.this.join();
|
|
|
+ EventReactor.super.dispose();
|
|
|
+ }
|
|
|
+ }.start();
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public String toString() {
|
|
|
+ return "EventReactor<" + name() + ">";
|
|
|
+ }
|
|
|
+
|
|
|
+}
|