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