CVS update: JGroups/src/org/jgroups/protocols DAISYCHAIN.java

"Bela Ban" <[email protected]>
Newsgroups gmane.comp.java.javagroups.cvs
Message-ID <[email protected]>
  User: belaban 
  Date: 10/08/13 15:18:16

  Modified:    src/org/jgroups/protocols DAISYCHAIN.java
  Log:
  added queues to implement fairness
  
  Revision  Changes    Path
  1.5       +61 -13    JGroups/src/org/jgroups/protocols/DAISYCHAIN.java
  
  Index: DAISYCHAIN.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/DAISYCHAIN.java,v
  retrieving revision 1.4
  retrieving revision 1.5
  diff -u -r1.4 -r1.5
  --- DAISYCHAIN.java	13 Aug 2010 13:07:39 -0000	1.4
  +++ DAISYCHAIN.java	13 Aug 2010 15:18:16 -0000	1.5
  @@ -1,17 +1,18 @@
   package org.jgroups.protocols;
   
   import org.jgroups.*;
  -import org.jgroups.annotations.Experimental;
  -import org.jgroups.annotations.MBean;
  -import org.jgroups.annotations.Property;
  -import org.jgroups.annotations.Unsupported;
  +import org.jgroups.annotations.*;
   import org.jgroups.stack.Protocol;
  +import org.jgroups.util.ConcurrentLinkedBlockingQueue;
   import org.jgroups.util.Util;
   
   import java.io.DataInputStream;
   import java.io.DataOutputStream;
   import java.io.IOException;
  +import java.util.concurrent.BlockingQueue;
   import java.util.concurrent.Executor;
  +import java.util.concurrent.locks.Lock;
  +import java.util.concurrent.locks.ReentrantLock;
   
   /**
    * Implementation of daisy chaining. Multicast messages are sent to our neighbor, which sends them to its neighbor etc.
  @@ -21,7 +22,7 @@
    * send another message. This leads to much better throughput, see the ref in the JIRA.<p/> 
    * JIRA: https://jira.jboss.org/browse/JGRP-1021
    * @author Bela Ban
  - * @version $Id: DAISYCHAIN.java,v 1.4 2010/08/13 13:07:39 belaban Exp $
  + * @version $Id: DAISYCHAIN.java,v 1.5 2010/08/13 15:18:16 belaban Exp $
    */
   @Experimental @Unsupported
   @MBean(description="Protocol just above the transport which disseminates multicasts via daisy chaining")
  @@ -37,8 +38,23 @@
       protected Address local_addr, next;
       protected int     view_size=0;
   
  +    protected final BlockingQueue<Message> send_queue=new ConcurrentLinkedBlockingQueue<Message>(1000);
  +    protected final BlockingQueue<Message> forward_queue=new ConcurrentLinkedBlockingQueue<Message>(1000);
  +    protected boolean forward=false; // flipped between true and false, to ensure fairness
  +    protected final Lock lock=new ReentrantLock();
   
   
  +    @ManagedAttribute
  +    public int msgs_forwarded=0;
  +
  +    @ManagedAttribute
  +    public int msgs_sent=0;
  +
  +    @ManagedAttribute
  +    public int getForwardQueueSize() {return forward_queue.size();}
  +
  +    @ManagedAttribute
  +    public int getSendQueueSize() {return send_queue.size();}
   
   
       public Object down(final Event evt) {
  @@ -60,6 +76,14 @@
                   copy.setDest(next);
                   copy.putHeader(getId(), hdr);
   
  +                try {
  +                    send_queue.put(copy);
  +                }
  +                catch(InterruptedException e) {
  +                    Thread.currentThread().interrupt();
  +                    return null;
  +                }
  +
                   if(loopback) {
                       if(log.isTraceEnabled()) log.trace(new StringBuilder("looping back message ").append(msg));
                       msg.setSrc(local_addr);
  @@ -73,9 +97,8 @@
                       });
                   }
   
  -                if(log.isTraceEnabled())
  -                    log.trace(local_addr + ": forwarding message with ttl=" + hdr.getTTL() + " to " + next);
  -                return down_prot.down(new Event(Event.MSG, copy)); // don't pass down
  +                return forward();
  +
   
               case Event.VIEW_CHANGE:
                   handleView((View)evt.getArg());
  @@ -109,9 +132,12 @@
                       Message copy=msg.copy(true);
                       copy.setDest(next);
                       copy.putHeader(getId(), new DaisyHeader(ttl));
  -                    if(log.isTraceEnabled())
  -                        log.trace(local_addr + ": forwarding message with ttl=" + ttl + " to " + next);
  -                    down_prot.down(new Event(Event.MSG, copy));
  +                    try {
  +                        forward_queue.put(copy);
  +                    }
  +                    catch(InterruptedException e) {
  +                    }
  +                    forward();
                   }
   
                   // 2. Pass up
  @@ -122,13 +148,35 @@
       }
   
   
  +    protected Object forward() {
  +        lock.lock();
  +        try {
  +            Message msg=forward? forward_queue.poll() : send_queue.poll();
  +            if(msg == null) {
  +                msg=forward? send_queue.poll() : forward_queue.poll();
  +                msgs_sent++;
  +            }
  +            else {
  +                msgs_forwarded++;
  +            }
  +            if(log.isTraceEnabled()) {
  +                DaisyHeader hdr=(DaisyHeader)msg.getHeader(getId());
  +                log.trace(local_addr + ": " + (forward? " forwarding" : " sending") + " message with ttl=" + hdr.getTTL() + " to " + next);
  +            }
  +            return down_prot.down(new Event(Event.MSG, msg));
  +        }
  +        finally {
  +            forward=!forward;
  +            lock.unlock();
  +        }
  +    }
       
   
       protected void handleView(View view) {
           view_size=view.size();
           next=Util.pickNext(view.getMembers(), local_addr);
  -        if(log.isTraceEnabled())
  -            log.trace("next=" + next);
  +        if(log.isDebugEnabled())
  +            log.debug("next=" + next);
       }
   
   
  
  
  

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

Make an app they can't live without
Enter the BlackBerry Developer Challenge
http://p.sf.net/sfu/RIM-dev2dev
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.