106 lines
2.3 KiB
C++
106 lines
2.3 KiB
C++
|
|
#include <cassert>
|
||
|
|
#include <thread>
|
||
|
|
#include <vector>
|
||
|
|
|
||
|
|
#include "log_level.hpp"
|
||
|
|
#include "message_queue.hpp"
|
||
|
|
#include "thread_test_utils.hpp"
|
||
|
|
|
||
|
|
namespace app::tests {
|
||
|
|
|
||
|
|
using namespace app;
|
||
|
|
|
||
|
|
void testPushAndPop() {
|
||
|
|
MessageQueue queue;
|
||
|
|
|
||
|
|
queue.push({logger::LogLevel::info(), "Hello"});
|
||
|
|
|
||
|
|
auto message = queue.pop();
|
||
|
|
|
||
|
|
assert(message.level == logger::LogLevel::info());
|
||
|
|
assert(message.message == "Hello");
|
||
|
|
}
|
||
|
|
|
||
|
|
void testFifoOrder() {
|
||
|
|
MessageQueue queue;
|
||
|
|
|
||
|
|
queue.push({logger::LogLevel::debug(), "first"});
|
||
|
|
queue.push({logger::LogLevel::info(), "second"});
|
||
|
|
queue.push({logger::LogLevel::error(), "third"});
|
||
|
|
|
||
|
|
assert(queue.pop().message == "first");
|
||
|
|
assert(queue.pop().message == "second");
|
||
|
|
assert(queue.pop().message == "third");
|
||
|
|
}
|
||
|
|
|
||
|
|
void testPopWaitsForMessage() {
|
||
|
|
MessageQueue queue;
|
||
|
|
|
||
|
|
bool received = false;
|
||
|
|
|
||
|
|
std::thread worker([&]() {
|
||
|
|
auto message = queue.pop();
|
||
|
|
|
||
|
|
assert(message.level == logger::LogLevel::error());
|
||
|
|
assert(message.message == "Delayed");
|
||
|
|
|
||
|
|
received = true;
|
||
|
|
});
|
||
|
|
|
||
|
|
// Даём потоку worker возможность дойти до queue.pop()
|
||
|
|
// и заблокироваться в ожидании сообщения.
|
||
|
|
std::this_thread::sleep_for(std::chrono::milliseconds(100));
|
||
|
|
|
||
|
|
assert(!received);
|
||
|
|
|
||
|
|
// Добавляем сообщение и тем самым разблокируем worker.
|
||
|
|
queue.push({logger::LogLevel::error(), "Delayed"});
|
||
|
|
|
||
|
|
// Ожидаем полного завершения потока worker.
|
||
|
|
worker.join();
|
||
|
|
|
||
|
|
assert(received);
|
||
|
|
}
|
||
|
|
|
||
|
|
void testConcurrentPushAndPop() {
|
||
|
|
MessageQueue queue;
|
||
|
|
|
||
|
|
constexpr std::size_t threadCount = 8;
|
||
|
|
constexpr std::size_t messagesPerThread = 1000;
|
||
|
|
|
||
|
|
std::vector<std::thread> producers;
|
||
|
|
|
||
|
|
for (std::size_t i = 0; i < threadCount; ++i) {
|
||
|
|
producers.emplace_back(
|
||
|
|
helpers::createProducer(queue, i, messagesPerThread));
|
||
|
|
}
|
||
|
|
const auto expectedMessages = threadCount * messagesPerThread;
|
||
|
|
|
||
|
|
for (std::size_t i = 0; i < expectedMessages; ++i) {
|
||
|
|
auto message = queue.pop();
|
||
|
|
|
||
|
|
assert(!message.message.empty());
|
||
|
|
}
|
||
|
|
|
||
|
|
for (auto &producer : producers) {
|
||
|
|
producer.join();
|
||
|
|
}
|
||
|
|
|
||
|
|
assert(queue.empty());
|
||
|
|
}
|
||
|
|
|
||
|
|
void runTests() {
|
||
|
|
testPushAndPop();
|
||
|
|
testFifoOrder();
|
||
|
|
testPopWaitsForMessage();
|
||
|
|
testConcurrentPushAndPop();
|
||
|
|
}
|
||
|
|
|
||
|
|
} // namespace app::tests
|
||
|
|
|
||
|
|
int main() {
|
||
|
|
app::tests::runTests();
|
||
|
|
|
||
|
|
return 0;
|
||
|
|
}
|