Ver código fonte

net.ranides.assira.events
- usunięcie method-level-lock, w efekcie better readility szczególnie podczas dziedziczenia
- "FAL": field access lock - trzyma lock tylko na czas zmiany fields, m.in. nie wysyła zdarzeń oraz nie uruchamia listenerów, jeśli trzyma jakiś lock

fix: EventDispatcher # add,remove,getEventListeners - FAL (performance)
fix: EventDispatcher # dispatchEvent - FAL (critical)
fix: EventDispatcher # dispose - FAL (critical)

fix: EventProactor # stop - lock na "executor" zamiast "this" (it's much safer, no conflict with ancestor)
fix: EventProactor # dispose - lock na "executor" zamiast "this" (it's much safer, no conflict with ancestor)

fix: EventReactor # dispose - remove lock

Ranides Atterwim 11 anos atrás
pai
commit
21300ea9f3

+ 70 - 47
src/main/java/net/ranides/assira/events/EventDispatcher.java

@@ -33,61 +33,76 @@ public class EventDispatcher implements EventRouter {
     private Object[] listeners = null;
 
     @Override
-    public synchronized <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener) {
-        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;
+    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 synchronized <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener) {
-        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;
+    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);
+            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;
                 }
-                listeners = realloc;
             }
         }
     }
 
     @Override
