Author: ataylor
Date: 2012-02-17 08:44:48 -0500 (Fri, 17 Feb 2012)
New Revision: 12134
Modified:
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/distribution/ClusterTestBase.java
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/failover/ClusterWithBackupFailoverTestBase.java
Log:
test fix to allow backups to announce before starting test
Modified:
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/distribution/ClusterTestBase.java
===================================================================
---
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/distribution/ClusterTestBase.java 2012-02-17
12:39:49 UTC (rev 12133)
+++
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/distribution/ClusterTestBase.java 2012-02-17
13:44:48 UTC (rev 12134)
@@ -17,6 +17,7 @@
import java.io.StringWriter;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
@@ -40,6 +41,8 @@
import org.hornetq.api.core.client.ClientSessionFactory;
import org.hornetq.api.core.client.HornetQClient;
import org.hornetq.api.core.client.ServerLocator;
+import org.hornetq.core.client.impl.Topology;
+import org.hornetq.core.client.impl.TopologyMember;
import org.hornetq.core.config.BroadcastGroupConfiguration;
import org.hornetq.core.config.ClusterConnectionConfiguration;
import org.hornetq.core.config.Configuration;
@@ -57,6 +60,7 @@
import org.hornetq.core.server.cluster.ClusterConnection;
import org.hornetq.core.server.cluster.ClusterManager;
import org.hornetq.core.server.cluster.RemoteQueueBinding;
+import org.hornetq.core.server.cluster.impl.ClusterConnectionImpl;
import org.hornetq.core.server.group.GroupingHandler;
import org.hornetq.core.server.group.impl.GroupingHandlerConfiguration;
import org.hornetq.core.server.impl.InVMNodeManager;
@@ -238,6 +242,64 @@
return consumers[node].consumer;
}
+
+ protected void waitForFailoverTopology(final int bNode, final int... nodes) throws
Exception
+ {
+ HornetQServer server = servers[bNode];
+
+ log.debug("waiting for " + nodes + " on the topology for server =
" + server);
+
+ long start = System.currentTimeMillis();
+
+ Set<ClusterConnection> ccs =
server.getClusterManager().getClusterConnections();
+
+ if (ccs.size() != 1)
+ {
+ throw new IllegalStateException("You need a single cluster connection on
this version of waitForTopology on ServiceTestBase");
+ }
+
+ boolean exists = false;
+
+ for (int node : nodes)
+ {
+ ClusterConnectionImpl clusterConnection = (ClusterConnectionImpl)
ccs.iterator().next();
+ Topology topology = clusterConnection.getTopology();
+ TransportConfiguration nodeConnector=
+
servers[node].getClusterManager().getClusterConnections().iterator().next().getConnector();
+ do
+ {
+ Collection<TopologyMember> members = topology.getMembers();
+ for (TopologyMember member : members)
+ {
+ if(member.getConnector().getA() != null &&
member.getConnector().getA().equals(nodeConnector))
+ {
+ exists = true;
+ break;
+ }
+ }
+ if(exists)
+ {
+ break;
+ }
+ Thread.sleep(10);
+ }
+ while (System.currentTimeMillis() - start < WAIT_TIMEOUT);
+ if(!exists)
+ {
+ String msg = "Timed out waiting for cluster topology of " + nodes
+
+ " (received " +
+ topology.getMembers().size() +
+ ") topology = " +
+ topology +
+ ")";
+
+ log.error(msg);
+
+ throw new Exception(msg);
+ }
+ }
+ }
+
protected void waitForMessages(final int node, final String address, final int count)
throws Exception
{
HornetQServer server = servers[node];
Modified:
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/failover/ClusterWithBackupFailoverTestBase.java
===================================================================
---
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/failover/ClusterWithBackupFailoverTestBase.java 2012-02-17
12:39:49 UTC (rev 12133)
+++
trunk/tests/integration-tests/src/test/java/org/hornetq/tests/integration/cluster/failover/ClusterWithBackupFailoverTestBase.java 2012-02-17
13:44:48 UTC (rev 12134)
@@ -67,6 +67,10 @@
startServers(3, 4, 5, 0, 1, 2);
+ waitForFailoverTopology(3, 0, 1, 2);
+ waitForFailoverTopology(4, 0, 1, 2);
+ waitForFailoverTopology(5, 0, 1, 2);
+
setupSessionFactory(0, 3, isNetty(), false);
setupSessionFactory(1, 4, isNetty(), false);
setupSessionFactory(2, 5, isNetty(), false);
@@ -89,9 +93,12 @@
send(2, QUEUES_TESTADDRESS, 10, false, null);
verifyReceiveRoundRobinInSomeOrder(true, 10, 0, 1, 2);
-
+ Thread.sleep(1000);
failNode(0);
+ waitForFailoverTopology(4, 3, 1, 2);
+ waitForFailoverTopology(5, 3, 1, 2);
+
// live nodes
waitForBindings(1, QUEUES_TESTADDRESS, 1, 1, true);
waitForBindings(2, QUEUES_TESTADDRESS, 1, 1, true);
@@ -117,6 +124,9 @@
failNode(1);
+ waitForFailoverTopology(5, 3, 4, 2);
+
+ Thread.sleep(1000);
// live nodes
waitForBindings(2, QUEUES_TESTADDRESS, 1, 1, true);
// activated backup nodes
@@ -140,6 +150,7 @@
failNode(2);
+ Thread.sleep(1000);
// activated backup nodes
waitForBindings(3, QUEUES_TESTADDRESS, 1, 1, true);
waitForBindings(4, QUEUES_TESTADDRESS, 1, 1, true);
@@ -293,6 +304,10 @@
setupSessionFactory(1, 4, isNetty(), false);
setupSessionFactory(2, 5, isNetty(), false);
+ waitForFailoverTopology(3, 0, 1, 2);
+ waitForFailoverTopology(4, 0, 1, 2);
+ waitForFailoverTopology(5, 0, 1, 2);
+
createQueue(0, QUEUES_TESTADDRESS, QUEUE_NAME, null, true);
createQueue(1, QUEUES_TESTADDRESS, QUEUE_NAME, null, true);
createQueue(2, QUEUES_TESTADDRESS, QUEUE_NAME, null, true);
@@ -314,6 +329,8 @@
failNode(0);
+ waitForFailoverTopology(4, 3, 1, 2);
+ waitForFailoverTopology(5, 3, 1, 2);
// live nodes
waitForBindings(1, QUEUES_TESTADDRESS, 1, 1, true);
waitForBindings(2, QUEUES_TESTADDRESS, 1, 1, true);
@@ -356,6 +373,7 @@
failNode(1);
+ waitForFailoverTopology(5, 2, 4);
// live nodes
waitForBindings(2, QUEUES_TESTADDRESS, 1, 1, true);
// activated backup nodes