jchannel.java
来自「JGRoups源码」· Java 代码 · 共 1,619 行 · 第 1/5 页
JAVA
1,619 行
View tmp=(View)evt.getArg(); if(tmp instanceof MergeView) my_view=new View(tmp.getVid(), tmp.getMembers()); else my_view=tmp; /* * Bela&Vladimir Oct 27th,2006 (JGroups 2.4)- we need to switch to * connected=true because client can invoke channel.getView() in * viewAccepted() callback invoked on this thread * (see Event.VIEW_CHANGE handling below) * * We do not set connect_promise because we want to wait for * CONNECT_OK and then return from user's JChannel.connect() call. * This is important since we have to wait for Event.UNBLOCK after * CONNECT_OK if blocks are turned on. See JChannel.connect() for * details. * */ if(connected == false) { connected=true; } // unblock queueing of messages due to previous BLOCK event: down(new Event(Event.STOP_QUEUEING)); break; case Event.CONFIG: HashMap config=(HashMap)evt.getArg(); if(config != null) { if(config.containsKey("state_transfer")) { state_transfer_supported=((Boolean)config.get("state_transfer")).booleanValue(); } if(config.containsKey("flush_supported")) { flush_supported=((Boolean)config.get("flush_supported")).booleanValue(); } } break; case Event.CONNECT_OK: connect_promise.setResult(evt.getArg()); break; case Event.SUSPEND_OK: flush_promise.setResult(Boolean.TRUE); break; case Event.DISCONNECT_OK: disconnect_promise.setResult(Boolean.TRUE); break; case Event.GET_STATE_OK: StateTransferInfo info=(StateTransferInfo)evt.getArg(); byte[] state=info.state; state_promise.setResult(state != null? Boolean.TRUE : Boolean.FALSE); if(up_handler != null) { up_handler.up(evt); return; } if(state != null) { String state_id=info.state_id; if(receiver != null) { if(receiver instanceof ExtendedReceiver && state_id!=null) ((ExtendedReceiver)receiver).setState(state_id, state); else receiver.setState(state); } else { try {mq.add(new Event(Event.STATE_RECEIVED, info));} catch(Exception e) {} } } break; case Event.STATE_TRANSFER_INPUTSTREAM: StateTransferInfo sti=(StateTransferInfo)evt.getArg(); InputStream is=sti.inputStream; //Oct 13,2006 moved to down() when Event.STATE_TRANSFER_INPUTSTREAM_CLOSED is received //state_promise.setResult(is != null? Boolean.TRUE : Boolean.FALSE); if(up_handler != null) { up_handler.up(evt); return; } if(is != null) { if(receiver instanceof ExtendedReceiver) { if(sti.state_id == null) ((ExtendedReceiver)receiver).setState(is); else ((ExtendedReceiver)receiver).setState(sti.state_id, is); } else { try { mq.add(new Event(Event.STATE_TRANSFER_INPUTSTREAM, sti)); } catch(Exception e) { } } } break; case Event.SET_LOCAL_ADDRESS: local_addr_promise.setResult(evt.getArg()); break; case Event.EXIT: handleExit(evt); return; // no need to pass event up; already done in handleExit() default: break; } // If UpHandler is installed, pass all events to it and return (UpHandler is e.g. a building block) if(up_handler != null) { up_handler.up(evt); if(type == Event.UNBLOCK){ flush_unblock_promise.setResult(Boolean.TRUE); } return; } switch(type) { case Event.MSG: if(receiver != null) { receiver.receive((Message)evt.getArg()); return; } break; case Event.VIEW_CHANGE: if(receiver != null) { receiver.viewAccepted((View)evt.getArg()); return; } break; case Event.SUSPECT: if(receiver != null) { receiver.suspect((Address)evt.getArg()); return; } break; case Event.GET_APPLSTATE: if(receiver != null) { StateTransferInfo info=(StateTransferInfo)evt.getArg(); byte[] tmp_state; String state_id=info.state_id; if(receiver instanceof ExtendedReceiver && state_id!=null) { tmp_state=((ExtendedReceiver)receiver).getState(state_id); } else { tmp_state=receiver.getState(); } returnState(tmp_state, state_id); return; } break; case Event.STATE_TRANSFER_OUTPUTSTREAM: if(receiver != null) { StateTransferInfo sti=(StateTransferInfo)evt.getArg(); OutputStream os=sti.outputStream; if(os != null && receiver instanceof ExtendedReceiver) { if(sti.state_id == null) ((ExtendedReceiver)receiver).getState(os); else ((ExtendedReceiver)receiver).getState(sti.state_id, os); } return; } break; case Event.BLOCK: if(!receive_blocks) { // discard if client has not set 'receiving blocks' to 'on' down(new Event(Event.BLOCK_OK)); down(new Event(Event.START_QUEUEING)); return; } if(receiver != null) { try { receiver.block(); } catch(Throwable t) { if(log.isErrorEnabled()) log.error("failed calling block() on Receiver", t); } finally { blockOk(); } return; } break; case Event.UNBLOCK: //discard if client has not set 'receiving blocks' to 'on' if(!receive_blocks) { return; } if(receiver instanceof ExtendedReceiver) { try { ((ExtendedReceiver)receiver).unblock(); } catch(Throwable t) { if(log.isErrorEnabled()) log.error("failed calling unblock() on Receiver", t); } finally{ flush_unblock_promise.setResult(Boolean.TRUE); } return; } break; default: break; } if(type == Event.MSG || type == Event.VIEW_CHANGE || type == Event.SUSPECT || type == Event.GET_APPLSTATE || type== Event.STATE_TRANSFER_OUTPUTSTREAM || type == Event.BLOCK || type == Event.UNBLOCK) { try { mq.add(evt); } catch(Exception e) { if(log.isErrorEnabled()) log.error("exception adding event " + evt + " to message queue", e); } } } /** * Sends a message through the protocol stack if the stack is available * @param evt the message to send down, encapsulated in an event */ public void down(Event evt) { if(evt == null) return; // handle setting of additional data (kludge, will be removed soon) if(evt.getType() == Event.CONFIG) { try { Map m=(Map)evt.getArg(); if(m != null && m.containsKey("additional_data")) { additional_data=(byte[])m.get("additional_data"); if(local_addr instanceof IpAddress) ((IpAddress)local_addr).setAdditionalData(additional_data); } } catch(Throwable t) { if(log.isErrorEnabled()) log.error("CONFIG event did not contain a hashmap: " + t); } } if(evt.getType() == Event.STATE_TRANSFER_INPUTSTREAM_CLOSED){ state_promise.setResult(Boolean.TRUE); } if(prot_stack != null) prot_stack.down(evt); else if(log.isErrorEnabled()) log.error("no protocol stack available"); } public String toString(boolean details) { StringBuffer sb=new StringBuffer(); sb.append("local_addr=").append(local_addr).append('\n'); sb.append("cluster_name=").append(cluster_name).append('\n'); sb.append("my_view=").append(my_view).append('\n'); sb.append("connected=").append(connected).append('\n'); sb.append("closed=").append(closed).append('\n'); if(mq != null) sb.append("incoming queue size=").append(mq.size()).append('\n'); if(details) { sb.append("receive_blocks=").append(receive_blocks).append('\n'); sb.append("receive_local_msgs=").append(receive_local_msgs).append('\n'); sb.append("auto_reconnect=").append(auto_reconnect).append('\n'); sb.append("auto_getstate=").append(auto_getstate).append('\n'); sb.append("state_transfer_supported=").append(state_transfer_supported).append('\n'); sb.append("props=").append(props).append('\n'); } return sb.toString(); } /* ----------------------------------- Private Methods ------------------------------------- */ protected final void init(ProtocolStackConfigurator configurator) throws ChannelException { if(log.isInfoEnabled()) log.info("JGroups version: " + Version.description); ConfiguratorFactory.substituteVariables(configurator); // replace vars with system props props=configurator.getProtocolStackString(); prot_stack=new ProtocolStack(this, props); try { prot_stack.setup(); // Setup protocol stack (create layers, queues between them } catch(Throwable e) { throw new ChannelException("unable to setup the protocol stack", e); } } /** * Initializes all variables. Used after <tt>close()</tt> or <tt>disconnect()</tt>, * to be ready for new <tt>connect()</tt> */ private void init() { local_addr=null; cluster_name=null; my_view=null; // changed by Bela Sept 25 2003 //if(mq != null && mq.closed()) // mq.reset(); connect_promise.reset(); disconnect_promise.reset(); connected=false; }
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?