Ranides Atterwim 3 years ago
parent
commit
79d72f3db1

+ 2 - 2
assira.core/src/main/java/net/ranides/assira/concurrent/TaskBuilder.java

@@ -292,7 +292,7 @@ public class TaskBuilder {
             if(router != null) {
             if(router != null) {
                 throw new IllegalStateException("Both router and router name is specified in invoker");
                 throw new IllegalStateException("Both router and router name is specified in invoker");
             }
             }
-            router = EventDispatcher.find(routerName).orElse(EventDispatcher.NULL);
+            router = EventDispatcher.find(routerName).orElse(null);
         }
         }
     }
     }
 
 
@@ -303,7 +303,7 @@ public class TaskBuilder {
         if(action != null) {
         if(action != null) {
             throw new IllegalStateException("Both action and event supplier is specified in invoker");
             throw new IllegalStateException("Both action and event supplier is specified in invoker");
         }
         }
-        if(router == EventDispatcher.NULL) {
+        if(router == null) {
             throw new IllegalStateException("Event supplier is defined, but event dispatcher not");
             throw new IllegalStateException("Event supplier is defined, but event dispatcher not");
         }
         }
         return (v) -> router.signalEvent( event.apply(v) );
         return (v) -> router.signalEvent( event.apply(v) );

+ 11 - 389
assira.core/src/main/java/net/ranides/assira/events/EventDispatcher.java

@@ -6,12 +6,6 @@
  */
  */
 package net.ranides.assira.events;
 package net.ranides.assira.events;
 
 
-import java.lang.ref.WeakReference;
-import java.util.*;
-import java.util.function.Supplier;
-
-import net.ranides.assira.collection.maps.OpenMap;
-import net.ranides.assira.text.StringTraits;
 import net.ranides.assira.trace.LoggerUtils;
 import net.ranides.assira.trace.LoggerUtils;
 import net.ranides.assira.trace.ThreadDump;
 import net.ranides.assira.trace.ThreadDump;
 import org.slf4j.Logger;
 import org.slf4j.Logger;
@@ -31,44 +25,16 @@ import org.slf4j.Logger;
  *
  *
  * @author ranides
  * @author ranides
  */
  */
-public class EventDispatcher implements EventRouter {
-
-    /**
-     * Dispatcher which ignores all events and all listeners.
-     *
-     * Please note, that adding listener won't cause any errors, but won't be reflects by "getEventListeners"
-     */
-    public static final EventDispatcher NULL = new EmptyDispatcher();
+public class EventDispatcher extends EventRouter {
 
 
     private static final Logger LOGGER = LoggerUtils.getLogger();
     private static final Logger LOGGER = LoggerUtils.getLogger();
-    
-    private static final Map<String, WeakReference<EventRouter>> DISPATCHERS = new OpenMap<>();
-
-    private Object[] listeners = null;
-    
-    private final String name;
 
 
-    /**
-     * Creates anonymous instant dispatcher.
-     * Anonymous dispatchers can't be find by using {@link EventDispatcher#find(String)}.
-     *
-     * Subclasses should use this constructor only inside protected or private constructors.
-     * Although anonymous dispatcher could be created directly, it is preferred not to do so.
-     */
-    protected EventDispatcher() {
-        this.name = "";
+    private EventDispatcher() {
+        super();
     }
     }
 
 
-    /**
-     * Creates named instant dispatcher.
-     * Subclasses should use this constructor only inside protected or private constructors.
-     *
-     * Correctly created dispatcher is always created by static method {@link #newInstance(String, Supplier)}.
-     *
-     * @param name name
-     */
-    protected EventDispatcher(String name) {
-        this.name = name;
+    private EventDispatcher(String name) {
+        super(name);
     }
     }
 
 
     /**
     /**
@@ -91,368 +57,24 @@ public class EventDispatcher implements EventRouter {
      * @return dispatcher
      * @return dispatcher
      */
      */
     public static EventDispatcher newInstance(String name) {
     public static EventDispatcher newInstance(String name) {
-        return newInstance(name, () -> new EventDispatcher(name));
-    }
-
-    /**
-     * Creates new dispatcher with provided name, using factory method.
-     *
-     * Registers provided dispatcher if it is not anonymous.
-     * Registered named dispatchers can be find by using {@link EventDispatcher#find(String)}.
-     *
-     * If there was already dispatcher with specified name,
-     * it won't create new one but it will throw IllegalArgumentException.
-     *
-     * WARNING:
-     * all dispatchers MUST BE constructed using this method.
-     * Especially subclasses MUST NOT construct objects without using it.
-     * This is the only method in whole codebase, which is allowed to create dispatchers.
-     *
-     * @param name name
-     * @param factory factory
-     * @param <R> dispatcher type
-     * @return dispatcher
-     */
-    public static <R extends EventRouter> R newInstance(String name, Supplier<R> factory) {
-        if(StringTraits.isEmpty(name)) {
-            return factory.get();
-        }
-        synchronized(DISPATCHERS) {
-            WeakReference<EventRouter> ref = DISPATCHERS.get(name);
-            if (ref != null && ref.get() != null) {
-                throw new IllegalArgumentException("EventDispatcher with name " + name + " already exists");
-            }
-            R that = factory.get();
-            DISPATCHERS.put(name, new WeakReference<>(that));
-            return that;
-        }
-    }
-
-    /**
-     * Finds named EventDispatcher.
-     * If there is no dispatcher with specified name, it returns none.
-     *
-     * @param name name
-     * @return optional
-     */
-    public static Optional<EventRouter> find(String name) {
-        synchronized (DISPATCHERS) {
-            WeakReference<EventRouter> ref = DISPATCHERS.get(name);
-            if(ref != null) {
-                return Optional.ofNullable(ref.get());
-            }
-            return Optional.empty();
-        }
+        return EventRouter.newInstance(name, () -> new EventDispatcher(name));
     }
     }
 
 
-    @Override
-    public String name() {
-        return name;
-    }
-
-    @Override
-    public <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener) {
-        // atomic operation, signal update
-        synchronized(this) {
-            if(null == listeners) {
-                listeners = new Object[]{event, listener};
-            } else {
-                int size = listeners.length;
-                Object[] realloc = new Object[size+2];
-                System.arraycopy(listeners, 0, realloc, 0, size);
-                realloc[size] = event;
-                realloc[size+1] = listener;
-                listeners = realloc;
-            }
-        }
-    }
-
-    @Override
-    public <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener) {
-        // atomic operation, signal update
-        synchronized(this) {
-            if(null == listeners) {
-                return;
-            }
-            int index = -1;
-            for (int i = listeners.length-2; i>=0; i-=2) {
-                if ((listeners[i]==event) && listeners[i+1].equals(listener) ) {
-                    index = i;
-                    break;
-                }
-            }
-
-            if (index != -1) {
-                if(listeners.length==2) {
-                    listeners = null;
-                } else {
-                    int size = listeners.length-2;
-                    Object[] realloc = new Object[size];
-                    System.arraycopy(listeners, 0, realloc, 0, index);
-                    if (index < size) {
-                        System.arraycopy(listeners, index+2, realloc, index, size-index);
-                    }
-                    listeners = realloc;
-                }
-            }
-        }
-    }
-
-    @Override
-    public <T extends Event> void removeEventListener(Class<T> event) {
-        // atomic operation, signal update
-        synchronized(this) {
-            if(null == listeners) {
-                return;
-            }
-            int found = 0;
-            for (int i = 0; i <= listeners.length; i += 2) {
-                if (listeners[i]==event) {
-                    found++;
-                }
-            }
-            if(found == listeners.length/2) {
-                listeners = null;
-            } else {
-                Object[] realloc = new Object[listeners.length - 2*found];
-
-                int index = 0;
-                for (int i = 0; i <= listeners.length; i += 2) {
-                    if (listeners[i]!=event) {
-                        realloc[index] = listeners[i];
-                        realloc[index+1] = listeners[i+1];
-                        index++;
-                    }
-                }
-                listeners = realloc;
-            }
-        }
-    }
-
-    @Override
-    public void removeAllEventListeners() {
-        // atomic operation, signal update
-        synchronized(this) {
-            listeners = null;
-        }
-    }
-
-    @Override
-    @SuppressWarnings("unchecked")
-    public Collection<EventBinding<?>> getEventListeners() {
-        Object[] array;
-        // read updated
-        synchronized(this) {
-            array = listeners;
-        }
-        // read update
-        if(null == array) {
-            return Collections.emptyList();
-        }
-        List<EventBinding<?>> list = new ArrayList<>(array.length/2);
-        for(int i=0, n=array.length; i<n; i+=2) {
-            list.add(new EventBinding.Immutable<>((Class)array[i], (EventListener)array[i+1]));
-        }
-        return list;
-    }
-    
-    @Override
-    public int getEventListenersCount() {
-        synchronized(this) {
-            return null == listeners ? 0 : listeners.length / 2;
-        }
-    }
-
-    /**
-     * Metoda blokująca, uruchamia kolejno obsługę zdarzenia u wszystkich
-     * zarejestrowanych obserwatorów, które na zdarzenie o określonym typie oczekują.
-     * Domyślna implementacja wywołuje tę metodę bezpośrednio w {@link #signalEvent}.
-     * Metody pochodne mogą użyć poniższej metody w zupełnie inny sposób i w
-     * zupełnie innym wątku, aby zaimplementować inną strategię przetwarzania zdarzeń.
-     *
-     * <div class="message-note">thread-safe method</div>
-     *
-     * @param event event
-     *
-     * @todo #85 and then document
-     */
-    protected void dispatchEvent(Event event) {
-        dispatchEvent(event, false);
-    }
-
-    @SuppressWarnings("unchecked")
-    protected void dispatchEvent(Event event, boolean direct) {
-        Object[] array;
-        // read updated
-        synchronized(this) {
-            array = this.listeners;
-        }
-        if(array == null) {
-            return;
-        }
-        Class<?> clazz = event.getClass();
-        for(int i=0, n=array.length; i<n; i+=2) {
-            if( ((Class)array[i]).isAssignableFrom(clazz) ) {
-                try {
-                    EventListener listener = (EventListener) array[i + 1];
-                    if(direct) {
-                        listener.handleEvent(event);
-                    } else {
-                        dispatchEvent(listener, event);
-                    }
-                } catch(Exception ex) {
-					String message = String.format("Exception thrown from event listener. Thread: %s. Listener: %s. Event: %s", 
-						Thread.currentThread().getName(), 
-						array[i+1], event
-					);
-					LOGGER.debug(message, ex);
-                    signalEvent(Events.failure(ex));
-                }
-            }
-        }
-    }
-
-    protected <T extends Event> void dispatchEvent(EventListener<? super T> listener, T event) {
-        listener.handleEvent(event);
-    }
-
-    /**
-     * {@inheritDoc}
-     * <p>
-     * Metoda czeka na obsłużenie komunikatu przez wszystkich zarejestrowanych obserwatorów.
-     * Obsługa komunikatu jest uruchamiana natychmiastowo, w bieżącym wątku
-     * (t.j. w tym samym, który zasygnalizował zdarzenie).
-     * </p>
-     * <p>Interesującym przypadkiem jest przekazanie Dispatcher'owi obserwatora,
-     * który jest {@code EventReactor}-em. Taki obserwator w bieżącym wątku
-     * obsłuży zdarzenie <i>w swoim imieniu</i>, ale ta obsługa, to nic więcej, niż zapis do
-     * wewnętrznej kolejki. Dokładniej, zgodnie z dokumentacją:</p>
-     *  <ul>
-     *      <li>ruszy metoda {@link EventReactor#handleEvent}</li>
-     *      <li>metoda wydeleguje obsługę do {@link EventReactor#signalEvent}</li>
-     *      <li>{@code signalEvent} wstawi komunikat do kolejki zdarzeń oczekujących</li>
-     *      <li>wstawiony komunikat zostanie przekazany do obserwatorów w zupełnie
-     *      innym wątku</li>
-     *      <li>w szczególności, czas obsługi jest całkowicie niezdefiniowany.
-     *          Może nastąpić przed zakończeniem metody signal z bieżącego wątku.
-     *          A może być odwrotnie - nastapić długo po jej zakończeniu</li>
-     *  </ul>
-     *
-     * <p>Podobny przypadek to {@code EventProactor}. Kolejność wywołań jest następująca:</p>
-     *  <ul>
-     *      <li>ruszy metoda {@link EventProactor#handleEvent}</li>
-     *      <li>metoda wydeleguje obsługę do {@link EventProactor#signalEvent}</li>
-     *      <li>metoda wydeleguje obsługę do {@link EventReactor#dispatchEvent}</li>
-     *      <li>{@code dispatchEvent} przekaże komunikat do ThreadPoolExecutor, który wywoła obsługę
-     *      w sobie tylko znany sposób.</li>
-     *      <li>Nie tylko czas obsługi jest całkowicie niezdefiniowany, ale dodatkowo odbiorcy mogą być uruchomieni
-     *      jednocześnie w kilku wątkach.</li>
-     *  </ul>
-     *
-     * <div class="message-note">thread-safe method</div>
-     * @param event
-     *
-     * @todo #85 and then document
-     */
     @Override
     @Override
     public boolean signalEvent(Event event) {
     public boolean signalEvent(Event event) {
 		ThreadDump.printIf(LOGGER.isTraceEnabled(), LOGGER::trace);
 		ThreadDump.printIf(LOGGER.isTraceEnabled(), LOGGER::trace);
-        dispatchEvent(event);
+        executeListenersNow(event);
         return true;
         return true;
     }
     }
 
 
-    /**
-     * Rozsyła otrzymany komunikat do wszystkich obserwatorów, delegując
-     * obsługę zdarzenia do metody {@link #signalEvent}
-     * <div class="message-note">thread-safe method</div>
-     * @param event event
-     */
     @Override
     @Override
-    public void handleEvent(Event event) {
-        signalEvent(event);
-    }
-
-    protected final void dispose(boolean direct) {
-        dispatchEvent(Events.dispose(this), direct);
-        // signal update
-        synchronized(this) {
-            this.listeners = null;
-            synchronized (DISPATCHERS) {
-                DISPATCHERS.remove(name());
-            }
-        }
+    public void dispose() {
+        executeDisposeNow();
     }
     }
 
 
     @Override
     @Override
-    public void dispose() {
-        dispose(true);
+    public String toString() {
+        return "EventDispatcher<" + name() + ">";
     }
     }
 
 
-    private static class EmptyDispatcher extends EventDispatcher {
-        @Override
-        public String name() {
-            return "$null";
-        }
-
-        @Override
-        public <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener) {
-            // do nothing
-        }
-
-        @Override
-        public <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener) {
-            // do nothing
-        }
-
-        @Override
-        public <T extends Event> void removeEventListener(Class<T> event) {
-            // do nothing
-        }
-
-        @Override
-        public void removeAllEventListeners() {
-            // do nothing
-        }
-
-        @Override
-        public Collection<EventBinding<?>> getEventListeners() {
-            return Collections.emptyList();
-        }
-
-        @Override
-        public int getEventListenersCount() {
-            return 0;
-        }
-
-        @Override
-        public boolean signalEvent(Event event) {
-            return false;
-        }
-
-        @Override
-        public void dispose() {
-            // do nothing
-        }
-
-        @Override
-        public void handleEvent(Event event) {
-            // do nothing
-        }
-
-        @Override
-        protected void dispatchEvent(Event event) {
-            // do nothing
-        }
-
-        @Override
-        protected void dispatchEvent(Event event, boolean direct) {
-            // do nothing
-        }
-
-        @Override
-        protected <T extends Event> void dispatchEvent(EventListener<? super T> listener, T event) {
-            // do nothing
-        }
-
-    }
 }
 }

