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 + -
显示快捷键?