Prechádzať zdrojové kódy

resolve #10

new: ActiveEvent , ActiveEventListener
new: TaskBroker
new: TaskRouter

change: Event & EventDispatcher: use of #dispatch/#process methods

fix: EventLock: wait methods
fix: ObserveLogger: we can control log level in unit tests
Mariusz Sieroń 6 rokov pred
rodič
commit
4dd94c85d1

+ 0 - 62
assira.drafts/src/main/java/net/ranides/assira/prototype/ActiveEvent.java

@@ -1,62 +0,0 @@
-package net.ranides.assira.prototype;
-
-import net.ranides.assira.events.Event;
-
-public abstract class ActiveEvent implements Event {
-
-    private final ActiveObject source;
-
-    private final String property;
-
-    private final Object value;
-
-    public ActiveEvent(ActiveObject source, String property, Object value) {
-        this.source = source;
-        this.property = property;
-        this.value = value;
-    }
-
-    public ActiveObject source() {
-        return source;
-    }
-
-    public String property() {
-        return property;
-    }
-
-    public Object value() {
-        return value;
-    }
-
-    class Create extends ActiveEvent {
-
-        public Create(ActiveObject source, String property, Object value) {
-            super(source, property, value);
-        }
-
-    }
-
-    class Update extends ActiveEvent {
-
-        private final Object previous;
-
-        public Update(ActiveObject source, String property, Object previous, Object value) {
-            super(source, property, value);
-            this.previous = previous;
-        }
-
-        public Object previous() {
-            return previous;
-        }
-
-    }
-
-    class Delete extends ActiveEvent {
-
-        public Delete(ActiveObject source, String property, Object value) {
-            super(source, property, value);
-        }
-
-    }
-
-}

+ 0 - 6
assira.drafts/src/main/java/net/ranides/assira/prototype/ActiveObject.java

@@ -1,6 +0,0 @@
-package net.ranides.assira.prototype;
-
-import net.ranides.assira.events.EventRouter;
-
-public interface ActiveObject extends LightObject, EventRouter {
-}

+ 80 - 0
assira.drafts/src/main/java/net/ranides/assira/prototype/LightEvents.java

@@ -0,0 +1,80 @@
+package net.ranides.assira.prototype;
+
+import lombok.experimental.UtilityClass;
+import net.ranides.assira.events.ActiveEvent;
+
+@UtilityClass
+public class LightEvents {
+
+    public static Create create(LightObject source, String property, Object value) {
+        return new Create(source, property, value);
+    }
+
+    public static Update update(LightObject source, String property, Object value) {
+        return new Update(source, property, value);
+    }
+
+    public static Delete delete(LightObject source, String property) {
+        return new Delete(source, property);
+    }
+
+    public static class LightEvent extends ActiveEvent {
+
+        private final LightObject source;
+
+        private final String property;
+
+        private final Object value;
+
+        LightEvent(LightObject source, String property, Object value) {
+            super();
+            this.source = source;
+            this.property = property;
+            this.value = value;
+        }
+
+        public LightObject source() {
+            return source;
+        }
+
+        public String property() {
+            return property;
+        }
+
+        public Object value() {
+            return value;
+        }
+    }
+
+    public static class Create extends LightEvent {
+
+        Create(LightObject source, String property, Object value) {
+            super(source, property, value);
+        }
+
+    }
+
+    public static class Update extends LightEvent {
+
+        private final Object previous;
+
+        Update(LightObject source, String property, Object value) {
+            super(source, property, value);
+            this.previous = source.get(property);
+        }
+
+        public Object previous() {
+            return previous;
+        }
+
+    }
+
+    public static class Delete extends LightEvent {
+
+        Delete(LightObject source, String property) {
+            super(source, property, source.get(property));
+        }
+
+    }
+
+}

+ 3 - 0
assira.drafts/src/main/java/net/ranides/assira/prototype/LightObject.java

@@ -1,4 +1,7 @@
 package net.ranides.assira.prototype;
 
 public interface LightObject {
+
+    Object get(String name);
+
 }

+ 26 - 5
assira.junit/src/main/java/org/slf4j/impl/ObserveLogger.java

