CVS update: JGroups/src/org/jgroups/protocols/pbcast NAKACK.java

"Bela Ban" <[email protected]> Tue, 21 Sep 2010 07:00:27 +0000
Newsgroups gmane.comp.java.javagroups.cvs
Message-ID <[email protected]>
  User: belaban 
  Date: 10/09/21 07:00:27

  Modified:    src/org/jgroups/protocols/pbcast NAKACK.java
  Log:
  - removed delivered_msgs
  - OOB messages are not removed (minus 1 lock) and we don't check for OOB messages availability (minus 1 lock): OOb messages are removed after CAS acquisition
  - JIRA: https://jira.jboss.org/browse/JGRP-1104
  
  Revision  Changes    Path
  1.259     +14 -65    JGroups/src/org/jgroups/protocols/pbcast/NAKACK.java
  
  Index: NAKACK.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/pbcast/NAKACK.java,v
  retrieving revision 1.258
  retrieving revision 1.259
  diff -u -r1.258 -r1.259
  --- NAKACK.java	20 Sep 2010 11:51:10 -0000	1.258
  +++ NAKACK.java	21 Sep 2010 07:00:27 -0000	1.259
  @@ -32,7 +32,7 @@
    * instead of the requester by setting use_mcast_xmit to true.
    *
    * @author Bela Ban
  - * @version $Id: NAKACK.java,v 1.258 2010/09/20 11:51:10 belaban Exp $
  + * @version $Id: NAKACK.java,v 1.259 2010/09/21 07:00:27 belaban Exp $
    */
   @MBean(description="Reliable transmission multipoint FIFO protocol")
   @DeprecatedProperty(names={"max_xmit_size", "eager_lock_release", "stats_list_size"})
  @@ -237,17 +237,14 @@
       protected final BoundedList<String> digest_history=new BoundedList<String>(10);
   
   
  -    /** <em>Regular</em> messages which have been added, but not removed */
  -    private final AtomicInteger undelivered_msgs=new AtomicInteger(0);
  -
   
       public NAKACK() {
       }
   
   
  -    @ManagedAttribute
  -    public int getUndeliveredMessages() {
  -        return undelivered_msgs.get();
  +    @Deprecated
  +    public static int getUndeliveredMessages() {
  +        return 0;
       }
   
       public long getXmitRequestsReceived() {return xmit_reqs_received;}
  @@ -391,7 +388,7 @@
        * @return
        * @deprecated removed in 2.6
        */
  -    public long getMaxXmitSize() {
  +    public static long getMaxXmitSize() {
           return -1;
       }
   
  @@ -659,8 +656,7 @@
               if(hdr == null)
                   break;  // pass up (e.g. unicast msg)
   
  -            // discard messages while not yet server (i.e., until JOIN has returned)
  -            if(!is_server) {
  +            if(!is_server) { // discard messages while not yet server (i.e., until JOIN has returned)
                   if(log.isTraceEnabled())
                       log.trace("message was discarded (not yet server)");
                   return null;
  @@ -749,8 +745,7 @@
               try { // incrementing seqno and adding the msg to sent_msgs needs to be atomic
                   msg_id=seqno +1;
                   msg.putHeader(this.id, NakAckHeader.createMessageHeader(msg_id));
  -                if(win.add(msg_id, msg) && !msg.isFlagSet(Message.OOB))
  -                    undelivered_msgs.incrementAndGet();
  +                win.add(msg_id, msg);
                   seqno=msg_id;
               }
               catch(Throwable t) {
  @@ -795,20 +790,14 @@
               if(leaving)
                   return;
               if(log.isWarnEnabled() && log_discard_msgs)
  -                log.warn(local_addr + ": dropped message from " + sender +
  -                        " (not in xmit_table), keys are " + xmit_table.keySet() +", view=" + view);
  +                log.warn(local_addr + ": dropped message from " + sender + " (not in table " + xmit_table.keySet() +"), view=" + view);
               return;
           }
   
           boolean loopback=local_addr.equals(sender);
  -        boolean added_to_window=false;
  -        boolean added=loopback || (added_to_window=win.add(hdr.seqno, msg));
  +        boolean added=loopback || win.add(hdr.seqno, msg);
   
  -        if(added_to_window && !msg.isFlagSet(Message.OOB))
  -            undelivered_msgs.incrementAndGet();
  -
  -        // message is passed up if OOB. Later, when remove() is called, we discard it. This affects ordering !
  -        // http://jira.jboss.com/jira/browse/JGRP-379
  +        // OOB msg is passed up. When removed, we discard it. Affects ordering: http://jira.jboss.com/jira/browse/JGRP-379
           if(msg.isFlagSet(Message.OOB)) {
               if(added) {
                   if(loopback)
  @@ -818,17 +807,6 @@
                           up_prot.up(new Event(Event.MSG, msg));
                   }
               }
  -            List<Message> msgs;
  -            while(!(msgs=win.removeOOBMessages()).isEmpty()) {
  -                for(Message tmp_msg: msgs) {
  -                    if(tmp_msg.setTransientFlagIfAbsent(Message.OOB_DELIVERED)) {
  -                        up_prot.up(new Event(Event.MSG, tmp_msg));
  -                    }
  -                }
  -            }
  -
  -            if(!(win.hasMessagesToRemove() && undelivered_msgs.get() > 0))
  -                return;
           }
   
           // Efficient way of checking whether another thread is already processing messages from 'sender'.
  @@ -840,17 +818,6 @@
               return;
           }
   
  -        // 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 NAKACK) says messages need to be delivered in the
  -        // order in which they were sent by the sender
  -        int num_regular_msgs_removed=0;
  -
  -        // 2nd line of defense: in case of an exception, remove() might not be called, therefore processing would never
  -        // be set back to false. If we get an exception and released_processing is not true, then we set
  -        // processing to false in the finally clause
           boolean released_processing=false;
           try {
               while(true) {
  @@ -862,22 +829,11 @@
                   }
   
                   for(final Message msg_to_deliver: msgs) {
  -
                       // discard OOB msg if it has already been delivered (http://jira.jboss.com/jira/browse/JGRP-379)
  -                    if(msg_to_deliver.isFlagSet(Message.OOB)) {
  -                        if(msg_to_deliver.setTransientFlagIfAbsent(Message.OOB_DELIVERED)) {
  -                            timer.execute(new Runnable() {
  -                                public void run() {
  -                                    up_prot.up(new Event(Event.MSG, msg_to_deliver));
  -                                }
  -                            });
  -                        }
  +                    if(msg_to_deliver.isFlagSet(Message.OOB) && !msg_to_deliver.setTransientFlagIfAbsent(Message.OOB_DELIVERED))
                           continue;
  -                    }
  -                    num_regular_msgs_removed++;
   
  -                    // Changed by bela Jan 29 2003: not needed (see above)
  -                    //msg_to_deliver.removeHeader(getName());
  +                    //msg_to_deliver.removeHeader(getName()); // Changed by bela Jan 29 2003: not needed (see above)
                       try {
                           up_prot.up(new Event(Event.MSG, msg_to_deliver));
                       }
  @@ -888,11 +844,6 @@
               }
           }
           finally {
  -            // We keep track of regular messages that we added, but couldn't remove (because of ordering).
  -            // When we have such messages pending, then even OOB threads will remove and process them
  -            // http://jira.jboss.com/jira/browse/JGRP-781
  -            undelivered_msgs.addAndGet(-num_regular_msgs_removed);
  -
               // processing is always set in win.remove(processing) above and never here ! This code is just a
               // 2nd line of defense should there be an exception before win.remove(processing) sets processing
               if(!released_processing)
  @@ -973,7 +924,7 @@
                   }
                   continue;
               }
  -            sendXmitRsp(xmit_requester, msg, i);
  +            sendXmitRsp(xmit_requester, msg);
           }
       }
   
  @@ -1011,9 +962,8 @@
        * to preserve the original message's properties, such as src, headers etc.
        * @param dest
        * @param msg
  -     * @param seqno
        */
  -    private void sendXmitRsp(Address dest, Message msg, long seqno) {
  +    private void sendXmitRsp(Address dest, Message msg) {
           Buffer buf;
           if(msg == null) {
               if(log.isErrorEnabled())
  @@ -1579,7 +1529,6 @@
               win.destroy();
           }
           xmit_table.clear();
  -        undelivered_msgs.set(0);
       }
   
   
  
  
  

------------------------------------------------------------------------------
Start uncovering the many advantages of virtual appliances
and start using them to simplify application deployment and
accelerate your shift to cloud computing.
http://p.sf.net/sfu/novell-sfdev2dev