+ 10 - 7
assira.core/src/main/java/net/ranides/assira/events/EventProactor.java

@@ -11,6 +11,7 @@ import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicInteger;
 import net.ranides.assira.trace.LoggerUtils;
 import net.ranides.assira.trace.LoggerUtils;
+import net.ranides.assira.trace.ThreadDump;
 import org.slf4j.Logger;
 import org.slf4j.Logger;
 
 
 /**
 /**
@@ -24,7 +25,7 @@ import org.slf4j.Logger;
  * <div class="message-note">thread-safe method</div>
  * <div class="message-note">thread-safe method</div>
  * @author ranides
  * @author ranides
  */
  */
-public class EventProactor extends EventDispatcher {
+public class EventProactor extends EventRouter {
 
 
     private static final Logger LOGGER = LoggerUtils.getLogger();
     private static final Logger LOGGER = LoggerUtils.getLogger();
 
 
@@ -109,8 +110,8 @@ public class EventProactor extends EventDispatcher {
             executor.shutdown();
             executor.shutdown();
         }
         }
         this.join();
         this.join();
-        dispatchEvent(Events.shutdown(EventProactor.this), true);
-        dispose(true);
+        executeListeners(Events.shutdown(EventProactor.this), EventListener::handleEvent);
+        executeDisposeNow();
         return this;
         return this;
     }
     }
 
 
