소스 검색

new: EventRouter # getEventListenersCount()
new: io.AsyncStream - come back after redesign
new: io.CheckedStream
remove: IORequest - read requests (because they are useless)

Ranides Atterwim 11 년 전
부모
커밋
cb8b3a59a7

+ 1 - 1
pom.xml

@@ -5,7 +5,7 @@
 
     <groupId>net.ranides</groupId>
     <artifactId>assira</artifactId>
-    <version>0.74.0</version>
+    <version>0.75.0</version>
     <packaging>jar</packaging>
 
     <name>assira</name>

+ 4 - 0
src/main/java/net/ranides/assira/config/AbstractConfiguration.java

@@ -138,5 +138,9 @@ public abstract class AbstractConfiguration extends AbstractBranch implements Co
         dispatcher.handleEvent(event);
     }
 
+    @Override
+    public int getEventListenersCount() {
+        return dispatcher.getEventListenersCount();
+    }
 
 }

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

@@ -93,7 +93,12 @@ public class EventDispatcher implements EventRouter {
         }
         return list;
     }
-
+    
+    @Override
+    @SuppressWarnings("unchecked")
+    public synchronized int getEventListenersCount() {
+        return null == listeners ? 0 : listeners.length / 2;
+    }
 
     /**
      * Metoda blokująca, uruchamia kolejno obsługę zdarzenia u wszystkich

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

@@ -52,6 +52,8 @@ public interface EventRouter extends EventListener<Event> {
     void removeAllEventListeners();
 
     Collection<EventBinding<?>> getEventListeners();
+    
+    int getEventListenersCount();
 
     /**
      * Sygnalizuje podane zdarzenie rozsyłając komunikat do wszystkich obserwatorów.

+ 353 - 0
src/main/java/net/ranides/assira/io/AsyncStream.java

@@ -0,0 +1,353 @@
+/*
+ * @author Ranides Atterwim <ranides@gmail.com>
+ * @copyright Ranides Atterwim
+ * @license WTFPL
+ * @url http://ranides.net/projects/assira
+ */
+package net.ranides.assira.io;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.util.Collection;
+import net.ranides.assira.annotations.Meta;
+import net.ranides.assira.events.Event;
+import net.ranides.assira.events.EventBinding;
+import net.ranides.assira.events.EventListener;
+import net.ranides.assira.events.EventReactor;
+import net.ranides.assira.events.EventRouter;
+import net.ranides.assira.events.ReflectEventListener;
+
+/**
+ *
+ * @author Ranides Atterwim <ranides@gmail.com>
+ */
+public final class AsyncStream {
+    
+    private AsyncStream() {
+        // utility class
+    }
+    
+    public static OStream wrap(OutputStream stream) {
+        return new OStream(stream);
+    }
+    
+    public static IStream wrap(InputStream stream) {
+        return new IStream(stream);
+    }
+    
+    public static class IStream extends InputStream implements EventRouter {
+
+        private final InputStream stream;
+        private final EventReactor router;
+
+        public IStream(InputStream stream) {
+            this.stream = stream;
+            this.router = EventReactor.newInstance("AsyncStream.IStream.Router", 0, 0);
+        }
+        
+        @Override
+        public <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener) {
+            router.addEventListener(event, listener);
+        }
+
+        @Override
+        public <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener) {
+            router.removeEventListener(event, listener);
+        }
+
+        @Override
+        public void removeAllEventListeners() {
+            router.removeAllEventListeners();
+        }
+
+        @Override
+        public Collection<EventBinding<?>> getEventListeners() {
+            return router.getEventListeners();
+        }
+
+        @Override
+        public int getEventListenersCount() {
+            return router.getEventListenersCount();
+        }
+
+        @Override
+        public boolean signalEvent(Event event) {
+            return router.signalEvent(event);
+        }
+
+        @Override
+        public void handleEvent(Event event) {
+            router.handleEvent(event);
+        }
+        
+        private boolean sendEvents() {
+            return getEventListenersCount() > 0;
+        }
+        
+        @Override
+        public void dispose() {
+            try {
+                close();
+                if(sendEvents()) {
+                    router.signalEvent(IOEvent.close());
+                }
+            } catch (IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+            } finally {
+                router.dispose();
+            }
+        }
+
+        @Override
+        public int read() throws IOException {
+            try {
+                int ret = stream.read();
+                if( sendEvents() ) {
+                    router.signalEvent(IOEvent.read(ret));
+                }
+                return ret;
+            } catch(IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+                throw ex;
+            }
+        }
+
+        @Override
+        public boolean markSupported() {
+            return stream.markSupported();
+        }
+
+        @Override
+        public synchronized void reset() throws IOException {
+            try {
+                stream.reset();
+                if( sendEvents() ) {
+                    router.signalEvent(IOEvent.reset());
+                }
+            } catch(IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+                throw ex;
+            } 
+        }
+
+        @Override
+        public synchronized void mark(int readlimit) {
+            stream.mark(readlimit);
+            if( sendEvents() ) {
+                router.signalEvent(IOEvent.mark(readlimit));
+            }
+        }
+
+        @Override
+        public void close() throws IOException {
+            try {
+                stream.close();
+                if( sendEvents() ) {
+                    router.signalEvent(IOEvent.close());
+                }
+            } catch(IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+                throw ex;
+            } finally {
+                router.dispose();
+            }
+        }
+
+        @Override
+        public int available() throws IOException {
+            return stream.available();
+        }
+
+        @Override
+        public long skip(long n) throws IOException {
+            try {
+                long ret = stream.skip(n);
+                if( sendEvents() ) {
+                    router.signalEvent(IOEvent.skip(ret));
+                }
+                return ret;
+            } catch(IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+                throw ex;
+            }
+        }
+
+        @Override
+        public int read(byte[] data, int off, int len) throws IOException {
+            try {
+                int ret = stream.read(data, off, len);
+                if( sendEvents() ) {
+                    router.signalEvent(IOEvent.read(data, off, ret));
+                }
+                return ret;
+            } catch(IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+                throw ex;
+            }
+        }
+
+        @Override
+        public int read(byte[] data) throws IOException {
+            try {
+                int ret = stream.read(data);
+                if( sendEvents() ) {
+                    router.signalEvent(IOEvent.read(data, 0, ret));
+                }
+                return ret;
+            } catch(IOException ex) {
+                router.signalEvent(IOEvent.failure(ex));
+                throw ex;
+            }
+        }
+        
+        
+        
+    }
+    
+    public static class OStream extends OutputStream implements EventRouter {
+        
+        private final OutputStream stream;
+        private final EventReactor worker;
+        
+        private OStream(OutputStream ostream) {
+            this.stream = ostream;
+            this.worker = EventReactor.newInstance("AsyncStream.OStream.Router", 0, 0);
+            this.worker.addEventListener(IORequest.class, new ASHandler(this));
+        }
+
+        @Override
+        public void write(int data) throws IOException {
+            worker.signalEvent(IORequest.write(data));
+        }
+
+        @Override
+        public void write(byte[] data, int off, int len) throws IOException {
+            try {
+                worker.signalEvent(IORequest.write(data, off, len));
+            } catch(Exception cause) {
+                worker.signalEvent(IOEvent.failure(new IOException(cause)));
+            }
+        }
+
+        @Override
+        public void write(byte[] data) throws IOException {
+            worker.signalEvent(IORequest.write(data));
+        }
+
+        @Override
+        public void close() throws IOException {
+            worker.signalEvent(IORequest.close());
+        }
+
+        @Override
+        public void flush() throws IOException {
+            worker.signalEvent(IORequest.flush());
+        }
+
+        @Override
+        public synchronized <T extends Event> void addEventListener(Class<T> event, EventListener<? super T> listener) {
+            worker.addEventListener(event, listener);
+        }
+
+        @Override
+        public synchronized <T extends Event> void removeEventListener(Class<T> event, EventListener<? super T> listener) {
+            worker.removeEventListener(event, listener);
+        }
+
+        @Override
+        public synchronized void removeAllEventListeners() {
+            worker.removeAllEventListeners();
+        }
+
+        @Override
+        public synchronized Collection<EventBinding<?>> getEventListeners() {
+            return worker.getEventListeners();
+        }
+
+        @Override
+        public synchronized int getEventListenersCount() {
+            return worker.getEventListenersCount();
+        }
+
+        @Override
+        public boolean signalEvent(Event event) {
+            return worker.signalEvent(event);
+        }
+        
+        @Override
+        public void handleEvent(Event event) {
+            worker.handleEvent(event);
+        }
+        
+        @Override
+        public void dispose() {
+            try {
+                stream.close();
+                if(sendEvents()) {
+                    worker.signalEvent(IOEvent.close());
+                }
+            } catch (IOException ex) {
+                worker.signalEvent(IOEvent.failure(ex));
+            } finally {
+                worker.dispose();
+            }
+        }
+
+        private boolean sendEvents() {
+            return worker.getEventListenersCount() > 1;
+        }
+        
+    }
+    
+    private static class ASHandler extends ReflectEventListener<Event> {
+        
+        private final OStream that;
+
+        public ASHandler(OStream that) {
+            this.that = that;
+        }
+
+        @Meta.EventHandler
+        void write(IORequest.WriteByte event) throws IOException {
+            that.stream.write(event.data());
+            if(that.sendEvents()) {
+                that.worker.signalEvent(IOEvent.write(event.data()));
+            }
+        }
+        
+        @Meta.EventHandler
+        void write(IORequest.WriteByteArray event) throws IOException {
+            that.stream.write(event.data());
+            if(that.sendEvents()) {
+                that.worker.signalEvent(IOEvent.write(event.data()));
+            }
+        }
+        
+        @Meta.EventHandler
+        void close(IORequest.Close event) throws IOException {
+            that.stream.close();
+            if(that.sendEvents()) {
+                that.worker.signalEvent(IOEvent.close());
+            }
+            that.worker.dispose();
+        }
+        
+        @Meta.EventHandler
+        void flush(IORequest.Flush event) throws IOException {
+            that.stream.flush();
+            if(that.sendEvents()) {
+                that.worker.signalEvent(IOEvent.flush());
+            }
+        }
+        
+        @Meta.ErrorHandler
+        void handle(IOException cause) {
+            if(that.sendEvents()) {
+                that.worker.signalEvent(IOEvent.failure(cause));
+            }
+        }
+        
+    }
+    
+}

