Sfoglia il codice sorgente

revert changes in "net.ranides.assira.events"

lesson for the future: don't change working code written 3 years ago
changes didn't fix anything, but introduced deadlocks. 
not investigated too much, just reverted
Ranides Atterwim 11 anni fa
parent
commit
62adf93950

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

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

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

@@ -67,8 +67,7 @@ public class EventProactor extends EventDispatcher {
     }
 
     public EventProactor stop() {
-        // read updated, signal update, in atomic
-        synchronized(executor) {
+        synchronized(this) {
             if( executor.isShutdown()) {
                 return this;
             }
@@ -97,16 +96,14 @@ public class EventProactor extends EventDispatcher {
     }
 
     @Override
-    public void dispose() {
-        // read updated, signal update, in atomic
-        synchronized(executor) {
-            if( executor.isShutdown()) {
-                return;
-            }
-            signalEvent(CoreEvent.stop(this));
-            executor.shutdown();
+    public synchronized void dispose() {
+        if( executor.isShutdown()) {
+            return;
         }
 
+        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 void dispose() {
+    public synchronized 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 {
 
     /**
-     * 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.
+     * 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
      * 
      */
-    @Ignore("manual slow")
+    @Ignore("unstable") // @todo (assira # 4) stabilize that unit test
     @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(250); // 1 pierwszy powinien długo czekać               
+                TimeUtils.sleep(50); // 1                
                 main.signalEvent(new WriteEvent(77));
-                TimeUtils.sleep(250); // 2
+                TimeUtils.sleep(50); // 2
                 main.signalEvent(new ReadEvent(55));
-                TimeUtils.sleep(250); // 3
+                TimeUtils.sleep(50); // 3
                 main.signalEvent(new ReadEvent(57));
-                TimeUtils.sleep(250); // 4
+                TimeUtils.sleep(50); // 4
                 main.signalEvent(new ReadEvent(59));
-                TimeUtils.sleep(250); // 5
+                TimeUtils.sleep(50); // 5
                 main.signalEvent(new Event(){});
-                TimeUtils.sleep(250); // 6
+                TimeUtils.sleep(50); // 6
                 main.signalEvent(new ReadEvent(44));
-                TimeUtils.sleep(250); // 7
+                TimeUtils.sleep(50); // 7
                 main.signalEvent(new WriteEvent(12));
-                TimeUtils.sleep(250); // 8
+                TimeUtils.sleep(50); // 8
                 main.signalEvent(new ReadEvent(77));
                 
                 main.stop();