@@ -126,8 +127,10 @@ public class EventProactor extends EventDispatcher {
     }
     }
 
 
     @Override
     @Override
-    protected <T extends Event> void dispatchEvent(final EventListener<? super T> listener, final T event) {
-        executor.execute(()->listener.handleEvent(event));
+    public boolean signalEvent(Event event) {
+        ThreadDump.printIf(LOGGER.isTraceEnabled(), LOGGER::trace);
+        executeListeners(event, (listener, e) -> executor.execute(()->listener.handleEvent(e)));
+        return true;
     }
     }
 
 
     @Override
     @Override
@@ -143,8 +146,8 @@ public class EventProactor extends EventDispatcher {
 
 
         new Thread(() -> {
         new Thread(() -> {
             EventProactor.this.join();
             EventProactor.this.join();
-            dispatchEvent(Events.shutdown(EventProactor.this), true);
-            dispose(true);
+            executeListeners(Events.shutdown(EventProactor.this), EventListener::handleEvent);
+            executeDisposeNow();
         }).start();
         }).start();
     }
     }
 
 

+ 10 - 12
assira.core/src/main/java/net/ranides/assira/events/EventReactor.java

@@ -23,7 +23,7 @@ import org.slf4j.Logger;
  * <div class="message-note">thread-safe method</div>
  * <div class="message-note">thread-safe method</div>
  * @author ranides
  * @author ranides
  */
  */