@@ -21,8 +21,29 @@ public final class ObserveLogger extends MarkerIgnoringBase {
 
     private static final long serialVersionUID = 7L;
 
+    private final int level;
+
     ObserveLogger(String name) {
         this.name = name;
+        this.level = levelToInt(System.getProperty("assira.junit.log", "debug"));
+    }
+
+    private static int levelToInt(String text) {
+        switch (text) {
+            case "trace":
+                return TRACE_INT;
+            case "debug":
+                return DEBUG_INT;
+            case "info":
+                return INFO_INT;
+            case "warn":
+            case "warning":
+                return WARN_INT;
+            case "error":
+                return ERROR_INT;
+            default:
+                return DEBUG_INT;
+        }
     }
 
     private void log(int level, String message, Throwable t) {
@@ -60,7 +81,7 @@ public final class ObserveLogger extends MarkerIgnoringBase {
 
     @Override
     public boolean isTraceEnabled() {
-        return true;
+        return level <= TRACE_INT;
     }
 
     @Override
@@ -90,7 +111,7 @@ public final class ObserveLogger extends MarkerIgnoringBase {
 
     @Override
     public boolean isDebugEnabled() {
-        return true;
+        return level <= DEBUG_INT;
     }
 
     @Override
@@ -120,7 +141,7 @@ public final class ObserveLogger extends MarkerIgnoringBase {
 
     @Override
     public boolean isInfoEnabled() {
-        return true;
+        return level <= INFO_INT;
     }
 
     @Override
@@ -150,7 +171,7 @@ public final class ObserveLogger extends MarkerIgnoringBase {
 
     @Override
     public boolean isWarnEnabled() {
-        return true;
+        return level <= WARN_INT;
     }
 
     @Override
@@ -180,7 +201,7 @@ public final class ObserveLogger extends MarkerIgnoringBase {
 
     @Override
     public boolean isErrorEnabled() {
-        return true;
+        return level <= ERROR_INT;
     }
 
     @Override

+ 62 - 0
assira/src/main/java/net/ranides/assira/events/ActiveEvent.java

@@ -0,0 +1,62 @@
+package net.ranides.assira.events;
+
+
+import java.util.concurrent.atomic.AtomicBoolean;
+
+public class ActiveEvent implements Event {
+
+    private final Object lock = new Object();
+
+    private int counter = 0;
+
+    private final AtomicBoolean canceled = new AtomicBoolean(false);
+
+    public ActiveEvent() {
+        // do nothing
+    }
+
+    public boolean cancel() {
+        return canceled.compareAndSet(false, true);
+    }
+
+    @Override
+    public void dispatch(EventListener<?> listener) {
+        if(!(listener instanceof ActiveEventListener)) {
+            return;
+        }
+        synchronized (lock) {
+            counter++;
+            lock.notify();
+        }
+    }
+
+    @Override
+    public void process(EventListener<?> listener) {
+        if(!(listener instanceof ActiveEventListener)) {
+            return;
+        }
+        synchronized (lock) {
+            counter--;
+            lock.notify();
+        }
+    }
+
+    public boolean canceled() {
+        return canceled.get();
+    }
+
+    public boolean processed() {
+        synchronized (lock) {
+            return counter == 0;
+        }
+    }
+
+    public boolean await() throws InterruptedException {
+        synchronized (lock) {
+            while (counter > 0) {
+                lock.wait();
+            }
+        }
+        return !canceled.get();
+    }
+}

+ 18 - 0
assira/src/main/java/net/ranides/assira/events/ActiveEventListener.java

@@ -0,0 +1,18 @@
+package net.ranides.assira.events;
+
+public interface ActiveEventListener<T extends ActiveEvent> extends EventListener<T> {
+
+    @Override
+    default void handleEvent(T event) {
+        try {
+            if (!event.canceled()) {
+                handleActiveEvent(event);
+            }
+        } finally {
+            event.process(this);
+        }
+
+    }
+
+    void handleActiveEvent(T event);
+}

+ 5 - 2
assira/src/main/java/net/ranides/assira/events/Event.java

@@ -12,6 +12,9 @@ package net.ranides.assira.events;
  * @author ranides
  */
 public interface Event {
-	
-	
+
+    default void dispatch(EventListener<?> listener) { }
+
+    default void process(EventListener<?> listener) { }
+
 }

+ 4 - 2
assira/src/main/java/net/ranides/assira/events/EventDispatcher.java

@@ -177,10 +177,12 @@ public class EventDispatcher implements EventRouter {
         for(int i=0, n=array.length; i<n; i+=2) {
             if( ((Class)array[i]).isAssignableFrom(clazz) ) {
                 try {
+                    EventListener listener = (EventListener) array[i + 1];
+                    event.dispatch(listener);
                     if(direct) {
-                        ((EventListener)array[i+1]).handleEvent(event);
+                        listener.handleEvent(event);
                     } else {
-                        dispatchEvent((EventListener)array[i+1], event);
+                        dispatchEvent(listener, event);
                     }
                 } catch(Exception ex) {
 					String message = String.format("Exception thrown from event listener. Thread: %s. Listener: %s. Event: %s", 

+ 63 - 29
assira/src/main/java/net/ranides/assira/events/EventLock.java

@@ -6,9 +6,12 @@
  */
 package net.ranides.assira.events;
 
+import java.util.Optional;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
+import java.util.function.Predicate;
+
 import net.ranides.assira.annotations.Meta;
 import net.ranides.assira.trace.LoggerUtils;
 import org.slf4j.Logger;
@@ -111,10 +114,10 @@ public abstract class EventLock<E extends Event> {
      * minie podany czas.
      * Zobacz opis klasy {@link EventLock} aby poznać szczegóły.
      * @param timeout
-     * @return oczekiwane zdarzenie
+     * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
      * @throws InterruptedException
      */
-    public abstract E waitForEvent(long timeout) throws InterruptedException;
+    public abstract Optional<E> waitForEvent(long timeout) throws InterruptedException;
 
     /**
      * Wersja oczekująca na zdarzenie konkretnego rodzaju. Zachowuje się identycznie
@@ -127,10 +130,12 @@ public abstract class EventLock<E extends Event> {
      * </p>
      * @param <T>
      * @param event
-     * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
+     * @return oczekiwane zdarzenie
      * @throws InterruptedException
      */
-    public abstract <T extends E> T waitForEvent(Class<T> event) throws InterruptedException;
+    public final <T extends E> T waitForEvent(Class<T> event) throws InterruptedException {
+        return event.cast(waitForEvent(event::isInstance));
+    }
 
     /**
      * Wersja oczekująca na zdarzenie konkretnego rodzaju. Zachowuje się identycznie
@@ -147,7 +152,36 @@ public abstract class EventLock<E extends Event> {
      * @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;
+    public final <T extends E> Optional<T> waitForEvent(Class<T> event, long timeout) throws InterruptedException {
+        return waitForEvent(event::isInstance, timeout).map(event::cast);
+    }
+
+    /**
+     * Wersja oczekująca na zdarzenie spełniające warunek. Zachowuje się identycznie
+     * jak {@link #waitForEvent()}, z tą różnicą, że reaguje na mniejszy zakres
+     * zdarzeń - zawężony tylko do tych spełniających warunek.
+     * <p>
+     * Metoda przydatna szczególnie wtedy, gdy chcemy oczekiwać na zdarzenie z konkretnym correlation id.
+     * </p>
+     * @param condition
+     * @return oczekiwane zdarzenie
+     * @throws InterruptedException
+     */
+    public abstract E waitForEvent(Predicate<E> condition) throws InterruptedException;
+
+    /**
+     * Wersja oczekująca na zdarzenie spełniające warunek. Zachowuje się identycznie
+     * jak {@link #waitForEvent(long timeout)}, z tą różnicą, że reaguje na mniejszy zakres
+     * zdarzeń - zawężony tylko do tych spełniających warunek.
+     * <p>
+     * Metoda przydatna szczególnie wtedy, gdy chcemy oczekiwać na zdarzenie z konkretnym correlation id.
+     * </p>
+     * @param condition
+     * @param timeout
+     * @return oczekiwane zdarzenie lub {@code null}, jeśli minął {@code timeout}
+     * @throws InterruptedException
+     */
+    public abstract Optional<E> waitForEvent(Predicate<E> condition, long timeout) throws InterruptedException;
 
 
 
@@ -266,7 +300,7 @@ public abstract class EventLock<E extends Event> {
         }
 
         @Override
-        public EE waitForEvent(long timeout) throws InterruptedException {
+        public Optional<EE> waitForEvent(long timeout) throws InterruptedException {
             long rem = TimeUnit.NANOSECONDS.convert(timeout, TimeUnit.MILLISECONDS);
             long end = System.nanoTime()+ rem;
             synchronized (this) {
@@ -276,32 +310,32 @@ public abstract class EventLock<E extends Event> {
                 }
                 EE result = hit;
                 hit = null;
-                return result;
+                return Optional.ofNullable(result);
             }
         }
 
         @Override
-        public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
+        public EE waitForEvent(Predicate<EE> condition) throws InterruptedException {
             synchronized (this) {
-                while( !event.isInstance(hit) ) { this.wait(); }
-                Event ret = hit;
+                while(hit==null || condition.test(hit) ) { this.wait(); }
+                EE ret = hit;
                 hit = null;
-                return event.cast(ret);
+                return ret;
             }
         }
-        
+
         @Override
-        public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
+        public Optional<EE> waitForEvent(Predicate<EE> condition, long timeout) throws InterruptedException {
             long rem = TimeUnit.NANOSECONDS.convert(timeout, TimeUnit.MILLISECONDS);
             long end = System.nanoTime()+ rem;
             synchronized (this) {
-                while( !event.isInstance(hit) && rem>0) {
+                while ((hit==null || !condition.test(hit)) && rem>0) {
                     this.waitNS(rem);
                     rem = end - System.nanoTime();
                 }
                 EE ret = hit;
                 hit = null;
-                return event.cast(ret);
+                return Optional.ofNullable(ret);
             }
         }
 
@@ -341,32 +375,32 @@ public abstract class EventLock<E extends Event> {
         }
 
         @Override
-        public EE waitForEvent(long timeout) throws InterruptedException {
-            return events.poll(timeout, TimeUnit.MILLISECONDS);
+        public Optional<EE> waitForEvent(long timeout) throws InterruptedException {
+            return Optional.ofNullable(events.poll(timeout, TimeUnit.MILLISECONDS));
         }
 
         @Override
-        public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
+        public EE waitForEvent(Predicate<EE> condition) throws InterruptedException {
             while(true) {
                 EE recent = events.take();
-                if( event.isInstance(recent) ) {
-                    return event.cast(recent);
+                if( condition.test(recent) ) {
+                    return recent;
                 }
             }
         }
 
         @Override
-        public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
+        public Optional<EE> waitForEvent(Predicate<EE> condition, long timeout) throws InterruptedException {
             long rem = TimeUnit.NANOSECONDS.convert(timeout, TimeUnit.MILLISECONDS);
             long end = System.nanoTime()+ rem;
             while(rem > 0) {
                 EE recent = events.poll(rem, TimeUnit.NANOSECONDS);
-                if( event.isInstance(recent) ) {
-                    return event.cast(recent);
+                if( recent !=null && condition.test(recent)) {
+                    return Optional.of(recent);
                 }
                 rem = end - System.nanoTime();
             }
-            return null;
+            return Optional.empty();
         }
 
         @Override
@@ -428,25 +462,25 @@ public abstract class EventLock<E extends Event> {
         }
 
         @Override
-        public <T extends EE> T waitForEvent(Class<T> event) throws InterruptedException {
+        public EE waitForEvent(Predicate<EE> condition) throws InterruptedException {
             try {
-                return super.waitForEvent(event);
+                return super.waitForEvent(condition);
             } finally {
                 release();
             }
         }
 
         @Override
-        public <T extends EE> T waitForEvent(Class<T> event, long timeout) throws InterruptedException {
+        public Optional<EE> waitForEvent(Predicate<EE> condition, long timeout) throws InterruptedException {
             try {
-                return super.waitForEvent(event, timeout);
+                return super.waitForEvent(condition, timeout);
             } finally {
                 release();
             }
         }
 
         @Override
-        public EE waitForEvent(long timeout) throws InterruptedException {
+        public Optional<EE> waitForEvent(long timeout) throws InterruptedException {
             try {
                 return super.waitForEvent(timeout);
             } finally {

+ 6 - 2
assira/src/main/java/net/ranides/assira/events/EventProactor.java

@@ -33,6 +33,10 @@ public class EventProactor extends EventDispatcher {
     private final AtomicInteger threads = new AtomicInteger(0);
     private final BlockingQueue<Runnable> queue = new LinkedBlockingQueue<>();
 
+    public static EventProactor newInstance(String name) {
+        return newInstance(name, Runtime.getRuntime().availableProcessors());
+    }
+
     public static EventProactor newInstance(String name, int size) {
         return new EventProactor(name, size);
     }
@@ -40,8 +44,8 @@ public class EventProactor extends EventDispatcher {
     protected EventProactor(final String name, final int size) {
         super(name);
         this.executor = new ThreadPoolExecutor(
-            size, size, 
-            0L, TimeUnit.MILLISECONDS, 
+            0, size,
+            1L, TimeUnit.MINUTES,
             queue, 
             this::newThread, 
             this::rejectedExecution

+ 36 - 0
assira/src/main/java/net/ranides/assira/events/TaskBroker.java

@@ -0,0 +1,36 @@
+package net.ranides.assira.events;
+
+import net.ranides.assira.functional.VarFunction;
+import net.ranides.assira.reflection.IClass;
+
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+public class TaskBroker<I> {
+
+    private final IClass<I> type;
+
+    private final TaskRouter router;
+
+    private final Map<Method, VarFunction<Object>> reflective;
+
+    public TaskBroker(IClass<I> type, I producer) {
+        this(type, producer, new TaskRouter(EventProactor.newInstance(producer.getClass().getName())));
+    }
+
+    public TaskBroker(IClass<I> type, I producer, TaskRouter router) {
+        this.type = type;
+        this.router = router;
+        this.reflective = type.methods().collect(Collectors.toMap(m -> m.reflective(), m -> m.bind(producer)));
+    }
+
+    public I client() {
+        ClassLoader c = type.raw().getClassLoader();
+        Class[] i = {type.raw()};
+        Object o = Proxy.newProxyInstance(c, i, (p, method, args) -> router.await(() -> reflective.get(method).apply(args) ));
+        return type.cast(o);
+    }
+
+}

+ 129 - 0
assira/src/main/java/net/ranides/assira/events/TaskRouter.java

@@ -0,0 +1,129 @@
+package net.ranides.assira.events;
+
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.function.Consumer;
+import java.util.function.Supplier;
+
+public class TaskRouter {
+
+    private static final Object SUCCESS = new Object();
+
+    private static final AtomicLong COUNTER = new AtomicLong(0);
+
+    private final long tid;
+
+    private final AtomicLong cid = new AtomicLong(0);
+
+    private final EventRouter producer;
+
+    private final EventRouter consumer;
+
+    private final ConcurrentHashMap<Long, Consumer<Object>> callbacks = new ConcurrentHashMap<>();
+
+    public TaskRouter(EventRouter producer) {
+        this(producer, producer);
+    }
+
+    public TaskRouter(EventRouter producer, EventRouter consumer) {
+        this.tid = COUNTER.incrementAndGet();
+        this.producer = producer;
+        this.consumer = consumer;
+
+        this.producer.addEventListener(RunnableEvent.class, event -> {
+            if(event.tid == tid) {
+                event.action.run();
+            }
+        });
+
+        this.producer.addEventListener(SupplierEvent.class, event -> {
+            if(event.tid == tid) {
+                consumer.signalEvent(new SupplierResponse(tid, event.cid, event.function.get()));
+            }
+        });
+
+        this.consumer.addEventListener(SupplierResponse.class, event -> {
+            if(event.tid == tid) {
+                Consumer<Object> callback = callbacks.remove(event.cid);
+                if(callback != null) {
+                    callback.accept(event.value);
+                }
+            }
+        });
+    }
+
+    public void schedule(Runnable action) {
+        producer.signalEvent(new RunnableEvent(tid, action));
+    }
+
+    public void schedule(Runnable action, Runnable callback) {
+        schedule(() -> { action.run(); return SUCCESS; }, (v) -> { callback.run(); });
+    }
+
+    @SuppressWarnings("unchecked")
+    public <R> void schedule(Supplier<R> action, Consumer<R> callback) {
+        long expect = cid.incrementAndGet();
+        callbacks.put(expect, (Consumer)callback);
+        producer.signalEvent(new SupplierEvent(tid, expect, action));
+    }
+
+    public boolean await(Runnable action) throws InterruptedException {
+        return SUCCESS == await(() -> { action.run(); return SUCCESS; });
+    }
+
+    @SuppressWarnings("unchecked")
+    public <R> R await(Supplier<R> action) throws InterruptedException {
+        long expect = cid.incrementAndGet();
+        EventLock<SupplierResponse> lock = EventLock.singleLock(SupplierResponse.class, consumer);
+
+        producer.signalEvent(new SupplierEvent(tid, expect, action));
+
+        return (R)lock.waitForEvent(e -> e.tid == tid && e.cid == expect).value;
+    }
+
+    private static class RunnableEvent implements Event {
+
+        private final long tid;
+
+        private final Runnable action;
+
+        public RunnableEvent(long tid, Runnable action) {
+            this.tid = tid;
+            this.action = action;
+        }
+    }
+
+    private static class SupplierEvent implements Event {
+
+        private final long tid;
+
+        private final long cid;
+
+        private final Supplier<?> function;
+
+        public SupplierEvent(long tid, long cid, Supplier<?> function) {
+            this.tid = tid;
+            this.cid = cid;
+            this.function = function;
+        }
+
+    }
+
+    private static class SupplierResponse implements Event {
+
+        private final long tid;
+
+        private final long cid;
+
+        private final Object value;
+
+        public SupplierResponse(long tid, long cid, Object value) {
+            this.tid = tid;
+            this.cid = cid;
+            this.value = value;
+        }
+
+    }
+
+
+}

+ 12 - 12
assira/src/test/java/net/ranides/assira/events/EventLockTest.java

@@ -83,16 +83,16 @@ public class EventLockTest {
         event = EventLock.singleLock(WriteEvent.class, main).waitForEvent(WriteEvent.class);
         assertEquals(17, event.handle());
         
-        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(300);
+        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(300).get();
         assertEquals(19, event.handle());
         
-        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(200);
+        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(200).orElse(null);
         assertNull(event);
         
-        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(ReadEvent.class, 600);
+        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(ReadEvent.class, 600).get();
         assertEquals(21, event.handle());
         
-        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(ReadEvent.class, 200);
+        event = EventLock.singleLock(IOEvent.class, main).waitForEvent(ReadEvent.class, 200).get();
         assertNull(event);
 
         assertEquals(0, main.getEventListenersCount());
@@ -154,16 +154,16 @@ public class EventLockTest {
         event = lock.waitForEvent(WriteEvent.class);
         assertEquals(17, event.handle());
         
-        event = (IOEvent)lock.waitForEvent(300);
+        event = (IOEvent)lock.waitForEvent(300).get();
         assertEquals(19, event.handle());
         
-        event = (IOEvent)lock.waitForEvent(200);
+        event = (IOEvent)lock.waitForEvent(200).get();
         assertNull(event);
         
-        event = lock.waitForEvent(ReadEvent.class, 600);
+        event = lock.waitForEvent(ReadEvent.class, 600).get();
         assertEquals(21, event.handle());
         
-        event = lock.waitForEvent(ReadEvent.class, 200);
+        event = lock.waitForEvent(ReadEvent.class, 200).get();
         assertNull(event);
 
         assertEquals(1, main.getEventListenersCount());
@@ -224,16 +224,16 @@ public class EventLockTest {
         event = lock.waitForEvent(WriteEvent.class);
         assertEquals(17, event.handle());
         
-        event = (IOEvent)lock.waitForEvent(500);
+        event = (IOEvent)lock.waitForEvent(500).get();
         assertEquals(19, event.handle());
         
-        event = (IOEvent)lock.waitForEvent(200);
+        event = (IOEvent)lock.waitForEvent(200).get();
         assertNull(event);
         
-        event = lock.waitForEvent(ReadEvent.class, 500);
+        event = lock.waitForEvent(ReadEvent.class, 500).get();
         assertEquals(21, event.handle());
         
-        event = lock.waitForEvent(ReadEvent.class, 100);
+        event = lock.waitForEvent(ReadEvent.class, 100).get();
         assertNull(event);
 
         assertEquals(1, main.getEventListenersCount());

+ 66 - 0
assira/src/test/java/net/ranides/assira/events/TaskBrokerTest.java

@@ -0,0 +1,66 @@
+package net.ranides.assira.events;
+
+import net.ranides.assira.reflection.IClass;
+import org.junit.Test;
+
+import static org.junit.Assert.assertEquals;
+
+public class TaskBrokerTest {
+
+    @Test
+    public void broker_reactor() {
+        testBroker(new TaskRouter(EventReactor.newInstance("name")));
+    }
+
+    @Test
+    public void broker_proactor() {
+        testBroker(new TaskRouter(EventProactor.newInstance("name")));
+    }
+
+    @Test
+    public void broker_dispatcher() {
+        testBroker(new TaskRouter(new EventDispatcher()));
+    }
+
+    protected void testBroker(TaskRouter router) {
+        TaskBroker<ICalc> broker = new TaskBroker<>(IClass.typeinfo(ICalc.class), new CCalc(), router);
+
+        ICalc client = broker.client();
+
+        for(int i=0; i<10; i++) {
+            assertEquals(11, client.add(4, 7));
+            assertEquals(2, client.mod(11, 3));
+            assertEquals(8, client.mul(2, 4));
+        }
+    }
+
+    private interface ICalc {
+
+        int add(int a, int b);
+
+        int mod(int a, int b);
+
+        int mul(int a, int b);
+
+    }
+
+    private static class CCalc implements ICalc {
+
+        @Override
+        public int add(int a, int b) {
+            return a + b;
+        }
+
+        @Override
+        public int mod(int a, int b) {
+            return a % b;
+        }
+
+        @Override
+        public int mul(int a, int b) {
+            return a * b;
+        }
+
+    }
+
+}

+ 3 - 6
pom.xml

@@ -15,6 +15,8 @@
         <maven.compiler.target>1.8</maven.compiler.target>
         <netbeans.hint.license>WTFPL</netbeans.hint.license>
         <assira.junit.debug>false</assira.junit.debug>
+        <assira.junit.log>debug</assira.junit.log>
+
 
         <org-netbeans-modules-editor-indent.CodeStyle.project.text-line-wrap>none</org-netbeans-modules-editor-indent.CodeStyle.project.text-line-wrap>
         <org-netbeans-modules-editor-indent.CodeStyle.project.indent-shift-width>4</org-netbeans-modules-editor-indent.CodeStyle.project.indent-shift-width>
@@ -181,6 +183,7 @@
                     <systemPropertyVariables>
                         <assira.junit.debug>${assira.junit.debug}</assira.junit.debug>
                         <assira.version>${project.version}</assira.version>
+                        <assira.junit.log>${assira.junit.log}</assira.junit.log>
                     </systemPropertyVariables>
                     <testFailureIgnore>false</testFailureIgnore>
                     <skip>false</skip>
@@ -250,12 +253,6 @@
                 <artifactId>slf4j-api</artifactId>
                 <version>1.7.13</version>
             </dependency>
-            <dependency>
-                <groupId>org.slf4j</groupId>
-                <artifactId>slf4j-nop</artifactId>
-                <version>1.7.16</version>
-                <scope>compile</scope>
-            </dependency>
             <dependency>
                 <groupId>org.slf4j</groupId>
                 <artifactId>slf4j-simple</artifactId>