[jboss-cvs] JBoss Messaging SVN: r4934 - in trunk: src/main/org/jboss/messaging/core/client/impl and 5 other directories.

jboss-cvs-commits at lists.jboss.org jboss-cvs-commits at lists.jboss.org
Thu Sep 11 05:44:31 EDT 2008


Author: timfox
Date: 2008-09-11 05:44:30 -0400 (Thu, 11 Sep 2008)
New Revision: 4934

Modified:
   trunk/src/main/org/jboss/messaging/core/client/ClientSessionFactory.java
   trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionFactoryImpl.java
   trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionImpl.java
   trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionPacketHandler.java
   trunk/src/main/org/jboss/messaging/core/remoting/RemotingService.java
   trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingConnectionImpl.java
   trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingServiceImpl.java
   trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java
   trunk/src/main/org/jboss/messaging/core/server/impl/ServerSessionImpl.java
   trunk/tests/jms-tests/src/org/jboss/test/messaging/jms/JMSTest.java
   trunk/tests/src/org/jboss/messaging/tests/integration/cluster/ReplicationTest.java
Log:
Fixed test etc


Modified: trunk/src/main/org/jboss/messaging/core/client/ClientSessionFactory.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/client/ClientSessionFactory.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/client/ClientSessionFactory.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -89,4 +89,6 @@
 
    void setCallTimeout(final long callTimeout);
    
+   boolean isFailedOver();
+   
 }

Modified: trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionFactoryImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionFactoryImpl.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionFactoryImpl.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -33,6 +33,7 @@
 import org.jboss.messaging.core.remoting.Channel;
 import org.jboss.messaging.core.remoting.ChannelHandler;
 import org.jboss.messaging.core.remoting.ConnectionRegistry;
+import org.jboss.messaging.core.remoting.FailureListener;
 import org.jboss.messaging.core.remoting.Packet;
 import org.jboss.messaging.core.remoting.RemotingConnection;
 import org.jboss.messaging.core.remoting.impl.ConnectionRegistryImpl;
@@ -53,7 +54,7 @@
  * @version <tt>$Revision: 3602 $</tt>
  *
  */
-public class ClientSessionFactoryImpl implements ClientSessionFactory
+public class ClientSessionFactoryImpl implements ClientSessionFactory, FailureListener
 {
    // Constants ------------------------------------------------------------------------------------
 
@@ -110,6 +111,8 @@
    private volatile boolean blockOnPersistentSend;
    
    private volatile boolean blockOnNonPersistentSend;
+   
+   private volatile boolean failedOver;
         
    // Static ---------------------------------------------------------------------------------------
    
@@ -127,7 +130,7 @@
                                    final boolean blockOnAcknowledge,
                                    final boolean blockOnNonPersistentSend,
                                    final boolean blockOnPersistentSend)
