feat(app): add concurrent message queue

This commit is contained in:
user
2026-07-20 23:00:56 +04:00
parent 31d1fbdb8a
commit a3cd178155
8 changed files with 298 additions and 24 deletions
+6 -24
View File
@@ -1,41 +1,23 @@
ROOT_DIR := $(abspath ..) CXXFLAGS := \
-I$(ROOT_DIR)/core/include \
CC := gcc -I$(ROOT_DIR)/logger/include
CXXFLAGS := -std=c++17 \
-Wall \
-Wextra \
-Werror \
-I$(ROOT_DIR)/core/include \
-I$(ROOT_DIR)/logger/include
LDFLAGS := -L$(ROOT_DIR)/build \ LDFLAGS := -L$(ROOT_DIR)/build \
-llogger \ -Wl,-rpath,'$$ORIGIN'
-lstdc++ \
-Wl,-rpath,$(ROOT_DIR)/build
BUILD_DIR := $(ROOT_DIR)/build LDLIBS += -llogger
TARGET := $(BUILD_DIR)/logger_app TARGET := $(BUILD_DIR)/logger_app
SRC := src/main.cpp SRC := src/main.cpp
.PHONY: all clean .PHONY: all clean
all: $(TARGET) all: $(TARGET)
$(TARGET): $(SRC) $(TARGET): $(SRC)
@mkdir -p $(BUILD_DIR) @mkdir -p $(BUILD_DIR)
$(CC) $(CXXFLAGS) $(LDFLAGS) $^ -o $@ $(LDLIBS)
$(CC) \
$(CXXFLAGS) \
$< \
-o $@ \
$(LDFLAGS)
clean: clean:
rm -f $(TARGET) rm -f $(TARGET)
+27
View File
@@ -0,0 +1,27 @@
#pragma once
#include <string>
#include "log_level.hpp"
namespace app {
/**
* @brief Сообщение, передаваемое в поток записи.
*
* Содержит текст сообщения и уровень логирования,
* с которым оно должно быть записано в журнал.
*/
struct LogMessage {
/**
* @brief Уровень логирования сообщения.
*/
logger::LogLevel level;
/**
* @brief Текст сообщения.
*/
std::string message;
};
} // namespace app
+70
View File
@@ -0,0 +1,70 @@
#pragma once
#include <condition_variable>
#include <mutex>
#include <queue>
#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<LogMessage> queue_;
/**
* @brief Мьютекс для синхронизации доступа к очереди.
*/
mutable std::mutex mutex_;
/**
* @brief Условная переменная для ожидания новых сообщений.
*/
std::condition_variable conditionVariable_;
};
} // namespace app
+31
View File
@@ -0,0 +1,31 @@
#include "message_queue.hpp"
namespace app {
void MessageQueue::push(LogMessage message) {
std::lock_guard<std::mutex> lock(mutex_);
queue_.push(std::move(message));
conditionVariable_.notify_one();
}
LogMessage MessageQueue::pop() {
std::unique_lock<std::mutex> 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<std::mutex> lock(mutex_);
return queue_.empty();
}
} // namespace app
+23
View File
@@ -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)
+23
View File
@@ -0,0 +1,23 @@
#include "thread_test_utils.hpp"
#include <cassert>
#include <cstddef>
#include <thread>
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
+13
View File
@@ -0,0 +1,13 @@
#pragma once
#include <cstddef>
#include <thread>
#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
+105
View File
@@ -0,0 +1,105 @@
#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;
}