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 10:23:53

  Modified:    src/org/jgroups/protocols DAISYCHAIN.java
  Log:
  ns
  
  Revision  Changes    Path
  1.2       +35 -10    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.1
  retrieving revision 1.2
  diff -u -r1.1 -r1.2
  --- DAISYCHAIN.java	13 Aug 2010 09:48:24 -0000	1.1
  +++ DAISYCHAIN.java	13 Aug 2010 10:23:52 -0000	1.2
  @@ -2,12 +2,14 @@
   
   import org.jgroups.*;
   import org.jgroups.annotations.MBean;
  +import org.jgroups.annotations.Property;
   import org.jgroups.stack.Protocol;
   import org.jgroups.util.Util;
   
   import java.io.DataInputStream;
   import java.io.DataOutputStream;
   import java.io.IOException;
  +import java.util.concurrent.Executor;
   
   /**
    * Implementation of daisy chaining. Multicast messages are sent to our neighbor, which sends them to its neighbor etc.
  @@ -17,14 +19,15 @@
    * 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.1 2010/08/13 09:48:24 belaban Exp $
  + * @version $Id: DAISYCHAIN.java,v 1.2 2010/08/13 10:23:52 belaban Exp $
    */
   @MBean(description="Protocol just above the transport which disseminates multicasts via daisy chaining")
   public class DAISYCHAIN extends Protocol {
   
   
       /* -----------------------------------------    Properties     -------------------------------------------------- */
  -
  +    @Property(description="Loop back multicast messages")
  +    boolean loopback=true;
   
       
       /* --------------------------------------------- Fields ------------------------------------------------------ */
  @@ -35,7 +38,7 @@
   
   
   
  -    public Object down(Event evt) {
  +    public Object down(final Event evt) {
           switch(evt.getType()) {
               case Event.MSG:
                   Message msg=(Message)evt.getArg();
  @@ -49,17 +52,35 @@
                   // we need to copy the message, as we cannot do a msg.setSrc(next): the next retransmission
                   // would use 'next' as destination  !
                   Message copy=msg.copy(true);
  -                DaisyHeader hdr=new DaisyHeader((short)view_size);
  +                short hdr_ttl=(short)(loopback? view_size -1 : view_size);
  +                DaisyHeader hdr=new DaisyHeader(hdr_ttl);
                   copy.setDest(next);
                   copy.putHeader(getId(), hdr);
  +
  +                if(loopback) {
  +                    if(log.isTraceEnabled()) log.trace(new StringBuilder("looping back message ").append(msg));
  +
  +                    Executor pool=msg.isFlagSet(Message.OOB)? getTransport().getOOBThreadPool()
  +                            : getTransport().getDefaultThreadPool();
  +                    pool.execute(new Runnable() {
  +                        public void run() {
  +                            up_prot.up(evt);
  +                        }
  +                    });
  +                }
  +
                   if(log.isTraceEnabled())
  -                    log.trace(local_addr + ": forwarding message (" + copy.getObject() + ") 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, copy)); // don't pass down
   
               case Event.VIEW_CHANGE:
                   handleView((View)evt.getArg());
                   break;
   
  +            case Event.TMP_VIEW:
  +                view_size=((View)evt.getArg()).size();
  +                break;
  +            
               case Event.SET_LOCAL_ADDRESS:
                   local_addr=(Address)evt.getArg();
                   break;
  @@ -79,13 +100,14 @@
                   // 1. forward the message to the next in line if ttl > 0
                   short ttl=hdr.getTTL();
                   if(log.isTraceEnabled())
  -                    log.trace(local_addr + ": received message (" + msg.getObject() + ") from " + msg.getSrc() + " with ttl=" + ttl);
  +                    log.trace(local_addr + ": received message from " + msg.getSrc() + " with ttl=" + ttl);
                   if(--ttl > 0) {
  -                    msg.setDest(next);
  -                    hdr.setTTL(ttl);
  +                    Message copy=msg.copy(true);
  +                    copy.setDest(next);
  +                    copy.putHeader(getId(), new DaisyHeader(ttl));
                       if(log.isTraceEnabled())
  -                        log.trace(local_addr + ": forwarding message (" +msg.getObject() + ") with ttl=" + ttl + " to " + next);
  -                    down_prot.down(evt);
  +                        log.trace(local_addr + ": forwarding message with ttl=" + ttl + " to " + next);
  +                    down_prot.down(new Event(Event.MSG, copy));
                   }
   
                   // 2. Pass up
  @@ -95,6 +117,9 @@
           return up_prot.up(evt);
       }
   
  +
  +    
  +
       protected void handleView(View view) {
           view_size=view.size();
           next=Util.pickNext(view.getMembers(), local_addr);
  
  
  

------------------------------------------------------------------------------
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.