Проблема разделяемого состояния

Типичный подход к многопоточности: разделяемые данные, защищённые мьютексами.

 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 — тот же паттерн, другая среда выполнения
  • Используйте ограниченные очереди и явно обрабатывайте противодавление

Что дальше