-
Notifications
You must be signed in to change notification settings - Fork 5
/
MultiClientTest.java
173 lines (151 loc) · 6.09 KB
/
MultiClientTest.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
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
/*
* Java Reliable Event Logging Protocol Library Server Implementation RLP-03
* Copyright (C) 2021 Suomen Kanuuna Oy
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program 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 Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*
*
* Additional permission under GNU Affero General Public License version 3
* section 7
*
* If you modify this Program, or any covered work, by linking or combining it
* with other code, such other code is not for that reason alone subject to any
* of the requirements of the GNU Affero GPL version 3 as long as this Program
* is the same Program as licensed from Suomen Kanuuna Oy without any additional
* modifications.
*
* Supplemented terms under GNU Affero General Public License version 3
* section 7
*
* Origin of the software must be attributed to Suomen Kanuuna Oy. Any modified
* versions must be marked as "Modified version of" The Program.
*
* Names of the licensors and authors may not be used for publicity purposes.
*
* No rights are granted for use of trade names, trademarks, or service marks
* which are in The Program if any.
*
* Licensee must indemnify licensors and authors for any liability that these
* contractual assumptions impose on licensors and authors.
*
* To the extent this program is licensed as part of the Commercial versions of
* Teragrep, the applicable Commercial License may apply to this file if you as
* a licensee so wish it.
*/
package com.teragrep.rlp_03;
import com.teragrep.rlp_01.RelpBatch;
import com.teragrep.rlp_01.RelpConnection;
import com.teragrep.rlp_03.config.Config;
import org.junit.jupiter.api.*;
import org.junit.jupiter.api.condition.DisabledOnJre;
import org.junit.jupiter.api.condition.JRE;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.Deque;
import java.util.Random;
import java.util.concurrent.ConcurrentLinkedDeque;
import java.util.concurrent.TimeoutException;
import java.util.function.Supplier;
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
@DisabledOnJre(JRE.JAVA_8)
public class MultiClientTest extends Thread{
private final String hostname = "localhost";
private Server server;
private static int port = 1239;
private final Deque<byte[]> messageList = new ConcurrentLinkedDeque<>();
@Test
public void testMultiClient() throws InterruptedException, IllegalStateException {
int n = 10;
Thread[] threads = new Thread[n];
for(int i=0; i<n; i++) {
Thread thread = new Thread(new MultiClientTest());
thread.start();
threads[i] = thread;
}
for (int i=0; i<n; i++) {
threads[i].join();
}
}
// testMultiClient executes this with new MultiClientTest() thread
public void run() {
Random random = new Random();
for(int i=0;i<3;i++) {
// Sleep to make the ordering unpredictable
try {
Thread.sleep(random.nextInt(100));
} catch (InterruptedException e) {
e.printStackTrace();
}
try {
testSendBatch();
testSendMessage();
} catch (IOException | TimeoutException | IllegalStateException e) {
e.printStackTrace();
}
}
}
@BeforeAll
public void init() throws IOException, InterruptedException {
Supplier<FrameProcessor> frameProcessorSupplier = new Supplier<FrameProcessor>() {
@Override
public FrameProcessor get() {
return new SyslogFrameProcessor((frame) -> messageList.add(frame.relpFrame().payload().toBytes()));
}
};
port = getPort();
Config config = new Config(port, 4);
ServerFactory serverFactory = new ServerFactory(config, frameProcessorSupplier);
server = serverFactory.create();
Thread serverThread = new Thread(server);
serverThread.start();
server.startup.waitForCompletion();
}
@AfterAll
public void cleanup() throws InterruptedException {
server.stop();
// 10 threads each: run 3 times 50 msgs of testSendBatch plus 3 times
// 1 msgs of testSendMessage (1530 messages total)
Assertions.assertEquals(10 * 3 * 50 + 10 * 3, messageList.size());
}
private synchronized int getPort() {
return ++port;
}
private void testSendMessage() throws IOException, TimeoutException {
RelpConnection relpSession = new RelpConnection();
relpSession.connect(hostname, port);
String msg = "<14>1 2020-05-15T13:24:03.603Z CFE-16 capsulated - - [CFE-16-metadata@48577 authentication_token=\"AUTH_TOKEN_11111\" channel=\"CHANNEL_11111\" time_source=\"generated\"][CFE-16-origin@48577] \"Hello, world!\"\n";
byte[] data = msg.getBytes(StandardCharsets.UTF_8);
RelpBatch batch = new RelpBatch();
long reqId = batch.insert(data);
relpSession.commit(batch);
// verify successful transaction
Assertions.assertTrue(batch.verifyTransaction(reqId));
relpSession.disconnect();
}
private void testSendBatch() throws IllegalStateException, IOException,
TimeoutException {
RelpConnection relpSession = new RelpConnection();
relpSession.connect(hostname, port);
String msg = "Hello, world!";
byte[] data = msg.getBytes(StandardCharsets.UTF_8);
int n = 50;
RelpBatch batch = new RelpBatch();
for (int i = 0; i < n; i++) {
batch.insert(data);
}
relpSession.commit(batch);
Assertions.assertTrue(batch.verifyTransactionAll());
relpSession.disconnect();
}
}