CVS update: JGroups/src/org/jgroups/protocols UFC.java MFC.java FlowControl.java

"Bela Ban" <[email protected]>
Newsgroups gmane.comp.java.javagroups.cvs
Message-ID <[email protected]>
  User: belaban 
  Date: 10/09/09 11:34:47

  Modified:    src/org/jgroups/protocols UFC.java MFC.java FlowControl.java
  Log:
  first draft of separate impls for unicast (UFC) and multicast (MFC) flow control (https://jira.jboss.org/browse/JGRP-1154)
  
  Revision  Changes    Path
  1.3       +64 -99    JGroups/src/org/jgroups/protocols/UFC.java
  
  Index: UFC.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/UFC.java,v
  retrieving revision 1.2
  retrieving revision 1.3
  diff -u -r1.2 -r1.3
  --- UFC.java	7 Sep 2010 10:38:21 -0000	1.2
  +++ UFC.java	9 Sep 2010 11:34:47 -0000	1.3
  @@ -4,6 +4,13 @@
   import org.jgroups.Event;
   import org.jgroups.Message;
   import org.jgroups.annotations.MBean;
  +import org.jgroups.annotations.ManagedAttribute;
  +import org.jgroups.annotations.ManagedOperation;
  +
  +import java.util.Iterator;
  +import java.util.Map;
  +import java.util.Vector;
  +import java.util.concurrent.ConcurrentHashMap;
   
   
   /**
  @@ -24,50 +31,62 @@
    * <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
    * <ol/>
    * @author Bela Ban
  - * @version $Id: UFC.java,v 1.2 2010/09/07 10:38:21 belaban Exp $
  + * @version $Id: UFC.java,v 1.3 2010/09/09 11:34:47 belaban Exp $
    */
   @MBean(description="Simple flow control protocol based on a credit system")
   public class UFC extends FlowControl {
  -
       
  +    /**
  +     * Map<Address,Long>: keys are members, values are credits left. For each send,
  +     * the number of credits is decremented by the message size
  +     */
  +    protected final Map<Address,Credit> sent=new ConcurrentHashMap<Address,Credit>(11);
   
  -    protected boolean handleMulticastMessage() {
  -        return false;
  -    }
   
   
  -    protected Credit createCredit(long credits) {
  -        return new UfcCredit(credits);
  +    @ManagedOperation(description="Print sender credits")
  +    public String printSenderCredits() {
  +        return printMap(sent);
       }
   
  -    public void unblock() {
  -        super.unblock();
  +    
  +    @ManagedOperation(description="Print credits")
  +    public String printCredits() {
  +        StringBuilder sb=new StringBuilder(super.printCredits());
  +        sb.append("\nsenders:\n").append(printMap(sent));
  +        return sb.toString();
       }
   
  -    public double getAverageTimeBlocked() {
  -        int    blockings=0;
  -        long   total_time_blocked=0;
  +    public Map<String, Object> dumpStats() {
  +        Map<String, Object> retval=super.dumpStats();
  +        retval.put("senders", printMap(sent));
  +        return retval;
  +    }
   
  -        for(Credit cred: sent.values()) {
  -            blockings+=((UfcCredit)cred).getNumBlockings();
  -            total_time_blocked=((UfcCredit)cred).getTotalBlockingTime();
  -        }
   
  -        return blockings > 0? total_time_blocked / (double)blockings : 0.0; // prevent div-by-zero
  +    protected boolean handleMulticastMessage() {
  +        return false;
       }
   
   
  +
  +    public void unblock() {
  +        super.unblock();
  +    }
  +
  +    @ManagedAttribute(description="Number of times flow control blocks sender")
       public int getNumberOfBlockings() {
           int retval=0;
           for(Credit cred: sent.values())
  -            retval+=((UfcCredit)cred).getNumBlockings();
  +            retval+=cred.getNumBlockings();
           return retval;
       }
   
  +    @ManagedAttribute(description="Total time (ms) spent in flow control block")
       public long getTotalTimeBlocked() {
           long retval=0;
           for(Credit cred: sent.values())
  -            retval+=((UfcCredit)cred).getTotalBlockingTime();
  +            retval+=cred.getTotalBlockingTime();
           return retval;
       }
   
  @@ -91,13 +110,7 @@
               return down_prot.down(evt);
           }
   
  -        if(ignore_synchronous_response && ignore_thread.get()) { // JGRP-465
  -            if(log.isTraceEnabled())
  -                log.trace("bypassing blocking to avoid deadlocking " + Thread.currentThread());
  -            return down_prot.down(evt);
  -        }
  -
  -        UfcCredit cred=(UfcCredit)sent.get(dest);
  +        Credit cred=sent.get(dest);
           if(cred == null) {
               log.error("destination " + dest + " not found; passing message down");
               return down_prot.down(evt);
  @@ -106,36 +119,47 @@
           long block_time=max_block_times != null? getMaxBlockTime(length) : max_block_time;
           
           while(running) {
  -            if(cred.decrementIfEnoughCredits(length, 0)) // timeout == 0: don't block
  -                break;
  -
  -            if(log.isTraceEnabled())
  -                log.trace("blocking for credits (for " + block_time + " ms)");
               boolean rc=cred.decrementIfEnoughCredits(length, block_time);
  -            if(rc && log.isTraceEnabled())
  -                log.trace("unblocking (received credits)");
  -            
               if(rc || !running || max_block_times != null)
                   break;
   
  -            if(cred.sendCreditRequest(System.currentTimeMillis()))
  -                sendCreditRequest(dest, cred.get());
  +            if(cred.needToSendCreditRequest())
  +                sendCreditRequest(dest, Math.max(0, max_credits - cred.get()));
           }
   
           // send message - either after regular processing, or after blocking (when enough credits available again)
           return down_prot.down(evt);
       }
  -    
   
   
  -    protected void handleCredit(Address sender, Number increase) {
  +    protected void handleViewChange(Vector<Address> mbrs) {
  +        super.handleViewChange(mbrs);
  +        if(mbrs == null) return;
  +
  +        // add members not in membership to received and sent hashmap (with full credits)
  +        for(int i=0; i < mbrs.size(); i++) {
  +            Address addr=mbrs.elementAt(i);
  +            if(!sent.containsKey(addr))
  +                sent.put(addr, new Credit(max_credits));
  +        }
  +
  +        // remove members that left
  +        for(Iterator<Address> it=sent.keySet().iterator(); it.hasNext();) {
  +            Address addr=it.next();
  +            if(!mbrs.contains(addr))
  +                it.remove(); // modified the underlying map
  +        }
  +    }
  +
  +
  +    protected void handleCredit(Address sender, long increase) {
           if(sender == null) return;
           StringBuilder sb=null;
   
           Credit cred=sent.get(sender);
           if(cred == null)
               return;
  -        long new_credit=Math.min(max_credits, cred.get() + increase.longValue());
  +        long new_credit=Math.min(max_credits, cred.get() + increase);
   
           if(log.isTraceEnabled()) {
               sb=new StringBuilder();
  @@ -143,67 +167,8 @@
                       .append(", new credits: ").append(new_credit);
               log.trace(sb);
           }
  -
  -        cred.increment(increase.longValue());
  +        cred.increment(increase);
       }
       
   
  -    protected class UfcCredit extends Credit {
  -        int num_blockings=0;
  -        long total_blocking_time=0;
  -        long last_credit_request=0;
  -
  -        protected UfcCredit(long credits) {
  -            super(credits);
  -        }
  -
  -        protected synchronized boolean decrementIfEnoughCredits(long credits, long timeout) {
  -            if(credits <= credits_left) {
  -                credits_left-=credits;
  -                return true;
  -            }
  -
  -            if(timeout <= 0)
  -                return false;
  -
  -            long start=System.currentTimeMillis();
  -            try {
  -                this.wait(timeout);
  -            }
  -            catch(InterruptedException e) {
  -            }
  -            finally {
  -                total_blocking_time+=System.currentTimeMillis() - start;
  -                num_blockings++;
  -            }
  -
  -            if(credits <= credits_left) {
  -                credits_left-=credits;
  -                return true;
  -            }
  -            return false;
  -        }
  -
  -        protected synchronized long increment(long credits) {
  -            long retval=super.increment(credits);
  -            notifyAll();
  -            return retval;
  -        }
  -
  -        protected synchronized boolean sendCreditRequest(long current_time) {
  -            if(current_time - last_credit_request >= max_block_time) {
  -                // we have to set this var now, because we release the lock below (for sending a credit request), so
  -                // all blocked threads would send a credit request, leading to a credit request storm
  -                last_credit_request=System.currentTimeMillis();
  -                return true;
  -            }
  -            return false;
  -        }
  -
  -        protected int getNumBlockings() {return num_blockings;}
  -
  -        protected long getTotalBlockingTime() {return total_blocking_time;}
  -    }
  -
  -
   }
  
  
  
  1.4       +81 -195   JGroups/src/org/jgroups/protocols/MFC.java
  
  Index: MFC.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/MFC.java,v
  retrieving revision 1.3
  retrieving revision 1.4
  diff -u -r1.3 -r1.4
  --- MFC.java	7 Sep 2010 10:38:32 -0000	1.3
  +++ MFC.java	9 Sep 2010 11:34:47 -0000	1.4
  @@ -3,15 +3,16 @@
   import org.jgroups.Address;
   import org.jgroups.Event;
   import org.jgroups.Message;
  -import org.jgroups.annotations.GuardedBy;
   import org.jgroups.annotations.MBean;
  +import org.jgroups.annotations.ManagedAttribute;
   import org.jgroups.annotations.ManagedOperation;
  +import org.jgroups.util.CreditMap;
   
  -import java.util.*;
  -import java.util.concurrent.TimeUnit;
  -import java.util.concurrent.locks.Condition;
  -import java.util.concurrent.locks.Lock;
  -import java.util.concurrent.locks.ReentrantLock;
  +import java.util.HashSet;
  +import java.util.List;
  +import java.util.Set;
  +import java.util.Vector;
  +import java.util.concurrent.atomic.AtomicLong;
   
   
   /**
  @@ -32,7 +33,7 @@
    * <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
    * <ol/>
    * @author Bela Ban
  - * @version $Id: MFC.java,v 1.3 2010/09/07 10:38:32 belaban Exp $
  + * @version $Id: MFC.java,v 1.4 2010/09/09 11:34:47 belaban Exp $
    */
   @MBean(description="Simple flow control protocol based on a credit system")
   public class MFC extends FlowControl {
  @@ -41,83 +42,62 @@
       
       /* --------------------------------------------- Fields ------------------------------------------------------ */
       
  -    
  -  
  -    /**
  -     * the lowest credits of any destination (sent_msgs)
  -     */
  -    @GuardedBy("lock")
  -    private long lowest_credit=max_credits;
  -
  -    /** Lock protecting sent credits table and some other vars (creditors for example) */
  -    private final Lock lock=new ReentrantLock();
   
  +    /** Maintains credits per member */
  +    protected CreditMap credits;
   
  -    /** List of members from whom we expect credits */
  -    @GuardedBy("lock")
  -    protected final Set<Address> creditors=new HashSet<Address>(11);
  +    /** Number of credits waiting to be decremented, used to determine how many credits to ask for in credit requests */
  +    @ManagedAttribute(description="The total number of bytes accumulated by messages which cannot be sent due " +
  +            "to insufficient credits",writable=false)
  +    protected final AtomicLong blocked_credits=new AtomicLong(0);
   
  +    /** Last time a credit request was sent. Used to prevent credit request storms */
  +    protected long last_credit_request=0;
   
  -    /** Mutex to block on down() */
  -    private final Condition credits_available=lock.newCondition();
      
   
  -    /**
  -     * Allows to unblock a blocked sender from an external program, e.g. JMX
  -     */
  +    /** Allows to unblock a blocked sender from an external program, e.g. JMX */
       @ManagedOperation(description="Unblock a sender")
       public void unblock() {
  -        lock.lock();
  -        try {
  -            if(log.isTraceEnabled())
  -                log.trace("unblocking the sender and replenishing all members");
  -
  -            for(Map.Entry<Address,Credit> entry: sent.entrySet())
  -                entry.getValue().set(max_credits);
  +        if(log.isTraceEnabled())
  +            log.trace("unblocking the sender and replenishing all members");
  +        credits.replenishAll();
  +    }
   
  -            lowest_credit=computeLowestCredit(sent);
  -            creditors.clear();
  -            credits_available.signalAll();
  -        }
  -        finally {
  -            lock.unlock();
  -        }
  +    @ManagedOperation(description="Print credits")
  +    public String printCredits() {
  +        return super.printCredits() + "\nsenders min credits: " + credits.getMinCredits();
       }
  -    
   
  -    public void init() throws Exception {
  -        super.init();
  -        lowest_credit=max_credits;
  +    @ManagedOperation(description="Print sender credits")
  +    public String printSenderCredits() {
  +        return credits.toString();
       }
   
  -    public void start() throws Exception {
  -        super.start();
  -        lowest_credit=max_credits;
  +    @ManagedAttribute(description="Number of times flow control blocks sender")
  +    public int getNumberOfBlockings() {
  +        return credits.getNumBlockings();
       }
   
  -    public void stop() {
  -        super.stop();
  -        lock.lock();
  -        try {
  -            running=false;
  -            ignore_thread.set(false);
  -            credits_available.signalAll(); // notify all threads waiting on the mutex that we are done
  -        }
  -        finally {
  -            lock.unlock();
  -        }
  +    @ManagedAttribute(description="Total time (ms) spent in flow control block")
  +    public long getTotalTimeBlocked() {
  +        return credits.getTotalBlockTime();
       }
   
       protected boolean handleMulticastMessage() {
           return true;
       }
   
  -    protected Credit createCredit(long credits) {
  -        return new MfcCredit(credits);
  +   
  +    public void init() throws Exception {
  +        super.init();
  +        credits=new CreditMap(max_credits);
       }
   
  -
  -
  +    public void stop() {
  +        super.stop();
  +        credits.clear();
  +    }
   
       protected Object handleDownMessage(final Event evt, final Message msg, int length) {
           Address dest=msg.getDest();
  @@ -126,133 +106,54 @@
               return down_prot.down(evt);
           }
   
  -        lock.lock();
  +        long block_time=max_block_times != null? getMaxBlockTime(length) : max_block_time;
           try {
  -            if(length > lowest_credit) { // then block and loop asking for credits until enough credits are available
  -                if(ignore_synchronous_response && ignore_thread.get()) { // JGRP-465
  -                    if(log.isTraceEnabled())
  -                        log.trace("bypassing blocking to avoid deadlocking " + Thread.currentThread());
  -                }
  -                else {
  -                    determineCreditors(length);
  -                    num_blockings++; // we count overall blockings, not blockings for *all* threads
  -                    if(log.isTraceEnabled())
  -                        log.trace("Starting blocking. lowest_credit=" + lowest_credit + "; msg length =" + length);
  -
  -                    long block_time=max_block_times != null? getMaxBlockTime(length) : max_block_time;
  -                    while(length > lowest_credit && running) {
  -                        try {
  -                            long start=System.currentTimeMillis();
  -                            boolean rc=credits_available.await(block_time, TimeUnit.MILLISECONDS);
  -                            total_time_blocking+=System.currentTimeMillis() - start;
  -                            if(length <= lowest_credit || rc || !running)
  -                                break;
  -
  -                            // if we use max_block_times, then we do *not* send credit requests, even if we run
  -                            // into timeouts: in this case, it is up to the receivers to send new credits
  -                            if(!rc && max_block_times != null)
  -                                break;
  -
  -                            long curr_time=System.currentTimeMillis();
  -                            long wait_time=curr_time - last_credit_request;
  -                            if(wait_time >= max_block_time) {
  -
  -                                // we have to set this var now, because we release the lock below (for sending a
  -                                // credit request), so all blocked threads would send a credit request, leading to
  -                                // a credit request storm
  -                                last_credit_request=curr_time;
  -
  -                                // we need to send the credit requests down *without* holding the lock, otherwise we might
  -                                // run into the deadlock described in http://jira.jboss.com/jira/browse/JGRP-292
  -                                Map<Address,Credit> sent_copy=new HashMap<Address,Credit>(sent);
  -                                sent_copy.keySet().retainAll(creditors);
  -                                lock.unlock();
  -                                try {
  -                                    for(Map.Entry<Address,Credit> entry: sent_copy.entrySet())
  -                                        sendCreditRequest(entry.getKey(), entry.getValue().get());
  -                                }
  -                                finally {
  -                                    lock.lock();
  -                                }
  -                            }
  -                        }
  -                        catch(InterruptedException e) {
  -                            // bela June 15 2007: don't interrupt the thread again, as this will trigger an infinite loop !!
  -                            // (http://jira.jboss.com/jira/browse/JGRP-536)
  -                            // Thread.currentThread().interrupt();
  -                        }
  -                    }
  +            if(length > 0 && max_block_times == null)
  +                this.blocked_credits.addAndGet(length);
  +            while(running) {
  +                boolean rc=credits.decrement(length, block_time);
  +                if(rc || max_block_times != null || !running)
  +                    break;
  +
  +                if(needToSendCreditRequest()) {
  +                    long credits_blocked=blocked_credits.get();
  +                    List<Address> targets=credits.getMembersWithInsufficientCredits(credits_blocked);
  +                    for(Address target: targets)
  +                        sendCreditRequest(target, credits_blocked);
                   }
               }
  -
  -            long tmp=decrementCredit(sent, length);
  -            if(tmp != -1)
  -                lowest_credit=Math.min(tmp, lowest_credit);
           }
           finally {
  -            lock.unlock();
  +            if(length > 0 && max_block_times == null)
  +                this.blocked_credits.getAndAdd(-length);
           }
   
  -        // send message - either after regular processing, or after blocking (when enough credits available again)
  +        // send message - either after regular processing, or after blocking (when enough credits are available again)
           return down_prot.down(evt);
       }
   
  -    /**
  -     * Adds members which have not enough credits to the creditors list. Called with lock held
  -     * @param length
  -     */
  -    protected void determineCreditors(int length) {
  -        for(Map.Entry<Address,Credit> entry: sent.entrySet()) {
  -            if(entry.getValue().get() <= length)
  -                creditors.add(entry.getKey());
  -        }
  -    }
  -
   
  -  
   
  -    /**
  -     * Decrements credits from a single member, or all members in sent_msgs, depending on whether it is a multicast
  -     * or unicast message. No need to acquire mutex (must already be held when this method is called)
  -     * @param dest
  -     * @param credits
  -     * @return The lowest number of credits left, or -1 if a unicast member was not found
  -     */
  -    protected long decrementCredit(Map<Address,Credit> map, long credits) {
  -        if(map.isEmpty())
  -            return -1;
  -        long lowest=max_credits;
  -        for(Credit cred: map.values())
  -            lowest=Math.min(cred.decrement(credits), lowest);
  -        return lowest;
  -    }
  -
  -    protected void handleCredit(Address sender, Number increase) {
  -        if(sender == null) return;
  -        StringBuilder sb=null;
   
  -        lock.lock();
  -        try {
  -            Credit cred=sent.get(sender);
  -            if(cred == null)
  -                return;
  -            long new_credit=Math.min(max_credits, cred.get() + increase.longValue());
  -
  -            if(log.isTraceEnabled()) {
  -                sb=new StringBuilder();
  -                sb.append("received " + increase + " credits from ").append(sender).append(", old credits: ").append(cred)
  -                        .append(", new credits: ").append(new_credit).append(".\nCreditors before are: ").append(creditors);
  -                log.trace(sb);
  -            }
  +    protected synchronized boolean needToSendCreditRequest() {
  +        long curr_time=System.currentTimeMillis();
  +        long wait_time=curr_time - last_credit_request;
  +        if(wait_time >= max_block_time) {
  +            last_credit_request=curr_time;
  +            return true;
  +        }
  +        return false;
  +    }
  +  
   
  -            cred.increment(increase.longValue());
   
  -            lowest_credit=computeLowestCredit(sent);
  -            if(!creditors.isEmpty() && creditors.remove(sender) && creditors.isEmpty())
  -                credits_available.signalAll();
  -        }
  -        finally {
  -            lock.unlock();
  +    protected void handleCredit(Address sender, long increase) {
  +        credits.replenish(sender, increase);
  +        if(log.isTraceEnabled()) {
  +            StringBuilder sb=new StringBuilder();
  +            sb.append("received " + increase + " credits from ").append(sender).append(", new credits for " + sender + " : ")
  +                    .append(credits.get(sender) + ", min_credits=" + credits.getMinCredits());
  +            log.trace(sb);
           }
       }
   
  @@ -260,32 +161,17 @@
       protected void handleViewChange(Vector<Address> mbrs) {
           super.handleViewChange(mbrs);
   
  -        lock.lock();
  -        try {
  -            // fixed http://jira.jboss.com/jira/browse/JGRP-754 (CCME)
  -            for(Iterator<Address> it=creditors.iterator(); it.hasNext();) {
  -                Address creditor=it.next();
  -                if(!mbrs.contains(creditor))
  -                    it.remove();
  -            }
  -
  -            if(log.isTraceEnabled()) log.trace("creditors are " + creditors);
  -            if(creditors.isEmpty()) {
  -                lowest_credit=computeLowestCredit(sent);
  -                credits_available.signalAll();
  -            }
  -        }
  -        finally {
  -            lock.unlock();
  +        Set<Address> keys=new HashSet<Address>(credits.keys());
  +        for(Address key: keys) {
  +            if(!mbrs.contains(key))
  +                credits.remove(key);
           }
  +
  +        for(Address key: mbrs)
  +            credits.putIfAbsent(key);
       }
   
  -    protected class MfcCredit extends Credit {
   
  -        protected MfcCredit(long credits) {
  -            super(credits);
  -        }
  -    }
   
   
   }
  
  
  
  1.5       +99 -173   JGroups/src/org/jgroups/protocols/FlowControl.java
  
  Index: FlowControl.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/FlowControl.java,v
  retrieving revision 1.4
  retrieving revision 1.5
  diff -u -r1.4 -r1.5
  --- FlowControl.java	7 Sep 2010 10:38:01 -0000	1.4
  +++ FlowControl.java	9 Sep 2010 11:34:47 -0000	1.5
  @@ -6,7 +6,6 @@
   import org.jgroups.View;
   import org.jgroups.annotations.*;
   import org.jgroups.stack.Protocol;
  -import org.jgroups.util.BoundedList;
   import org.jgroups.util.Util;
   
   import java.util.*;
  @@ -31,7 +30,7 @@
    * <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
    * <ol/>
    * @author Bela Ban
  - * @version $Id: FlowControl.java,v 1.4 2010/09/07 10:38:01 belaban Exp $
  + * @version $Id: FlowControl.java,v 1.5 2010/09/09 11:34:47 belaban Exp $
    */
   @MBean(description="Simple flow control protocol based on a credit system")
   public abstract class FlowControl extends Protocol {
  @@ -88,10 +87,9 @@
       protected long min_credits=0;
       
       /**
  -     * Whether an up thread that comes back down should be allowed to
  -     * bypass blocking if all credits are exhausted. Avoids JGRP-465.
  -     * Set to false by default in 2.5 because we have OOB messages for credit replenishments - this flag should not be set
  -     * to true if the concurrent stack is used
  +     * Whether an up thread that comes back down should be allowed to bypass blocking if all credits are exhausted.
  +     * Avoids JGRP-465. Set to false by default in 2.5 because we have OOB messages for credit replenishments -
  +     * this flag should not be set to true if the concurrent stack is used
        */
       @Property(description="Does not block a down message if it is a result of handling an up message in the" +
               "same thread. Fixes JGRP-928")
  @@ -101,27 +99,12 @@
       
       
       /* ---------------------------------------------   JMX      ------------------------------------------------------ */
  -    
  -    
  -    protected int num_blockings=0;
  -    protected int num_credit_requests_received=0, num_credit_requests_sent=0;
  -    protected int num_credit_responses_sent=0, num_credit_responses_received=0;
  -    protected long total_time_blocking=0;
  +    protected int  num_credit_requests_received=0, num_credit_requests_sent=0;
  +    protected int  num_credit_responses_sent=0, num_credit_responses_received=0;
  +
   
  -    // protected final BoundedList<Long> last_blockings=new BoundedList<Long>(50);
  -    
  -    
  -    
       /* --------------------------------------------- Fields ------------------------------------------------------ */
  -    
  -    
  -    /**
  -     * Map<Address,Long>: keys are members, values are credits left. For each send, the
  -     * number of credits is decremented by the message size. A HashMap rather than a ConcurrentHashMap is
  -     * currently used as there might be null values
  -     */
  -    @GuardedBy("lock")
  -    protected final Map<Address,Credit> sent=new ConcurrentHashMap<Address,Credit>(11);
  +   
   
       /**
        * Keeps track of credits / member at the receiver's side. Keys are members, values are credits left (in bytes).
  @@ -132,11 +115,6 @@
       protected final Map<Address,Credit> received=new ConcurrentHashMap<Address,Credit>(11);
   
   
  -  
  -    
  -    /** Peers who have asked for credit that we didn't have */
  -    protected final Set<Address> pending_requesters=new HashSet<Address>(11);
  -
       /**
        * Whether FlowControl is still running, this is set to false when the protocol terminates (on stop())
        */
  @@ -159,16 +137,10 @@
           }
       };   
   
  -    /** Last time a credit request was sent. Used to prevent credit request storms */
  -    @GuardedBy("lock")
  -    protected long last_credit_request=0;   
   
       public void resetStats() {
           super.resetStats();
  -        num_blockings=0;
           num_credit_responses_sent=num_credit_responses_received=num_credit_requests_received=num_credit_requests_sent=0;
  -        total_time_blocking=0;
  -        // last_blockings.clear();
       }
   
       public long getMaxCredits() {
  @@ -195,10 +167,8 @@
           this.min_credits=min_credits;
       }
   
  -    @ManagedAttribute(description="Number of times flow control blocks sender")
  -    public int getNumberOfBlockings() {
  -        return num_blockings;
  -    }
  +
  +    public abstract int getNumberOfBlockings();
   
       public long getMaxBlockTime() {
           return max_block_time;
  @@ -259,14 +229,13 @@
           return sb.toString();
       }
       
  -    @ManagedAttribute(description="Total time (ms) spent in flow control block")
  -    public long getTotalTimeBlocked() {
  -        return total_time_blocking;
  -    }
  +
  +    public abstract long getTotalTimeBlocked();
   
       @ManagedAttribute(description="Average time spent in a flow control block")
       public double getAverageTimeBlocked() {
  -        return num_blockings == 0? 0.0 : total_time_blocking / (double)num_blockings;
  +        long number_of_blockings=getNumberOfBlockings();
  +        return number_of_blockings == 0? 0.0 : getTotalTimeBlocked() / (double)number_of_blockings;
       }
   
       @ManagedAttribute(description="Number of credit requests received")
  @@ -288,36 +257,27 @@
       public int getNumberOfCreditResponsesSent() {
           return num_credit_responses_sent;
       }
  -
  -    @ManagedOperation(description="Print sender credits")
  -    public String printSenderCredits() {
  -        return printMap(sent);
  -    }
  +    
  +    public abstract String printSenderCredits();
   
       @ManagedOperation(description="Print receiver credits")
       public String printReceiverCredits() {
           return printMap(received);
       }
   
  -    @ManagedOperation(description="Print credits")
  +
       public String printCredits() {
           StringBuilder sb=new StringBuilder();
  -        sb.append("senders:\n").append(printMap(sent)).append("\n\nreceivers:\n").append(printMap(received));
  +        sb.append("receivers:\n").append(printMap(received));
           return sb.toString();
       }
   
       public Map<String, Object> dumpStats() {
           Map<String, Object> retval=super.dumpStats();      
  -        retval.put("senders", printMap(sent));
           retval.put("receivers", printMap(received));
           return retval;
       }
   
  -    //@ManagedOperation(description="Print last blocking times")
  -    //public String showLastBlockingTimes() {
  -      //  return last_blockings.toString();
  -    //}
  -
   
       protected long getMaxBlockTime(long length) {
           if(max_block_times == null)
  @@ -339,9 +299,7 @@
        */
       protected abstract boolean handleMulticastMessage();
   
  -    protected abstract Credit createCredit(long credits);
  -
  -    protected abstract void handleCredit(Address sender, Number increase);
  +    protected abstract void handleCredit(Address sender, long increase);
   
   
       /**
  @@ -394,6 +352,13 @@
                   int length=msg.getLength();
                   if(length == 0)
                       break;
  +
  +                if(ignore_synchronous_response && ignore_thread.get()) { // JGRP-465
  +                    if(log.isTraceEnabled())
  +                        log.trace("bypassing flow control because of synchronous response " + Thread.currentThread());
  +                    break;
  +                }
  +
                   return handleDownMessage(evt, msg, length);
               case Event.CONFIG:
                   handleConfigEvent((Map<String,Object>)evt.getArg()); 
  @@ -430,14 +395,14 @@
                       switch(hdr.type) {
                           case FcHeader.REPLENISH:
                               num_credit_responses_received++;
  -                            handleCredit(msg.getSrc(), (Number)msg.getObject());
  +                            handleCredit(msg.getSrc(), (Long)msg.getObject());
                               break;
                           case FcHeader.CREDIT_REQUEST:
                               num_credit_requests_received++;
                               Address sender=msg.getSrc();
  -                            Long sent_credits=(Long)msg.getObject();
  -                            if(sent_credits != null)
  -                                handleCreditRequest(received, sender, sent_credits.longValue());
  +                            Long requested_credits=(Long)msg.getObject();
  +                            if(requested_credits != null)
  +                                handleCreditRequest(received, sender, requested_credits.longValue());
                               break;
                           default:
                               log.error("header type " + hdr.type + " not known");
  @@ -459,10 +424,8 @@
                   finally {
                       if(ignore_synchronous_response)
                           ignore_thread.set(false); // need to revert because the thread is placed back into the pool
  -                    if(new_credits > 0) {
  -                        if(log.isTraceEnabled()) log.trace("sending " + new_credits + " credits to " + sender);
  +                    if(new_credits > 0)
                           sendCredit(sender, new_credits);
  -                    }
                   }
   
               case Event.VIEW_CHANGE:
  @@ -498,39 +461,6 @@
   
   
   
  -
  -
  -
  -    /**
  -     * Decrements credits from a single member, or all members in sent_msgs, depending on whether it is a multicast
  -     * or unicast message. No need to acquire mutex (must already be held when this method is called)
  -     * @param dest
  -     * @param credits
  -     * @return The lowest number of credits left, or -1 if a unicast member was not found
  -     */
  -    protected long decrementCredit(Map<Address,Credit> map, Address dest, long credits) {
  -        if(dest == null || dest.isMulticastAddress()) {
  -            if(map.isEmpty())
  -                return -1;
  -            long lowest=max_credits;
  -            for(Credit cred: map.values())
  -                lowest=Math.min(cred.decrement(credits), lowest);
  -            return lowest;
  -        }
  -        else {
  -            Credit cred=map.get(dest);
  -            if(cred != null)
  -                return cred.decrement(credits);
  -            }
  -        return -1;
  -    }
  -
  -
  -
  -
  -   
  -
  -
       /**
        * Check whether sender has enough credits left. If not, send it some more
        * @param map The hashmap to use
  @@ -548,7 +478,7 @@
               return 0;
   
           if(log.isTraceEnabled())
  -            log.trace("sender " + sender + " minus " + length + " credits, " + (cred.get() - length) + " remaining");
  +            log.trace(sender + " used " + length + " credits, " + (cred.get() - length) + " remaining");
   
           return cred.decrementAndGet(length);
       }
  @@ -556,65 +486,27 @@
       /**
        * @param map The map to modify
        * @param sender The sender who requests credits
  -     * @param left_credits Number of bytes that the sender has left to send messages to us
  +     * @param requested_credits Number of bytes that the sender has left to send messages to us
        */
  -    protected void handleCreditRequest(Map<Address,Credit> map, Address sender, long left_credits) {
  -        if(sender == null) return;
  -        long credit_response=0;
  -        Credit cred=map.get(sender);
  -
  -        long old_credit=cred != null? cred.get() : 0;
  -        if(old_credit > 0)
  -            credit_response=Math.min(max_credits, max_credits - old_credit);
  +    protected void handleCreditRequest(Map<Address,Credit> map, Address sender, long requested_credits) {
  +        Credit cred;
  +        if(sender == null || (cred=map.get(sender)) == null)
  +            return;
   
  +        long credit_response=Math.min(max_credits, Math.min(requested_credits, max_credits - cred.get()));
           if(credit_response > 0) {
               if(log.isTraceEnabled())
                   log.trace("received credit request from " + sender + ": sending " + credit_response + " credits");
  -            if(cred != null)
  -                cred.set(max_credits);
  -            else
  -                map.put(sender, createCredit(max_credits));
  -            pending_requesters.remove(sender);
  -        }
  -        else {
  -            if(pending_requesters.contains(sender)) {
  -                // a sender might have negative credits, e.g. -20000. If we subtracted -20000 from max_credits,
  -                // we'd end up with max_credits + 20000, and send too many credits back. So if the sender's
  -                // credits is negative, we simply send max_credits back
  -                long credits_left=Math.max(0, left_credits);
  -                credit_response=max_credits - credits_left;
  -                // credit_response = max_credits;
  -                if(cred != null)
  -                    cred.set(max_credits);
  -                else
  -                    map.put(sender, createCredit(max_credits));
  -                pending_requesters.remove(sender);
  -                if(log.isWarnEnabled())
  -                    log.warn("Received two credit requests from " + sender +
  -                            " without any intervening messages; sending " + credit_response + " credits");
  -            }
  -            else {
  -                pending_requesters.add(sender);
  -                if(log.isTraceEnabled())
  -                    log.trace("received credit request from " + sender + " but have no credits available");
  -            }
  -        }
  -
  -
  -        if(credit_response > 0)
  +            cred.set(max_credits);
               sendCredit(sender, credit_response);
  +        }
       }
   
   
  -    protected void sendCredit(Address dest, long credit) {
  +    protected void sendCredit(Address dest, long credits) {
           if(log.isTraceEnabled())
  -            log.trace("replenishing " + dest + " with " + credit	+ " credits");
  -        Number number;
  -        if(credit < Integer.MAX_VALUE)
  -            number=(int)credit;
  -        else
  -            number=credit;
  -        Message msg=new Message(dest, null, number);
  +            if(log.isTraceEnabled()) log.trace("sending " + credits + " credits to " + dest);
  +        Message msg=new Message(dest, null, new Long(credits));
           msg.setFlag(Message.OOB);
           msg.putHeader(this.id, REPLENISH_HDR);
           down_prot.down(new Event(Event.MSG, msg));
  @@ -625,12 +517,12 @@
        * We cannot send this request as OOB messages, as the credit request needs to queue up behind the regular messages;
        * if a receiver cannot process the regular messages, that is a sign that the sender should be throttled !
        * @param dest The member to which we send the credit request
  -     * @param credits_left The number of bytes (of credits) left for dest
  +     * @param credits_needed The number of bytes (of credits) left for dest
        */
  -    protected void sendCreditRequest(final Address dest, Long credits_left) {
  +    protected void sendCreditRequest(final Address dest, Long credits_needed) {
           if(log.isTraceEnabled())
               log.trace("sending credit request to " + dest);
  -        Message msg=new Message(dest, null, credits_left);
  +        Message msg=new Message(dest, null, credits_needed);
           msg.putHeader(this.id, CREDIT_REQUEST_HDR);
           down_prot.down(new Event(Event.MSG, msg));
           num_credit_requests_sent++;
  @@ -647,9 +539,7 @@
           for(int i=0; i < mbrs.size(); i++) {
               addr=mbrs.elementAt(i);
               if(!received.containsKey(addr))
  -                received.put(addr, createCredit(max_credits));
  -            if(!sent.containsKey(addr))
  -                sent.put(addr, createCredit(max_credits));
  +                received.put(addr, new Credit(max_credits));
           }
           // remove members that left
           for(Iterator<Address> it=received.keySet().iterator(); it.hasNext();) {
  @@ -657,17 +547,9 @@
               if(!mbrs.contains(addr))
                   it.remove();
           }
  -
  -        // remove members that left
  -        for(Iterator<Address> it=sent.keySet().iterator(); it.hasNext();) {
  -            addr=it.next();
  -            if(!mbrs.contains(addr))
  -                it.remove(); // modified the underlying map
  -        }
  -
  -      
       }
   
  +    
       protected static long computeLowestCredit(Map<Address,Credit> m) {
           Collection<Credit> credits=m.values();
           return Collections.min(credits).get();
  @@ -683,13 +565,46 @@
   
   
   
  -    protected abstract class Credit implements Comparable {
  +    protected class Credit implements Comparable {
           protected long credits_left;
  +        protected int  num_blockings=0;
  +        protected long total_blocking_time=0;
  +        protected long last_credit_request=0;
   
  +        
           protected Credit(long credits) {
               this.credits_left=credits;
           }
   
  +
  +        protected synchronized boolean decrementIfEnoughCredits(long credits, long timeout) {
  +            if(credits <= credits_left) {
  +                credits_left-=credits;
  +                return true;
  +            }
  +
  +            if(timeout <= 0)
  +                return false;
  +
  +            long start=System.currentTimeMillis();
  +            try {
  +                this.wait(timeout);
  +            }
  +            catch(InterruptedException e) {
  +            }
  +            finally {
  +                total_blocking_time+=System.currentTimeMillis() - start;
  +                num_blockings++;
  +            }
  +
  +            if(credits <= credits_left) {
  +                credits_left-=credits;
  +                return true;
  +            }
  +            return false;
  +        }
  +
  +
           protected synchronized long decrementAndGet(long credits) {
               credits_left=Math.max(0, credits_left - credits);
               long credit_response=max_credits - credits_left;
  @@ -700,19 +615,30 @@
               return 0;
           }
   
  -        protected synchronized long decrement(long credits) {
  -            return credits_left=Math.max(0, credits_left - credits);
  +
  +        protected synchronized long increment(long credits) {
  +            long retval=credits_left=Math.min(max_credits, credits_left + credits);
  +            notifyAll();
  +            return retval;
           }
   
  -        
  +        protected synchronized boolean needToSendCreditRequest() {
  +            long current_time=System.currentTimeMillis();
  +            if(current_time - last_credit_request >= max_block_time) {
  +                last_credit_request=current_time;
  +                return true;
  +            }
  +            return false;
  +        }
  +
  +        protected int getNumBlockings() {return num_blockings;}
  +
  +        protected long getTotalBlockingTime() {return total_blocking_time;}
   
           protected synchronized long get() {return credits_left;}
   
           protected synchronized void set(long new_credits) {credits_left=Math.min(max_credits, new_credits);}
   
  -        protected synchronized long increment(long credits) {
  -            return credits_left=Math.min(max_credits, credits_left + credits);
  -        }
   
           public String toString() {
               return String.valueOf(credits_left);
  
  
  

------------------------------------------------------------------------------
This SF.net Dev2Dev email is sponsored by:

Show off your parallel programming skills.
Enter the Intel(R) Threading Challenge 2010.
http://p.sf.net/sfu/intel-thread-sfd
lmpx.com only provides a reader for public news (NNTP) servers. It is not affiliated with the servers or forums shown here and is not responsible for the content of articles, which is written by their respective authors.