fd_sock.java
来自「JGRoups源码」· Java 代码 · 共 1,269 行 · 第 1/4 页
JAVA
1,269 行
void teardownPingSocket() { synchronized(sock_mutex) { if(ping_sock != null) { try { ping_sock.shutdownInput(); ping_sock.close(); } catch(Exception ex) { } ping_sock=null; } Util.close(ping_input); ping_input=null; } } /** * Determines coordinator C. If C is null and we are the first member, return. Else loop: send GET_CACHE message * to coordinator and wait for GET_CACHE_RSP response. Loop until valid response has been received. */ void getCacheFromCoordinator() { Address coord; int attempts=num_tries; Message msg; FdHeader hdr; Hashtable result; get_cache_promise.reset(); while(attempts > 0) { if((coord=determineCoordinator()) != null) { if(coord.equals(local_addr)) { // we are the first member --> empty cache if(log.isDebugEnabled()) log.debug("first member; cache is empty"); return; } hdr=new FdHeader(FdHeader.GET_CACHE); hdr.mbr=local_addr; msg=new Message(coord, null, null); msg.putHeader(name, hdr); passDown(new Event(Event.MSG, msg)); result=(Hashtable) get_cache_promise.getResult(get_cache_timeout); if(result != null) { cache.putAll(result); // replace all entries (there should be none !) in cache with the new values if(trace) log.trace("got cache from " + coord + ": cache is " + cache); return; } else { if(log.isErrorEnabled()) log.error("received null cache; retrying"); } } Util.sleep(get_cache_retry_timeout); --attempts; } } /** * Sends a SUSPECT message to all group members. Only the coordinator (or the next member in line if the coord * itself is suspected) will react to this message by installing a new view. To overcome the unreliability * of the SUSPECT message (it may be lost because we are not above any retransmission layer), the following scheme * is used: after sending the SUSPECT message, it is also added to the broadcast task, which will periodically * re-send the SUSPECT until a view is received in which the suspected process is not a member anymore. The reason is * that - at one point - either the coordinator or another participant taking over for a crashed coordinator, will * react to the SUSPECT message and issue a new view, at which point the broadcast task stops. */ void broadcastSuspectMessage(Address suspected_mbr) { Message suspect_msg; FdHeader hdr; if(suspected_mbr == null) return; if(trace) log.trace("suspecting " + suspected_mbr + " (own address is " + local_addr + ')'); // 1. Send a SUSPECT message right away; the broadcast task will take some time to send it (sleeps first) hdr=new FdHeader(FdHeader.SUSPECT); hdr.mbrs=new Vector(1); hdr.mbrs.addElement(suspected_mbr); suspect_msg=new Message(); suspect_msg.putHeader(name, hdr); passDown(new Event(Event.MSG, suspect_msg)); // 2. Add to broadcast task and start latter (if not yet running). The task will end when // suspected members are removed from the membership bcast_task.addSuspectedMember(suspected_mbr); if(stats) { num_suspect_events++; suspect_history.add(suspected_mbr); } } void broadcastWhoHasSockMessage(Address mbr) { Message msg; FdHeader hdr; if(local_addr != null && mbr != null) if(log.isDebugEnabled()) log.debug("[" + local_addr + "]: who-has " + mbr); msg=new Message(); // bcast msg hdr=new FdHeader(FdHeader.WHO_HAS_SOCK); hdr.mbr=mbr; msg.putHeader(name, hdr); passDown(new Event(Event.MSG, msg)); } /** Sends or broadcasts a I_HAVE_SOCK response. If 'dst' is null, the reponse will be broadcast, otherwise it will be unicast back to the requester */ void sendIHaveSockMessage(Address dst, Address mbr, IpAddress addr) { Message msg=new Message(dst, null, null); FdHeader hdr=new FdHeader(FdHeader.I_HAVE_SOCK); hdr.mbr=mbr; hdr.sock_addr=addr; msg.putHeader(name, hdr); passDown(new Event(Event.MSG, msg)); } /** Attempts to obtain the ping_addr first from the cache, then by unicasting q request to <code>mbr</code>, then by multicasting a request to all members. */ private IpAddress fetchPingAddress(Address mbr) { IpAddress ret; Message ping_addr_req; FdHeader hdr; if(mbr == null) { if(log.isErrorEnabled()) log.error("mbr == null"); return null; } // 1. Try to get from cache. Add a little delay so that joining mbrs can send their socket address before // we ask them to do so ret=(IpAddress)cache.get(mbr); if(ret != null) { return ret; } Util.sleep(300); if((ret=(IpAddress)cache.get(mbr)) != null) return ret; // 2. Try to get from mbr ping_addr_promise.reset(); ping_addr_req=new Message(mbr, null, null); // unicast hdr=new FdHeader(FdHeader.WHO_HAS_SOCK); hdr.mbr=mbr; ping_addr_req.putHeader(name, hdr); passDown(new Event(Event.MSG, ping_addr_req)); if(!running) return null; ret=(IpAddress)ping_addr_promise.getResult(3000); if(ret != null) { return ret; } // 3. Try to get from all members ping_addr_req=new Message(null); // multicast hdr=new FdHeader(FdHeader.WHO_HAS_SOCK); hdr.mbr=mbr; ping_addr_req.putHeader(name, hdr); passDown(new Event(Event.MSG, ping_addr_req)); ret=(IpAddress) ping_addr_promise.getResult(3000); return ret; } Address determinePingDest() { Address tmp; if(pingable_mbrs == null || pingable_mbrs.size() < 2 || local_addr == null) return null; for(int i=0; i < pingable_mbrs.size(); i++) { tmp=(Address) pingable_mbrs.elementAt(i); if(local_addr.equals(tmp)) { if(i + 1 >= pingable_mbrs.size()) return (Address) pingable_mbrs.elementAt(0); else return (Address) pingable_mbrs.elementAt(i + 1); } } return null; } Address determineCoordinator() { return members.size() > 0 ? (Address) members.elementAt(0) : null; } static String signalToString(int signal) { switch(signal) { case NORMAL_TERMINATION: return "NORMAL_TERMINATION"; case ABNORMAL_TERMINATION: return "ABNORMAL_TERMINATION"; case INTERRUPT: return "INTERRUPT"; default: return "n/a"; } } /* ------------------------------- End of Private Methods ------------------------------------ */ public static class FdHeader extends Header implements Streamable { public static final byte SUSPECT=10; public static final byte WHO_HAS_SOCK=11; public static final byte I_HAVE_SOCK=12; public static final byte GET_CACHE=13; // sent by joining member to coordinator public static final byte GET_CACHE_RSP=14; // sent by coordinator to joining member in response to GET_CACHE byte type=SUSPECT; Address mbr=null; // set on WHO_HAS_SOCK (requested mbr), I_HAVE_SOCK IpAddress sock_addr; // set on I_HAVE_SOCK // Hashtable<Address,IpAddress> Hashtable cachedAddrs=null; // set on GET_CACHE_RSP Vector mbrs=null; // set on SUSPECT (list of suspected members) public FdHeader() { } // used for externalization public FdHeader(byte type) { this.type=type; } public FdHeader(byte type, Address mbr) { this.type=type; this.mbr=mbr; } public FdHeader(byte type, Vector mbrs) { this.type=type; this.mbrs=mbrs; } public FdHeader(byte type, Hashtable cachedAddrs) { this.type=type; this.cachedAddrs=cachedAddrs; } public String toString() { StringBuffer sb=new StringBuffer(); sb.append(type2String(type)); if(mbr != null) sb.append(", mbr=").append(mbr); if(sock_addr != null) sb.append(", sock_addr=").append(sock_addr); if(cachedAddrs != null) sb.append(", cache=").append(cachedAddrs); if(mbrs != null) sb.append(", mbrs=").append(mbrs); return sb.toString(); } public static String type2String(byte type) { switch(type) { case SUSPECT: return "SUSPECT"; case WHO_HAS_SOCK: return "WHO_HAS_SOCK"; case I_HAVE_SOCK: return "I_HAVE_SOCK"; case GET_CACHE: return "GET_CACHE"; case GET_CACHE_RSP: return "GET_CACHE_RSP"; default: return "unknown type (" + type + ')'; } } public void writeExternal(ObjectOutput out) throws IOException { out.writeByte(type); out.writeObject(mbr); out.writeObject(sock_addr); out.writeObject(cachedAddrs); out.writeObject(mbrs); } public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException { type=in.readByte(); mbr=(Address) in.readObject(); sock_addr=(IpAddress) in.readObject(); cachedAddrs=(Hashtable) in.readObject(); mbrs=(Vector) in.readObject(); } public long size() { long retval=Global.BYTE_SIZE; // type retval+=Util.size(mbr); retval+=Util.size(sock_addr); retval+=Global.INT_SIZE; // cachedAddrs size Map.Entry entry; Address key; IpAddress val; if(cachedAddrs != null) { for(Iterator it=cachedAddrs.entrySet().iterator(); it.hasNext();) { entry=(Map.Entry)it.next(); if((key=(Address)entry.getKey()) != null) retval+=Util.size(key); retval+=Global.BYTE_SIZE; // presence for val if((val=(IpAddress)entry.getValue()) != null) retval+=val.size(); } }
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?