From a3cd178155e8385e38f5c1c173a13df78792e7e8 Mon Sep 17 00:00:00 2001 From: user Date: Mon, 20 Jul 2026 23:00:56 +0400 Subject: [PATCH] feat(app): add concurrent message queue --- app/Makefile | 30 ++----- app/include/log_message.hpp | 27 ++++++ app/include/message_queue.hpp | 70 ++++++++++++++++ app/src/message_queue.cpp | 31 +++++++ tests/app/Makefile | 23 ++++++ tests/app/helpers/thread_test_utils.cpp | 23 ++++++ tests/app/helpers/thread_test_utils.hpp | 13 +++ tests/app/message_queue_test.cpp | 105 ++++++++++++++++++++++++ 8 files changed, 298 insertions(+), 24 deletions(-) create mode 100644 app/include/log_message.hpp create mode 100644 app/include/message_queue.hpp create mode 100644 app/src/message_queue.cpp create mode 100644 tests/app/Makefile create mode 100644 tests/app/helpers/thread_test_utils.cpp create mode 100644 tests/app/helpers/thread_test_utils.hpp create mode 100644 tests/app/message_queue_test.cpp diff --git a/app/Makefile b/app/Makefile index 502fd32..d81ac72 100644 --- a/app/Makefile +++ b/app/Makefile @@ -1,41 +1,23 @@ -ROOT_DIR := $(abspath ..) - -CC := gcc - -CXXFLAGS := -std=c++17 \ - -Wall \ - -Wextra \ - -Werror \ - -I$(ROOT_DIR)/core/include \ - -I$(ROOT_DIR)/logger/include +CXXFLAGS := \ + -I$(ROOT_DIR)/core/include \ + -I$(ROOT_DIR)/logger/include LDFLAGS := -L$(ROOT_DIR)/build \ - -llogger \ - -lstdc++ \ - -Wl,-rpath,$(ROOT_DIR)/build + -Wl,-rpath,'$$ORIGIN' -BUILD_DIR := $(ROOT_DIR)/build +LDLIBS += -llogger TARGET := $(BUILD_DIR)/logger_app SRC := src/main.cpp - .PHONY: all clean - all: $(TARGET) - $(TARGET): $(SRC) @mkdir -p $(BUILD_DIR) - - $(CC) \ - $(CXXFLAGS) \ - $< \ - -o $@ \ - $(LDFLAGS) - + $(CC) $(CXXFLAGS) $(LDFLAGS) $^ -o $@ $(LDLIBS) clean: rm -f $(TARGET) diff --git a/app/include/log_message.hpp b/app/include/log_message.hpp new file mode 100644 index 0000000..caf2b32 --- /dev/null +++ b/app/include/log_message.hpp @@ -0,0 +1,27 @@ +#pragma once + +#include + +#include "log_level.hpp" + +namespace app { + +/** + * @brief Сообщение, передаваемое в поток записи. + * + * Содержит текст сообщения и уровень логирования, + * с которым оно должно быть записано в журнал. + */ +struct LogMessage { + /** + * @brief Уровень логирования сообщения. + */ + logger::LogLevel level; + + /** + * @brief Текст сообщения. + */ + std::string message; +}; + +} // namespace app diff --git a/app/include/message_queue.hpp b/app/include/message_queue.hpp new file mode 100644 index 0000000..57e8f09 --- /dev/null +++ b/app/include/message_queue.hpp @@ -0,0 +1,70 @@ +#pragma once + +#include +#include +#include + +#include "log_message.hpp" + +namespace app { + +/** + * @brief Потокобезопасная очередь сообщений. + * + * Используется для передачи сообщений от потока, + * принимающего ввод пользователя, к потоку, + * выполняющему запись в журнал. + */ +class MessageQueue { +public: + /** + * @brief Добавляет сообщение в очередь. + * + * Метод является потокобезопасным. + * + * @param message Сообщение для передачи. + */ + void push(LogMessage message); + + /** + * @brief Извлекает сообщение из очереди. + * + * Если очередь пуста, метод блокирует вызывающий + * поток до появления нового сообщения. + * + * Метод является потокобезопасным. + * + * @return Следующее сообщение из очереди. + */ + [[nodiscard]] + LogMessage pop(); + + /** + * @brief Проверяет, пуста ли очередь. + * + * Метод является потокобезопасным. + * + * @return true, если очередь не содержит сообщений. + * @return false, если очередь содержит хотя бы одно сообщение. + */ + [[nodiscard]] + bool empty() const; + +private: + /** + * @brief Очередь сообщений. + */ + std::queue queue_; + + /** + * @brief Мьютекс для синхронизации доступа к очереди. + */ + mutable std::mutex mutex_; + + /** + * @brief Условная переменная для ожидания новых сообщений. + */ + std::condition_variable conditionVariable_; +}; + +} // namespace app diff --git a/app/src/message_queue.cpp b/app/src/message_queue.cpp new file mode 100644 index 0000000..a513c8f --- /dev/null +++ b/app/src/message_queue.cpp @@ -0,0 +1,31 @@ +#include "message_queue.hpp" + +namespace app { + +void MessageQueue::push(LogMessage message) { + std::lock_guard lock(mutex_); + + queue_.push(std::move(message)); + + conditionVariable_.notify_one(); +} + +LogMessage MessageQueue::pop() { + std::unique_lock lock(mutex_); + + conditionVariable_.wait(lock, [this]() { return !queue_.empty(); }); + + auto message = std::move(queue_.front()); + + queue_.pop(); + + return message; +} + +bool MessageQueue::empty() const { + std::lock_guard lock(mutex_); + + return queue_.empty(); +} + +} // namespace app diff --git a/tests/app/Makefile b/tests/app/Makefile new file mode 100644 index 0000000..05a58c6 --- /dev/null +++ b/tests/app/Makefile @@ -0,0 +1,23 @@ +CXXFLAGS += \ + -I$(ROOT_DIR)/app/include \ + -I$(ROOT_DIR)/logger/include \ + -Ihelpers + +TARGET := $(BUILD_DIR)/app_test + +SRC := \ + $(ROOT_DIR)/app/src/message_queue.cpp \ + helpers/thread_test_utils.cpp \ + message_queue_test.cpp + +.PHONY: all clean + +all: $(TARGET) + $(TARGET) + +$(TARGET): $(SRC) + @mkdir -p $(BUILD_DIR) + $(CC) $(CXXFLAGS) $^ -o $@ $(LDLIBS) + +clean: + rm -f $(TARGET) diff --git a/tests/app/helpers/thread_test_utils.cpp b/tests/app/helpers/thread_test_utils.cpp new file mode 100644 index 0000000..d31be07 --- /dev/null +++ b/tests/app/helpers/thread_test_utils.cpp @@ -0,0 +1,23 @@ +#include "thread_test_utils.hpp" + +#include +#include +#include + +namespace app::tests::helpers { + +std::thread createProducer(MessageQueue &queue, std::size_t threadIndex, + std::size_t messagesCount) { + return std::thread([&queue, threadIndex, messagesCount]() { + for (std::size_t messageIndex = 0; messageIndex < messagesCount; + ++messageIndex) { + queue.push({ + logger::LogLevel::info(), + "Thread " + std::to_string(threadIndex) + ", message " + + std::to_string(messageIndex), + }); + } + }); +} + +} // namespace app::tests::helpers diff --git a/tests/app/helpers/thread_test_utils.hpp b/tests/app/helpers/thread_test_utils.hpp new file mode 100644 index 0000000..34ae30f --- /dev/null +++ b/tests/app/helpers/thread_test_utils.hpp @@ -0,0 +1,13 @@ +#pragma once + +#include +#include + +#include "message_queue.hpp" + +namespace app::tests::helpers { + +std::thread createProducer(MessageQueue &queue, std::size_t threadIndex, + std::size_t messagesCount); + +} // namespace app::tests::helpers diff --git a/tests/app/message_queue_test.cpp b/tests/app/message_queue_test.cpp new file mode 100644 index 0000000..7598541 --- /dev/null +++ b/tests/app/message_queue_test.cpp @@ -0,0 +1,105 @@ +#include +#include +#include + +#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 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; +}