Parcourir la source

poprawiony EventProactor

EventReactor - obiekt NIE jest dispose'owany automatycznie po zatrzymaniu, dublowało to DispatchEvent'y
Ranides Atterwim il y a 12 ans
Parent
commit
a8a7e50ec6

+ 13 - 15
src/main/java/net/ranides/assira/events/EventProactor.java

@@ -18,14 +18,11 @@ import net.ranides.assira.trace.AdvLogger;
 import net.ranides.assira.trace.LoggerUtils;
 
 /**
- * Kiepski EventProactor.
- * Obsługuje każdy EVENT w osobnym wątku.
- * A ten wątek obsłuje każdego LISTENERA w osobnym wątku.
- *
- * Obie strategie powinno się dać włączyć / wyłączyć.
- * Ale przede wszystkim to trzeba ograniczyć liczbę tworzonych i niszczonych wątków
- * bo to wydajnościowo jest chyba zarżnięciem.
+ * EventProactor, obsługuje każdego listenera i każde zdarzenie asynchronicznie.
+ * Korzysta z puli wątków, żeby nie powodować eksplozji wątków.
  *
+ * Standardowy ThreadPoolExecutor ma w sumie głupią tę strategię zarządzania pulą,
+ * napisać *kiedyś* własną wersję, albo przemyśleć konfigurację.
  *
  * <div class="message-note">thread-safe method</div>
  * @author ranides
@@ -40,28 +37,29 @@ public class EventProactor extends EventDispatcher {
     private final ThreadPoolExecutor executor;
     private final String name;
     private final AtomicInteger threads = new AtomicInteger(0);
-    private final BlockingQueue<Runnable> tasks = new LinkedBlockingQueue<Runnable>();
+    private final BlockingQueue<Runnable> queue = new LinkedBlockingQueue<Runnable>();
     private final RejectedExecutionHandler handler = new RejectedExecutionHandler() {
         @Override
         public void rejectedExecution(Runnable task, ThreadPoolExecutor executor) {
-            LOGGER.fdebug("ignored task: %s", task);
+            LOGGER.fwarn("ignored task: %s", task);
         }
     };
     ThreadFactory factory = new ThreadFactory() {
         @Override
         public Thread newThread(Runnable runnable) {
-            return new Thread(runnable, name + "-" + threads.getAndIncrement());
+            String tid = name + "-" + threads.getAndIncrement();
+            LOGGER.fdebug("new thread: %s", tid);
+            return new Thread(runnable, tid);
         }
     };
 
-
-    public static EventProactor newInstance(String name, int min, int max) {
-        return new EventProactor(name, min, max);
+    public static EventProactor newInstance(String name, int size) {
+        return new EventProactor(name, size);
     }
 
-    protected EventProactor(final String name, final int min, final int max) {
+    protected EventProactor(final String name, final int size) {
         this.name = name;
-        this.executor = new ThreadPoolExecutor(min, max, Integer.MAX_VALUE, TimeUnit.MILLISECONDS, tasks, factory, handler);
+        this.executor = new ThreadPoolExecutor(size, size, 0L, TimeUnit.MILLISECONDS, queue, factory, handler);
     }
 
     public EventProactor start() {

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

@@ -97,7 +97,7 @@ public class EventReactor extends EventDispatcher {
                 }
                 if(!interrupted) {
                     dispatchEvent(CoreEvent.shutdown(EventReactor.this));
-                    dispatchEvent(CoreEvent.dispose(EventReactor.this));
+                    // dispatchEvent(CoreEvent.dispose(EventReactor.this));
                 }
             }
         };
@@ -114,7 +114,7 @@ public class EventReactor extends EventDispatcher {
     }
 
     /**
-     * Wyłącza {@code EventRouter}. Metoda blokująca - czeka, aż EventRouter
+     * Wyłącza {@code EventRouter}. Metoda blokująca - czeka, aż EventReactor
      * poprawnie zwolni wszystkie zasoby.
      * <div class="message-note">thread-safe method</div>
      * @return

+ 50 - 31
src/test/java/net/ranides/assira/events/EventProactorTest.java

@@ -8,14 +8,11 @@ package net.ranides.assira.events;
 
 import java.util.LinkedList;
 import java.util.List;
-import net.ranides.assira.collection.ArrayUtils;
 import net.ranides.assira.collection.SetUtils;
 import net.ranides.assira.collection.map.LazyMap;
-import net.ranides.assira.math.Randomizer;
 import net.ranides.assira.time.TimeUtils;
 import net.ranides.assira.trace.LoggerUtils;
 import static org.junit.Assert.*;
-import org.junit.Ignore;
 import org.junit.Test;
 
 
@@ -31,7 +28,7 @@ public class EventProactorTest {
     }
 
     private static class HitMap extends LazyMap<String, List<Integer>>  {
-        private int counter = 0;
+        private int counter;
 
         @Override
         public List<Integer> apply(String source) {
@@ -48,7 +45,10 @@ public class EventProactorTest {
 
             @Override
             public void handleEvent(T event) {
-                get( name ).add(counter++);
+                // "handleEvent" can be invoked by many threads - have to be thread-safe
+                synchronized(this) {
+                    get( name ).add(counter++);
+                }
             }
         }
 
@@ -57,27 +57,45 @@ public class EventProactorTest {
     @Test
     public void testPropagation() {
         HitMap hits = new HitMap();
-        EventProactor unknown = EventProactor.newInstance("unknown", 1, 2);
-        EventProactor main = EventProactor.newInstance("name", 2, 4);
+        EventProactor unknown = EventProactor.newInstance("unknown", 2);
+        EventProactor main = EventProactor.newInstance("name", 10);
 
         main.addEventListener(Event.class, hits.new HitListener<Event>("global") );
         main.addEventListener(IOEvent.class, hits.new HitListener<IOEvent>("io") );
         main.addEventListener(ReadEvent.class, hits.new HitListener<ReadEvent>("read") );
         main.addEventListener(WriteEvent.class, hits.new HitListener<WriteEvent>("write") );
+        main.addEventListener(FloodEvent.class, hits.new HitListener<FloodEvent>("flood") );
         main.addEventListener(CoreEvent.Shutdown.class, hits.new HitListener<CoreEvent.Shutdown>("shutdown") );
         main.addEventListener(CoreEvent.Stop.class, hits.new HitListener<CoreEvent.Stop>("exit") );
+        main.addEventListener(FloodEvent.class, new EventListener<FloodEvent>(){
+
+            @Override
+            public void handleEvent(FloodEvent event) {
+                throw new UnsupportedOperationException("Not supported yet."); //To change body of generated methods, choose Tools | Templates.
+            }
+
+        });
+                        // we consume thread, but not lock handler
+                TimeUtils.sleep(25);
+
 
         main.signalEvent(new WriteEvent());
         main.signalEvent(new ReadEvent());
         main.signalEvent(new Event(){});
+
+        for(int i=0; i<100; i++) {
+            main.signalEvent(new FloodEvent());
+        }
+
         main.signalEvent(CoreEvent.stop(unknown));
         main.signalEvent(CoreEvent.stop(main) );
 
         main.stop();
 
-        assertEquals(SetUtils.asHashSet("global", "io", "write", "read", "shutdown", "exit"), hits.keySet());
+        assertEquals(SetUtils.asHashSet("global", "io", "write", "read", "shutdown", "exit", "flood"), hits.keySet());
 
-        assertEquals(8, hits.get("global").size());
+        assertEquals(108, hits.get("global").size());
+        assertEquals(100, hits.get("flood").size());
         assertEquals(2, hits.get("io").size());
         assertEquals(1, hits.get("write").size());
         assertEquals(1, hits.get("read").size());
@@ -87,28 +105,26 @@ public class EventProactorTest {
         unknown.dispose();
     }
 
-    @Ignore
-    @Test
-    public void testStop() {
-        HitMap hits = new HitMap();
-        EventReactor unknown = EventReactor.newInstance("unknown", 32, 1000);
-        EventReactor main = EventReactor.newInstance("name", 32, 1000);
-
-        main.addEventListener(Event.class, hits.new HitListener<Event>("global") );
-        main.addEventListener(CoreEvent.Shutdown.class, hits.new HitListener<CoreEvent.Shutdown>("shutdown") );
-        main.addEventListener(CoreEvent.Stop.class, hits.new HitListener<CoreEvent.Stop>("exit") );
-
-        main.signalEvent(CoreEvent.stop(unknown) );
-
-        main.stop();
-
-        assertEquals(SetUtils.asHashSet("global", "exit", "shutdown"), hits.keySet());
-        assertEquals(ArrayUtils.asList(0,2,4,6), hits.get("global"));
-        assertEquals(ArrayUtils.asList(1,3), hits.get("exit"));
-        assertEquals(ArrayUtils.asList(5), hits.get("shutdown"));
-
-        unknown.dispose();
-    }
+//    @Test
+//    public void testStop() {
+//        HitMap hits = new HitMap();
+//        EventReactor unknown = EventReactor.newInstance("unknown", 32, 1000);
+//        EventReactor main = EventReactor.newInstance("name", 32, 1000);
+//
+//        main.addEventListener(Event.class, hits.new HitListener<Event>("global") );
+//        main.addEventListener(CoreEvent.Shutdown.class, hits.new HitListener<CoreEvent.Shutdown>("shutdown") );
+//        main.addEventListener(CoreEvent.Stop.class, hits.new HitListener<CoreEvent.Stop>("exit") );
+//
+//        main.signalEvent(CoreEvent.stop(unknown) );
+//        main.stop();
+//
+//        assertEquals(SetUtils.asHashSet("global", "exit", "shutdown"), hits.keySet());
+//        assertEquals(ArrayUtils.asList(0,2,4,6), hits.get("global"));
+//        assertEquals(ArrayUtils.asList(1,3), hits.get("exit"));
+//        assertEquals(ArrayUtils.asList(5), hits.get("shutdown"));
+//
+//        unknown.dispose();
+//    }
 
     @SuppressWarnings("PMD")
     class IOEvent implements Event { }
@@ -118,4 +134,7 @@ public class EventProactorTest {
 
     @SuppressWarnings("PMD")
     class WriteEvent extends IOEvent { }
+
+    @SuppressWarnings("PMD")
+    class FloodEvent implements Event { }
 }

+ 2 - 2
src/test/java/net/ranides/assira/events/EventReactorTest.java

@@ -71,7 +71,7 @@ public class EventReactorTest {
         main.stop();
 
         assertEquals(SetUtils.asHashSet("global", "io", "write", "read", "exit", "shutdown"), hits.keySet());
-        assertEquals(ArrayUtils.asList(0,3,6,7,9,11,13,14), hits.get("global"));
+        assertEquals(ArrayUtils.asList(0,3,6,7,9,11,13), hits.get("global"));
         assertEquals(ArrayUtils.asList(1,4), hits.get("io"));
         assertEquals(ArrayUtils.asList(2), hits.get("write"));
         assertEquals(ArrayUtils.asList(5), hits.get("read"));
@@ -96,7 +96,7 @@ public class EventReactorTest {
         main.stop();
 
         assertEquals(SetUtils.asHashSet("global", "exit", "shutdown"), hits.keySet());
-        assertEquals(ArrayUtils.asList(0,2,4,6,7), hits.get("global"));
+        assertEquals(ArrayUtils.asList(0,2,4,6), hits.get("global"));
         assertEquals(ArrayUtils.asList(1,3), hits.get("exit"));
         assertEquals(ArrayUtils.asList(5), hits.get("shutdown"));