messagedispatcher.java

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

JAVA
945
字号
        public void stopInternal() {            // do nothing, DON'T REMOVE !!!!        }        protected void receiveUpEvent(Event evt) {        }        protected void receiveDownEvent(Event evt) {        }        /**         * Called by request correlator when message was not generated by it. We handle it and call the message         * listener's corresponding methods         */        public void passUp(Event evt) {            switch(evt.getType()) {                case Event.MSG:                    if(msg_listener != null) {                        msg_listener.receive((Message) evt.getArg());                    }                    break;                case Event.GET_APPLSTATE: // reply with GET_APPLSTATE_OK                    StateTransferInfo info=(StateTransferInfo)evt.getArg();                    String state_id=info.state_id;                    byte[] tmp_state=null;                    if(msg_listener != null) {                        try {                            if(msg_listener instanceof ExtendedMessageListener && state_id!=null) {                                tmp_state=((ExtendedMessageListener)msg_listener).getState(state_id);                            }                            else {                                tmp_state=msg_listener.getState();                            }                        }                        catch(Throwable t) {                            this.log.error("failed getting state from message listener (" + msg_listener + ')', t);                        }                    }                    channel.returnState(tmp_state, state_id);                    break;                case Event.GET_STATE_OK:                    if(msg_listener != null) {                        try {                            info=(StateTransferInfo)evt.getArg();                            String id=info.state_id;                            if(msg_listener instanceof ExtendedMessageListener && id!=null) {                                ((ExtendedMessageListener)msg_listener).setState(id, info.state);                            }                            else {                                msg_listener.setState(info.state);                            }                        }                        catch(ClassCastException cast_ex) {                            if(this.log.isErrorEnabled())                                this.log.error("received SetStateEvent, but argument " +                                        evt.getArg() + " is not serializable. Discarding message.");                        }                    }                    break;                case Event.STATE_TRANSFER_OUTPUTSTREAM:                    if(msg_listener != null) {                        StateTransferInfo sti=(StateTransferInfo)evt.getArg();                        OutputStream os=sti.outputStream;                        if(os != null && msg_listener instanceof ExtendedMessageListener) {                            if(sti.state_id == null)                                ((ExtendedMessageListener)msg_listener).getState(os);                            else                                ((ExtendedMessageListener)msg_listener).getState(sti.state_id, os);                        }                        return;                    }    				break;                case Event.STATE_TRANSFER_INPUTSTREAM:                    if(msg_listener != null) {                    	StateTransferInfo sti=(StateTransferInfo)evt.getArg();                        InputStream is=sti.inputStream;                        if(is!=null && msg_listener instanceof ExtendedMessageListener) {                            if(sti.state_id == null)                                ((ExtendedMessageListener)msg_listener).setState(is);                            else                                ((ExtendedMessageListener)msg_listener).setState(sti.state_id, is);                        }                    }        			break;                case Event.VIEW_CHANGE:                    View v=(View) evt.getArg();                    Vector new_mbrs=v.getMembers();                    setMembers(new_mbrs);                    if(membership_listener != null) {                        membership_listener.viewAccepted(v);                    }                    break;                case Event.SET_LOCAL_ADDRESS:                    if(log.isTraceEnabled())                        log.trace("setting local_addr (" + local_addr + ") to " + evt.getArg());                    local_addr=(Address)evt.getArg();                    break;                case Event.SUSPECT:                    if(membership_listener != null) {                        membership_listener.suspect((Address) evt.getArg());                    }                    break;                case Event.BLOCK:                    if(membership_listener != null) {                        membership_listener.block();                        channel.blockOk();                    }                    break;                case Event.UNBLOCK:                   if(membership_listener instanceof ExtendedMembershipListener) {                      ((ExtendedMembershipListener)membership_listener).unblock();                   }                   break;                }        }        public void passDown(Event evt) {            down(evt);        }        /**         * Called by channel (we registered before) when event is received. This is the UpHandler interface.         */        public void up(Event evt) {            if(corr != null) {                corr.receive(evt); // calls passUp()            }            else {                if(log.isErrorEnabled()) { //Something is seriously wrong, correlator should not be null since latch is not locked!                    log.error("correlator is null, event will be ignored (evt=" + evt + ")");                }            }        }        public void down(Event evt) {            if(channel != null) {                channel.down(evt);            }            else                if(this.log.isWarnEnabled()) {                    this.log.warn("channel is null, discarding event " + evt);                }        }        /* ----------------------- End of Protocol Interface ------------------------ */    }    class TransportAdapter implements Transport {        public void send(Message msg) throws Exception {            if(channel != null) {                channel.send(msg);            }            else                if(adapter != null) {                    try {                        if(id != null) {                            adapter.send(id, msg);                        }                        else {                            adapter.send(msg);                        }                    }                    catch(Throwable ex) {                        if(log.isErrorEnabled()) {                            log.error("exception=" + Util.print(ex));                        }                    }                }                else {                    if(log.isErrorEnabled()) {                        log.error("channel == null");                    }                }        }        public Object receive(long timeout) throws Exception {            return null;        }    }    class PullPushHandler implements ExtendedMessageListener, MembershipListener {        /* ------------------------- MessageListener interface ---------------------- */        public void receive(Message msg) {            boolean pass_up=true;            if(corr != null) {                pass_up=corr.receiveMessage(msg);            }            if(pass_up) {   // pass on to MessageListener                if(msg_listener != null) {                    msg_listener.receive(msg);                }            }        }        public byte[] getState() {            return msg_listener != null ? msg_listener.getState() : null;        }        public byte[] getState(String state_id) {            if(msg_listener == null) return null;            if(msg_listener instanceof ExtendedMessageListener && state_id!=null) {                return ((ExtendedMessageListener)msg_listener).getState(state_id);            }            else {                return msg_listener.getState();            }        }        public void setState(byte[] state) {            if(msg_listener != null) {                msg_listener.setState(state);            }        }        public void setState(String state_id, byte[] state) {            if(msg_listener != null) {                if(msg_listener instanceof ExtendedMessageListener && state_id!=null) {                    ((ExtendedMessageListener)msg_listener).setState(state_id, state);                }                else {                    msg_listener.setState(state);                }            }        }        public void getState(OutputStream ostream) {            if (msg_listener instanceof ExtendedMessageListener) {                ((ExtendedMessageListener) msg_listener).getState(ostream);            }        }        public void getState(String state_id, OutputStream ostream) {            if (msg_listener instanceof ExtendedMessageListener && state_id!=null) {                ((ExtendedMessageListener) msg_listener).getState(state_id,ostream);            }        }        public void setState(InputStream istream) {            if (msg_listener instanceof ExtendedMessageListener) {                ((ExtendedMessageListener) msg_listener).setState(istream);            }        }        public void setState(String state_id, InputStream istream) {            if (msg_listener instanceof ExtendedMessageListener && state_id != null) {                ((ExtendedMessageListener) msg_listener).setState(state_id,istream);            }        }        /*		 * --------------------- End of MessageListener interface		 * -------------------		 */        /* ------------------------ MembershipListener interface -------------------- */        public void viewAccepted(View v) {            if(corr != null) {                corr.receiveView(v);            }            Vector new_mbrs=v.getMembers();            setMembers(new_mbrs);            if(membership_listener != null) {                membership_listener.viewAccepted(v);            }        }        public void suspect(Address suspected_mbr) {            if(corr != null) {                corr.receiveSuspect(suspected_mbr);            }            if(membership_listener != null) {                membership_listener.suspect(suspected_mbr);            }        }        public void block() {            if(membership_listener != null) {                membership_listener.block();            }        }        /* --------------------- End of MembershipListener interface ---------------- */        // @todo: receive SET_LOCAL_ADDR event and call corr.setLocalAddress(addr)    }}

⌨️ 快捷键说明

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