-public class EventReactor extends EventDispatcher {
+public class EventReactor extends EventRouter {
     
     
     private static final Logger LOGGER = LoggerUtils.getLogger();
     private static final Logger LOGGER = LoggerUtils.getLogger();
 
 
@@ -112,23 +112,21 @@ public class EventReactor extends EventDispatcher {
 
 
             @Override
             @Override
             public void run() {
             public void run() {
-                boolean interrupted = true;
                 try {
                 try {
                     while (true) {
                     while (true) {
                         Event event = events.take();
                         Event event = events.take();
                         while(joiner.isJoinable(event, events.peek()) ) {
                         while(joiner.isJoinable(event, events.peek()) ) {
                             event = joiner.join(event, events.poll());
                             event = joiner.join(event, events.poll());
                         }
                         }
-                        dispatchEvent(event);
-                        if( dispatchExit(event) ) { interrupted=false; break; }
+                        executeListenersNow(event);
+                        if (isExitEvent(event)) {
+                            break;
+                        }
                     }
                     }
                 } catch (InterruptedException _e) {
                 } catch (InterruptedException _e) {
-                    dispatchEvent(Events.interrupt(EventReactor.this));
-                    interrupted = false;
-                }
-                if(!interrupted) {
-                    dispatchEvent(Events.shutdown(EventReactor.this));
+                    executeListenersNow(Events.interrupt(EventReactor.this));
                 }
                 }
+                executeListenersNow(Events.shutdown(EventReactor.this));
             }
             }
         };
         };
     }
     }
