flushtest.java

来自「JGRoups源码」· Java 代码 · 共 572 行 · 第 1/2 页

JAVA
572
字号

        System.out.println("=== fetching the state ====");
        c2.getState(null, 10000);
        // Util.sleep(2000);
        if(sendMessages){
            c1.send(new Message());
            c1.send(new Message());
            c2.send(new Message());
            c2.send(new Message());
        }

        checkNonStateTransferMemberSequence(receiver3);
        checkBlockStateUnBlockSequence(receiver);
        checkBlockStateUnBlockSequence(receiver2);

        c2.close();
        c2=null;
        // Util.sleep(2000);

        c2=createChannel();
        receiver2=new MyReceiver(c2,"c2");
        c2.setReceiver(receiver2);
        c2.connect("bla");
        // Util.sleep(2000);
        if(sendMessages){
            c1.send(new Message());
        }

        checkEventSequence(receiver);
        checkEventSequence(receiver2);
        checkEventSequence(receiver3);

        System.out.println("=== fetching the state ====");
        c3.getState(null, 10000);
        if(sendMessages){
            c1.send(new Message());
            c2.send(new Message());
        }
        // Util.sleep(2000);

        checkNonStateTransferMemberSequence(receiver2);
        checkBlockStateUnBlockSequence(receiver);
        checkBlockStateUnBlockSequence(receiver3);
    }



    private void checkBlockStateUnBlockSequence(MyReceiver receiver) {
        List events = receiver.getEvents();
        String name = receiver.getName();
        assertNotNull(events);
        assertEquals("Should have three events [block,get|setstate,unblock] but " + name + " has "
                + events, 3, events.size());
        Object obj=events.remove(0);
        assertTrue(name, obj instanceof BlockEvent);
        obj=events.remove(0);
        assertTrue(name, obj instanceof GetStateEvent || obj instanceof SetStateEvent);
        obj=events.remove(0);
        assertTrue(name, obj instanceof UnblockEvent);
        receiver.clear();
    }

    private void checkEventSequence(MyReceiver receiver) {
       List events = receiver.getEvents();
       String eventString = "[" + receiver.getName() + ",events:" + events;
       assertNotNull(events);
       int size = events.size();
       for (int i = 0; i < size; i++)
       {
          Object event = events.get(i);
          if(event instanceof BlockEvent)
          {
             if(i+1<size)
             {
                assertTrue("After Block should be View " + eventString,events.get(i+1) instanceof View);
             }
             if(i!=0)
             {
                assertTrue("Before Block should be Unblock " + eventString,events.get(i-1) instanceof UnblockEvent);
             }
          }
          if(event instanceof View)
          {
             if(i+1<size)
             {
                assertTrue("After View should be Unblock " + eventString,events.get(i+1) instanceof UnblockEvent);
             }
             assertTrue("Before View should be Block " +eventString,events.get(i-1) instanceof BlockEvent);
          }
          if(event instanceof UnblockEvent)
          {
             if(i+1<size)
             {
                assertTrue("After UnBlock should be Block " + eventString,events.get(i+1) instanceof BlockEvent);
             }
             assertTrue("Before UnBlock should be View "+eventString,events.get(i-1) instanceof View);
          }


       }
       receiver.clear();
   }

    private void checkNonStateTransferMemberSequence(MyReceiver receiver) {
        List events = receiver.getEvents();
        assertNotNull(events);
        assertEquals("Should have two events [block,unblock] but " + receiver.getName() + " has "
                        + events, 2, events.size());

        Object obj = events.remove(0);
        assertTrue(obj instanceof BlockEvent);
        obj = events.remove(0);
        assertTrue(obj instanceof UnblockEvent);
        receiver.clear();
    }

    private Channel createChannel() throws ChannelException {
        Channel ret=new JChannel(CONFIG);
        ret.setOpt(Channel.BLOCK, Boolean.TRUE);
        Protocol flush=((JChannel)ret).getProtocolStack().findProtocol("FLUSH");
        if(flush != null) {
            Properties p=new Properties();
            p.setProperty("timeout", "0");
            flush.setProperties(p);

            // send timeout up and down the stack, so other protocols can use the same value too
            Map map = new HashMap();
            map.put("flush_timeout", new Long(0));
            flush.passUp(new Event(Event.CONFIG, map));
            flush.passDown(new Event(Event.CONFIG, map));
        }
        return ret;
    }


    public static Test suite() {
        return new TestSuite(FlushTest.class);
    }

    public static void main(String[] args) {
        junit.textui.TestRunner.run(FlushTest.suite());
    }

    private static class MyReceiver extends ExtendedReceiverAdapter {
        List events;
        String name;
        boolean verbose = true;
        Channel channel = null;

        public MyReceiver(Channel ch,String name) {
           this.name=name;
           channel = ch;
           events=Collections.synchronizedList(new LinkedList());
        }
        
        public MyReceiver(String name) {
            this.name=name;
            events=Collections.synchronizedList(new LinkedList());
        }
        
        public String getLocalAddress()
        {
           String address = "";
           if (channel != null)
           {
              address = channel.getLocalAddress().toString();
           }
           return address;
        }

        public String getName()
        {
            return name;
        }

        public void clear() {
            events.clear();
        }

        public List getEvents() {return new LinkedList(events);}

        public void block() {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() +"]: BLOCK");
            events.add(new BlockEvent());
        }

        public void unblock() {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() + "]: UNBLOCK");
            events.add(new UnblockEvent());
        }

        public void viewAccepted(View new_view) {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() + "]: " + new_view);
            events.add(new_view);
        }

        public byte[] getState() {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() + "]: GetStateEvent");
            events.add(new GetStateEvent(null, null));
            return new byte[]{'b', 'e', 'l', 'a'};
        }

        public void setState(byte[] state) {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() + "]: SetStateEvent");
            events.add(new SetStateEvent(null, null));
        }

        public void getState(OutputStream ostream) {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() + "]: GetStateEvent streamed");
            events.add(new GetStateEvent(null, null));
            byte [] payload = new byte[]{'b', 'e', 'l', 'a'};
            try{
               ostream.write(payload);
            }
            catch (IOException e){              
               e.printStackTrace();
            }
            finally{
               Util.close(ostream);
            }
        }

        public void setState(InputStream istream) {
            if(verbose)
                System.out.println("[" + name + ":" +getLocalAddress() + "]: SetStateEvent streamed");
            events.add(new SetStateEvent(null, null));
            byte [] payload = new byte[4];
            try{
               istream.read(payload);
            }
            catch (IOException e){              
               e.printStackTrace();
            }
            finally{
               Util.close(istream);
            }
        }
    }

    private static class MySimpleReplier extends ExtendedReceiverAdapter {
        Channel channel;
        boolean handle_requests=false;

        public MySimpleReplier(Channel channel, boolean handle_requests) {
            this.channel=channel;
            this.handle_requests=handle_requests;
        }

        public void receive(Message msg) {
            Message reply=new Message(msg.getSrc());
            try {
                System.out.print("-- MySimpleReplier[" + channel.getLocalAddress() + "]: received message from " + msg.getSrc());
                if(handle_requests) {
                    System.out.println(", sending reply");
                    channel.send(reply);
                }
                else
                    System.out.println("\n");
            }
            catch(Exception e) {
                e.printStackTrace();
            }
        }

        public void viewAccepted(View new_view) {
            System.out.println("-- MySimpleReplier[" + channel.getLocalAddress() + "]: viewAccepted(" + new_view + ")");
        }

        public void block() {
            System.out.println("-- MySimpleReplier[" + channel.getLocalAddress() + "]: block()");
        }

        public void unblock() {
            System.out.println("-- MySimpleReplier[" + channel.getLocalAddress() + "]: unblock()");
        }


    }
}

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?