unicast.java

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

JAVA
600
字号
                }                catch(Throwable t) { // eat the exception, don't pass it up the stack                    if(warn) {                        log.warn("failure passing message down", t);                    }                }                msg=null;                return; // we already passed the msg down            case Event.VIEW_CHANGE:  // remove connections to peers that are not members anymore !                Vector new_members=((View)evt.getArg()).getMembers();                Vector left_members;                synchronized(members) {                    left_members=Util.determineLeftMembers(members, new_members);                    members.clear();                    if(new_members != null)                        members.addAll(new_members);                }                // Remove all connections for members that left between the current view and the new view                // See DESIGN for details                boolean rc;                if(use_gms && left_members.size() > 0) {                    Object mbr;                    for(int i=0; i < left_members.size(); i++) {                        mbr=left_members.elementAt(i);                        rc=removeConnection(mbr); // adds to previous_members                        if(rc && trace)                            log.trace("removed " + mbr + " from connection table, member(s) " + left_members + " left");                    }                }                // code by Matthias Weber May 23 2006                for(Enumeration e=previous_members.elements(); e.hasMoreElements();) {                    Object mbr=e.nextElement();                    if(members.contains(mbr)) {                        if(previous_members.removeElement(mbr) != null) {                            if(trace)                                log.trace("removing " + mbr + " from previous_members as result of VIEW_CHANGE event, " +                                        "previous_members=" + previous_members);                        }                    }                }                break;            case Event.ENABLE_UNICASTS_TO:                Object member=evt.getArg();                previous_members.removeElement(member);                if(trace)                    log.trace("removing " + member + " from previous_members as result of ENABLE_UNICAST_TO event, " +                            "previous_members=" + previous_members);                break;        }        passDown(evt);          // Pass on to the layer below us    }    /** Removes and resets from connection table (which is already locked). Returns true if member was found, otherwise false */    private boolean removeConnection(Object mbr) {        Entry entry;        synchronized(connections) {            entry=(Entry)connections.remove(mbr);            if(!previous_members.contains(mbr))                previous_members.add(mbr);        }        if(entry != null) {            entry.reset();            if(trace)                log.trace(local_addr + ": removed connection for dst " + mbr);            return true;        }        else            return false;    }    private void removeAllConnections() {        Entry entry;        synchronized(connections) {            for(Iterator it=connections.values().iterator(); it.hasNext();) {                entry=(Entry)it.next();                entry.reset();            }            connections.clear();        }    }    /** Called by AckSenderWindow to resend messages for which no ACK has been received yet */    public void retransmit(long seqno, Message msg) {        Object  dst=msg.getDest();        // bela Dec 23 2002:        // this will remove a member on a MERGE request, e.g. A and B merge: when A sends the unicast        // request to B and there's a retransmit(), B will be removed !        //          if(use_gms && !members.contains(dst) && !prev_members.contains(dst)) {        //        //                  if(warn) log.warn("UNICAST.retransmit()", "seqno=" + seqno + ":  dest " + dst +        //                             " is not member any longer; removing entry !");        //              synchronized(connections) {        //                  removeConnection(dst);        //              }        //              return;        //          }        if(trace)            log.trace("[" + local_addr + "] --> XMIT(" + dst + ": #" + seqno + ')');        if(Global.copy)            passDown(new Event(Event.MSG, msg.copy()));        else            passDown(new Event(Event.MSG, msg));        num_xmit_requests_received++;    }    /**     * Check whether the hashtable contains an entry e for <code>sender</code> (create if not). If     * e.received_msgs is null and <code>first</code> is true: create a new AckReceiverWindow(seqno) and     * add message. Set e.received_msgs to the new window. Else just add the message.     * @return boolean True if we can send an ack, false otherwise     */    private boolean handleDataReceived(Object sender, long seqno, Message msg) {        if(trace)            log.trace(new StringBuffer().append(local_addr).append(" <-- DATA(").append(sender).append(": #").append(seqno));        if(previous_members.contains(sender)) {            // we don't want to see messages from departed members            if(seqno > DEFAULT_FIRST_SEQNO) {                if(trace)                    log.trace("discarding message " + seqno + " from previous member " + sender);                return false; // don't ack this message so the sender keeps resending it !            }            if(trace)                log.trace("removed " + sender + " from previous_members as we received a message from it");            previous_members.removeElement(sender);        }        Entry    entry;        synchronized(connections) {            entry=(Entry)connections.get(sender);            if(entry == null) {                entry=new Entry();                connections.put(sender, entry);                if(trace)                    log.trace(local_addr + ": created new connection for dst " + sender);            }            if(entry.received_msgs == null)                entry.received_msgs=new AckReceiverWindow(DEFAULT_FIRST_SEQNO);        }        entry.received_msgs.add(seqno, msg); // entry.received_msgs is guaranteed to be non-null if we get here        num_msgs_received++;        num_bytes_received+=msg.getLength();        // Try to remove (from the AckReceiverWindow) as many messages as possible as pass them up        Message  m;        // Prevents concurrent passing up of messages by different threads (http://jira.jboss.com/jira/browse/JGRP-198);        // this is all the more important once we have a threadless stack (http://jira.jboss.com/jira/browse/JGRP-181),        // where lots of threads can come up to this point concurrently, but only 1 is allowed to pass at a time        // We *can* deliver messages from *different* senders concurrently, e.g. reception of P1, Q1, P2, Q2 can result in        // delivery of P1, Q1, Q2, P2: FIFO (implemented by UNICAST) says messages need to be delivered only in the        // order in which they were sent by their senders        synchronized(entry) {            while((m=entry.received_msgs.remove()) != null)                passUp(new Event(Event.MSG, m));        }        return true; // msg was successfully received - send an ack back to the sender    }    /** Add the ACK to hashtable.sender.sent_msgs */    private void handleAckReceived(Object sender, long seqno) {        Entry           entry;        AckSenderWindow win;        if(trace)            log.trace(new StringBuffer().append(local_addr).append(" <-- ACK(").append(sender).                      append(": #").append(seqno).append(')'));        synchronized(connections) {            entry=(Entry)connections.get(sender);        }        if(entry == null || entry.sent_msgs == null)            return;        win=entry.sent_msgs;        win.ack(seqno); // removes message from retransmission        num_acks_received++;    }    private void sendAck(Address dst, long seqno) {        Message ack=new Message(dst);        ack.putHeader(name, new UnicastHeader(UnicastHeader.ACK, seqno));        if(trace)            log.trace(new StringBuffer().append(local_addr).append(" --> ACK(").append(dst).                      append(": #").append(seqno).append(')'));        passDown(new Event(Event.MSG, ack));        num_acks_sent++;    }    public static class UnicastHeader extends Header implements Streamable {        public static final byte DATA=0;        public static final byte ACK=1;        byte    type=DATA;        long    seqno=0;        static final long serialized_size=Global.BYTE_SIZE + Global.LONG_SIZE;        public UnicastHeader() {} // used for externalization        public UnicastHeader(byte type, long seqno) {            this.type=type;            this.seqno=seqno;        }        public String toString() {            return "[UNICAST: " + type2Str(type) + ", seqno=" + seqno + ']';        }        public static String type2Str(byte t) {            switch(t) {                case DATA: return "DATA";                case ACK: return "ACK";                default: return "<unknown>";            }        }        public final long size() {            return serialized_size;        }        public void writeExternal(ObjectOutput out) throws IOException {            out.writeByte(type);            out.writeLong(seqno);        }        public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException {            type=in.readByte();            seqno=in.readLong();        }        public void writeTo(DataOutputStream out) throws IOException {            out.writeByte(type);            out.writeLong(seqno);        }        public void readFrom(DataInputStream in) throws IOException, IllegalAccessException, InstantiationException {            type=in.readByte();            seqno=in.readLong();        }    }    private static final class Entry {        AckReceiverWindow  received_msgs=null;  // stores all msgs rcvd by a certain peer in seqno-order        AckSenderWindow    sent_msgs=null;      // stores (and retransmits) msgs sent by us to a certain peer        long               sent_msgs_seqno=DEFAULT_FIRST_SEQNO;   // seqno for msgs sent by us        void reset() {            if(sent_msgs != null)                sent_msgs.reset();            if(received_msgs != null)                received_msgs.reset();            sent_msgs_seqno=DEFAULT_FIRST_SEQNO;        }        public String toString() {            StringBuffer sb=new StringBuffer();            if(sent_msgs != null)                sb.append("sent_msgs=").append(sent_msgs).append('\n');            if(received_msgs != null)                sb.append("received_msgs=").append(received_msgs).append('\n');            return sb.toString();        }    }}

⌨️ 快捷键说明

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