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/24 12:15:41

  Modified:    src/org/jgroups/protocols DAISYCHAIN.java
  Log:
  removed sent and forward queue
  
  Revision  Changes    Path
  1.10      +7 -44     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.9
  retrieving revision 1.10
  diff -u -r1.9 -r1.10
  --- DAISYCHAIN.java	24 Aug 2010 10:07:53 -0000	1.9
  +++ DAISYCHAIN.java	24 Aug 2010 12:15:41 -0000	1.10
  @@ -8,8 +8,6 @@
   import java.io.DataInputStream;
   import java.io.DataOutputStream;
   import java.io.IOException;
  -import java.util.Queue;
  -import java.util.concurrent.ConcurrentLinkedQueue;
   import java.util.concurrent.Executor;
   
   /**
  @@ -20,7 +18,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.9 2010/08/24 10:07:53 belaban Exp $
  + * @version $Id: DAISYCHAIN.java,v 1.10 2010/08/24 12:15:41 belaban Exp $
    */
   @Experimental @Unsupported
   @MBean(description="Protocol just above the transport which disseminates multicasts via daisy chaining")
  @@ -31,18 +29,9 @@
       @Property(description="Loop back multicast messages")
       boolean loopback=true;
   
  -    @Property(description="Max number of messages in the send queue. The adder will block until more space is available")
  -    int send_queue_max_size=1000;
  -
  -    @Property(description="Max number of messages in the forward queue. The adder will block until more space is available")
  -    int forward_queue_max_size=1000;
  -
       /* --------------------------------------------- Fields ------------------------------------------------------ */
       protected Address              local_addr, next;
       protected int                  view_size=0;
  -    protected final Queue<Message> send_queue=new ConcurrentLinkedQueue<Message>();
  -    protected final Queue<Message> forward_queue=new ConcurrentLinkedQueue<Message>();
  -    protected volatile boolean     forward=false; // flipped between true and false, to ensure fairness
       protected Executor             default_pool=null;
       protected Executor             oob_pool=null;
   
  @@ -53,11 +42,6 @@
       @ManagedAttribute
       public int msgs_sent=0;
   
  -    @ManagedAttribute
  -    public int getForwardQueueSize() {return forward_queue.size();}
  -
  -    @ManagedAttribute
  -    public int getSendQueueSize() {return send_queue.size();}
   
   
       public void init() throws Exception {
  @@ -83,7 +67,6 @@
                   DaisyHeader hdr=new DaisyHeader(hdr_ttl);
                   copy.setDest(next);
                   copy.putHeader(getId(), hdr);
  -                send_queue.offer(copy);
   
                   if(loopback) {
                       if(log.isTraceEnabled()) log.trace(new StringBuilder("looping back message ").append(msg));
  @@ -97,7 +80,8 @@
                       });
                   }
   
  -                return forward();
  +                msgs_sent++;
  +                return forward(copy);
   
   
               case Event.VIEW_CHANGE:
  @@ -132,8 +116,8 @@
                       Message copy=msg.copy(true);
                       copy.setDest(next);
                       copy.putHeader(getId(), new DaisyHeader(ttl));
  -                    forward_queue.offer(copy);
  -                    forward();
  +                    forward(copy);
  +                    msgs_forwarded++;
                   }
   
                   // 2. Pass up
  @@ -144,34 +128,13 @@
       }
   
   
  -    protected Object forward() {
  -        Message        msg=null;
  -        Queue<Message> queue=null;
  -        String         tmp=null;
  -
  -        while(!(send_queue.isEmpty() && forward_queue.isEmpty())) {
  -            tmp=forward? " forwarding" : " sending";
  -            queue=forward? forward_queue : send_queue;
  -            if(queue.isEmpty()) {
  -                queue=forward? send_queue : forward_queue;
  -                msgs_sent++;
  -            }
  -            else {
  -                msgs_forwarded++;
  -            }
  -
  -            msg=queue.poll();
  -            if(msg != null)
  -                break;
  -            forward=!forward;
  -        }
  -
  +    protected Object forward(Message msg) {
           if(msg == null)
               return null;
   
           if(log.isTraceEnabled()) {
               DaisyHeader hdr=(DaisyHeader)msg.getHeader(getId());
  -            log.trace(local_addr + ": " + tmp + " message with ttl=" + hdr.getTTL() + " to " + next);
  +            log.trace(local_addr + ": forwarding message with ttl=" + hdr.getTTL() + " to " + next);
           }
           return down_prot.down(new Event(Event.MSG, msg));
       }
  
  
  

------------------------------------------------------------------------------
Sell apps to millions through the Intel(R) Atom(Tm) Developer Program
Be part of this innovative community and reach millions of netbook users 
worldwide. Take advantage of special opportunities to increase revenue and 
speed time-to-market. Join now, and jumpstart your future.
http://p.sf.net/sfu/intel-atom-d2d
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.