@@ -163,7 +161,7 @@ public class EventReactor extends EventDispatcher {
     public EventReactor stop() {
     public EventReactor stop() {
         signalEvent(exit);
         signalEvent(exit);
         join();
         join();
-        super.dispose();
+        executeDisposeNow();
         return this;
         return this;
     }
     }
 
 
@@ -181,7 +179,7 @@ public class EventReactor extends EventDispatcher {
     }
     }
 
 
     @SuppressWarnings("PMD.CompareObjectsWithEquals")
     @SuppressWarnings("PMD.CompareObjectsWithEquals")
-    private boolean dispatchExit(Event event) {
+    private boolean isExitEvent(Event event) {
         if(exit == event) { return true; }
         if(exit == event) { return true; }
         return (event instanceof Events.Dispose) && this == ((Events.Dispose)event).router();
         return (event instanceof Events.Dispose) && this == ((Events.Dispose)event).router();
     }
     }
@@ -219,7 +217,7 @@ public class EventReactor extends EventDispatcher {
         signalEvent(exit);
         signalEvent(exit);
         new Thread(() -> {
         new Thread(() -> {
             EventReactor.this.join();
             EventReactor.this.join();
-            EventReactor.super.dispose();
+            EventReactor.this.executeDisposeNow();
         }).start();
         }).start();
     }
     }
 
 

