/
AllIn.java
88 lines (75 loc) · 2.79 KB
/
AllIn.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
package shoe;
import org.apache.activemq.broker.BrokerService;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import javax.jms.*;
import static org.hamcrest.CoreMatchers.is;
import static org.junit.Assert.assertThat;
public class AllIn {
public static final String BROKER_URL = "tcp://localhost:61616";
public static final String QUEUE_NAME = "test";
public static final int MESSAGES_TO_SEND = 10;
private BrokerService broker;
private volatile boolean shutdown;
private int messagesReceived;
@Before
public void startBroker() throws Exception {
broker = new BrokerService();
broker.addConnector(BROKER_URL);
broker.start();
}
@Test
public void smoke() throws Exception {
produce();
consume();
waitForConsumptionToComplete();
assertThat(messagesReceived, is(MESSAGES_TO_SEND + 1));
}
private void waitForConsumptionToComplete() throws InterruptedException {
while(!shutdown) {
Thread.sleep(10);
}
}
private void consume() throws JMSException {
ConnectionFactory factory = new ActiveMQConnectionFactory(BROKER_URL);
Connection connection = factory.createConnection();
connection.start();
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Destination destination = session.createQueue(QUEUE_NAME);
MessageConsumer consumer = session.createConsumer(destination);
consumer.setMessageListener(new SimpleConsumer());
}
private void produce() throws JMSException {
ConnectionFactory factory = new ActiveMQConnectionFactory(BROKER_URL);
Connection connection = factory.createConnection();
connection.start();
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
Destination destination = session.createQueue(QUEUE_NAME);
MessageProducer producer = session.createProducer(destination);
for(int i = 1; i <= MESSAGES_TO_SEND; ++i) {
Message message = session.createTextMessage("Hello World!" + i);
producer.send(message);
}
producer.send(session.createTextMessage("terminate"));
}
@After
public void stopBroker() throws Exception {
broker.stop();
}
class SimpleConsumer implements MessageListener {
@Override
public void onMessage(Message message) {
try {
String body = ((TextMessage)message).getText();
++messagesReceived;
if("terminate".equals(body)) {
shutdown = true;
}
} catch (JMSException e) {
e.printStackTrace();
}
}
}
}