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

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

  Modified:    src/org/jgroups/protocols RELAY.java
  Log:
  added async prop
  
  Revision  Changes    Path
  1.2       +49 -39    JGroups/src/org/jgroups/protocols/RELAY.java
  
  Index: RELAY.java
  ===================================================================
  RCS file: /cvsroot/javagroups/JGroups/src/org/jgroups/protocols/RELAY.java,v
  retrieving revision 1.1
  retrieving revision 1.2
  diff -u -r1.1 -r1.2
  --- RELAY.java	18 Aug 2010 09:46:12 -0000	1.1
  +++ RELAY.java	18 Aug 2010 10:02:15 -0000	1.2
  @@ -5,6 +5,7 @@
   import org.jgroups.annotations.Property;
   import org.jgroups.annotations.Unsupported;
   import org.jgroups.stack.Protocol;
  +import org.jgroups.util.TimeScheduler;
   import org.jgroups.util.Util;
   
   import java.io.DataInputStream;
  @@ -15,7 +16,7 @@
    * Simple relaying protocol: RELAY is added to the top of the stack, creates a channel to a bridge cluster,
    * and - if coordinator - relays all multicast messages via the bridge cluster to the remote relay.  
    * @author Bela Ban
  - * @version $Id: RELAY.java,v 1.1 2010/08/18 09:46:12 belaban Exp $
  + * @version $Id: RELAY.java,v 1.2 2010/08/18 10:02:15 belaban Exp $
    */
   @Experimental @Unsupported
   public class RELAY extends Protocol {
  @@ -26,35 +27,28 @@
       @Property(description="Name of the bridge cluster")
       protected String cluster_name="bridge-cluster";
   
  +    @Property(description="If true, messages are relayed asynchronously, ie. via submission of a task to the timer thread pool")
  +    protected boolean async=false;
  +
       protected Address local_addr;
   
       protected volatile boolean is_coord=false;
   
       protected JChannel ch;
   
  +    protected TimeScheduler timer;
   
  -    public Object down(Event evt) {
  -        switch(evt.getType()) {
   
  -            case Event.MSG:
  -                if(is_coord && ch != null) {
  -                    Message msg=(Message)evt.getArg();
  -                    Address dest=msg.getDest();
  -                    if(dest == null || dest.isMulticastAddress()) {
  -                        Message copy=msg.copy(true, false);
  -                        copy.setSrc(local_addr);
  -                        try {
  -                            if(log.isTraceEnabled())
  -                                log.trace("down(): relaying message from " + copy.getSrc());
  -                            ch.send(copy);
  -                        }
  -                        catch(Throwable e) {
  -                            log.error("failed forwarding message " + copy, e);
  -                        }
  -                    }
  -                }
  -                break;
  +    public void init() throws Exception {
  +        timer=getTransport().getTimer();
  +    }
  +
  +    public void stop() {
  +        Util.close(ch);
  +    }
   
  +    public Object down(Event evt) {
  +        switch(evt.getType()) {
               case Event.VIEW_CHANGE:
                   handleView((View)evt.getArg());
                   break;
  @@ -67,28 +61,21 @@
                   local_addr=(Address)evt.getArg();
                   break;
           }
  -
  -
  -
           return down_prot.down(evt);
       }
   
  +
       public Object up(Event evt) {
           switch(evt.getType()) {
               case Event.MSG:
                   if(is_coord && ch != null) {
                       Message msg=(Message)evt.getArg();
                       Address dest=msg.getDest();
  -                    if((dest == null || dest.isMulticastAddress()) && !local_addr.equals(msg.getSrc())) {
  -                        Message copy=msg.copy(true, false);
  -                        try {
  -                            if(log.isTraceEnabled())
  -                                log.trace("up(): relaying message from " + copy.getSrc());
  -                            ch.send(copy);
  -                        }
  -                        catch(Throwable e) {
  -                            log.error("failed forwarding message " + copy, e);
  -                        }
  +                    RelayHeader hdr=(RelayHeader)msg.getHeader(getId());
  +                    boolean multicast=dest == null || dest.isMulticastAddress();
  +
  +                    if(multicast && hdr == null) {
  +                        relay(msg);
                       }
                   }
                   break;
  @@ -96,14 +83,10 @@
                   handleView((View)evt.getArg());
                   break;
           }
  -
  -
           return up_prot.up(evt);
       }
   
  -    public void stop() {
  -        Util.close(ch);
  -    }
  +
   
   
       protected void handleView(View view) {
  @@ -134,10 +117,37 @@
           }
       }
   
  +    protected void relay(Message msg) {
  +        final Message copy=msg.copy(true, false);
  +        if(async) {
  +            timer.execute(new Runnable() {
  +                public void run() {
  +                    _relay(copy);
  +                }
  +            });
  +        }
  +        else {
  +            _relay(copy);
  +        }
  +    }
  +
  +
  +    protected void _relay(Message msg) {
  +        try {
  +            if(log.isTraceEnabled())
  +                log.trace("relaying message from " + msg.getSrc());
  +            ch.send(msg);
  +        }
  +        catch(Throwable e) {
  +            log.error("failed relaying message " + msg, e);
  +        }
  +    }
  +
   
       protected class Receiver extends ReceiverAdapter {
           public void receive(Message msg) {
               Message copy=msg.copy(true, false); // copy the payload and everything else but the headers
  +            copy.putHeader(getId(), new RelayHeader());
               copy.setSrc(local_addr);
               if(log.isTraceEnabled())
                   log.trace("received msg from " + msg.getSrc() + ", passing down the stack with dest=" +
  
  
  

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