-
Notifications
You must be signed in to change notification settings - Fork 1
/
AbstractTramCommandTest.java
executable file
·63 lines (47 loc) · 1.88 KB
/
AbstractTramCommandTest.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
package io.eventuate.tram.examples.basic.commands;
import io.eventuate.tram.commands.common.ReplyMessageHeaders;
import io.eventuate.tram.commands.producer.CommandProducer;
import io.eventuate.tram.messaging.common.Message;
import io.eventuate.tram.messaging.consumer.MessageConsumer;
import org.junit.jupiter.api.Test;
import javax.inject.Inject;
import java.util.Collections;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
public abstract class AbstractTramCommandTest {
@Inject
private CommandProducer commandProducer;
@Inject
private AbstractTramCommandTestConfig config;
@Inject
private MessageConsumer messageConsumer;
private BlockingQueue<Message> queue = new LinkedBlockingDeque<>();
@Test
public void shouldInvokeCommand() throws InterruptedException {
subscribeToReplyChannel();
String commandId = sendCommand();
assertReplyReceived(commandId);
}
private void assertReplyReceived(String commandId) throws InterruptedException {
Message m = queue.poll(30, TimeUnit.SECONDS);
System.out.println("Got message = " + m);
assertNotNull(m);
assertEquals(commandId, m.getRequiredHeader(ReplyMessageHeaders.IN_REPLY_TO));
}
private String sendCommand() {
return commandProducer.send(config.getCommandChannel(),
new DoSomethingCommand(),
config.getReplyChannel(),
Collections.emptyMap());
}
private void subscribeToReplyChannel() {
String subscriberId = "subscriberId" + config.getUniqueId();
messageConsumer.subscribe(subscriberId, Collections.singleton(config.getReplyChannel()), this::handleMessage);
}
private void handleMessage(Message message) {
queue.add(message);
}
}