+ 135 - 0
src/main/java/net/ranides/assira/io/CheckedStream.java

@@ -0,0 +1,135 @@
+/*
+ * @author Ranides Atterwim <ranides@gmail.com>
+ * @copyright Ranides Atterwim
+ * @license WTFPL
+ * @url http://ranides.net/projects/assira
+ */
+package net.ranides.assira.io;
+
+import java.io.FilterInputStream;
+import java.io.FilterOutputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+
+/**
+ *
+ * @author Ranides Atterwim <ranides@gmail.com>
+ */
+public final class CheckedStream {
+    
+    private CheckedStream() {
+        // utility class
+    }
+    
+    public static OutputStream wrap(OutputStream ostream) {
+        return ostream instanceof OStream ? (OStream)ostream : new OStream(ostream);
+    }
+    
+    public static InputStream wrap(InputStream istream) {
+        return istream instanceof IStream ? (IStream)istream : new IStream(istream);
+    }
+    
+    private static final class OStream extends FilterOutputStream {
+        
+        private boolean closed = false;
+
+        public OStream(OutputStream out) {
+            super(out);
+        }
+        
+        @Override
+        public synchronized void write(byte[] data, int off, int len) throws IOException {
+            checkstate();
+            super.write(data, off, len);
+        }
+
+        @Override
+        public synchronized void write(int bytedata) throws IOException {
+            checkstate();
+            super.write(bytedata);
+        }
+        
+        @Override
+        public void close() throws IOException {
+            checkstate();
+            closed = true;
+        }
+
+        @Override
+        public String toString() {
+            return out.toString();
+        }
+
+        private void checkstate() throws IOException {
+            if(closed) {
+                throw new IOException("already closed");
+            }
+        }
+
+    }
+    
+    private static final class IStream extends FilterInputStream {
+        
+        private boolean closed = false;
+
+        public IStream(InputStream istream) {
+            super(istream);
+        }
+        
+        @Override
+        public synchronized void reset() throws IOException {
+            checkstate();
+            super.reset();
+        }
+
+        @Override
+        public synchronized int available() throws IOException {
+            checkstate();
+            return super.available();
+        }
+
+        @Override
+        public synchronized long skip(long n) throws IOException {
+            checkstate();
+            return super.skip(n);
+        }
+
+        @Override
+        public synchronized int read(byte[] target, int off, int len) throws IOException {
+            checkstate();
+            return super.read(target, off, len);
+        }
+
+        @Override
+        public synchronized int read() throws IOException {
+            checkstate();
+            return super.read();
+        }
+
+        @Override
+        public int read(byte[] target) throws IOException {
+            checkstate();
+            return super.read(target);
+        }
+        
+        @Override
+        public synchronized void close() throws IOException {
+            checkstate();
+            closed = true;
+        }
+        
+        @Override
+        public String toString() {
+            return in.toString();
+        }
+
+        private void checkstate() throws IOException {
+            if(closed) {
+                throw new IOException("already closed");
+            }
+        }
+
+    }
+    
+}