+ 252 - 28
assira.core/src/main/java/net/ranides/assira/events/EventRouter.java

@@ -6,7 +6,15 @@
  */
  */
 package net.ranides.assira.events;
 package net.ranides.assira.events;
 
 
-import java.util.Collection;
+import net.ranides.assira.collection.maps.OpenMap;
+import net.ranides.assira.text.StringTraits;
+import net.ranides.assira.trace.LoggerUtils;
+import org.slf4j.Logger;
+
+import java.lang.ref.WeakReference;
+import java.util.*;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.BiConsumer;
 import java.util.function.Supplier;
 import java.util.function.Supplier;
 
 
 /**
 /**
@@ -25,7 +33,109 @@ import java.util.function.Supplier;
  *
  *
  * @author ranides
  * @author ranides
  */
  */
-public interface EventRouter extends EventListener<Event> {
+public abstract class EventRouter implements EventListener<Event> {
+
+    private static final Logger LOGGER = LoggerUtils.getLogger();
+
+    private static final Map<String, WeakReference<EventRouter>> DISPATCHERS = new OpenMap<>();
+
+    private final AtomicReference<Object[]> listeners = new AtomicReference<>(null);
+
+    private final String name;
+
+    /**
+     * Creates anonymous router.
+     * Anonymous dispatchers can't be find by using {@link EventDispatcher#find(String)}.
+     *
+     * Subclasses should use this constructor only inside protected or private constructors.
+     * Although anonymous routers could be created directly, it is preferred not to do so.
+     */
+    protected EventRouter() {
+        this.name = "";
+    }
+
+    /**
+     * Creates new dispatcher with provided name, using factory method.
+     *
+     * Registers provided dispatcher if it is not anonymous.
+     * Registered named dispatchers can be find by using {@link EventDispatcher#find(String)}.
+     *
+     * If there was already dispatcher with specified name,
+     * it won't create new one but it will throw IllegalArgumentException.
+     *
+     * WARNING:
+     * all dispatchers MUST BE constructed using this method.
+     * Especially subclasses MUST NOT construct objects without using it.
+     * This is the only method in whole codebase, which is allowed to create dispatchers.
+     *
+     * @param name name
+     * @param factory factory
+     * @param <R> dispatcher type
+     * @return dispatcher
+     */
+    protected static <R extends EventRouter> R newInstance(String name, Supplier<R> factory) {
+        if(StringTraits.isEmpty(name)) {
+            return factory.get();
+        }
+        synchronized(DISPATCHERS) {
+            WeakReference<EventRouter> ref = DISPATCHERS.get(name);
+            if (ref != null && ref.get() != null) {
+                throw new IllegalArgumentException("EventDispatcher with name " + name + " already exists");
+            }
+            R that = factory.get();
+            DISPATCHERS.put(name, new WeakReference<>(that));
+            return that;
+        }
+    }
+
+    /**
+     * Finds named EventDispatcher.
+     * If there is no dispatcher with specified name, it returns none.
+     *
+     * @param name name
+     * @return optional
+     */
+    public static Optional<EventRouter> find(String name) {
+        synchronized (DISPATCHERS) {
+            WeakReference<EventRouter> ref = DISPATCHERS.get(name);
+            if(ref != null) {
+                return Optional.ofNullable(ref.get());
+            }
+            return Optional.empty();
+        }
+    }
+
+    /**
+     * Creates named router.
+     * Subclasses should use this constructor only inside protected or private constructors.
+     *
+     * Correctly created dispatcher is always created by static method {@code newInstance(String, Supplier)}.
+     *
+     * @param name name
+     */
+    protected EventRouter(String name) {
+        this.name = name;
+    }
+
+    /**
+     * Sends provided event to all listeners.
+     * Returns false if router was unable to schedule event processing and discared it.
+     *
+     * Returns true if event was scheduled for processing.
+     * It does not mean, that event was successfully processed by listeners.
+     *
+     * @param event event
+     * @return boolean
+     */
+    public abstract boolean signalEvent(Event event);
+
+    /**
+     * Stops this event router and release all resources used by it.
+     * It sends {@link Events.Dispose} event to all listeners.
+     *
+     * <div class="message-note">thread-safe method</div>
+     */
+    public abstract void dispose();
 
 
     /**
     /**
      * Event router can have a name, usable for debug, logging and dependency management.
      * Event router can have a name, usable for debug, logging and dependency management.
@@ -37,8 +147,10 @@ public interface EventRouter extends EventListener<Event> {
      *
      *
      * @return String
      * @return String
      */
      */
-    String name();
-    
+    public final String name() {
+        return name;
+    }
+
     /**
     /**
      * Registers listener which will receive all events of specified types (and descendant types).
      * Registers listener which will receive all events of specified types (and descendant types).
      *
      *
@@ -48,7 +160,22 @@ public interface EventRouter extends EventListener<Event> {
      * @param listener listener
      * @param listener listener
      * @param <T> T
      * @param <T> T
      */
      */
-    <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener);
+    public final <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener) {
+        // atomic write operation
+        synchronized(this) {
+            Object[] array = listeners.get();
+            if(null == array) {
+                listeners.set(new Object[]{event, listener});
+            } else {
+                int size = array.length;
+                Object[] realloc = new Object[size+2];
+                System.arraycopy(array, 0, realloc, 0, size);
+                realloc[size] = event;
+                realloc[size+1] = listener;
+                listeners.set(realloc);
+            }
+        }
+    }
 
 
     /**
     /**
      * Removes listener, which was previously registered for specified type.
      * Removes listener, which was previously registered for specified type.
@@ -61,7 +188,36 @@ public interface EventRouter extends EventListener<Event> {
      * @param listener listener
      * @param listener listener
      * @param <T> T
      * @param <T> T
      */
      */
