Browse Source

fix: EventReactor - support for unbounded queue
fix: LoggerUtils - remove previous appenders in #reset methods
unit test for AsyncStream

Ranides Atterwim 11 năm trước cách đây
mục cha
commit
27ea10146d

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

@@ -8,6 +8,8 @@
 package net.ranides.assira.events;
 
 import java.util.concurrent.ArrayBlockingQueue;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
 import net.ranides.assira.trace.AdvLogger;
 import net.ranides.assira.trace.LoggerUtils;
@@ -34,7 +36,7 @@ public class EventReactor extends EventDispatcher {
     
     private static final AdvLogger LOGGER = LoggerUtils.getLogger();
 
-    private final ArrayBlockingQueue<Event> events;
+    private final BlockingQueue<Event> events;
     private final long maxtime;
     private final String name;
     private final Thread distributor;
@@ -80,7 +82,7 @@ public class EventReactor extends EventDispatcher {
      */
     protected EventReactor(final String name, int size, long maxtime, final EventJoiner joiner) {
         this.maxtime = maxtime;
-        this.events = new ArrayBlockingQueue<>(size, true);
+        this.events = createQueue(size);
         this.name = name;
         this.exit = CoreEvent.stop(this);
         this.distributor = new Thread(name) {
@@ -108,6 +110,14 @@ public class EventReactor extends EventDispatcher {
             }
         };
     }
+    
+    private static BlockingQueue<Event> createQueue(int size) {
+        if( 0==size || Integer.MAX_VALUE==size) {
+            return new LinkedBlockingQueue<>();
+        } else {
+            return new ArrayBlockingQueue<>(size, true);
+        }
+    }
 
     /**
      * Uruchamia {@code EventRouter}. Metoda nieblokująca.

+ 6 - 3
src/main/java/net/ranides/assira/trace/LoggerUtils.java

@@ -19,10 +19,12 @@ import net.ranides.assira.system.RuntimeUtils;
  * @author ranides
  */
 public final class LoggerUtils {
-
+    
     private static final Map<String, AdvLogger> CACHE = new HashMap<>();
 
-    private static final org.apache.log4j.Appender APPENDER = new org.apache.log4j.ConsoleAppender(new org.apache.log4j.PatternLayout("%d %p [%c %M] %m%n"), "System.out");
+    public static final org.apache.log4j.Layout LAYOUT = new org.apache.log4j.PatternLayout("%d %p [%c %M] %m%n");
+    
+    public static final org.apache.log4j.Appender PRINTER = new org.apache.log4j.ConsoleAppender(LAYOUT, "System.out");
 
     private LoggerUtils() { }
 
@@ -52,11 +54,12 @@ public final class LoggerUtils {
     }
 
     public static org.apache.log4j.Logger resetLogger4j(org.apache.log4j.Level level) {
-        return resetLogger4j(level, APPENDER);
+        return resetLogger4j(level, PRINTER);
     }
 
     public static org.apache.log4j.Logger resetLogger4j(org.apache.log4j.Level level, org.apache.log4j.Appender appender) {
         org.apache.log4j.Logger logger = org.apache.log4j.Logger.getRootLogger();
+        logger.removeAllAppenders();
         logger.addAppender(appender);
         logger.setLevel(level);
         return logger;

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

@@ -0,0 +1,107 @@
+/*
+ * @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.OutputStream;
+import java.io.StringWriter;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import net.ranides.assira.text.TextEncoding;
+import net.ranides.assira.text.TextOutputStream;
+import net.ranides.assira.time.TimeUtils;
+import net.ranides.assira.trace.LoggerUtils;
+import org.apache.log4j.Level;
+import org.apache.log4j.WriterAppender;
+import org.junit.Test;
+import static org.junit.Assert.*;
+
+/**
+ *
+ * @author msieron
+ */
+
+
+public class AsyncStreamTest {
+
+    @Test
+    public void testWriteLogger() throws IOException {
+        StringWriter errorsWriter = new StringWriter();
+        LoggerUtils.resetLogger4j(Level.ALL, new WriterAppender(LoggerUtils.LAYOUT, errorsWriter));
+        
+        OutputStream ostream = new TSG();
+        OutputStream astream = AsyncStream.create(ostream, "logger");
+        
+        astream.write(65);
+        astream.write("qwerty".getBytes(TextEncoding.NIO_UTF8));
+        astream.write("HELLO".getBytes(TextEncoding.NIO_UTF8), 1, 3);
+        astream.close();
+        astream.close();
+        
+        TimeUtils.sleep(100);
+//        System.out.println( "value=" + ostream.toString() );
+//        System.out.println( "errors=" + errorsWriter.toString() );
+        assertEquals("AqwertyELL", ostream.toString());
+        // we don't check if ALL errors are reported 
+        assertTrue(errorsWriter.toString().contains("already closed") );
+        
+        LoggerUtils.resetLogger4j(false);
+    }
+    
+    @Test
+    public void testWriteListener() throws IOException {
+        LoggerUtils.resetLogger4j(false);
+        
+        final ConcurrentLinkedQueue<String> out = new ConcurrentLinkedQueue<>();
+
+        OutputStream ostream = new TSG();
+        OutputStream astream = AsyncStream.create(new AsyncStream.OutputListener(ostream) {
+            @Override
+            public void handleEvent(IOEvent event) {
+                out.add(event.toString());
+            }
+        });
+        
+        astream.write(65);
+        astream.write("qwerty".getBytes(TextEncoding.NIO_UTF8));
+        astream.write("HELLO".getBytes(TextEncoding.NIO_UTF8), 1, 3);
+        astream.close();
+        astream.close();
+        
+        TimeUtils.sleep(100);
+//        System.out.println( "value=" + ostream.toString() );
+//        System.out.println( "events=" + out.toString() );
+        assertEquals("AqwertyELL", ostream.toString());
+        List<String> expected = Arrays.asList(
+            "IOEvent.WriteByte: 65", 
+            "IOEvent.WriteByteArray: 6", 
+            "IOEvent.WriteByteArray: 3", 
+            "IOEvent.Close", 
+            "IOEvent.Failure: java.io.IOException: already closed", 
+            "IOEvent.Failure: java.io.IOException: already closed"
+        );
+        assertEquals(expected, new ArrayList<>(out));
+    }
+    
+    private static class TSG extends TextOutputStream {
+        
+        private boolean opened = true;
+
+        @Override
+        public void close() throws IOException {
+            if( !opened ) {
+                throw new IOException("already closed");
+            }
+            opened = false;
+            super.close();
+        }
+        
+    }
+    
+}