- Added test case to check if any message dispatch are sent over the bridge even when there is no consumer on the other end

- Added better setup for some to the test case

git-svn-id: https://svn.apache.org/repos/asf/incubator/activemq/trunk@367501 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
Adrian T. Co 2006-01-10 06:14:27 +00:00
parent aa7435064a
commit 5eff976d28
4 changed files with 167 additions and 4 deletions

View File

@ -22,15 +22,15 @@ import org.apache.activemq.network.DemandForwardingBridge;
import org.apache.activemq.transport.TransportFactory;
import java.util.List;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.ArrayList;
import java.net.URI;
/**
* @version $Revision: 1.1.1.1 $
*/
public class MultiBrokersMultiClientsUsingTcpTest extends MultiBrokersMultiClientsTest {
protected List bridges = new ArrayList();
protected List bridges;
protected void bridgeAllBrokers(String groupName) throws Exception {
for (int i=1; i<=BROKER_COUNT; i++) {
@ -77,5 +77,7 @@ public class MultiBrokersMultiClientsUsingTcpTest extends MultiBrokersMultiClien
for (Iterator i=brokers.values().iterator(); i.hasNext();) {
((BrokerItem)i.next()).broker.addConnector("tcp://localhost:" + (61616 + j++));
}
bridges = new ArrayList();
}
}

View File

@ -29,7 +29,7 @@ import java.net.URI;
* @version $Revision: 1.1.1.1 $
*/
public class ThreeBrokerQueueNetworkUsingTcpTest extends ThreeBrokerQueueNetworkTest {
protected List bridges = new ArrayList();
protected List bridges;
protected void bridgeBrokers(BrokerService localBroker, BrokerService remoteBroker) throws Exception {
List remoteTransports = remoteBroker.getTransportConnectors();
@ -57,4 +57,10 @@ public class ThreeBrokerQueueNetworkUsingTcpTest extends ThreeBrokerQueueNetwork
MAX_SETUP_TIME = 2000;
}
public void setUp() throws Exception {
super.setUp();
bridges = new ArrayList();
}
}

View File

@ -29,7 +29,7 @@ import java.net.URI;
* @version $Revision: 1.1.1.1 $
*/
public class ThreeBrokerTopicNetworkUsingTcpTest extends ThreeBrokerTopicNetworkTest {
protected List bridges = new ArrayList();
protected List bridges;
protected void bridgeBrokers(BrokerService localBroker, BrokerService remoteBroker) throws Exception {
List remoteTransports = remoteBroker.getTransportConnectors();
@ -57,4 +57,10 @@ public class ThreeBrokerTopicNetworkUsingTcpTest extends ThreeBrokerTopicNetwork
MAX_SETUP_TIME = 2000;
}
public void setUp() throws Exception {
super.setUp();
bridges = new ArrayList();
}
}

View File

@ -0,0 +1,149 @@
/**
*
* Copyright 2005-2006 The Apache Software Foundation
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.activemq.usecases;
import org.apache.activemq.broker.BrokerService;
import org.apache.activemq.broker.TransportConnector;
import org.apache.activemq.network.DemandForwardingBridge;
import org.apache.activemq.transport.TransportFactory;
import org.apache.activemq.JmsMultipleBrokersTestSupport;
import org.apache.activemq.command.Command;
import org.apache.activemq.util.MessageIdList;
import javax.jms.Destination;
import javax.jms.MessageConsumer;
import java.util.List;
import java.util.ArrayList;
import java.net.URI;
import edu.emory.mathcs.backport.java.util.concurrent.atomic.AtomicInteger;
/**
* @version $Revision: 1.1.1.1 $
*/
public class TwoBrokerMessageNotSentToRemoteWhenNoConsumerTest extends JmsMultipleBrokersTestSupport {
protected List bridges;
protected AtomicInteger msgDispatchCount;
/**
* BrokerA -> BrokerB
*/
public void testRemoteBrokerHasConsumer() throws Exception {
// Setup broker networks
bridgeBrokers("BrokerA", "BrokerB");
startAllBrokers();
// Setup destination
Destination dest = createDestination("TEST.FOO", true);
// Setup consumers
MessageConsumer clientA = createConsumer("BrokerA", dest);
MessageConsumer clientB = createConsumer("BrokerB", dest);
// Send messages
sendMessages("BrokerA", dest, 10);
// Get message count
MessageIdList msgsA = getConsumerMessages("BrokerA", clientA);
MessageIdList msgsB = getConsumerMessages("BrokerB", clientB);
msgsA.waitForMessagesToArrive(10);
msgsB.waitForMessagesToArrive(10);
assertEquals(10, msgsA.getMessageCount());
assertEquals(10, msgsB.getMessageCount());
// Check that 10 message dispatch commands are send over the network
assertEquals(10, msgDispatchCount.get());
}
/**
* BrokerA -> BrokerB
*/
public void testRemoteBrokerHasNoConsumer() throws Exception {
// Setup broker networks
bridgeBrokers("BrokerA", "BrokerB");
startAllBrokers();
// Setup destination
Destination dest = createDestination("TEST.FOO", true);
// Setup consumers
MessageConsumer clientA = createConsumer("BrokerA", dest);
// Send messages
sendMessages("BrokerA", dest, 10);
// Get message count
MessageIdList msgsA = getConsumerMessages("BrokerA", clientA);
msgsA.waitForMessagesToArrive(10);
assertEquals(10, msgsA.getMessageCount());
// Check that no message dispatch commands are send over the network
assertEquals(0, msgDispatchCount.get());
}
protected void bridgeBrokers(BrokerService localBroker, BrokerService remoteBroker) throws Exception {
List remoteTransports = remoteBroker.getTransportConnectors();
List localTransports = localBroker.getTransportConnectors();
URI remoteURI, localURI;
if (!remoteTransports.isEmpty() && !localTransports.isEmpty()) {
remoteURI = ((TransportConnector)remoteTransports.get(0)).getConnectUri();
localURI = ((TransportConnector)localTransports.get(0)).getConnectUri();
// Ensure that we are connecting using tcp
if (remoteURI.toString().startsWith("tcp:") && localURI.toString().startsWith("tcp:")) {
DemandForwardingBridge bridge = new DemandForwardingBridge(TransportFactory.connect(localURI),
TransportFactory.connect(remoteURI)) {
protected void serviceLocalCommand(Command command) {
if (command.isMessageDispatch()) {
// Keep track of the number of message dispatches through the bridge
msgDispatchCount.incrementAndGet();
}
super.serviceLocalCommand(command);
}
};
bridge.setClientId(localBroker.getBrokerName() + "_to_" + remoteBroker.getBrokerName());
bridges.add(bridge);
bridge.start();
} else {
throw new Exception("Remote broker or local broker is not using tcp connectors");
}
} else {
throw new Exception("Remote broker or local broker has no registered connectors.");
}
MAX_SETUP_TIME = 2000;
}
public void setUp() throws Exception {
super.setAutoFail(true);
super.setUp();
createBroker(new URI("broker:(tcp://localhost:61616)/BrokerA?persistent=false&useJmx=false"));
createBroker(new URI("broker:(tcp://localhost:61617)/BrokerB?persistent=false&useJmx=false"));
bridges = new ArrayList();
msgDispatchCount = new AtomicInteger(0);
}
}