-    <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener);
+    public final <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener) {
+        // atomic write operation
+        synchronized(this) {
+            Object[] array = listeners.get();
+            if(null == array) {
+                return;
+            }
+            int index = -1;
+            for (int i = array.length-2; i>=0; i-=2) {
+                if ((array[i]==event) && array[i+1].equals(listener) ) {
+                    index = i;
+                    break;
+                }
+            }
+
+            if (index != -1) {
+                if(array.length==2) {
+                    listeners.set(null);
+                } else {
+                    int size = array.length-2;
+                    Object[] realloc = new Object[size];
+                    System.arraycopy(array, 0, realloc, 0, index);
+                    if (index < size) {
+                        System.arraycopy(array, index+2, realloc, index, size-index);
+                    }
+                    listeners.set(realloc);
+                }
+            }
+        }
+    }
 
 
     /**
     /**
      * Removes all listeners, which were previously registered for specified type.
      * Removes all listeners, which were previously registered for specified type.
@@ -71,14 +227,45 @@ public interface EventRouter extends EventListener<Event> {
      * @param event event
      * @param event event
      * @param <T> T
      * @param <T> T
      */
      */
-    <T extends Event> void removeEventListener(Class<T> event);
+    public final <T extends Event> void removeEventListener(Class<T> event) {
+        // atomic write operation
+        synchronized(this) {
+            Object[] array = listeners.get();
+            if(null == array) {
+                return;
+            }
+            int found = 0;
+            for (int i = 0; i <= array.length; i += 2) {
+                if (array[i]==event) {
+                    found++;
+                }
+            }
+            if(found == array.length/2) {
+                listeners.set(null);
+            } else {
+                Object[] realloc = new Object[array.length - 2*found];
+
+                int index = 0;
+                for (int i = 0; i <= array.length; i += 2) {
+                    if (array[i]!=event) {
+                        realloc[index] = array[i];
+                        realloc[index+1] = array[i+1];
+                        index++;
+                    }
+                }
+                listeners.set(realloc);
+            }
+        }
+    }
 
 
     /**
     /**
      * Removes all registered listeners.
      * Removes all registered listeners.
      *
      *
      * <div class="message-note">thread-safe method</div>
      * <div class="message-note">thread-safe method</div>
      */
      */
