diff --git a/app/include/message_queue.hpp b/app/include/message_queue.hpp index 57e8f09..887cc29 100644 --- a/app/include/message_queue.hpp +++ b/app/include/message_queue.hpp @@ -2,6 +2,7 @@ #include #include +#include #include #include "log_message.hpp" @@ -22,6 +23,8 @@ public: * * Метод является потокобезопасным. * + * После остановки очереди новые сообщения не принимаются. + * * @param message Сообщение для передачи. */ void push(LogMessage message); @@ -29,15 +32,28 @@ public: /** * @brief Извлекает сообщение из очереди. * - * Если очередь пуста, метод блокирует вызывающий - * поток до появления нового сообщения. + * Если очередь пуста, метод блокирует вызывающий поток + * до появления нового сообщения или до остановки очереди. + * + * После остановки очереди и обработки всех оставшихся сообщений + * возвращается std::nullopt. * * Метод является потокобезопасным. * - * @return Следующее сообщение из очереди. + * @return Следующее сообщение или std::nullopt, если очередь остановлена. */ [[nodiscard]] - LogMessage pop(); + std::optional pop(); + + /** + * @brief Останавливает очередь. + * + * Пробуждает все ожидающие потоки. После вызова stop() + * новые сообщения не принимаются. + * + * Метод является потокобезопасным. + */ + void stop(); /** * @brief Проверяет, пуста ли очередь. @@ -56,6 +72,11 @@ private: */ std::queue queue_; + /** + * @brief Признак остановки очереди. + */ + bool stopped_{false}; + /** * @brief Мьютекс для синхронизации доступа к очереди. */ diff --git a/app/src/message_queue.cpp b/app/src/message_queue.cpp index a513c8f..a6cc2c5 100644 --- a/app/src/message_queue.cpp +++ b/app/src/message_queue.cpp @@ -5,15 +5,24 @@ namespace app { void MessageQueue::push(LogMessage message) { std::lock_guard lock(mutex_); + if (stopped_) { + return; + } + queue_.push(std::move(message)); conditionVariable_.notify_one(); } -LogMessage MessageQueue::pop() { +std::optional MessageQueue::pop() { std::unique_lock 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()); @@ -22,6 +31,16 @@ LogMessage MessageQueue::pop() { return message; } +void MessageQueue::stop() { + { + std::lock_guard lock(mutex_); + + stopped_ = true; + } + + conditionVariable_.notify_all(); +} + bool MessageQueue::empty() const { std::lock_guard lock(mutex_); diff --git a/tests/app/message_queue_test.cpp b/tests/app/message_queue_test.cpp index 7598541..71133fc 100644 --- a/tests/app/message_queue_test.cpp +++ b/tests/app/message_queue_test.cpp @@ -17,8 +17,9 @@ void testPushAndPop() { auto message = queue.pop(); - assert(message.level == logger::LogLevel::info()); - assert(message.message == "Hello"); + assert(message.has_value()); + assert(message->level == logger::LogLevel::info()); + assert(message->message == "Hello"); } void testFifoOrder() { @@ -28,9 +29,9 @@ void testFifoOrder() { 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"); + assert(queue.pop()->message == "first"); + assert(queue.pop()->message == "second"); + assert(queue.pop()->message == "third"); } void testPopWaitsForMessage() { @@ -41,8 +42,9 @@ void testPopWaitsForMessage() { std::thread worker([&]() { auto message = queue.pop(); - assert(message.level == logger::LogLevel::error()); - assert(message.message == "Delayed"); + assert(message.has_value()); + assert(message->level == logger::LogLevel::error()); + assert(message->message == "Delayed"); received = true; }); @@ -74,12 +76,14 @@ void testConcurrentPushAndPop() { 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()); + assert(message.has_value()); + assert(!message->message.empty()); } for (auto &producer : producers) { @@ -89,11 +93,44 @@ void testConcurrentPushAndPop() { 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() { testPushAndPop(); testFifoOrder(); testPopWaitsForMessage(); testConcurrentPushAndPop(); + testStopReturnsNullopt(); + testStopReturnsRemainingMessages(); } } // namespace app::tests