-   {      
+   {           
       this.connectorFactory = instantiateConnectorFactory(connectorConfig.getFactoryClassName());
       this.transportParams = connectorConfig.getParams();
       if (backupConfig != null)
@@ -149,7 +152,7 @@
    
    public ClientSessionFactoryImpl(final TransportConfiguration connectorConfig,
                                    final TransportConfiguration backupConfig)
-   {      
+   {            
       this.connectorFactory = instantiateConnectorFactory(connectorConfig.getFactoryClassName());
       this.transportParams = connectorConfig.getParams();
       if (backupConfig != null)
@@ -337,6 +340,11 @@
    {
       this.callTimeout = callTimeout;
    }
+   
+   public boolean isFailedOver()
+   {
+      return failedOver;
+   }
             
    // Public ---------------------------------------------------------------------------------------
    
@@ -351,6 +359,21 @@
 
    // Private --------------------------------------------------------------------------------------
    
+   private void handleFailover(final MessagingException me)
+   {
+      log.info(this + " Factory Failure has been detected, initiating failover");
+      if (backupConnectorFactory == null)
+      {
+         throw new IllegalStateException("Cannot fail-over if backup connector factory is null");
+      }
+                  
+      this.connectorFactory = backupConnectorFactory;
+      this.transportParams = backupTransportParams;
+      
+      this.backupConnectorFactory = null;
+      this.backupTransportParams = null;               
+   }
+   
    private ConnectorFactory instantiateConnectorFactory(final String connectorFactoryClassName)
    {
       ClassLoader loader = Thread.currentThread().getContextClassLoader();
@@ -378,9 +401,11 @@
       {
          remotingConnection = connectionRegistry.getConnection(connectorFactory, transportParams,
                                                                pingPeriod, callTimeout);
-         
+                           
          if (backupConnectorFactory != null)
          {
+            remotingConnection.addFailureListener(this);
+            
             backupConnection = connectionRegistry.getConnection(backupConnectorFactory, backupTransportParams,
                      pingPeriod, callTimeout);
          }
@@ -449,5 +474,9 @@
    }
 
    
-   // Inner Classes --------------------------------------------------------------------------------
+   public void connectionFailed(final MessagingException me)
+   {
+      handleFailover(me);
+   }
+   
 }

Modified: trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionImpl.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionImpl.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -1076,7 +1076,7 @@
  
    private void handleFailover(final MessagingException me)
    {      
-      log.info("Failure has been detected, initiating failover");
+      log.info("Session Failure has been detected, initiating failover");
       
       channel.lock();           
       
@@ -1084,6 +1084,11 @@
       {
          Packet request = new ReattachSessionMessage(channel.getID(), channel.getLastReceivedCommandID());
          
+         //This is necessary for invm since the replicating connection will be the same connection
+         //as the original replicating connection since the key is the same in the registry, and that connection
+         //won't have any resend buffer etc
+         backupConnection.setBackup(false);
+         
          Channel channel1 = backupConnection.getChannel(1, false, -1);
          
          ReattachSessionResponseMessage response = (ReattachSessionResponseMessage)channel1.sendBlocking(request);             

Modified: trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionPacketHandler.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionPacketHandler.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/client/impl/ClientSessionPacketHandler.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -47,17 +47,14 @@
    private final ClientSessionInternal clientSession;
 
    public ClientSessionPacketHandler(final ClientSessionInternal clientSesssion)
-   {  
-     // log.info("creating clientsessionpacketHandler " + System.identityHashCode(this));
+   {     
       this.clientSession = clientSesssion;
    }
 
    public void handlePacket(final Packet packet)
    {
       byte type = packet.getType();
-      
-     // log.info(System.identityHashCode(this) + "handling packet");
-      
+       
       try
       {
          switch (type)

Modified: trunk/src/main/org/jboss/messaging/core/remoting/RemotingService.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/remoting/RemotingService.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/remoting/RemotingService.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -44,4 +44,6 @@
    void addInterceptor(Interceptor interceptor);
    
    boolean removeInterceptor(Interceptor interceptor);
+   
+   void setBackup(boolean backup);
 }

Modified: trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingConnectionImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingConnectionImpl.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingConnectionImpl.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -1,23 +1,23 @@
 /*
- * JBoss, Home of Professional Open Source
- * Copyright 2005-2008, Red Hat Middleware LLC, and individual contributors
- * by the @authors tag. See the copyright.txt in the distribution for a
- * full listing of individual contributors.
- *
- * This is free software; you can redistribute it and/or modify it
- * under the terms of the GNU Lesser General Public License as
- * published by the Free Software Foundation; either version 2.1 of
- * the License, or (at your option) any later version.
- *
- * This software is distributed in the hope that it will be useful,
- * but WITHOUT ANY WARRANTY; without even the implied warranty of
- * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
- * Lesser General Public License for more details.
- *
- * You should have received a copy of the GNU Lesser General Public
- * License along with this software; if not, write to the Free
- * Software Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA
- * 02110-1301 USA, or see the FSF site: http://www.fsf.org.
+ * JBoss, Home of Professional Open Source Copyright 2005-2008, Red Hat
+ * Middleware LLC, and individual contributors by the @authors tag. See the
+ * copyright.txt in the distribution for a full listing of individual
+ * contributors.
+ * 
+ * This is free software; you can redistribute it and/or modify it under the
+ * terms of the GNU Lesser General Public License as published by the Free
+ * Software Foundation; either version 2.1 of the License, or (at your option)
+ * any later version.
+ * 
+ * This software is distributed in the hope that it will be useful, but WITHOUT
+ * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS
+ * FOR A PARTICULAR PURPOSE. See the GNU Lesser General Public License for more
+ * details.
+ * 
+ * You should have received a copy of the GNU Lesser General Public License
+ * along with this software; if not, write to the Free Software Foundation,
+ * Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA, or see the FSF
+ * site: http://www.fsf.org.
  */
 
 package org.jboss.messaging.core.remoting.impl;
@@ -81,8 +81,10 @@
 import static org.jboss.messaging.core.remoting.impl.wireformat.PacketImpl.SESS_XA_SUSPEND;
 
 import java.util.ArrayList;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentLinkedQueue;
 import java.util.concurrent.Executor;
@@ -165,13 +167,15 @@
  * @version <tt>$Revision$</tt> $Id: RemotingConnectionImpl.java 4633
  *          2008-07-04 11:43:34Z timfox $
  */
-public class RemotingConnectionImpl extends AbstractBufferHandler implements RemotingConnection
+public class RemotingConnectionImpl extends AbstractBufferHandler implements
+         RemotingConnection
 {
    // Constants
    // ------------------------------------------------------------------------------------
 
-   private static final Logger log = Logger.getLogger(RemotingConnectionImpl.class);
-   
+   private static final Logger log = Logger
+            .getLogger(RemotingConnectionImpl.class);
+
    private static final float EXPIRE_FACTOR = 1.5f;
 
    // Static
@@ -181,63 +185,62 @@
    // -----------------------------------------------------------------------------------
 
    private final Connection transportConnection;
-   
+
    private final Map<Long, ChannelImpl> channels = new ConcurrentHashMap<Long, ChannelImpl>();
 
-   private final List<FailureListener> failureListeners = new ArrayList<FailureListener>();
+   private final Set<FailureListener> failureListeners = new HashSet<FailureListener>();
 
    private final long blockingCallTimeout;
-   
+
    private final ExecutorFactory executorFactory;
 
    private Runnable pinger;
-   
+
    private final List<Interceptor> interceptors;
-   
+
    private ScheduledFuture<?> future;
 
    private boolean firstTime = true;
 
    private volatile boolean gotPong;
-  
+
    private volatile boolean destroyed;
-   
+
    private long expirePeriod;
-   
+
    private volatile boolean stopPinging;
-   
+
    private volatile long expireTime = -1;
-         
+
    private final Channel pingChannel;
-   
+
    private final RemotingConnection replicatingConnection;
-   
+
    private volatile boolean backup;
-   
+
    private final boolean client;
-   
+
    private boolean writePackets;
-      
+
    // Constructors
    // ---------------------------------------------------------------------------------
 
    private final long pingPeriod;
-   
+
    private final ScheduledExecutorService pingExecutor;
-   
-   public RemotingConnectionImpl(final Connection transportConnection,               
-                                 final long blockingCallTimeout, final long pingPeriod,
-                                 final ExecutorService handlerExecutor,
-                                 final ScheduledExecutorService pingExecutor,
-                                 final List<Interceptor> interceptors,
-                                 final RemotingConnection replicatingConnection,
-                                 final boolean client)
-                                 
+
+   public RemotingConnectionImpl(final Connection transportConnection,
+            final long blockingCallTimeout, final long pingPeriod,
+            final ExecutorService handlerExecutor,
+            final ScheduledExecutorService pingExecutor,
+            final List<Interceptor> interceptors,
+            final RemotingConnection replicatingConnection, final boolean client)
+
    {
       this.transportConnection = transportConnection;
 
       this.blockingCallTimeout = blockingCallTimeout;
-      
+
       if (handlerExecutor != null)
       {
          this.executorFactory = new OrderedExecutorFactory(handlerExecutor);
@@ -248,42 +251,42 @@
       }
 
       this.interceptors = interceptors;
-      
+
       this.replicatingConnection = replicatingConnection;
-      
+
       this.client = client;
-      
+
       this.writePackets = client || !backup;
-      
+
       this.pingPeriod = pingPeriod;
-      
-      this.pingExecutor = pingExecutor;     
-                  
-      //Channel zero is reserved for pinging
+
+      this.pingExecutor = pingExecutor;
+
+      // Channel zero is reserved for pinging
       pingChannel = getChannel(0, false, -1);
-      
+
       ChannelHandler ppHandler = new PingPongHandler();
-      
-      pingChannel.setHandler(ppHandler);           
+
+      pingChannel.setHandler(ppHandler);
    }
-   
+
    public void startPinger()
    {
       if (pingPeriod != -1)
-      {   
+      {
          pinger = new Pinger();
-   
-         expirePeriod = (long)(EXPIRE_FACTOR * pingPeriod);
-         
+
+         expirePeriod = (long) (EXPIRE_FACTOR * pingPeriod);
+
          future = pingExecutor.scheduleWithFixedDelay(pinger, 0, pingPeriod,
-                                                      TimeUnit.MILLISECONDS);
+                  TimeUnit.MILLISECONDS);
       }
       else
       {
          pinger = null;
       }
    }
-   
+
    // RemotingConnection implementation
    // ------------------------------------------------------------
 
@@ -291,30 +294,31 @@
    {
       return transportConnection.getID();
    }
-   
-   public synchronized Channel getChannel(final long channelID, final boolean ordered,
-                                          final int packetConfirmationBatchSize)
-   {      
+
+   public synchronized Channel getChannel(final long channelID,
+            final boolean ordered, final int packetConfirmationBatchSize)
+   {
       ChannelImpl channel = channels.get(channelID);
-      
+
       if (channel == null)
       {
-         channel = new ChannelImpl(this, channelID, ordered, packetConfirmationBatchSize);
-         
+         channel = new ChannelImpl(this, channelID, ordered,
+                  packetConfirmationBatchSize);
+
          channels.put(channelID, channel);
       }
-      
+
       return channel;
    }
-   
-   //This is a bit hacky - can we somehow do this in the constructor?
+
+   // This is a bit hacky - can we somehow do this in the constructor?
    public void setBackup(final boolean backup)
-   {      
+   {
       this.backup = backup;
-      
+
       this.writePackets = client || !backup;
    }
-   
+
    public boolean isBackup()
    {
       return backup;
@@ -323,13 +327,15 @@
    public synchronized void addFailureListener(final FailureListener listener)
    {
       checkDestroyed();
-      
+
       if (listener == null)
       {
          throw new IllegalStateException("FailureListener cannot be null");
       }
 
+     // log.info("Adding failure listener " + listener);
       failureListeners.add(listener);
+    //  log.info("There are now " + failureListeners.size());
    }
 
    public synchronized boolean removeFailureListener(final FailureListener listener)
@@ -345,21 +351,29 @@
    public MessagingBuffer createBuffer(final int size)
    {
       checkDestroyed();
-      
+
       return transportConnection.createBuffer(size);
    }
 
    public synchronized void fail(final MessagingException me)
    {
+      if (destroyed)
+      {
+         return;
+      }
+      
       log.warn(me.getMessage());
 
       destroy();
 
       // Then call the listeners
-      for (FailureListener listener : new ArrayList<FailureListener>(failureListeners))
+      Set<FailureListener> listenersClone = new HashSet<FailureListener>(failureListeners);
+     // log.info(" client " + client + " backup " + backup + " There are " + listenersClone.size() + " listeners");
+      for (FailureListener listener: listenersClone)
       {
          try
          {
+         //   log.info("*** calling failed on " + listener);
             listener.connectionFailed(me);
          }
          catch (Throwable t)
@@ -380,42 +394,46 @@
 
       if (future != null)
       {
-         future.cancel(false);                  
+         future.cancel(false);
       }
-      
+
       pingChannel.close();
-            
-      channels.clear();
 
       destroyed = true;
 
       // We close the underlying transport connection
       transportConnection.close();
    }
-   
+
    public boolean isExpired(final long now)
    {
       return expireTime != -1 && now >= expireTime;
    }
-   
+
    /* For testing only */
    public void stopPingingAfterOne()
    {
       stopPinging = true;
    }
 
-   // Buffer Handler implementation ----------------------------------------------------
-   
-   public void bufferReceived(final Object connectionID, final MessagingBuffer buffer)
+   // Buffer Handler implementation
+   // ----------------------------------------------------
+
+   public void bufferReceived(final Object connectionID,
+            final MessagingBuffer buffer)
    {
-      //checkDestroyed();
-      
+      if (destroyed)
+      {
+         // Ignore packets that might come in after connection is destroyed
+         return;
+      }
+
       final Packet packet = decode(buffer);
-      
+
       long channelID = packet.getChannelID();
-      
+
       ChannelImpl channel = channels.get(channelID);
-      
+
       if (channel == null)
       {
          if (packet.getType() == PacketImpl.SESS_PACKETS_CONFIRMED)
@@ -430,14 +448,15 @@
             return;
          }
          else
-         {                        
-            throw new IllegalArgumentException("Cannot handle packet " + packet + " no channel is registered with id " + channelID);
+         {
+            throw new IllegalArgumentException("Cannot handle packet " + packet
+                     + " no channel is registered with id " + channelID);
          }
       }
-            
-      channel.handlePacket(packet);            
+
+      channel.handlePacket(packet);
    }
-        
+
    // Package protected
    // ----------------------------------------------------------------------------
 
@@ -449,24 +468,23 @@
 
    private void checkDestroyed()
    {
-      if (destroyed)
-      {
-         throw new IllegalStateException("Connection is destroyed");
-      }
+      if (destroyed) { throw new IllegalStateException(System
+               .identityHashCode(this)
+               + " Connection is destroyed"); }
    }
-   
+
    private void doWrite(final Packet packet)
-   {      
+   {
       checkDestroyed();
-      
-      MessagingBuffer buffer = transportConnection.createBuffer(PacketImpl.INITIAL_BUFFER_SIZE);
 
+      MessagingBuffer buffer = transportConnection
+               .createBuffer(PacketImpl.INITIAL_BUFFER_SIZE);
+
       packet.encode(buffer);
 
       transportConnection.write(buffer);
    }
 
-   
    private Packet decode(final MessagingBuffer in)
    {
       byte packetType = in.getByte();
@@ -476,7 +494,7 @@
       switch (packetType)
       {
          case PING:
-         {            
+         {
             packet = new Ping();
             break;
          }
@@ -706,7 +724,7 @@
             break;
          }
          case SESS_START:
-         {            
+         {
             packet = new PacketImpl(PacketImpl.SESS_START);
             break;
          }
@@ -788,100 +806,110 @@
 
    // Inner classes
    // --------------------------------------------------------------------------------
+
+   //FIXME - improve locking on this class
    
-   //Needs to be static so we can re-assign it to another remotingconnection
+   // Needs to be static so we can re-assign it to another remotingconnection
    private static class ChannelImpl implements Channel
    {
       private final long id;
-      
+
       private final Executor executor;
-      
+
       private ChannelHandler handler;
-      
+
       private Packet response;
-      
+
       private final java.util.Queue<Packet> resendCache;
-      
+
       private final int packetConfirmationBatchSize;
-      
+
       private volatile int firstStoredCommandID;
-      
+
       private volatile int lastReceivedCommandID = -1;
-      
+
       private volatile int nextConfirmation;
-      
-      private final Channel replicatingChannel;
-            
+
+      private Channel replicatingChannel;
+
       private final ReadWriteLock lock = new ReentrantReadWriteLock(true);
-      
+
       private volatile RemotingConnectionImpl connection;
-      
-      private ChannelImpl(final RemotingConnectionImpl connection, final long id, final boolean ordered, final int packetConfirmationBatchSize)
-      {                  
+
+      private ChannelImpl(final RemotingConnectionImpl connection,
+               final long id, final boolean ordered,
+               final int packetConfirmationBatchSize)
+      {
          this.connection = connection;
-         
+
          this.id = id;
-         
+
          if (ordered && connection.executorFactory != null)
-         {              
-            executor = connection.executorFactory.getExecutor();            
+         {
+            executor = connection.executorFactory.getExecutor();
          }
          else
          {
             executor = null;
-         }                  
-         
+         }
+
          this.packetConfirmationBatchSize = packetConfirmationBatchSize;
-                  
-         if (packetConfirmationBatchSize != -1 && (connection.client && !connection.backup || !connection.client && connection.replicatingConnection == null))
+
+         if (packetConfirmationBatchSize != -1
+                  && (connection.client && !connection.backup || !connection.client
+                           && connection.replicatingConnection == null))
          {
             resendCache = new ConcurrentLinkedQueue<Packet>();
-            
+
             this.nextConfirmation = packetConfirmationBatchSize - 1;
          }
          else
          {
             resendCache = null;
          }
-         
+
          if (connection.replicatingConnection != null)
          {
-            replicatingChannel = connection.replicatingConnection.getChannel(id, ordered, -1);
-            
-            replicatingChannel.setHandler(new ReplicatedPacketsConfirmedChannelHandler());
+            // log.info("Getting replicating channel");
+            replicatingChannel = connection.replicatingConnection.getChannel(
+                     id, ordered, -1);
+
+            replicatingChannel
+                     .setHandler(new ReplicatedPacketsConfirmedChannelHandler());
          }
          else
          {
             replicatingChannel = null;
-         }         
+         }
       }
-      
+
       public long getID()
       {
          return id;
       }
-      
+
       public int getLastReceivedCommandID()
       {
          return lastReceivedCommandID;
       }
-         
+
       public void send(final Packet packet)
       {
          lock.readLock().lock();
-         
+
          try
          {
             packet.setChannelID(id);
-               
+
             if (resendCache != null)
             {
                addToCache(packet);
             }
-                 
-            if (connection.writePackets || packet.getType() == PacketImpl.SESS_PACKETS_CONFIRMED
+
+            if (connection.writePackets
+                     || packet.getType() == PacketImpl.SESS_PACKETS_CONFIRMED
                      || packet.getType() == PacketImpl.SESS_REPLICATE_DELIVERY_RESP)
-            {                
+            {
                connection.doWrite(packet);
             }
          }
@@ -890,66 +918,69 @@
             lock.readLock().unlock();
          }
       }
-      
+
       private final Object blockingLock = new Object();
 
-      public synchronized Packet sendBlocking(final Packet packet) throws MessagingException
+      public Packet sendBlocking(final Packet packet)
+               throws MessagingException
       {
          lock.readLock().lock();
-         
+
          try
          {
-            //For now we only allow one blocking request-response at a time per channel
-            //We can relax this but it will involve some kind of correlation id
+            // For now we only allow one blocking request-response at a time per
+            // channel
+            // We can relax this but it will involve some kind of correlation id
             synchronized (blockingLock)
             {
-               response = null;
-                        
-               packet.setChannelID(id);
-      
-               if (resendCache != null)
+               synchronized (this)
                {
-                  addToCache(packet);
-               }
-               
-               connection.doWrite(packet);
-               
-               long toWait = connection.blockingCallTimeout;
-               
-               long start = System.currentTimeMillis();
-      
-               while (response == null && toWait > 0)
-               {
-                  try
+                  response = null;
+   
+                  packet.setChannelID(id);
+   
+                  if (resendCache != null)
                   {
-                     wait(toWait);
+                     addToCache(packet);
                   }
-                  catch (InterruptedException e)
+   
+                  connection.doWrite(packet);
+   
+                  long toWait = connection.blockingCallTimeout;
+   
+                  long start = System.currentTimeMillis();
+   
+                  while (response == null && toWait > 0)
                   {
+                     try
+                     {
+                        wait(toWait);
+                     }
+                     catch (InterruptedException e)
+                     {
+                     }
+   
+                     long now = System.currentTimeMillis();
+   
+                     toWait -= now - start;
+   
+                     start = now;
                   }
-      
-                  long now = System.currentTimeMillis();
-      
-                  toWait -= now - start;
-      
-                  start = now;
+   
+                  if (response == null) { throw new IllegalStateException(
+                           "Timed out waiting for response"); }
+   
+                  if (response.getType() == PacketImpl.EXCEPTION)
+                  {
+                     MessagingExceptionMessage mem = (MessagingExceptionMessage) response;
+   
+                     throw mem.getException();
+                  }
+                  else
+                  {
+                     return response;
+                  }
                }
-               
-               if (response == null)
-               {
-                  throw new IllegalStateException("Timed out waiting for response");
-               }
-               
-               if (response.getType() == PacketImpl.EXCEPTION)
-               {
-                  MessagingExceptionMessage mem = (MessagingExceptionMessage)response;
-                  
-                  throw mem.getException();
-               }
-               else
-               {
-                  return response;
-               }
             }
          }
          finally
@@ -960,101 +991,104 @@
 
       public void setHandler(final ChannelHandler handler)
       {
+         // log.info("client " + connection.client + " backup " +
+         // connection.backup + " setting handler " +
+         // System.identityHashCode(handler));
          this.handler = handler;
       }
-      
+
       public void close()
       {
-         if (!connection.destroyed && connection.channels.remove(id) == null)
-         {
-            throw new IllegalArgumentException("Cannot find channel with id " + id + " to close");
-         }         
-         
+         if (!connection.destroyed && connection.channels.remove(id) == null) { throw new IllegalArgumentException(
+                  "Cannot find channel with id " + id + " to close"); }
+
          if (replicatingChannel != null)
          {
             replicatingChannel.close();
          }
-         
+
          if (resendCache != null)
          {
-//            log.info(System.identityHashCode(this) +  " backup:" + backup
-//                     + " client:" + client + " replicatingconn:" + replicatingConnection +
-//                     " pcbs:" + packetConfirmationBatchSize + " channelid:" + id + 
-//            " at close resend cache size is " + this.resendCache.size());
+            // log.info(System.identityHashCode(this) + " backup:" + backup
+            // + " client:" + client + " replicatingconn:" +
+            // replicatingConnection +
+            // " pcbs:" + packetConfirmationBatchSize + " channelid:" + id +
+            // " at close resend cache size is " + this.resendCache.size());
          }
       }
-      
+
       public Channel getReplicatingChannel()
       {
          return replicatingChannel;
       }
-      
+
       public void transferConnection(final RemotingConnection newConnection)
       {
          if (executor != null)
          {
-            //First wait for anything in the executor to complete
+            // First wait for anything in the executor to complete
             Future future = new Future();
-            
+
             executor.execute(future);
-            
+
             boolean ok = future.await(10000);
-            
-            if (!ok)
-            {
-               throw new IllegalStateException("Timed out waiting for executor to complete");
-            }
+
+            if (!ok) { throw new IllegalStateException(
+                     "Timed out waiting for executor to complete"); }
          }
-         
-         RemotingConnectionImpl rnewConnection = (RemotingConnectionImpl)newConnection;
-         
+
+         RemotingConnectionImpl rnewConnection = (RemotingConnectionImpl) newConnection;
+
          connection.channels.remove(id);
-         
+
          rnewConnection.channels.put(id, this);
-         
-         connection = rnewConnection;         
+
+         connection = rnewConnection;
+
+         replicatingChannel = null;
       }
 
       public int replayCommands(final int otherLastReceivedCommandID)
-      {         
+      {
          clearUpTo(otherLastReceivedCommandID);
-         
+
          Packet packet = null;
-         
+
          int count = 0;
-         
+
          while ((packet = resendCache.poll()) != null)
          {
             connection.doWrite(packet);
-            
+
             count++;
          }
 
          return this.lastReceivedCommandID;
       }
-      
+
       public void lock()
       {
          lock.writeLock().lock();
       }
-      
+
       public void unlock()
       {
          lock.writeLock().unlock();
       }
-      
+
       private void handlePacket(final Packet packet)
-      {                                          
-       //  log.info("handling packet client " + connection.client + " backup " + connection.backup);
+      {
+         // log.info("handling packet client " + connection.client + " backup "
+         // + connection.backup);
          if (packet.getType() == PacketImpl.SESS_PACKETS_CONFIRMED)
          {
             if (resendCache != null)
             {
-               final PacketsConfirmedMessage msg = (PacketsConfirmedMessage)packet;
-                              
+               final PacketsConfirmedMessage msg = (PacketsConfirmedMessage) packet;
+
                if (executor == null)
-               {                  
-                  clearUpTo(msg.getCommandID());  
+               {
+                  clearUpTo(msg.getCommandID());
                }
                else
                {
@@ -1062,10 +1096,10 @@
                   {
                      public void run()
                      {
-                        clearUpTo(msg.getCommandID());  
+                        clearUpTo(msg.getCommandID());
                      }
                   });
-               }              
+               }
             }
             else if (connection.replicatingConnection != null)
             {
@@ -1075,55 +1109,58 @@
             {
                handler.handlePacket(packet);
             }
-            
+
             return;
-         }  
+         }
          else
          {
-            if (replicatingChannel != null && packet.getType() != PacketImpl.PING)
-            {            
+            if (replicatingChannel != null
+                     && packet.getType() != PacketImpl.PING)
+            {
                replicatingChannel.send(packet);
             }
-                                               
+
             if (connection.interceptors != null)
             {
                for (Interceptor interceptor : connection.interceptors)
                {
                   try
                   {
-                     boolean callNext = interceptor.intercept(packet, connection);
-                     
+                     boolean callNext = interceptor.intercept(packet,
+                              connection);
+
                      if (!callNext)
                      {
-                        //abort
-                                       
+                        // abort
+
                         return;
                      }
                   }
                   catch (Throwable e)
                   {
-                     log.warn("Failure in calling interceptor: " + interceptor, e);
+                     log.warn("Failure in calling interceptor: " + interceptor,
+                              e);
                   }
                }
             }
-            
+
             if (packet.isResponse())
             {
                synchronized (this)
                {
                   response = packet;
-                  
-                  checkConfirmation(packet);                  
-                  
-                  notify();                                          
+
+                  checkConfirmation(packet);
+
+                  notify();
                }
-            }      
+            }
             else if (handler != null)
-            {              
+            {
                if (executor == null)
                {
                   checkConfirmation(packet);
-                                    
+
                   handler.handlePacket(packet);
                }
                else
@@ -1132,88 +1169,86 @@
                   {
                      public void run()
                      {
-                        checkConfirmation(packet);                        
-                        
+                        checkConfirmation(packet);
+
                         handler.handlePacket(packet);
                      }
                   });
-               }            
-            } 
+               }
+            }
             else
             {
-               checkConfirmation(packet);               
+               checkConfirmation(packet);
             }
-         }            
-      }   
-         
+         }
+      }
+
       private void checkConfirmation(final Packet packet)
-      {        
+      {
          if (packet.isUsesConfirmations() && resendCache != null)
-         {            
+         {
             lastReceivedCommandID++;
-                 
+
             if (lastReceivedCommandID == nextConfirmation)
             {
-               Packet confirmed = new PacketsConfirmedMessage(lastReceivedCommandID);
-               
+               Packet confirmed = new PacketsConfirmedMessage(
+                        lastReceivedCommandID);
+
                nextConfirmation += packetConfirmationBatchSize;
-               
+
                confirmed.setChannelID(id);
-                               
+
                connection.doWrite(confirmed);
-            }                        
-         }         
-      }      
-      
+            }
+         }
+      }
+
       private void addToCache(final Packet packet)
-      {               
-         resendCache.add(packet);    
+      {
+         resendCache.add(packet);
       }
-      
+
       private void clearUpTo(final int lastReceivedCommandID)
-      {                          
+      {
          int numberToClear = 1 + lastReceivedCommandID - firstStoredCommandID;
-         
-         if (numberToClear == -1)
-         {
-            throw new IllegalArgumentException("Invalid lastReceivedCommandID: " + lastReceivedCommandID);
-         }
-         
+
+         if (numberToClear == -1) { throw new IllegalArgumentException(
+                  "Invalid lastReceivedCommandID: " + lastReceivedCommandID); }
+
          for (int i = 0; i < numberToClear; i++)
          {
             Packet packet = resendCache.poll();
-            
-            if (packet == null)
-            {
-               throw new IllegalStateException("Can't find packet to clear");
-            }
+
+            if (packet == null) { throw new IllegalStateException(
+                     "Can't find packet to clear"); }
          }
 
          firstStoredCommandID += numberToClear;
       }
-      
-      private class ReplicatedPacketsConfirmedChannelHandler implements ChannelHandler
+
+      private class ReplicatedPacketsConfirmedChannelHandler implements
+               ChannelHandler
       {
          public void handlePacket(final Packet packet)
          {
             if (packet.getType() == SESS_PACKETS_CONFIRMED)
-            {               
-               //Send it straight back to the client
+            {
+               // Send it straight back to the client
                connection.doWrite(packet);
             }
             else if (packet.getType() == PacketImpl.SESS_REPLICATE_DELIVERY_RESP)
             {
-               //Send it straight to the server handler
+               // Send it straight to the server handler
                handler.handlePacket(packet);
             }
             else
             {
                throw new IllegalArgumentException("Invalid packet " + packet);
             }
-         }         
+         }
       }
-   }   
-      
+   }
+
    private class Pinger implements Runnable
    {
       public synchronized void run()
@@ -1221,31 +1256,32 @@
          if (!firstTime && !gotPong)
          {
             // Error - didn't get pong back
-            MessagingException me = new MessagingException(MessagingException.NOT_CONNECTED,
-                                                           "Did not receive pong from server");
+            MessagingException me = new MessagingException(
+                     MessagingException.NOT_CONNECTED,
+                     "Did not receive pong from server");
 
             fail(me);
          }
 
          gotPong = false;
-         
+
          firstTime = false;
 
          // Send ping
          Packet ping = new Ping(expirePeriod);
-         
+
          pingChannel.send(ping);
       }
    }
-   
+
    private class PingPongHandler implements ChannelHandler
    {
       public void handlePacket(final Packet packet)
       {
          byte type = packet.getType();
-         
+
          if (type == PONG)
-         {            
+         {
             gotPong = true;
 
             if (stopPinging)
@@ -1255,11 +1291,12 @@
          }
          else if (type == PING)
          {
-            expireTime = System.currentTimeMillis() + ((Ping)packet).getExpirePeriod();
-            
-            //Parameter is placeholder for future
+            expireTime = System.currentTimeMillis()
+                     + ((Ping) packet).getExpirePeriod();
+
+            // Parameter is placeholder for future
             Packet pong = new Pong(-1);
-            
+
             pingChannel.send(pong);
          }
          else

Modified: trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingServiceImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingServiceImpl.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/remoting/impl/RemotingServiceImpl.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -91,7 +91,7 @@
    
    private final BufferHandler bufferHandler = new DelegatingBufferHandler();
    
-   private final boolean backup;
+   private volatile boolean backup;
    
    private volatile MessagingServer server;
 
@@ -227,6 +227,11 @@
    {
       this.server = server;
    }
+   
+   public void setBackup(final boolean backup)
+   {
+      this.backup = backup;
+   }
 
    // ConnectionLifeCycleListener implementation -----------------------------------
 

Modified: trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/server/impl/MessagingServerImpl.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -387,8 +387,17 @@
          throw new IllegalArgumentException("Cannot find session with id " + sessionID + " to reattach");
       }
       
+      //This is necessary for invm since the replicating connection will be the same connection
+      //as the original replicating connection since the key is the same in the registry, and that connection
+      //won't have any resend buffer etc
+      connection.setBackup(false);
+      
       postOffice.setBackup(false);
       
+      configuration.setBackup(false);
+      
+      remotingService.setBackup(false);
+      
       //Reconnect the channel to the new connection
       session.transferConnection(connection);
       

Modified: trunk/src/main/org/jboss/messaging/core/server/impl/ServerSessionImpl.java
===================================================================
--- trunk/src/main/org/jboss/messaging/core/server/impl/ServerSessionImpl.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/src/main/org/jboss/messaging/core/server/impl/ServerSessionImpl.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -1267,6 +1267,9 @@
  
       remotingConnection.removeFailureListener(this);
       
+      //Destroy the old connection
+    //  remotingConnection.destroy();
+      
       remotingConnection = newConnection;
       
       remotingConnection.addFailureListener(this);

Modified: trunk/tests/jms-tests/src/org/jboss/test/messaging/jms/JMSTest.java
===================================================================
--- trunk/tests/jms-tests/src/org/jboss/test/messaging/jms/JMSTest.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/tests/jms-tests/src/org/jboss/test/messaging/jms/JMSTest.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -325,7 +325,7 @@
       {
 	      conn = cf.createConnection();
 
-	      Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE);
+	      final Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE);
 
 	      final MessageConsumer cons = session.createConsumer(queue1);
 
@@ -340,11 +340,14 @@
 	         {
 	            try
 	            {
-	               Message m = cons.receive(5000);
-	               if (m != null)
+	               synchronized (session)
 	               {
-	                  message.set(m);
-	                  latch.countDown();
+   	               Message m = cons.receive(5000);
+   	               if (m != null)
+   	               {
+   	                  message.set(m);
+   	                  latch.countDown();
+   	               }
 	               }
 	            }
 	            catch(Exception e)
@@ -355,13 +358,16 @@
 	         }
 	      }, "Receiving Thread").start();
 
-	      MessageProducer prod = session.createProducer(queue1);
-	      prod.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
+         synchronized (session)
+         {
+   	      MessageProducer prod = session.createProducer(queue1);
+   	      prod.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
+   
+   	      TextMessage m = session.createTextMessage("message one");
+   
+   	      prod.send(m);
+         }
 
-	      TextMessage m = session.createTextMessage("message one");
-
-	      prod.send(m);
-
 	      boolean gotMessage = latch.await(5000, TimeUnit.MILLISECONDS);
 	      assertTrue(gotMessage);
 	      TextMessage rm = (TextMessage) message.get();

Modified: trunk/tests/src/org/jboss/messaging/tests/integration/cluster/ReplicationTest.java
===================================================================
--- trunk/tests/src/org/jboss/messaging/tests/integration/cluster/ReplicationTest.java	2008-09-11 08:31:24 UTC (rev 4933)
+++ trunk/tests/src/org/jboss/messaging/tests/integration/cluster/ReplicationTest.java	2008-09-11 09:44:30 UTC (rev 4934)
@@ -121,12 +121,112 @@
       backupService.stop();
    }
    
-   public void testFailover() throws Exception
+   public void testFailoverSameConnectionFactory() throws Exception
    {             
       final SimpleString QUEUE = new SimpleString("CoreClientTestQueue");
       
       Configuration backupConf = new ConfigurationImpl();      
       backupConf.setSecurityEnabled(false);        
+      backupConf.setPacketConfirmationBatchSize(10);
+      Map<String, Object> backupParams = new HashMap<String, Object>();
+      backupParams.put(TransportConstants.SERVER_ID_PROP_NAME, 1);
+      backupConf.getAcceptorConfigurations().add(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMAcceptorFactory", backupParams));
+      backupConf.setBackup(true);                  
+      MessagingService backupService = MessagingServiceImpl.newNullStorageMessagingServer(backupConf);              
+      backupService.start();
+            
+      Configuration liveConf = new ConfigurationImpl();      
+      liveConf.setSecurityEnabled(false);    
+      liveConf.setPacketConfirmationBatchSize(10);
+      liveConf.getAcceptorConfigurations().add(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMAcceptorFactory"));
+      liveConf.setBackupConnectorConfiguration(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMConnectorFactory", backupParams));
+      MessagingService liveService = MessagingServiceImpl.newNullStorageMessagingServer(liveConf);              
+      liveService.start();
+            
+      ClientSessionFactory sf =
+         new ClientSessionFactoryImpl(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMConnectorFactory"),
+                  new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMConnectorFactory", backupParams));
+
+      ClientSession session = sf.createSession(false, true, true, -1, false);
+                  
+      session.createQueue(QUEUE, QUEUE, null, false, false);
+       
+      ClientProducer producer = session.createProducer(QUEUE);     
+      
+      final int numMessages = 1000;
+      
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage message = session.createClientMessage(JBossTextMessage.TYPE, false, 0,
+               System.currentTimeMillis(), (byte) 1);         
+         message.putIntProperty(new SimpleString("blah"), i);
+         message.getBody().putString("testINVMCoreClient");
+         message.getBody().flip();  
+         producer.send(message);
+      }
+      
+      RemotingConnection conn = ((ClientSessionImpl)session).getConnection();
+      
+      //Simulate failure on connection
+      conn.fail(new MessagingException(MessagingException.NOT_CONNECTED));
+                      
+      ClientConsumer consumer = session.createConsumer(QUEUE);
+      
+      session.start();
+      
+      for (int i = 0; i < numMessages / 2; i++)
+      {
+         ClientMessage message2 = consumer.receive();
+
+         assertEquals("testINVMCoreClient", message2.getBody().getString());
+         
+         session.acknowledge();
+         
+         //log.info("got message " + message2.getProperty(new SimpleString("blah")));
+      }
+
+      session.close();
+                  
+      log.info("** creating new one");
+      
+//      sf =
+//         new ClientSessionFactoryImpl(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMConnectorFactory", backupParams));
+            
+      session = sf.createSession(false, true, true, -1, false);
+      
+      consumer = session.createConsumer(QUEUE);
+      
+      session.start();
+      
+      for (int i = 0; i < numMessages / 2; i++)
+      {
+         ClientMessage message2 = consumer.receive();
+
+         assertEquals("testINVMCoreClient", message2.getBody().getString());
+         
+         session.acknowledge();
+         
+        // log.info("got message " + message2.getProperty(new SimpleString("blah")));
+      }
+      
+      ClientMessage message3 = consumer.receive(1000);
+      
+      assertNull(message3);
+      
+      liveService.stop();
+      backupService.stop();
+      
+  //    todo - do we need to failover connection factories too?????
+               
+               
+   }
+   
+   public void testFailoverChangeConnectionFactory() throws Exception
+   {             
+      final SimpleString QUEUE = new SimpleString("CoreClientTestQueue");
+      
+      Configuration backupConf = new ConfigurationImpl();      
+      backupConf.setSecurityEnabled(false);        
       backupConf.setPacketConfirmationBatchSize(1);
       Map<String, Object> backupParams = new HashMap<String, Object>();
       backupParams.put(TransportConstants.SERVER_ID_PROP_NAME, 1);
@@ -186,7 +286,9 @@
       }
 
       session.close();
-            
+                  
+      log.info("** creating new one");
+      
       sf =
          new ClientSessionFactoryImpl(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.invm.InVMConnectorFactory", backupParams));
             
@@ -218,7 +320,107 @@
                
                
    }
+   
+   public void testFailoverNetty() throws Exception
+   {             
+      final SimpleString QUEUE = new SimpleString("CoreClientTestQueue");
+      
+      Configuration backupConf = new ConfigurationImpl();      
+      backupConf.setSecurityEnabled(false);        
+      backupConf.setPacketConfirmationBatchSize(1);
+      Map<String, Object> backupParams = new HashMap<String, Object>();
+      backupParams.put(org.jboss.messaging.core.remoting.impl.netty.TransportConstants.PORT_PROP_NAME, 7654);
+      backupConf.getAcceptorConfigurations().add(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.netty.NettyAcceptorFactory", backupParams));
+      backupConf.setBackup(true);                  
+      MessagingService backupService = MessagingServiceImpl.newNullStorageMessagingServer(backupConf);              
+      backupService.start();
+            
+      Configuration liveConf = new ConfigurationImpl();      
+      liveConf.setSecurityEnabled(false);    
+      liveConf.setPacketConfirmationBatchSize(1);
+      liveConf.getAcceptorConfigurations().add(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.netty.NettyAcceptorFactory"));
+      liveConf.setBackupConnectorConfiguration(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.netty.NettyConnectorFactory", backupParams));
+      MessagingService liveService = MessagingServiceImpl.newNullStorageMessagingServer(liveConf);              
+      liveService.start();
+            
+      ClientSessionFactory sf =
+         new ClientSessionFactoryImpl(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.netty.NettyConnectorFactory"),
+                  new TransportConfiguration("org.jboss.messaging.core.remoting.impl.netty.NettyConnectorFactory", backupParams));
 
+      ClientSession session = sf.createSession(false, true, true, -1, false);
+                  
+      session.createQueue(QUEUE, QUEUE, null, false, false);
+       
+      ClientProducer producer = session.createProducer(QUEUE);     
+      
+      final int numMessages = 10;
+      
+      for (int i = 0; i < numMessages; i++)
+      {
+         ClientMessage message = session.createClientMessage(JBossTextMessage.TYPE, false, 0,
+               System.currentTimeMillis(), (byte) 1);         
+         message.putIntProperty(new SimpleString("blah"), i);
+         message.getBody().putString("testINVMCoreClient");
+         message.getBody().flip();  
+         producer.send(message);
+      }
+      
+      RemotingConnection conn = ((ClientSessionImpl)session).getConnection();
+      
+      //Simulate failure on connection
+      conn.fail(new MessagingException(MessagingException.NOT_CONNECTED));
+                      
+      ClientConsumer consumer = session.createConsumer(QUEUE);
+      
+      session.start();
+      
+      for (int i = 0; i < numMessages / 2; i++)
+      {
+         ClientMessage message2 = consumer.receive();
+
+         assertEquals("testINVMCoreClient", message2.getBody().getString());
+         
+         session.acknowledge();
+         
+         log.info("got message " + message2.getProperty(new SimpleString("blah")));
+      }
+
+      session.close();
+                  
+      log.info("** creating new one");
+      
+      sf =
+         new ClientSessionFactoryImpl(new TransportConfiguration("org.jboss.messaging.core.remoting.impl.netty.NettyConnectorFactory", backupParams));
+            
+      session = sf.createSession(false, true, true, -1, false);
+      
+      consumer = session.createConsumer(QUEUE);
+      
+      session.start();
+      
+      for (int i = 0; i < numMessages / 2; i++)
+      {
+         ClientMessage message2 = consumer.receive();
+
+         assertEquals("testINVMCoreClient", message2.getBody().getString());
+         
+         session.acknowledge();
+         
+         log.info("got message " + message2.getProperty(new SimpleString("blah")));
+      }
+      
+      ClientMessage message3 = consumer.receive(1000);
+      
+      assertNull(message3);
+      
+      liveService.stop();
+      backupService.stop();
+      
+  //    todo - do we need to failover connection factories too?????
+               
+               
+   }
+
    // Package protected ---------------------------------------------
 
    // Protected -----------------------------------------------------




More information about the jboss-cvs-commits mailing list