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