-    public synchronized void removeAllEventListeners() {
-        listeners = null;
+    public void removeAllEventListeners() {
+        // atomic operation, signal update
+        synchronized(this) {
+            listeners = null;
+        }
     }
 
     @Override
     @SuppressWarnings("unchecked")
-    public synchronized Collection<EventBinding<?>> getEventListeners() {
-        if(null == listeners) {
+    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<>();
-        for(int i=0, n=listeners.length; i<n; i+=2) {
-            list.add(EventBinding.$make(listeners[i], listeners[i+1]));
+        for(int i=0, n=array.length; i<n; i+=2) {
+            list.add(EventBinding.$make(array[i], array[i+1]));
         }
         return list;
     }
@@ -103,28 +118,33 @@ public class EventDispatcher implements EventRouter {
      * <div class="message-note">thread-safe method</div>
      * @param event
      */
-    protected synchronized void dispatchEvent(Event event) {
+    protected void dispatchEvent(Event event) {
         dispatchEvent(event, false);
     }
 
     @SuppressWarnings("unchecked")
     @edu.umd.cs.findbugs.annotations.SuppressWarnings({"REC_CATCH_EXCEPTION"})
-    protected synchronized void dispatchEvent(Event event, boolean direct) {
-        if(listeners == null) {
+    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=listeners.length; i<n; i+=2) {
-            if( ((Class)listeners[i]).isAssignableFrom(clazz) ) {
+        for(int i=0, n=array.length; i<n; i+=2) {
+            if( ((Class)array[i]).isAssignableFrom(clazz) ) {
                 try {
                     if(direct) {
-                        ((EventListener)listeners[i+1]).handleEvent(event);
+                        ((EventListener)array[i+1]).handleEvent(event);
                     } else {
-                        dispatchEvent((EventListener)listeners[i+1], event);
+                        dispatchEvent((EventListener)array[i+1], event);
                     }
                 } catch(Exception ex) {
                     LOGGER.xdebug("Exception thrown from event listener. Thread: %s. Listener: %s. Event: %s",
-                        ex, Thread.currentThread().getName(), listeners[i+1], event
+                        ex, Thread.currentThread().getName(), array[i+1], event
                     );
                     signalEvent(CoreEvent.failure(this, ex));
                 }
@@ -178,13 +198,16 @@ public class EventDispatcher implements EventRouter {
         signalEvent(event);
     }
 
-    protected synchronized void dispose(boolean direct) {
+    protected void dispose(boolean direct) {
         dispatchEvent(CoreEvent.dispose(this), direct);
-        this.listeners = null;
+        // signal update
+        synchronized(this) {
+            this.listeners = null;
+        }
     }
 
     @Override
-    public synchronized void dispose() {
+    public void dispose() {
         dispose(true);
     }
 

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

@@ -67,7 +67,8 @@ public class EventProactor extends EventDispatcher {
     }
 
     public EventProactor stop() {
-        synchronized(this) {
+        // read updated, signal update, in atomic
+        synchronized(executor) {
             if( executor.isShutdown()) {
                 return this;
             }
@@ -96,14 +97,16 @@ public class EventProactor extends EventDispatcher {
     }
 
     @Override
-    public synchronized void dispose() {
-        if( executor.isShutdown()) {
-            return;
+    public void dispose() {
+        // read updated, signal update, in atomic
+        synchronized(executor) {
+            if( executor.isShutdown()) {
+                return;
+            }
+            signalEvent(CoreEvent.stop(this));
+            executor.shutdown();
         }
 
-        signalEvent(CoreEvent.stop(this));
-        executor.shutdown();
-
         new Thread(){
             @Override
             public void run() {

+ 1 - 1
src/main/java/net/ranides/assira/events/EventReactor.java

@@ -187,7 +187,7 @@ public class EventReactor extends EventDispatcher {
      * <div class="message-note">thread-safe method</div>
      */
     @Override
-    public synchronized void dispose() {
+    public void dispose() {
         signalEvent(exit);
         new Thread(){
             @Override

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

@@ -22,13 +22,13 @@ import org.junit.Test;
 public class EventLockTest {
 
     /**
-     * chciałoby się tu mieć wynik pomiaru czasu na poziomie 400ms
-     * wychodzi ponad 450 ... ;(
-     * uwaga! test się wysypie na zajętym systemie. 
-     * opóźnienie 50ms może być za małe, żeby uniknąć race condition
+     * MUSIMY tu mieć duże opóźnienia, jeśli jednocześnie chcemy sprawdzić
+     * naprawdę EventLock oraz uniknąć race-condition. Moglibyśmy wprowadzić
+     * jakąś synchronizację... ale wtedy z definicji byśmy testowali tę 
+     * smart-synchronizację, a nie nasz EventLock.
      * 
      */
-    @Ignore("unstable") // @todo (assira # 4) stabilize that unit test
+    @Ignore("manual slow")
     @Test
     public void testEventLockRace() throws InterruptedException {
         final EventReactor main = EventReactor.newInstance("name", 32, 1000);
@@ -39,21 +39,21 @@ public class EventLockTest {
             public void run() {
                 // sleep unikają race condition (event wysłany przed rozpoczęciem waitForEvent)
                 
-                TimeUtils.sleep(50); // 1                
+                TimeUtils.sleep(250); // 1 pierwszy powinien długo czekać               
                 main.signalEvent(new WriteEvent(77));
-                TimeUtils.sleep(50); // 2
+                TimeUtils.sleep(250); // 2
                 main.signalEvent(new ReadEvent(55));
-                TimeUtils.sleep(50); // 3
+                TimeUtils.sleep(250); // 3
                 main.signalEvent(new ReadEvent(57));
-                TimeUtils.sleep(50); // 4
+                TimeUtils.sleep(250); // 4
                 main.signalEvent(new ReadEvent(59));
-                TimeUtils.sleep(50); // 5
+                TimeUtils.sleep(250); // 5
                 main.signalEvent(new Event(){});
-                TimeUtils.sleep(50); // 6
+                TimeUtils.sleep(250); // 6
                 main.signalEvent(new ReadEvent(44));
-                TimeUtils.sleep(50); // 7
+                TimeUtils.sleep(250); // 7
                 main.signalEvent(new WriteEvent(12));
-                TimeUtils.sleep(50); // 8
+                TimeUtils.sleep(250); // 8
                 main.signalEvent(new ReadEvent(77));
                 
                 main.stop();