+ 0 - 149
src/main/java/net/ranides/assira/io/IORequest.java

@@ -26,7 +26,6 @@ public abstract class IORequest implements Event, Serializable {
     private static final class Li { // NOPMD lazy init idiom
         static final Close CLOSE = new Close();
         static final Flush FLUSH = new Flush();
-        static final Reset RESET = new Reset();
     }
 
     public static Close close() {
@@ -37,31 +36,6 @@ public abstract class IORequest implements Event, Serializable {
         return Li.FLUSH;
     }
 
-    public static Skip skip(long delta) {
-        return new Skip(delta);
-    }
-
-    public static Reset reset() {
-        return Li.RESET;
-    }
-
-    public static Mark mark(int limit) {
-        return new Mark(limit);
-    }
-
-    public static ReadByte read() {
-        return new ReadByte();
-    }
-
-    public static ReadByteArray read(int count) {
-        return new ReadByteArray(count);
-    }
-
-    @Deprecated
-    public static ReadByteBuffer read(byte[] buffer, int offset, int count) {
-        return new ReadByteBuffer(buffer, offset, count);
-    }
-
     public static WriteByte write(int value) {
         return new WriteByte(value);
     }
@@ -74,8 +48,6 @@ public abstract class IORequest implements Event, Serializable {
         return new WriteByteArray(Arrays.copyOfRange(array, offset, offset + length));
     }
 
-/* ************************************************************************** */
-
 
     public static class Close extends IORequest {
 
@@ -105,127 +77,6 @@ public abstract class IORequest implements Event, Serializable {
         }
     }
 
-    public static class Skip extends IORequest {
-
-        private static final long serialVersionUID = 1L;
-
-        private final long delta;
-
-        protected Skip(long delta) {
-            this.delta = delta;
-        }
-
-        public long delta() {
-            return delta;
-        }
-
-        @Override
-        public String toString() {
-            return "IORequest.Skip";
-        }
-    }
-
-    public static class Reset extends IORequest {
-
-        private static final long serialVersionUID = 1L;
-
-        protected Reset() {
-            // do nothing
-        }
-
-        @Override
-        public String toString() {
-            return "IORequest.Reset";
-        }
-    }
-
-    public static class Mark extends IORequest {
-
-        private static final long serialVersionUID = 1L;
-
-        private final int limit;
-
-        protected Mark(int limit) {
-            this.limit = limit;
-        }
-
-        public int limit() {
-            return limit;
-        }
-
-        @Override
-        public String toString() {
-            return "IORequest.Mark";
-        }
-    }
-
-    public static class ReadByte extends IORequest {
-
-        private static final long serialVersionUID = 1L;
-
-        protected ReadByte() {
-            // do nothing
-        }
-
-        @Override
-        public String toString() {
-            return "IORequest.ReadByte";
-        }
-    }
-
-    public static class ReadByteArray extends IORequest {
-
-        private static final long serialVersionUID = 1L;
-
-        private final int count;
-
-        protected ReadByteArray(int count) {
-            this.count = count;
-        }
-
-        public final int count() {
-            return count;
-        }
-
-        @Override
-        public String toString() {
-            return "IORequest.ReadByteArray: " + count;
-        }
-    }
-
-    @Deprecated
-    public static class ReadByteBuffer extends IORequest {
-
-        private static final long serialVersionUID = 1L;
-
-        private final byte[] buffer;
-        private final int offset;
-        private final int count;
-
-        protected ReadByteBuffer(byte[] buffer, int offset, int count) {
-            this.buffer = buffer;
-            this.offset = offset;
-            this.count = count;
-        }
-
-        public byte[] buffer() {
-            return buffer;
-        }
-
-        public int offset() {
-            return offset;
-        }
-
-        public final int count() {
-            return count;
-        }
-
-        @Override
-        public String toString() {
-            return "IORequest.ReadByteBuffer: " + count;
-        }
-    }
-
     public static class WriteByte extends IORequest {
 
         private static final long serialVersionUID = 1L;

+ 268 - 0
src/test/java/net/ranides/assira/io/AsyncStreamTest.java

@@ -0,0 +1,268 @@
+/*
+ * @author Ranides Atterwim <ranides@gmail.com>
+ * @copyright Ranides Atterwim
+ * @license WTFPL
+ * @url http://ranides.net/projects/assira
+ */
+package net.ranides.assira.io;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import net.ranides.assira.collection.SetUtils;
+import net.ranides.assira.events.Event;
+import net.ranides.assira.events.EventListener;
+import net.ranides.assira.text.TextEncoding;
+import net.ranides.assira.text.TextInputStream;
+import net.ranides.assira.text.TextOutputStream;
+import net.ranides.assira.time.TimeUtils;
+import org.junit.Test;
+import static org.junit.Assert.*;
+
+/**
+ *
+ * @author Ranides Atterwim <ranides@gmail.com>
+ */
+public class AsyncStreamTest {
+
+    @Test
+    public void testWriter() throws IOException {
+        OutputStream ostream = CheckedStream.wrap(new TextOutputStream());
+        AsyncStream.OStream astream = AsyncStream.wrap(ostream);
+        EventCollector<IOEvent> events1 = new EventCollector<>();
+        EventCollector<IORequest> events2 = new EventCollector<>();
+        
+        astream.addEventListener(IOEvent.class, events1);
+        astream.addEventListener(IORequest.class, events2);
+        
+        astream.write('H');
+        astream.write("ello ".getBytes(TextEncoding.NIO_UTF8));
+        astream.write("...world!...".getBytes(TextEncoding.NIO_UTF8), 3, 6);
+        astream.write("???".getBytes(TextEncoding.NIO_UTF8), 100, 100);
+        astream.flush();
+        astream.close();
+        
+        // @todo (assira # 8) issue: exception will be lost
+        // CStream won't write anything but exception will be lost
+        astream.write('?');
+        
+        TimeUtils.sleep(250);
+        assertEquals("Hello world!",ostream.toString());
+        
+        // @todo (assira # 8) issue: exception will be lost
+        astream.close();
+        
+        TimeUtils.sleep(250);
+        assertEquals("Hello world!",ostream.toString());
+        
+        assertTrue(events1.eq(SetUtils.asHashSet(
+            "IOEvent.WriteByte: 72",
+            "IOEvent.WriteByteArray: 5",
+            "IOEvent.WriteByteArray: 6",
+            "IOEvent.Failure: java.io.IOException: java.lang.ArrayIndexOutOfBoundsException",
+            "IOEvent.Flush",
+            "IOEvent.Close"
+        )));
+        assertTrue(events2.eq(Arrays.asList(
+            "IORequest.WriteByte: 72",
+            "IORequest.WriteByteArray: 5",
+            "IORequest.WriteByteArray: 6",
+            "IORequest.Flush",
+            "IORequest.Close",
+            "IORequest.WriteByte: 63"
+        )));
+    }
+    
+    @Test
+    public void testWriterDispose() throws IOException {
+        
+        OutputStream ostream = CheckedStream.wrap(new TextOutputStream());
+        AsyncStream.OStream astream = AsyncStream.wrap(ostream);
+        EventCollector<IOEvent> events1 = new EventCollector<>();
+        EventCollector<IORequest> events2 = new EventCollector<>();
+        
+        astream.addEventListener(IOEvent.class, events1);
+        astream.addEventListener(IORequest.class, events2);
+        
+        astream.dispose();
+        
+        // @todo (assira # 8) issue: exception will be lost
+        TimeUtils.sleep(250);
+        astream.dispose();
+        
+        TimeUtils.sleep(250);
+
+        assertTrue(events1.eq(SetUtils.asHashSet(
+            "IOEvent.Close"
+        )));
+        assertTrue(events2.eq(Collections.<String>emptyList()));
+        
+    }
+    
+    @Test
+    public void testWriterRouter() throws IOException {
+        
+        OutputStream ostream = CheckedStream.wrap(new TextOutputStream());
+        AsyncStream.OStream astream = AsyncStream.wrap(ostream);
+        EventListener<Event> e1 = new EventCollector<>();
+        EventListener<Event> e2 = new EventCollector<>();
+        
+        int prev = astream.getEventListenersCount();
+        
+        astream.addEventListener(IOEvent.class, e1);
+        astream.addEventListener(IORequest.class, e2);
+        assertEquals(prev+2, astream.getEventListenersCount());
+        
+        astream.removeEventListener(IOEvent.class, e2);
+        astream.removeEventListener(IORequest.class, e1);
+        assertEquals(prev+2, astream.getEventListenersCount());
+        
+        astream.removeEventListener(IOEvent.class, e1);
+        assertEquals(prev+1, astream.getEventListenersCount());
+        astream.removeEventListener(IORequest.class, e2);
+        assertEquals(prev+0, astream.getEventListenersCount());
+        
+        // @todo (assira # 8) issue: AsyncStream will be broken
+        // "hidden" listener is required to implement functionality
+        astream.removeAllEventListeners();
+        assertEquals(0, astream.getEventListenersCount());
+        
+        astream.close();
+    }
+    
+    @Test
+    public void testReader() throws IOException {
+        InputStream istream = CheckedStream.wrap(new TextInputStream("Hello world"));
+        AsyncStream.IStream astream = AsyncStream.wrap(istream);
+        EventCollector<IOEvent> events1 = new EventCollector<>();
+        
+        astream.addEventListener(IOEvent.class, events1);
+
+        assertEquals(11, astream.available());
+        assertTrue(astream.markSupported());
+        astream.mark(8);
+        assertEquals('H', astream.read());
+        astream.reset();
+        
+        byte[] target= new byte[6];
+        assertEquals(6, astream.read(target));
+        assertArrayEquals("Hello ".getBytes(), target);
+        
+        assertEquals(2, astream.skip(2));
+        assertEquals(2, astream.read(target, 1, 2));
+        assertArrayEquals("Hrllo ".getBytes(), target);
+        
+        astream.close();
+        // @todo (assira # 8) issue: exceptions will be lost
+        try {
+            astream.close();
+            fail("IOException expected");
+        } catch(IOException ex) { assertTrue(true); }
+        
+        try {
+            astream.read();
+            fail("IOException expected");
+        } catch(IOException ex) { assertTrue(true); }
+        
+        try {
+            astream.read(target);
+            fail("IOException expected");
+        } catch(IOException ex) { assertTrue(true); }
+        
+        try {
+            astream.read(target,1,2);
+            fail("IOException expected");
+        } catch(IOException ex) { assertTrue(true); }
+        
+        try {
+            astream.reset();
+            fail("IOException expected");
+        } catch(IOException ex) { assertTrue(true); }
+        
+        try {
+            astream.skip(5);
+            fail("IOException expected");
+        } catch(IOException ex) { assertTrue(true); }
+        
+        assertTrue(events1.eq(Arrays.asList(
+            "IOEvent.Mark", 
+            "IOEvent.ReadByte: 72", 
+            "IOEvent.Reset", 
+            "IOEvent.ReadByteArray: 6:48 65 6C 6C 6F 20", 
+            "IOEvent.Skip", 
+            "IOEvent.ReadByteArray: 2:72 6C", 
+            "IOEvent.Close"
+        )));
+    }
+    
+    @Test
+    public void testReaderDispose() throws IOException {
+        InputStream istream = CheckedStream.wrap(new TextInputStream("Hello world"));
+        AsyncStream.IStream astream = AsyncStream.wrap(istream);
+        EventCollector<IOEvent> events1 = new EventCollector<>();
+        
+        astream.addEventListener(IOEvent.class, events1);
+        
+        astream.dispose();
+        astream.dispose();
+        TimeUtils.sleep(250);
+
+        assertTrue(events1.eq(SetUtils.asHashSet(
+            "IOEvent.Close"
+        )));
+    }
+    
+    @Test
+    public void testReaderRouter() throws IOException {
+        InputStream istream = CheckedStream.wrap(new TextInputStream("Hello world"));
+        AsyncStream.IStream astream = AsyncStream.wrap(istream);
+        EventCollector<IOEvent> e1 = new EventCollector<>();
+        EventCollector<IOEvent> e2 = new EventCollector<>();
+        
+        astream.addEventListener(IOEvent.Close.class, e1);
+        astream.addEventListener(IOEvent.Read.class, e2);
+        assertEquals(2, astream.getEventListenersCount());
+        
+        astream.removeEventListener(IOEvent.Read.class, e1);
+        astream.removeEventListener(IOEvent.Close.class, e2);
+        assertEquals(2, astream.getEventListenersCount());
+        
+        astream.removeEventListener(IOEvent.Close.class, e1);
+        assertEquals(1, astream.getEventListenersCount());
+        astream.removeEventListener(IOEvent.Read.class, e2);
+        assertEquals(0, astream.getEventListenersCount());
+        
+        astream.removeAllEventListeners();
+        assertEquals(0, astream.getEventListenersCount());
+        
+        astream.close();
+    }
+    
+    private static final class EventCollector<E extends Event> implements EventListener<E> {
+        
+        private final ConcurrentLinkedQueue<String> events = new ConcurrentLinkedQueue<>();
+
+        @Override
+        public void handleEvent(E event) {
+            events.add(event.toString());
+        }
+        
+        public boolean eq(Set<String> collection) {
+            return collection.equals(new HashSet<>(events));
+        }
+        
+        public boolean eq(List<String> collection) {
+            return collection.equals(new ArrayList<>(events));
+        }
+        
+    }
+    
+    
+}