-    void removeAllEventListeners();
+    public final void removeAllEventListeners() {
+        listeners.set(null);
+    }
 
 
     /**
     /**
      * Returns collection with all registered bindings.
      * Returns collection with all registered bindings.
@@ -88,36 +275,73 @@ public interface EventRouter extends EventListener<Event> {
      *
      *
      * @return collection
      * @return collection
      */
      */
-    Collection<EventBinding<?>> getEventListeners();
+    @SuppressWarnings("unchecked")
+    public final Collection<EventBinding<?>> getEventListeners() {
+        Object[] array = listeners.get();
+
+        if (null == array) {
+            return Collections.emptyList();
+        }
+        List<EventBinding<?>> list = new ArrayList<>(array.length / 2);
+        for (int i = 0, n = array.length; i < n; i += 2) {
+            list.add(new EventBinding.Immutable<>((Class) array[i], (EventListener) array[i + 1]));
+        }
+        return list;
+    }
 
 
     /**
     /**
      * Returns number of registered event listeners.
      * Returns number of registered event listeners.
      *
      *
      * @return int
      * @return int
      */
      */
-    int getEventListenersCount();
-
-    @Override
-    void handleEvent(Event event);
+    public final int getEventListenersCount() {
+        Object[] array = listeners.get();
+        return null == array ? 0 : array.length / 2;
+    }
 
 
     /**
     /**
-     * Sends provided event to all listeners.
-     * Returns false if router was unable to schedule event processing and discared it.
-     *
-     * Returns true if event was scheduled for processing.
-     * It does not mean, that event was successfully processed by listeners.
-     *
+     * Rozsyła otrzymany komunikat do wszystkich obserwatorów, delegując
+     * obsługę zdarzenia do metody {@link #signalEvent}
+     * <div class="message-note">thread-safe method</div>
      * @param event event
      * @param event event
-     * @return boolean
      */
      */
-    boolean signalEvent(Event event);
+    @Override
+    public final void handleEvent(Event event) {
+        signalEvent(event);
+    }
 
 
-    /**
-     * Stops this event router and release all resources used by it.
-     * It sends {@link Events.Dispose} event to all listeners.
-     *
-     * <div class="message-note">thread-safe method</div>
-     */
-    void dispose();
+    protected final void executeDisposeNow() {
+        executeListenersNow(Events.dispose(this));
+        listeners.set(null);
+        synchronized (DISPATCHERS) {
+            DISPATCHERS.remove(name());
+        }
+    }
+
+    protected final void executeListenersNow(Event event) {
+        executeListeners(event, EventListener::handleEvent);
+    }
 
 
+    protected final void executeListeners(Event event, BiConsumer<EventListener<Event>, Event> executor) {
+        Object[] array = listeners.get();
+        if(array == null) {
+            return;
+        }
+        Class<?> clazz = event.getClass();
+        for(int i=0, n=array.length; i<n; i+=2) {
+            if( ((Class)array[i]).isAssignableFrom(clazz) ) {
+                try {
+                    EventListener listener = (EventListener) array[i + 1];
+                    executor.accept(listener, event);
+                } catch(Exception ex) {
+                    String message = String.format("Exception thrown from event listener. Thread: %s. Listener: %s. Event: %s",
+                        Thread.currentThread().getName(),
+                        array[i+1], event
+                    );
+                    LOGGER.debug(message, ex);
+                    signalEvent(Events.failure(ex));
+                }
+            }
+        }
+    }
 }
 }

+ 1 - 1
assira.core/src/test/java/net/ranides/assira/events/EventObserverTest.java

@@ -47,7 +47,7 @@ public class EventObserverTest {
             }
             }
         });
         });
 
 
-        System.out.println("creation2: " + tm.time());
+        System.out.println("creation: " + tm.time());
 
 
         observer.handleEvent(new MyEvent("first"));
         observer.handleEvent(new MyEvent("first"));
         observer.handleEvent(new MyEvent("second"));
         observer.handleEvent(new MyEvent("second"));