feat(app): add graceful shutdown support to MessageQueue
This commit is contained in:
@@ -2,6 +2,7 @@
|
|||||||
|
|
||||||
#include <condition_variable>
|
#include <condition_variable>
|
||||||
#include <mutex>
|
#include <mutex>
|
||||||
|
#include <optional>
|
||||||
#include <queue>
|
#include <queue>
|
||||||
|
|
||||||
#include "log_message.hpp"
|
#include "log_message.hpp"
|
||||||
@@ -22,6 +23,8 @@ public:
|
|||||||
*
|
*
|
||||||
* Метод является потокобезопасным.
|
* Метод является потокобезопасным.
|
||||||
*
|
*
|
||||||
|
* После остановки очереди новые сообщения не принимаются.
|
||||||
|
*
|
||||||
* @param message Сообщение для передачи.
|
* @param message Сообщение для передачи.
|
||||||
*/
|
*/
|
||||||
void push(LogMessage message);
|
void push(LogMessage message);
|
||||||
@@ -29,15 +32,28 @@ public:
|
|||||||
/**
|
/**
|
||||||
* @brief Извлекает сообщение из очереди.
|
* @brief Извлекает сообщение из очереди.
|
||||||
*
|
*
|
||||||
* Если очередь пуста, метод блокирует вызывающий
|
* Если очередь пуста, метод блокирует вызывающий поток
|
||||||
* поток до появления нового сообщения.
|
* до появления нового сообщения или до остановки очереди.
|
||||||
|
*
|
||||||
|
* После остановки очереди и обработки всех оставшихся сообщений
|
||||||
|
* возвращается std::nullopt.
|
||||||
*
|
*
|
||||||
* Метод является потокобезопасным.
|
* Метод является потокобезопасным.
|
||||||
*
|
*
|
||||||
* @return Следующее сообщение из очереди.
|
* @return Следующее сообщение или std::nullopt, если очередь остановлена.
|
||||||
*/
|
*/
|
||||||
[[nodiscard]]
|
[[nodiscard]]
|
||||||
LogMessage pop();
|
std::optional<LogMessage> pop();
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @brief Останавливает очередь.
|
||||||
|
*
|
||||||
|
* Пробуждает все ожидающие потоки. После вызова stop()
|
||||||
|
* новые сообщения не принимаются.
|
||||||
|
*
|
||||||
|
* Метод является потокобезопасным.
|
||||||
|
*/
|
||||||
|
void stop();
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @brief Проверяет, пуста ли очередь.
|
* @brief Проверяет, пуста ли очередь.
|
||||||
@@ -56,6 +72,11 @@ private:
|
|||||||
*/
|
*/
|
||||||
std::queue<LogMessage> queue_;
|
std::queue<LogMessage> queue_;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @brief Признак остановки очереди.
|
||||||
|
*/
|
||||||
|
bool stopped_{false};
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @brief Мьютекс для синхронизации доступа к очереди.
|
* @brief Мьютекс для синхронизации доступа к очереди.
|
||||||
*/
|
*/
|
||||||
|
|||||||
@@ -5,15 +5,24 @@ namespace app {
|
|||||||
void MessageQueue::push(LogMessage message) {
|
void MessageQueue::push(LogMessage message) {
|
||||||
std::lock_guard<std::mutex> lock(mutex_);
|
std::lock_guard<std::mutex> lock(mutex_);
|
||||||
|
|
||||||
|
if (stopped_) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
queue_.push(std::move(message));
|
queue_.push(std::move(message));
|
||||||
|
|
||||||
conditionVariable_.notify_one();
|
conditionVariable_.notify_one();
|
||||||
}
|
}
|
||||||
|
|
||||||
LogMessage MessageQueue::pop() {
|
std::optional<LogMessage> MessageQueue::pop() {
|
||||||
std::unique_lock<std::mutex> lock(mutex_);
|
std::unique_lock<std::mutex> lock(mutex_);
|
||||||
|
|
||||||
conditionVariable_.wait(lock, [this]() { return !queue_.empty(); });
|
conditionVariable_.wait(lock,
|
||||||
|
[this]() { return stopped_ || !queue_.empty(); });
|
||||||
|
|
||||||
|
if (queue_.empty()) {
|
||||||
|
return std::nullopt;
|
||||||
|
}
|
||||||
|
|
||||||
auto message = std::move(queue_.front());
|
auto message = std::move(queue_.front());
|
||||||
|
|
||||||
@@ -22,6 +31,16 @@ LogMessage MessageQueue::pop() {
|
|||||||
return message;
|
return message;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void MessageQueue::stop() {
|
||||||
|
{
|
||||||
|
std::lock_guard<std::mutex> lock(mutex_);
|
||||||
|
|
||||||
|
stopped_ = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
conditionVariable_.notify_all();
|
||||||
|
}
|
||||||
|
|
||||||
bool MessageQueue::empty() const {
|
bool MessageQueue::empty() const {
|
||||||
std::lock_guard<std::mutex> lock(mutex_);
|
std::lock_guard<std::mutex> lock(mutex_);
|
||||||
|
|
||||||
|
|||||||
@@ -17,8 +17,9 @@ void testPushAndPop() {
|
|||||||
|
|
||||||
auto message = queue.pop();
|
auto message = queue.pop();
|
||||||
|
|
||||||
assert(message.level == logger::LogLevel::info());
|
assert(message.has_value());
|
||||||
assert(message.message == "Hello");
|
assert(message->level == logger::LogLevel::info());
|
||||||
|
assert(message->message == "Hello");
|
||||||
}
|
}
|
||||||
|
|
||||||
void testFifoOrder() {
|
void testFifoOrder() {
|
||||||
@@ -28,9 +29,9 @@ void testFifoOrder() {
|
|||||||
queue.push({logger::LogLevel::info(), "second"});
|
queue.push({logger::LogLevel::info(), "second"});
|
||||||
queue.push({logger::LogLevel::error(), "third"});
|
queue.push({logger::LogLevel::error(), "third"});
|
||||||
|
|
||||||
assert(queue.pop().message == "first");
|
assert(queue.pop()->message == "first");
|
||||||
assert(queue.pop().message == "second");
|
assert(queue.pop()->message == "second");
|
||||||
assert(queue.pop().message == "third");
|
assert(queue.pop()->message == "third");
|
||||||
}
|
}
|
||||||
|
|
||||||
void testPopWaitsForMessage() {
|
void testPopWaitsForMessage() {
|
||||||
@@ -41,8 +42,9 @@ void testPopWaitsForMessage() {
|
|||||||
std::thread worker([&]() {
|
std::thread worker([&]() {
|
||||||
auto message = queue.pop();
|
auto message = queue.pop();
|
||||||
|
|
||||||
assert(message.level == logger::LogLevel::error());
|
assert(message.has_value());
|
||||||
assert(message.message == "Delayed");
|
assert(message->level == logger::LogLevel::error());
|
||||||
|
assert(message->message == "Delayed");
|
||||||
|
|
||||||
received = true;
|
received = true;
|
||||||
});
|
});
|
||||||
@@ -74,12 +76,14 @@ void testConcurrentPushAndPop() {
|
|||||||
producers.emplace_back(
|
producers.emplace_back(
|
||||||
helpers::createProducer(queue, i, messagesPerThread));
|
helpers::createProducer(queue, i, messagesPerThread));
|
||||||
}
|
}
|
||||||
|
|
||||||
const auto expectedMessages = threadCount * messagesPerThread;
|
const auto expectedMessages = threadCount * messagesPerThread;
|
||||||
|
|
||||||
for (std::size_t i = 0; i < expectedMessages; ++i) {
|
for (std::size_t i = 0; i < expectedMessages; ++i) {
|
||||||
auto message = queue.pop();
|
auto message = queue.pop();
|
||||||
|
|
||||||
assert(!message.message.empty());
|
assert(message.has_value());
|
||||||
|
assert(!message->message.empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
for (auto &producer : producers) {
|
for (auto &producer : producers) {
|
||||||
@@ -89,11 +93,44 @@ void testConcurrentPushAndPop() {
|
|||||||
assert(queue.empty());
|
assert(queue.empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void testStopReturnsNullopt() {
|
||||||
|
MessageQueue queue;
|
||||||
|
|
||||||
|
queue.stop();
|
||||||
|
|
||||||
|
auto message = queue.pop();
|
||||||
|
|
||||||
|
assert(!message.has_value());
|
||||||
|
}
|
||||||
|
|
||||||
|
void testStopReturnsRemainingMessages() {
|
||||||
|
MessageQueue queue;
|
||||||
|
|
||||||
|
queue.push({logger::LogLevel::info(), "first"});
|
||||||
|
queue.push({logger::LogLevel::error(), "second"});
|
||||||
|
|
||||||
|
queue.stop();
|
||||||
|
|
||||||
|
auto first = queue.pop();
|
||||||
|
auto second = queue.pop();
|
||||||
|
auto third = queue.pop();
|
||||||
|
|
||||||
|
assert(first.has_value());
|
||||||
|
assert(first->message == "first");
|
||||||
|
|
||||||
|
assert(second.has_value());
|
||||||
|
assert(second->message == "second");
|
||||||
|
|
||||||
|
assert(!third.has_value());
|
||||||
|
}
|
||||||
|
|
||||||
void runTests() {
|
void runTests() {
|
||||||
testPushAndPop();
|
testPushAndPop();
|
||||||
testFifoOrder();
|
testFifoOrder();
|
||||||
testPopWaitsForMessage();
|
testPopWaitsForMessage();
|
||||||
testConcurrentPushAndPop();
|
testConcurrentPushAndPop();
|
||||||
|
testStopReturnsNullopt();
|
||||||
|
testStopReturnsRemainingMessages();
|
||||||
}
|
}
|
||||||
|
|
||||||
} // namespace app::tests
|
} // namespace app::tests
|
||||||
|
|||||||
Reference in New Issue
Block a user