-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathexample.cpp
More file actions
126 lines (104 loc) · 3.52 KB
/
Copy pathexample.cpp
File metadata and controls
126 lines (104 loc) · 3.52 KB
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
#include <assert.h>
#include <thread>
#include "LWMessageQueue.h"
#include "Message.h"
const uint32_t numChannels = 2;
const uint32_t numMessages = 1000;
const uint32_t numMessagesPerThread = numMessages * 2;
const uint32_t queueSize = 2048;
using MessageQueue = LWMessageQueue::LWMessageQueue<queueSize, numChannels, MessageUnion>;
void inputThread0Run(MessageQueue::ThreadChannelInput inChannel) {
for (uint32_t i = 0; i < numMessages; ++i) {
{
Message1 message;
message.value = 17;
message.anotherValue = 4711;
assert(!inChannel.isFull());
inChannel.pushMessage(message);
}
{
Message2 message;
message.value = 0;
message.anotherValue = 0;
message.moreValues[0] = 0;
message.moreValues[1] = 0;
assert(!inChannel.isFull());
inChannel.pushMessage(message);
}
}
fprintf(stdout, "Input thread 0 done, sent %u messages.\n", numMessagesPerThread);
}
void inputThread1Run(MessageQueue::ThreadChannelInput inChannel) {
for (uint32_t i = 0; i < numMessages; ++i) {
{
Message1 message;
message.value = 17;
message.anotherValue = 4711;
assert(!inChannel.isFull());
inChannel.pushMessage(message);
}
{
Message2 message;
message.value = 1;
message.anotherValue = 1;
message.moreValues[0] = 1;
message.moreValues[1] = 1;
assert(!inChannel.isFull());
inChannel.pushMessage(message);
}
}
fprintf(stdout, "Input thread 1 done, sent %u messages.\n", numMessagesPerThread);
}
void verifyMessage(const uint32_t channelIndex, const MessageQueue::MessageContainer& messageContainer) {
if (messageContainer.isOfType<Message1>()) {
const Message1& message = messageContainer.getMessage<Message1>();
assert(message.value == 17);
assert(message.anotherValue == 4711);
} else if (messageContainer.isOfType<Message2>()) {
const Message2& message = messageContainer.getMessage<Message2>();
assert(message.value == channelIndex);
assert(message.anotherValue == channelIndex);
assert(message.moreValues[0] == channelIndex);
assert(message.moreValues[1] == channelIndex);
} else {
assert(false);
}
}
void outputThreadRun(MessageQueue* messageQueue) {
assert(messageQueue != nullptr);
MessageQueue::ThreadChannelOutput channel0 = messageQueue->getThreadChannelOutput(0);
MessageQueue::ThreadChannelOutput channel1 = messageQueue->getThreadChannelOutput(1);
uint32_t receivedMessages = 0;
const uint32_t totalTestWantedMessages = numMessagesPerThread * numChannels;
while (receivedMessages < totalTestWantedMessages) {
// Channel 0
{
const uint32_t pendingMessages = channel0.getNumMessages();
for (uint32_t messageIndex = 0; messageIndex < pendingMessages; ++messageIndex) {
MessageQueue::MessageContainer messageContainer = channel0.popMessage();
++receivedMessages;
verifyMessage(0, messageContainer);
}
}
// Channel 1
{
const uint32_t pendingMessages = channel1.getNumMessages();
for (uint32_t messageIndex = 0; messageIndex < pendingMessages; ++messageIndex) {
MessageQueue::MessageContainer messageContainer = channel1.popMessage();
++receivedMessages;
verifyMessage(1, messageContainer);
}
}
}
fprintf(stdout, "Output thread done, received %u messages\n", receivedMessages);
}
int main(int, char**) {
MessageQueue messageQueue;
std::thread outputThread(outputThreadRun, &messageQueue);
std::thread inputThread0(inputThread0Run, messageQueue.getThreadChannelInput(0));
std::thread inputThread1(inputThread1Run, messageQueue.getThreadChannelInput(1));
inputThread0.join();
inputThread1.join();
outputThread.join();
return 0;
}