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/13 09:08:37

  Modified:    src/org/jgroups/protocols UFC.java MFC.java FlowControl.java
  Log:
  ns
  
  Revision  Changes    Path
  1.4       +11 -20    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.3
  retrieving revision 1.4
  diff -u -r1.3 -r1.4
  --- UFC.java	9 Sep 2010 11:34:47 -0000	1.3
  +++ UFC.java	13 Sep 2010 09:08:36 -0000	1.4
  @@ -31,7 +31,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: UFC.java,v 1.3 2010/09/09 11:34:47 belaban Exp $
  + * @version $Id: UFC.java,v 1.4 2010/09/13 09:08:36 belaban Exp $
    */
   @MBean(description="Simple flow control protocol based on a credit system")
   public class UFC extends FlowControl {
  @@ -92,27 +92,22 @@
   
       public void stop() {
           super.stop();
  -        for(final Credit cred: sent.values()) {
  -            synchronized(cred) {
  -                cred.set(max_credits);
  -                cred.notifyAll();
  -            }
  -        }
  +        for(Credit cred: sent.values())
  +            cred.set(max_credits);
       }
   
   
   
   
  -    protected Object handleDownMessage(final Event evt, final Message msg, int length) {
  -        Address dest=msg.getDest();
  -        if(dest == null || dest.isMulticastAddress()) {
  +    protected Object handleDownMessage(final Event evt, final Message msg, Address dest, int length) {
  +        if(dest == null || dest.isMulticastAddress()) { // 2nd line of defense, not really needed
               log.error(getClass().getSimpleName() + " doesn't handle multicast messages; passing message down");
               return down_prot.down(evt);
           }
   
           Credit cred=sent.get(dest);
           if(cred == null) {
  -            log.error("destination " + dest + " not found; passing message down");
  +            log.warn("destination " + dest + " not found; passing message down");
               return down_prot.down(evt);
           }
   
  @@ -137,8 +132,7 @@
           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);
  +        for(Address addr: mbrs) {
               if(!sent.containsKey(addr))
                   sent.put(addr, new Credit(max_credits));
           }
  @@ -153,16 +147,13 @@
   
   
       protected void handleCredit(Address sender, long increase) {
  -        if(sender == null) return;
  -        StringBuilder sb=null;
  -
  -        Credit cred=sent.get(sender);
  -        if(cred == null)
  +        Credit cred;
  +        if(sender == null || (cred=sent.get(sender)) == null || increase <= 0)
               return;
  -        long new_credit=Math.min(max_credits, cred.get() + increase);
   
  +        long new_credit=Math.min(max_credits, cred.get() + increase);
           if(log.isTraceEnabled()) {
  -            sb=new StringBuilder();
  +            StringBuilder sb=new StringBuilder();
               sb.append("received " + increase + " credits from ").append(sender).append(", old credits: ").append(cred)
                       .append(", new credits: ").append(new_credit);
               log.trace(sb);
  
  
  
  1.6       +3 -4      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.5
  retrieving revision 1.6
  diff -u -r1.5 -r1.6
  --- MFC.java	10 Sep 2010 10:55:18 -0000	1.5
  +++ MFC.java	13 Sep 2010 09:08:36 -0000	1.6
  @@ -33,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.5 2010/09/10 10:55:18 belaban Exp $
  + * @version $Id: MFC.java,v 1.6 2010/09/13 09:08:36 belaban Exp $
    */
   @MBean(description="Simple flow control protocol based on a credit system")
   public class MFC extends FlowControl {
  @@ -95,9 +95,8 @@
           credits.clear();
       }
   
  -    protected Object handleDownMessage(final Event evt, final Message msg, int length) {
  -        Address dest=msg.getDest();
  -        if(dest != null && !dest.isMulticastAddress()) {
  +    protected Object handleDownMessage(final Event evt, final Message msg, Address dest, int length) {
  +        if(dest != null && !dest.isMulticastAddress()) { // 2nd line of defense, not really needed
               log.error(getClass().getSimpleName() + " doesn't handle unicast messages; passing message down");
               return down_prot.down(evt);
           }
  
  
  
  1.7       +31 -66    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.6
  retrieving revision 1.7
  diff -u -r1.6 -r1.7
  --- FlowControl.java	10 Sep 2010 10:53:42 -0000	1.6
  +++ FlowControl.java	13 Sep 2010 09:08:36 -0000	1.7
  @@ -16,21 +16,10 @@
    * Simple flow control protocol based on a credit system. Each sender has a number of credits (bytes
    * to send). When the credits have been exhausted, the sender blocks. Each receiver also keeps track of
    * how many credits it has received from a sender. When credits for a sender fall below a threshold,
  - * the receiver sends more credits to the sender. Works for both unicast and multicast messages.
  - * <p/>
  - * Note that this protocol must be located towards the top of the stack, or all down_threads from JChannel to this
  - * protocol must be set to false ! This is in order to block JChannel.send()/JChannel.down().
  - * <br/>This is the second simplified implementation of the same model. The algorithm is sketched out in
  - * doc/FlowControl.txt
  - * <br/>
  - * Changes (Brian) April 2006:
  - * <ol>
  - * <li>Receivers now send credits to a sender when more than min_credits have been received (rather than when min_credits
  - * are left)
  - * <li>Receivers don't send the full credits (max_credits), but rather the actual number of bytes received
  - * <ol/>
  + * the receiver sends more credits to the sender.
  + * 
    * @author Bela Ban
  - * @version $Id: FlowControl.java,v 1.6 2010/09/10 10:53:42 belaban Exp $
  + * @version $Id: FlowControl.java,v 1.7 2010/09/13 09:08:36 belaban Exp $
    */
   @MBean(description="Simple flow control protocol based on a credit system")
   public abstract class FlowControl extends Protocol {
  @@ -65,9 +54,6 @@
        */
       protected Map<Long,Long> max_block_times=null;
   
  -    /** Keeps track of the end time after which a message should not get blocked anymore */
  -    // protected static final ThreadLocal<Long> end_time=new ThreadLocal<Long>();
  -
   
       /**
        * If we've received (min_threshold * max_credits) bytes from P, we send more credits to P. Example: if
  @@ -107,20 +93,17 @@
      
   
       /**
  -     * Keeps track of credits / member at the receiver's side. Keys are members, values are credits left (in bytes).
  -     * For each receive, the credits for the sender are decremented by the size of the received message.
  -     * When the credits fall below the threshold, we refill and send a REPLENISH message to the sender.
  -     * The sender blocks until REPLENISH message is received.
  +     * Keeps track of credits per member at the receiver. For each message, the credits for the sender are decremented
  +     * by the size of the received message. When the credits fall below the threshold, we refill and send a REPLENISH
  +     * message to the sender.
        */
       protected final Map<Address,Credit> received=new ConcurrentHashMap<Address,Credit>(11);
   
   
  -    /**
  -     * Whether FlowControl is still running, this is set to false when the protocol terminates (on stop())
  -     */
  +    /** Whether FlowControl is still running, this is set to false when the protocol terminates (on stop()) */
       protected volatile boolean running=true;
   
  -
  +    
       protected boolean frag_size_received=false;
   
      
  @@ -167,7 +150,6 @@
           this.min_credits=min_credits;
       }
   
  -
       public abstract int getNumberOfBlockings();
   
       public long getMaxBlockTime() {
  @@ -288,7 +270,6 @@
               if(length <= entry.getKey())
                   break;
           }
  -
           return retval != null? retval : 0;
       }
   
  @@ -303,9 +284,9 @@
   
   
       /**
  -     * Allows to unblock a blocked sender from an external program, e.g. JMX
  +     * Allows to unblock all blocked senders from an external program, e.g. JMX
        */
  -    @ManagedOperation(description="Unblock a sender")
  +    @ManagedOperation(description="Unblocks all senders")
       public void unblock() {
           ;
       }
  @@ -322,9 +303,7 @@
               log.warn("No fragmentation protocol was found. When flow control is used, we recommend " +
                       "a fragmentation protocol, due to http://jira.jboss.com/jira/browse/JGRP-590");
           }
  -
           running=true;
  -
       }
   
       public void stop() {
  @@ -358,11 +337,12 @@
                           log.trace("bypassing flow control because of synchronous response " + Thread.currentThread());
                       break;
                   }
  +                return handleDownMessage(evt, msg, dest, length);
   
  -                return handleDownMessage(evt, msg, length);
               case Event.CONFIG:
                   handleConfigEvent((Map<String,Object>)evt.getArg()); 
                   break;
  +            
               case Event.VIEW_CHANGE:
                   handleViewChange(((View)evt.getArg()).getMembers());
                   break;
  @@ -376,9 +356,6 @@
           switch(evt.getType()) {
   
               case Event.MSG:
  -
  -                // JGRP-465. We only deal with msgs to avoid having to use a concurrent collection; ignore views,
  -                // suspicions, etc which can come up on unusual threads.
                   Message msg=(Message)evt.getArg();
                   if(msg.isFlagSet(Message.NO_FC))
                       break;
  @@ -457,7 +434,7 @@
       }
   
       
  -    protected abstract Object handleDownMessage(final Event evt, final Message msg, int length);
  +    protected abstract Object handleDownMessage(final Event evt, final Message msg, Address dest, int length);
   
   
   
  @@ -470,16 +447,11 @@
        * @return long Number of credits to be sent. Greater than 0 if credits needs to be sent, 0 otherwise
        */
       protected long adjustCredit(Map<Address,Credit> map, Address sender, int length) {
  -        if(sender == null || length == 0)
  -            return 0;
  -
  -        Credit cred=map.get(sender);
  -        if(cred == null)
  +        Credit cred;
  +        if(sender == null || length == 0 || (cred=map.get(sender)) == null)
               return 0;
  -
           if(log.isTraceEnabled())
               log.trace(sender + " used " + length + " credits, " + (cred.get() - length) + " remaining");
  -
           return cred.decrementAndGet(length);
       }
   
  @@ -530,30 +502,23 @@
   
   
       protected void handleViewChange(Vector<Address> mbrs) {
  -        Address addr;
           if(mbrs == null) return;
           if(log.isTraceEnabled()) log.trace("new membership: " + mbrs);
   
  -
           // add members not in membership to received and sent hashmap (with full credits)
  -        for(int i=0; i < mbrs.size(); i++) {
  -            addr=mbrs.elementAt(i);
  +        for(Address addr: mbrs) {
               if(!received.containsKey(addr))
                   received.put(addr, new Credit(max_credits));
           }
           // remove members that left
           for(Iterator<Address> it=received.keySet().iterator(); it.hasNext();) {
  -            addr=it.next();
  +            Address addr=it.next();
               if(!mbrs.contains(addr))
                   it.remove();
           }
       }
   
  -    
  -    protected static long computeLowestCredit(Map<Address,Credit> m) {
  -        Collection<Credit> credits=m.values();
  -        return Collections.min(credits).get();
  -    }
  +
   
       protected static String printMap(Map<Address,Credit> m) {
           StringBuilder sb=new StringBuilder();
  @@ -565,7 +530,7 @@
   
   
   
  -    protected class Credit implements Comparable {
  +    protected class Credit {
           protected long credits_left;
           protected int  num_blockings=0;
           protected long total_blocking_time=0;
  @@ -578,10 +543,8 @@
   
   
           protected synchronized boolean decrementIfEnoughCredits(long credits, long timeout) {
  -            if(credits <= credits_left) {
  -                credits_left-=credits;
  +            if(decrement(credits))
                   return true;
  -            }
   
               if(timeout <= 0)
                   return false;
  @@ -597,6 +560,11 @@
                   num_blockings++;
               }
   
  +            return decrement(credits);
  +        }
  +
  +        
  +        protected boolean decrement(long credits) {
               if(credits <= credits_left) {
                   credits_left-=credits;
                   return true;
  @@ -616,10 +584,9 @@
           }
   
   
  -        protected synchronized long increment(long credits) {
  -            long retval=credits_left=Math.min(max_credits, credits_left + credits);
  +        protected synchronized void increment(long credits) {
  +            credits_left=Math.min(max_credits, credits_left + credits);
               notifyAll();
  -            return retval;
           }
   
           protected synchronized boolean needToSendCreditRequest() {
  @@ -637,17 +604,15 @@
   
           protected synchronized long get() {return credits_left;}
   
  -        protected synchronized void set(long new_credits) {credits_left=Math.min(max_credits, new_credits);}
  -
  +        protected synchronized void set(long new_credits) {
  +            credits_left=Math.min(max_credits, new_credits);
  +            notifyAll();
  +        }
   
           public String toString() {
               return String.valueOf(credits_left);
           }
   
  -        public int compareTo(Object o) {
  -            Credit other=(Credit)o;
  -            return credits_left < other.credits_left ? -1 : credits_left > other.credits_left ? 1 : 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
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.