|
|
@@ -1,457 +1,456 @@
|
|
|
-/*
|
|
|
- * @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.BlockingQueue;
|
|
|
-import java.util.concurrent.LinkedBlockingQueue;
|
|
|
-import java.util.concurrent.TimeUnit;
|
|
|
-import net.ranides.assira.annotations.Meta;
|
|
|
-import net.ranides.assira.trace.LoggerUtils;
|
|
|
-import org.slf4j.Logger;
|
|
|
-
|
|
|
-/**
|
|
|
- * Klasa pozwalająca na synchroniczną obsługę zdarzeń. Pozwala wstrzymać wykonanie
|
|
|
- * bieżącego wątku i oczekiwać, aż do wskazanego {@link EventRouter}'a dotrze
|
|
|
- * zdarzenie określonego typu.
|
|
|
- * <p>
|
|
|
- * Zależnie od rodzaju implementacji, klasa może oferować usługi o różnym stopniu
|
|
|
- * dokładności i bezpieczeństwa.
|
|
|
- * </p>
|
|
|
- * <p>Najsilniejsze gwarancje, to możliwość obsługi wszystkich zdarzeń określonego
|
|
|
- * rodzaju bez ryzyka pominięcia żadnego. Nawet jeśli jedno lub więcej zdarzeń
|
|
|
- * nastąpiło przed (lub pomiędzy) wywołaniem metody {@link #waitForEvent}, to zostanie
|
|
|
- * ono dostarczone. Gwarantowana jest również prawidłowa kolejność dostarczenia
|
|
|
- * komunikatów przez kolejne wywołania metod blokujących. Implementacje "silne"
|
|
|
- * mają duże wymagania pamięciowe, w ekstremalnym przypadku mogą spowodować
|
|
|
- * przepełnienie sterty, jeśli program zostanie zalany komunikatami.
|
|
|
- * </p>
|
|
|
- * <p>
|
|
|
- * "Słaby {@code EventLock}" daje gwarancje, że metody {@link #waitForEvent} nie
|
|
|
- * będą blokować, jeśli pomiędzy ich wywołaniami (albo przed wywołaniem) wystąpiło
|
|
|
- * obserwowane zdarzenie. Nie dają jednak gwarancji dostarczenia wszystkich
|
|
|
- * komunikatów - część z nich może być utracona/pominięta, w przypadku wysokiego
|
|
|
- * obciążenia. Słaby {@code EventLock} udostępnia informację o ilości zgubionych
|
|
|
- * komunikatów. Jego zaletą jest znacznie mniejsze zużycie pamięci oraz brak podatności
|
|
|
- * na przeciążenie aplikacji.
|
|
|
- * </p>
|
|
|
- * <p>
|
|
|
- * "Niebezpieczny {@code EventLock}" nie oferuje żadnych gwarancji. Co więcej:
|
|
|
- * metody {@link #waitForEvent} zawsze blokują, a jako wynik zwracają tylko i
|
|
|
- * wyłącznie wydarzenie, które nastąpiło w czasie blokady. Z tego powodu jest
|
|
|
- * silnie podatny na <i>race condition</i>. Jego zaletą jest niemal całkowity
|
|
|
- * brak obciążenia pamięci oraz procesora. Z tego powodu może być z powodzeniem
|
|
|
- * używany do obserwacji zdarzeń przychodzących w bardzo dużo odstępach (względem czasu
|
|
|
- * obsługi) bez obaw o wydajność aplikacji.
|
|
|
- * </p>
|
|
|
- * @author ranides
|
|
|
- *
|
|
|
- * @todo (assira #1) EventStream zamiast EventLock
|
|
|
- */
|
|
|
-public abstract class EventLock<E extends Event> {
|
|
|
-
|
|
|
- private static final int LOG_BATCH = 100;
|
|
|
-
|
|
|
- /**
|
|
|
- * Zwraca liczbę komunikatów zgubionych przez cały czas istnienia obiektu.
|
|
|
- * @return
|
|
|
- */
|
|
|
- public int discarded() {
|
|
|
- throw new UnsupportedOperationException();
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Zwraca liczbę wszystkich otrzymanych komunikatów przez cały czas istnienia obiektu
|
|
|
- * (w tym komuniaktów zgubionych).
|
|
|
- * @return
|
|
|
- */
|
|
|
- public int counter() {
|
|
|
- throw new UnsupportedOperationException();
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Zwraca liczbę aktualnie oczekujących komunikatów.
|
|
|
- * @return
|
|
|
- */
|
|
|
- public int pending() {
|
|
|
- throw new UnsupportedOperationException();
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Niszczy obiekt {@code EventLock} i zwalnia wszystkie zasoby.
|
|
|
- */
|
|
|
- public abstract void release();
|
|
|
-
|
|
|
- /**
|
|
|
- * Resetuje stan obiektu - jeśli {@code EventLock} przechowywał w kolejce
|
|
|
- * nieobsłużone zdarzenia - są one usuwane. Wywołanie funkcji {@code reset}
|
|
|
- * pomiędzy dwoma wywołaniami metody {@code wait} może spowodować, że część zdarzeń,
|
|
|
- * którymi {@code EventLock} był zainteresowany zostanie utracona.
|
|
|
- */
|
|
|
- public abstract void reset();
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda blokująca: wstrzymuje wykonanie bieżącego wątku do momentu, aż
|
|
|
- * do obserwowanego {@code EventRouter}'a dotrze dowolne zdarzenie oczekiwane
|
|
|
- * przez {@code EventLock} lub wystąpi wyjątek {@code InterruptedException}.
|
|
|
- * Zobacz opis klasy {@link EventLock} aby poznać szczegóły.
|
|
|
- * @return oczekiwane zdarzenie
|
|
|
- * @throws InterruptedException
|
|
|
- */
|
|
|
- public abstract E waitForEvent() throws InterruptedException;
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda blokująca: wstrzymuje wykonanie bieżącego wątku do momentu, aż
|
|
|
- * do obserwowanego {@code EventRouter}'a dotrze dowolne zdarzenie oczekiwane
|
|
|
- * przez {@code EventLock}, wystąpi wyjątek {@code InterruptedException}, lub
|
|
|
- * minie podany czas.
|
|
|
- * Zobacz opis klasy {@link EventLock} aby poznać szczegóły.
|
|
|
- * @param timeout
|
|
|
- * @return oczekiwane zdarzenie
|
|
|
- * @throws InterruptedException
|
|
|
- */
|
|
|
- public abstract E waitForEvent(long timeout) throws InterruptedException;
|
|
|
-
|
|
|
- /**
|
|
|
- * Wersja oczekująca na zdarzenie konkretnego rodzaju. Zachowuje się identycznie
|
|
|
- * jak {@link #waitForEvent()}, z tą różnicą, że reaguje na mniejszy zakres
|
|
|
- * zdarzeń - zawężony tylko do podanej klasy oraz jej pochodnych.
|
|
|
- * <p>
|
|
|
- * Metoda przydatna szczególnie wtedy, gdy utworzony obiekt {@code EventLock}
|
|
|
- * reaguje na bardzo ogólną klasę zdarzeń, a użytkownik w danym momencie chce
|
|
|
- * wstrzymać wykonanie w oczekiwaniu na zdarzenie bardziej szczegółowe.
|
|
|
- * </p>
|
|
|
- * @param <T>
|
|
|
- * @param event
|
|
|
- * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
|
|
|
- * @throws InterruptedException
|
|
|
- */
|
|
|
- public abstract <T extends E> T waitForEvent(Class<T> event) throws InterruptedException;
|
|
|
-
|
|
|
- /**
|
|
|
- * Wersja oczekująca na zdarzenie konkretnego rodzaju. Zachowuje się identycznie
|
|
|
- * jak {@link #waitForEvent(long timeout)}, z tą różnicą, że reaguje na mniejszy zakres
|
|
|
- * zdarzeń - zawężony tylko do podanej klasy oraz jej pochodnych.
|
|
|
- * <p>
|
|
|
- * Metoda przydatna szczególnie wtedy, gdy utworzony obiekt {@code EventLock}
|
|
|
- * reaguje na bardzo ogólną klasę zdarzeń, a użytkownik w danym momencie chce
|
|
|
- * wstrzymać wykonanie w oczekiwaniu na zdarzenie bardziej szczegółowe.
|
|
|
- * </p>
|
|
|
- * @param <T>
|
|
|
- * @param event
|
|
|
- * @param timeout
|
|
|
- * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
|
|
|
- * @throws InterruptedException
|
|
|
- */
|
|
|
- public abstract <T extends E> T waitForEvent(Class<T> event, long timeout) throws InterruptedException;
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda tworzy silny {@code EventLock} - szczegóły "silnego kontraktu"
|
|
|
- * są w opisie klasy {@link EventLock}).
|
|
|
- * <p>
|
|
|
- * Jeśli obiekt nie jest już potrzebny, musi zostać zwolniony za pomocą metody {@link #release()}
|
|
|
- * </p>
|
|
|
- * @param event
|
|
|
- * @param router
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static <EE extends Event> EventLock<EE> lock(Class<EE> event, EventRouter router) {
|
|
|
- return lock(event, router, Integer.MAX_VALUE);
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda tworzy słaby {@code EventLock} - szczegóły "słabego kontraktu"
|
|
|
- * są w opisie klasy {@link EventLock}).
|
|
|
- * <p>
|
|
|
- * Jeśli obiekt nie jest już potrzebny, musi zostać zwolniony za pomocą metody {@link #release()}
|
|
|
- * </p>
|
|
|
- * @param event
|
|
|
- * @param router
|
|
|
- * @param capacity
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static <EE extends Event> EventLock<EE> lock(Class<EE> event, EventRouter router, int capacity) {
|
|
|
- QueEventLock<EE> lock = new QueEventLock<>(event, router, capacity);
|
|
|
- router.addEventListener(event, lock);
|
|
|
- return lock;
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda tworzy silny {@code EventLock}, który może zostać wykorzystany tylko raz
|
|
|
- * - szczegóły "silnego kontraktu" są w opisie klasy {@link EventLock}).
|
|
|
- * <p>
|
|
|
- * Po zakończeniu metody {@code waitForEvent} obiekt sam automatycznie
|
|
|
- * zwalnia wszystkie zasoby i przestaje być funcjonalny.
|
|
|
- * </p>
|
|
|
- * @param event
|
|
|
- * @param router
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static <EE extends Event> EventLock<EE> singleLock(Class<EE> event, EventRouter router) {
|
|
|
- return singleLock(event, router, Integer.MAX_VALUE);
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda tworzy słaby {@code EventLock}, który może zostać wykorzystany tylko raz
|
|
|
- * - szczegóły "słabego kontraktu" są w opisie klasy {@link EventLock}).
|
|
|
- * <p>
|
|
|
- * Po zakończeniu metody {@code waitForEvent} obiekt sam automatycznie
|
|
|
- * zwalnia wszystkie zasoby i przestaje być funcjonalny.
|
|
|
- * </p>
|
|
|
- * @param event
|
|
|
- * @param router
|
|
|
- * @return
|
|
|
- */
|
|
|
- public static <EE extends Event> EventLock<EE> singleLock(Class<EE> event, EventRouter router, int capacity) {
|
|
|
- SingleLock<EE> lock = new SingleLock<>(event, router, capacity);
|
|
|
- router.addEventListener(event, lock);
|
|
|
- return lock;
|
|
|
- }
|
|
|
-
|
|
|
- /**
|
|
|
- * Metoda tworzy niebezpieczny {@code EventLock} - szczegóły "niebezpiecznego kontraktu"
|
|
|
- * są w opisie klasy {@link EventLock}).
|
|
|
- * <p>
|
|
|
- * Jeśli obiekt nie jest już potrzebny, musi zostać zwolniony za pomocą metody {@link #release()}
|
|
|
- * </p>
|
|
|
- * @param event
|
|
|
- * @param router
|
|
|
- * @return
|
|
|
- */
|
|
|
- @Meta.Unsafe
|
|
|
- public static <EE extends Event> EventLock<EE> unsafeLock(Class<EE> event, EventRouter router) {
|
|
|
- UnsafeEventLock<EE> lock = new UnsafeEventLock<>(event, router);
|
|
|
- router.addEventListener(event, lock);
|
|
|
- return lock;
|
|
|
- }
|
|
|
-
|
|
|
- private static class UnsafeEventLock<EE extends Event> extends EventLock<EE> implements EventListener<EE> {
|
|
|
-
|
|
|
- private final Class<? extends EE> observed;
|
|
|
- private final EventRouter router;
|
|
|
- private volatile EE hit;
|
|
|
-
|
|
|
- public UnsafeEventLock(Class<? extends EE> event, EventRouter router) {
|
|
|
- super();
|
|
|
- this.observed = event;
|
|
|
- this.router = router;
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void handleEvent(EE event) {
|
|
|
- synchronized (this) {
|
|
|
- hit = event;
|
|
|
- this.notifyAll();
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public EE waitForEvent() throws InterruptedException {
|
|
|
- synchronized (this) {
|
|
|
- while (hit==null) { this.wait(); }
|
|
|
- EE result = hit;
|
|
|
- hit = null;
|
|
|
- return result;
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public EE waitForEvent(long timeout) throws InterruptedException {
|
|
|
- long end = System.currentTimeMillis() + timeout;
|
|
|
- long rem = timeout;
|
|
|
- synchronized (this) {
|
|
|
- while (hit==null && rem>0) {
|
|
|
- this.wait(rem);
|
|
|
- rem = end - System.currentTimeMillis();
|
|
|
- }
|
|
|
- EE result = hit;
|
|
|
- hit = null;
|
|
|
- return result;
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
|
|
|
- synchronized (this) {
|
|
|
- while( !event.isInstance(hit) ) { this.wait(); }
|
|
|
- Event ret = hit;
|
|
|
- hit = null;
|
|
|
- return event.cast(ret);
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
|
|
|
- long end = System.currentTimeMillis() + timeout;
|
|
|
- long rem = timeout;
|
|
|
- synchronized (this) {
|
|
|
- while( !event.isInstance(hit) && rem>0) {
|
|
|
- this.wait(rem);
|
|
|
- rem = end - System.currentTimeMillis();
|
|
|
- }
|
|
|
- EE ret = hit;
|
|
|
- hit = null;
|
|
|
- return event.cast(ret);
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void reset() {
|
|
|
- hit = null;
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void release() {
|
|
|
- router.removeEventListener(observed, this);
|
|
|
- }
|
|
|
-
|
|
|
- }
|
|
|
-
|
|
|
- private static class QueEventLock<EE extends Event> extends EventLock<EE> implements EventListener<EE> {
|
|
|
-
|
|
|
- private static final Logger LOGGER = LoggerUtils.getLogger();
|
|
|
-
|
|
|
- private final Class<? extends EE> observed;
|
|
|
- private final EventRouter router;
|
|
|
- private final BlockingQueue<EE> events;
|
|
|
- private int discarded;
|
|
|
- private int processed;
|
|
|
-
|
|
|
- public QueEventLock(Class<? extends EE> event, EventRouter router, int capacity) {
|
|
|
- this.observed = event;
|
|
|
- this.router = router;
|
|
|
- this.processed = 0;
|
|
|
- this.discarded = 0;
|
|
|
- this.events = new LinkedBlockingQueue<>(capacity);
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public EE waitForEvent() throws InterruptedException {
|
|
|
- return events.take();
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public EE waitForEvent(long timeout) throws InterruptedException {
|
|
|
- return events.poll(timeout, TimeUnit.MILLISECONDS);
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
|
|
|
- while(true) {
|
|
|
- EE recent = events.take();
|
|
|
- if( event.isInstance(recent) ) {
|
|
|
- return event.cast(recent);
|
|
|
- }
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
|
|
|
- long end = System.currentTimeMillis() + timeout;
|
|
|
- long rem = timeout;
|
|
|
- while(rem > 0) {
|
|
|
- EE recent = events.poll(rem, TimeUnit.MILLISECONDS);
|
|
|
- if( event.isInstance(recent) ) {
|
|
|
- return event.cast(recent);
|
|
|
- }
|
|
|
- rem = end - System.currentTimeMillis();
|
|
|
- }
|
|
|
- return null;
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void release() {
|
|
|
- LOGGER.trace("overall processed={}", processed);
|
|
|
- if( discarded > 0 ) {
|
|
|
- LOGGER.warn("overall discarded={}",discarded);
|
|
|
- }
|
|
|
- router.removeEventListener(observed, this);
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void reset() {
|
|
|
- events.clear();
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public void handleEvent(EE event) {
|
|
|
- processed++;
|
|
|
- while( !events.offer(event) ) {
|
|
|
- Event ignored = events.poll();
|
|
|
- discarded++;
|
|
|
- if( 0 == (processed % LOG_BATCH) ) {
|
|
|
- LOGGER.warn("events flood. processed={} discarded={} event={}", processed, discarded, ignored);
|
|
|
- }
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public int discarded() {
|
|
|
- return discarded;
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public int counter() {
|
|
|
- return processed;
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public int pending() {
|
|
|
- return events.size();
|
|
|
- }
|
|
|
-
|
|
|
- }
|
|
|
-
|
|
|
- private static class SingleLock<EE extends Event> extends QueEventLock<EE> {
|
|
|
-
|
|
|
- public <T extends EE> SingleLock(Class<T> event, EventRouter router, int capacity) {
|
|
|
- super(event, router, capacity);
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public EE waitForEvent() throws InterruptedException {
|
|
|
- try {
|
|
|
- return super.waitForEvent();
|
|
|
- } finally {
|
|
|
- release();
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
|
|
|
- try {
|
|
|
- return super.waitForEvent(event);
|
|
|
- } finally {
|
|
|
- release();
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
|
|
|
- try {
|
|
|
- return super.waitForEvent(event, timeout);
|
|
|
- } finally {
|
|
|
- release();
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- @Override
|
|
|
- public EE waitForEvent(long timeout) throws InterruptedException {
|
|
|
- try {
|
|
|
- return super.waitForEvent(timeout);
|
|
|
- } finally {
|
|
|
- release();
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
-
|
|
|
-
|
|
|
- }
|
|
|
-
|
|
|
-}
|
|
|
+/*
|
|
|
+ * @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.BlockingQueue;
|
|
|
+import java.util.concurrent.LinkedBlockingQueue;
|
|
|
+import java.util.concurrent.TimeUnit;
|
|
|
+import net.ranides.assira.annotations.Meta;
|
|
|
+import net.ranides.assira.trace.LoggerUtils;
|
|
|
+import org.slf4j.Logger;
|
|
|
+
|
|
|
+/**
|
|
|
+ * Klasa pozwalająca na synchroniczną obsługę zdarzeń. Pozwala wstrzymać wykonanie
|
|
|
+ * bieżącego wątku i oczekiwać, aż do wskazanego {@link EventRouter}'a dotrze
|
|
|
+ * zdarzenie określonego typu.
|
|
|
+ * <p>
|
|
|
+ * Zależnie od rodzaju implementacji, klasa może oferować usługi o różnym stopniu
|
|
|
+ * dokładności i bezpieczeństwa.
|
|
|
+ * </p>
|
|
|
+ * <p>Najsilniejsze gwarancje, to możliwość obsługi wszystkich zdarzeń określonego
|
|
|
+ * rodzaju bez ryzyka pominięcia żadnego. Nawet jeśli jedno lub więcej zdarzeń
|
|
|
+ * nastąpiło przed (lub pomiędzy) wywołaniem metody {@link #waitForEvent}, to zostanie
|
|
|
+ * ono dostarczone. Gwarantowana jest również prawidłowa kolejność dostarczenia
|
|
|
+ * komunikatów przez kolejne wywołania metod blokujących. Implementacje "silne"
|
|
|
+ * mają duże wymagania pamięciowe, w ekstremalnym przypadku mogą spowodować
|
|
|
+ * przepełnienie sterty, jeśli program zostanie zalany komunikatami.
|
|
|
+ * </p>
|
|
|
+ * <p>
|
|
|
+ * "Słaby {@code EventLock}" daje gwarancje, że metody {@link #waitForEvent} nie
|
|
|
+ * będą blokować, jeśli pomiędzy ich wywołaniami (albo przed wywołaniem) wystąpiło
|
|
|
+ * obserwowane zdarzenie. Nie dają jednak gwarancji dostarczenia wszystkich
|
|
|
+ * komunikatów - część z nich może być utracona/pominięta, w przypadku wysokiego
|
|
|
+ * obciążenia. Słaby {@code EventLock} udostępnia informację o ilości zgubionych
|
|
|
+ * komunikatów. Jego zaletą jest znacznie mniejsze zużycie pamięci oraz brak podatności
|
|
|
+ * na przeciążenie aplikacji.
|
|
|
+ * </p>
|
|
|
+ * <p>
|
|
|
+ * "Niebezpieczny {@code EventLock}" nie oferuje żadnych gwarancji. Co więcej:
|
|
|
+ * metody {@link #waitForEvent} zawsze blokują, a jako wynik zwracają tylko i
|
|
|
+ * wyłącznie wydarzenie, które nastąpiło w czasie blokady. Z tego powodu jest
|
|
|
+ * silnie podatny na <i>race condition</i>. Jego zaletą jest niemal całkowity
|
|
|
+ * brak obciążenia pamięci oraz procesora. Z tego powodu może być z powodzeniem
|
|
|
+ * używany do obserwacji zdarzeń przychodzących w bardzo dużo odstępach (względem czasu
|
|
|
+ * obsługi) bez obaw o wydajność aplikacji.
|
|
|
+ * </p>
|
|
|
+ * @author ranides
|
|
|
+ *
|
|
|
+ */
|
|
|
+public abstract class EventLock<E extends Event> {
|
|
|
+
|
|
|
+ private static final int LOG_BATCH = 100;
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Zwraca liczbę komunikatów zgubionych przez cały czas istnienia obiektu.
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public int discarded() {
|
|
|
+ throw new UnsupportedOperationException();
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Zwraca liczbę wszystkich otrzymanych komunikatów przez cały czas istnienia obiektu
|
|
|
+ * (w tym komuniaktów zgubionych).
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public int counter() {
|
|
|
+ throw new UnsupportedOperationException();
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Zwraca liczbę aktualnie oczekujących komunikatów.
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public int pending() {
|
|
|
+ throw new UnsupportedOperationException();
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Niszczy obiekt {@code EventLock} i zwalnia wszystkie zasoby.
|
|
|
+ */
|
|
|
+ public abstract void release();
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Resetuje stan obiektu - jeśli {@code EventLock} przechowywał w kolejce
|
|
|
+ * nieobsłużone zdarzenia - są one usuwane. Wywołanie funkcji {@code reset}
|
|
|
+ * pomiędzy dwoma wywołaniami metody {@code wait} może spowodować, że część zdarzeń,
|
|
|
+ * którymi {@code EventLock} był zainteresowany zostanie utracona.
|
|
|
+ */
|
|
|
+ public abstract void reset();
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda blokująca: wstrzymuje wykonanie bieżącego wątku do momentu, aż
|
|
|
+ * do obserwowanego {@code EventRouter}'a dotrze dowolne zdarzenie oczekiwane
|
|
|
+ * przez {@code EventLock} lub wystąpi wyjątek {@code InterruptedException}.
|
|
|
+ * Zobacz opis klasy {@link EventLock} aby poznać szczegóły.
|
|
|
+ * @return oczekiwane zdarzenie
|
|
|
+ * @throws InterruptedException
|
|
|
+ */
|
|
|
+ public abstract E waitForEvent() throws InterruptedException;
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda blokująca: wstrzymuje wykonanie bieżącego wątku do momentu, aż
|
|
|
+ * do obserwowanego {@code EventRouter}'a dotrze dowolne zdarzenie oczekiwane
|
|
|
+ * przez {@code EventLock}, wystąpi wyjątek {@code InterruptedException}, lub
|
|
|
+ * minie podany czas.
|
|
|
+ * Zobacz opis klasy {@link EventLock} aby poznać szczegóły.
|
|
|
+ * @param timeout
|
|
|
+ * @return oczekiwane zdarzenie
|
|
|
+ * @throws InterruptedException
|
|
|
+ */
|
|
|
+ public abstract E waitForEvent(long timeout) throws InterruptedException;
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Wersja oczekująca na zdarzenie konkretnego rodzaju. Zachowuje się identycznie
|
|
|
+ * jak {@link #waitForEvent()}, z tą różnicą, że reaguje na mniejszy zakres
|
|
|
+ * zdarzeń - zawężony tylko do podanej klasy oraz jej pochodnych.
|
|
|
+ * <p>
|
|
|
+ * Metoda przydatna szczególnie wtedy, gdy utworzony obiekt {@code EventLock}
|
|
|
+ * reaguje na bardzo ogólną klasę zdarzeń, a użytkownik w danym momencie chce
|
|
|
+ * wstrzymać wykonanie w oczekiwaniu na zdarzenie bardziej szczegółowe.
|
|
|
+ * </p>
|
|
|
+ * @param <T>
|
|
|
+ * @param event
|
|
|
+ * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
|
|
|
+ * @throws InterruptedException
|
|
|
+ */
|
|
|
+ public abstract <T extends E> T waitForEvent(Class<T> event) throws InterruptedException;
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Wersja oczekująca na zdarzenie konkretnego rodzaju. Zachowuje się identycznie
|
|
|
+ * jak {@link #waitForEvent(long timeout)}, z tą różnicą, że reaguje na mniejszy zakres
|
|
|
+ * zdarzeń - zawężony tylko do podanej klasy oraz jej pochodnych.
|
|
|
+ * <p>
|
|
|
+ * Metoda przydatna szczególnie wtedy, gdy utworzony obiekt {@code EventLock}
|
|
|
+ * reaguje na bardzo ogólną klasę zdarzeń, a użytkownik w danym momencie chce
|
|
|
+ * wstrzymać wykonanie w oczekiwaniu na zdarzenie bardziej szczegółowe.
|
|
|
+ * </p>
|
|
|
+ * @param <T>
|
|
|
+ * @param event
|
|
|
+ * @param timeout
|
|
|
+ * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
|
|
|
+ * @throws InterruptedException
|
|
|
+ */
|
|
|
+ public abstract <T extends E> T waitForEvent(Class<T> event, long timeout) throws InterruptedException;
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda tworzy silny {@code EventLock} - szczegóły "silnego kontraktu"
|
|
|
+ * są w opisie klasy {@link EventLock}).
|
|
|
+ * <p>
|
|
|
+ * Jeśli obiekt nie jest już potrzebny, musi zostać zwolniony za pomocą metody {@link #release()}
|
|
|
+ * </p>
|
|
|
+ * @param event
|
|
|
+ * @param router
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static <EE extends Event> EventLock<EE> lock(Class<EE> event, EventRouter router) {
|
|
|
+ return lock(event, router, Integer.MAX_VALUE);
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda tworzy słaby {@code EventLock} - szczegóły "słabego kontraktu"
|
|
|
+ * są w opisie klasy {@link EventLock}).
|
|
|
+ * <p>
|
|
|
+ * Jeśli obiekt nie jest już potrzebny, musi zostać zwolniony za pomocą metody {@link #release()}
|
|
|
+ * </p>
|
|
|
+ * @param event
|
|
|
+ * @param router
|
|
|
+ * @param capacity
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static <EE extends Event> EventLock<EE> lock(Class<EE> event, EventRouter router, int capacity) {
|
|
|
+ QueEventLock<EE> lock = new QueEventLock<>(event, router, capacity);
|
|
|
+ router.addEventListener(event, lock);
|
|
|
+ return lock;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda tworzy silny {@code EventLock}, który może zostać wykorzystany tylko raz
|
|
|
+ * - szczegóły "silnego kontraktu" są w opisie klasy {@link EventLock}).
|
|
|
+ * <p>
|
|
|
+ * Po zakończeniu metody {@code waitForEvent} obiekt sam automatycznie
|
|
|
+ * zwalnia wszystkie zasoby i przestaje być funcjonalny.
|
|
|
+ * </p>
|
|
|
+ * @param event
|
|
|
+ * @param router
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static <EE extends Event> EventLock<EE> singleLock(Class<EE> event, EventRouter router) {
|
|
|
+ return singleLock(event, router, Integer.MAX_VALUE);
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda tworzy słaby {@code EventLock}, który może zostać wykorzystany tylko raz
|
|
|
+ * - szczegóły "słabego kontraktu" są w opisie klasy {@link EventLock}).
|
|
|
+ * <p>
|
|
|
+ * Po zakończeniu metody {@code waitForEvent} obiekt sam automatycznie
|
|
|
+ * zwalnia wszystkie zasoby i przestaje być funcjonalny.
|
|
|
+ * </p>
|
|
|
+ * @param event
|
|
|
+ * @param router
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ public static <EE extends Event> EventLock<EE> singleLock(Class<EE> event, EventRouter router, int capacity) {
|
|
|
+ SingleLock<EE> lock = new SingleLock<>(event, router, capacity);
|
|
|
+ router.addEventListener(event, lock);
|
|
|
+ return lock;
|
|
|
+ }
|
|
|
+
|
|
|
+ /**
|
|
|
+ * Metoda tworzy niebezpieczny {@code EventLock} - szczegóły "niebezpiecznego kontraktu"
|
|
|
+ * są w opisie klasy {@link EventLock}).
|
|
|
+ * <p>
|
|
|
+ * Jeśli obiekt nie jest już potrzebny, musi zostać zwolniony za pomocą metody {@link #release()}
|
|
|
+ * </p>
|
|
|
+ * @param event
|
|
|
+ * @param router
|
|
|
+ * @return
|
|
|
+ */
|
|
|
+ @Meta.Unsafe
|
|
|
+ public static <EE extends Event> EventLock<EE> unsafeLock(Class<EE> event, EventRouter router) {
|
|
|
+ UnsafeEventLock<EE> lock = new UnsafeEventLock<>(event, router);
|
|
|
+ router.addEventListener(event, lock);
|
|
|
+ return lock;
|
|
|
+ }
|
|
|
+
|
|
|
+ private static class UnsafeEventLock<EE extends Event> extends EventLock<EE> implements EventListener<EE> {
|
|
|
+
|
|
|
+ private final Class<? extends EE> observed;
|
|
|
+ private final EventRouter router;
|
|
|
+ private volatile EE hit;
|
|
|
+
|
|
|
+ public UnsafeEventLock(Class<? extends EE> event, EventRouter router) {
|
|
|
+ super();
|
|
|
+ this.observed = event;
|
|
|
+ this.router = router;
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void handleEvent(EE event) {
|
|
|
+ synchronized (this) {
|
|
|
+ hit = event;
|
|
|
+ this.notifyAll();
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public EE waitForEvent() throws InterruptedException {
|
|
|
+ synchronized (this) {
|
|
|
+ while (hit==null) { this.wait(); }
|
|
|
+ EE result = hit;
|
|
|
+ hit = null;
|
|
|
+ return result;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public EE waitForEvent(long timeout) throws InterruptedException {
|
|
|
+ long end = System.currentTimeMillis() + timeout;
|
|
|
+ long rem = timeout;
|
|
|
+ synchronized (this) {
|
|
|
+ while (hit==null && rem>0) {
|
|
|
+ this.wait(rem);
|
|
|
+ rem = end - System.currentTimeMillis();
|
|
|
+ }
|
|
|
+ EE result = hit;
|
|
|
+ hit = null;
|
|
|
+ return result;
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
|
|
|
+ synchronized (this) {
|
|
|
+ while( !event.isInstance(hit) ) { this.wait(); }
|
|
|
+ Event ret = hit;
|
|
|
+ hit = null;
|
|
|
+ return event.cast(ret);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
|
|
|
+ long end = System.currentTimeMillis() + timeout;
|
|
|
+ long rem = timeout;
|
|
|
+ synchronized (this) {
|
|
|
+ while( !event.isInstance(hit) && rem>0) {
|
|
|
+ this.wait(rem);
|
|
|
+ rem = end - System.currentTimeMillis();
|
|
|
+ }
|
|
|
+ EE ret = hit;
|
|
|
+ hit = null;
|
|
|
+ return event.cast(ret);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void reset() {
|
|
|
+ hit = null;
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void release() {
|
|
|
+ router.removeEventListener(observed, this);
|
|
|
+ }
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+ private static class QueEventLock<EE extends Event> extends EventLock<EE> implements EventListener<EE> {
|
|
|
+
|
|
|
+ private static final Logger LOGGER = LoggerUtils.getLogger();
|
|
|
+
|
|
|
+ private final Class<? extends EE> observed;
|
|
|
+ private final EventRouter router;
|
|
|
+ private final BlockingQueue<EE> events;
|
|
|
+ private int discarded;
|
|
|
+ private int processed;
|
|
|
+
|
|
|
+ public QueEventLock(Class<? extends EE> event, EventRouter router, int capacity) {
|
|
|
+ this.observed = event;
|
|
|
+ this.router = router;
|
|
|
+ this.processed = 0;
|
|
|
+ this.discarded = 0;
|
|
|
+ this.events = new LinkedBlockingQueue<>(capacity);
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public EE waitForEvent() throws InterruptedException {
|
|
|
+ return events.take();
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public EE waitForEvent(long timeout) throws InterruptedException {
|
|
|
+ return events.poll(timeout, TimeUnit.MILLISECONDS);
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
|
|
|
+ while(true) {
|
|
|
+ EE recent = events.take();
|
|
|
+ if( event.isInstance(recent) ) {
|
|
|
+ return event.cast(recent);
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
|
|
|
+ long end = System.currentTimeMillis() + timeout;
|
|
|
+ long rem = timeout;
|
|
|
+ while(rem > 0) {
|
|
|
+ EE recent = events.poll(rem, TimeUnit.MILLISECONDS);
|
|
|
+ if( event.isInstance(recent) ) {
|
|
|
+ return event.cast(recent);
|
|
|
+ }
|
|
|
+ rem = end - System.currentTimeMillis();
|
|
|
+ }
|
|
|
+ return null;
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void release() {
|
|
|
+ LOGGER.trace("overall processed={}", processed);
|
|
|
+ if( discarded > 0 ) {
|
|
|
+ LOGGER.warn("overall discarded={}",discarded);
|
|
|
+ }
|
|
|
+ router.removeEventListener(observed, this);
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void reset() {
|
|
|
+ events.clear();
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public void handleEvent(EE event) {
|
|
|
+ processed++;
|
|
|
+ while( !events.offer(event) ) {
|
|
|
+ Event ignored = events.poll();
|
|
|
+ discarded++;
|
|
|
+ if( 0 == (processed % LOG_BATCH) ) {
|
|
|
+ LOGGER.warn("events flood. processed={} discarded={} event={}", processed, discarded, ignored);
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public int discarded() {
|
|
|
+ return discarded;
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public int counter() {
|
|
|
+ return processed;
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public int pending() {
|
|
|
+ return events.size();
|
|
|
+ }
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+ private static class SingleLock<EE extends Event> extends QueEventLock<EE> {
|
|
|
+
|
|
|
+ public <T extends EE> SingleLock(Class<T> event, EventRouter router, int capacity) {
|
|
|
+ super(event, router, capacity);
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public EE waitForEvent() throws InterruptedException {
|
|
|
+ try {
|
|
|
+ return super.waitForEvent();
|
|
|
+ } finally {
|
|
|
+ release();
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
|
|
|
+ try {
|
|
|
+ return super.waitForEvent(event);
|
|
|
+ } finally {
|
|
|
+ release();
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
|
|
|
+ try {
|
|
|
+ return super.waitForEvent(event, timeout);
|
|
|
+ } finally {
|
|
|
+ release();
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ @Override
|
|
|
+ public EE waitForEvent(long timeout) throws InterruptedException {
|
|
|
+ try {
|
|
|
+ return super.waitForEvent(timeout);
|
|
|
+ } finally {
|
|
|
+ release();
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+
|
|
|
+ }
|
|
|
+
|
|
|
+}
|