Проблема разделяемого состояния
Типичный подход к многопоточности: разделяемые данные, защищённые мьютексами.
1class SensorManager {
2 std::mutex mtx_;
3 SensorData latest_; // разделяемое — оба потока обращаются
4 bool updated_ = false;
5public:
6 void updateFromSensor() { // вызывается из потока сенсора
7 std::lock_guard lk(mtx_);
8 latest_ = readSensor();
9 updated_ = true;
10 }
11
12 SensorData getLatest() { // вызывается из потока дисплея
13 std::lock_guard lk(mtx_);
14 updated_ = false;
15 return latest_;
16 }
17};
Для двух потоков это работает, но по мере роста числа вызывающих конкуренция за блокировку растёт, а рассуждения о том, какой поток владеет каким состоянием, усложняются. Каждый метод — потенциальное место дедлока.
Паттерн «активный объект» устраняет разделяемое состояние: объект владеет собственным потоком и обрабатывает все вызовы как сообщения из очереди. Ни один внешний поток никогда не трогает его внутренности напрямую.
Активный объект
1#include <thread>
2#include <queue>
3#include <mutex>
4#include <condition_variable>
5#include <functional>
6#include <atomic>
7
8class ActiveObject {
9public:
10 ActiveObject() : running_(true), thread_([this] { run(); }) {}
11
12 ~ActiveObject() {
13 post([this] { running_ = false; }); // «таблетка яда»
14 thread_.join();
15 }
16
17 // Поставить сообщение (callable) во внутреннюю очередь — неблокирующий
18 void post(std::function<void()> msg) {
19 {
20 std::lock_guard<std::mutex> lk(mtx_);
21 queue_.push(std::move(msg));
22 }
23 cv_.notify_one();
24 }
25
26private:
27 void run() {
28 while (running_) {
29 std::function<void()> msg;
30 {
31 std::unique_lock<std::mutex> lk(mtx_);
32 cv_.wait(lk, [this] { return !queue_.empty(); });
33 msg = std::move(queue_.front());
34 queue_.pop();
35 }
36 msg(); // выполнять в этом потоке — внешняя блокировка не нужна
37 }
38 }
39
40 std::atomic<bool> running_;
41 std::thread thread_;
42 std::mutex mtx_;
43 std::condition_variable cv_;
44 std::queue<std::function<void()>> queue_;
45};
Всё состояние живёт внутри объекта. Поток выполняет сообщения последовательно — два сообщения никогда не выполняются одновременно, поэтому внутренний мьютекс для состояния не нужен.
Конкретный активный сенсор
1class SensorService : private ActiveObject {
2 // Состояние — доступно только из внутреннего потока
3 SensorData latest_;
4 std::vector<std::function<void(SensorData)>> subscribers_;
5
6public:
7 // Публичный API — вызывается из любого потока, отправляет во внутреннюю очередь
8 void requestRead() {
9 post([this] { doRead(); });
10 }
11
12 void subscribe(std::function<void(SensorData)> cb) {
13 post([this, cb = std::move(cb)] mutable {
14 subscribers_.push_back(std::move(cb));
15 });
16 }
17
18 // Возвращает future — вызывающий может ждать результата
19 std::future<SensorData> getLatest() {
20 auto promise = std::make_shared<std::promise<SensorData>>();
21 auto future = promise->get_future();
22 post([this, promise] {
23 promise->set_value(latest_);
24 });
25 return future;
26 }
27
28private:
29 // Выполняется только во внутреннем потоке
30 void doRead() {
31 latest_ = readSensor(); // нет блокировки — только этот поток обращается к latest_
32 for (auto& cb : subscribers_)
33 cb(latest_);
34 }
35};
Использование из нескольких потоков:
1SensorService sensor;
2
3// Поток дисплея — подписывается на обновления
4sensor.subscribe([](SensorData d) {
5 display.update(d.temperature); // вызывается в потоке сенсора — осторожно!
6});
7
8// Поток управления — периодические чтения
9while (true) {
10 sensor.requestRead();
11 std::this_thread::sleep_for(std::chrono::milliseconds(100));
12}
13
14// Поток запроса — ждать значения
15auto future = sensor.getLatest();
16auto data = future.get(); // блокирует до ответа активного объекта
Проблема владения callback
В примере с подписчиком, callback display.update выполняется в потоке сенсора,
а не в потоке дисплея. Если у display есть собственный поток и собственное
состояние, это нарушение потокобезопасности.
Правильное решение: сделать дисплей тоже активным объектом, а callback-и публиковать в очередь дисплея:
1class DisplayService : private ActiveObject {
2public:
3 void update(SensorData d) {
4 post([this, d] { doUpdate(d); });
5 }
6private:
7 void doUpdate(SensorData d) {
8 lcd.draw(d.temperature); // выполняется только в потоке DisplayService
9 }
10};
11
12DisplayService display;
13
14sensor.subscribe([&display](SensorData d) {
15 display.update(d); // публикует в очередь дисплея — неблокирующий, потокобезопасный
16});
Каждый активный объект владеет своим состоянием и обрабатывает сообщения из своего потока.
Коммуникация между ними всегда через post — никогда через прямые вызовы.
Вариант для FreeRTOS
В FreeRTOS паттерн «активный объект» напрямую отображается на задачу + очередь:
1class ActiveSensorTask {
2 QueueHandle_t queue_;
3 TaskHandle_t task_;
4
5 struct Message {
6 enum class Type { Read, Subscribe, Stop } type;
7 // payload...
8 };
9
10public:
11 ActiveSensorTask() {
12 queue_ = xQueueCreate(16, sizeof(Message));
13 xTaskCreate(taskFn, "Sensor", 512, this, 2, &task_);
14 }
15
16 void requestRead() {
17 Message msg{Message::Type::Read};
18 xQueueSend(queue_, &msg, 0);
19 }
20
21 void requestStop() {
22 Message msg{Message::Type::Stop};
23 xQueueSend(queue_, &msg, portMAX_DELAY);
24 }
25
26private:
27 static void taskFn(void* param) {
28 static_cast<ActiveSensorTask*>(param)->run();
29 }
30
31 void run() {
32 Message msg;
33 while (true) {
34 if (xQueueReceive(queue_, &msg, portMAX_DELAY) == pdTRUE) {
35 if (msg.type == Message::Type::Stop) break;
36 if (msg.type == Message::Type::Read) doRead();
37 }
38 }
39 vTaskDelete(nullptr);
40 }
41
42 void doRead() {
43 SensorData d = readSensor();
44 // уведомить другие задачи через их очереди
45 }
46};
Структура идентична: задача владеет своим состоянием, получает типизированные сообщения из очереди, обрабатывает их последовательно. Мьютекс для внутреннего состояния не нужен.
Ограниченная очередь и противодавление
std::queue в базовой реализации растёт без ограничений — медленный потребитель может
исчерпать память. Используйте ограниченную очередь с post, сигнализирующим о противодавлении:
1bool post(std::function<void()> msg) {
2 std::lock_guard<std::mutex> lk(mtx_);
3 if (queue_.size() >= maxQueueSize_) return false; // отклонить
4 queue_.push(std::move(msg));
5 cv_.notify_one();
6 return true;
7}
В FreeRTOS xQueueSend с timeout 0 возвращает errQUEUE_FULL — вызывающий
решает: сбросить, повторить или эскалировать.
Активный объект vs модель акторов
Паттерн «активный объект» — это C++-версия модели акторов (Erlang, Akka). Каждый объект — актор: изолированное состояние, коммуникация через обмен сообщениями.
Разница: акторы обычно имеют среду выполнения, маршрутизирующую сообщения между объектами по имени или адресу. Активные объекты соединяют свои очереди явно — проще, без фреймворка, но менее динамично.
Для небольших встраиваемых систем паттерн «активный объект» с задачами и очередями FreeRTOS даёт преимущества изоляции акторов без фреймворка.
Итоги
- Активный объект = один поток + одна очередь + всё состояние внутри
- Внешние потоки публикуют сообщения (callable или типизированные структуры) — никогда не вызывают методы напрямую на разделяемом состоянии
- Внутренний поток выполняет сообщения последовательно — внутренний мьютекс не нужен
- Callback-и, публикующие в очередь другого активного объекта, сохраняют инвариант: нет доступа к состоянию между потоками
- Напрямую отображается на задачу + очередь FreeRTOS — тот же паттерн, другая среда выполнения
- Используйте ограниченные очереди и явно обрабатывайте противодавление
Что дальше
- FreeRTOS: задачи и очереди — реализация этого паттерна на FreeRTOS
- Lock-free очереди — замена очереди с мьютексом на lock-free SPSC
- Паттерн «наблюдатель» — callback-и между активными объектами