Быстрая очередь сообщений (FMQ)

Инфраструктура удаленного вызова процедур (RPC) HIDL использует механизмы связывания, а это значит, что вызовы требуют дополнительных ресурсов, операций ядра и могут запускать действия планировщика. Однако в случаях, когда данные должны передаваться между процессами с меньшими накладными расходами и без участия ядра, используется система Fast Message Queue (FMQ).

FMQ создает очереди сообщений с нужными свойствами. Вы можете отправить объект MQDescriptorSync или MQDescriptorUnsync через вызов HIDL RPC, и этот объект будет использоваться процессом-получателем для доступа к очереди сообщений.

Типы очередей

В Android поддерживаются два типа очередей (варианты):

  • Несинхронизированные очереди могут переполняться и иметь много читателей. Каждый читатель должен вовремя считывать данные, иначе они будут потеряны.
  • Синхронизированные очереди не могут переполняться и могут иметь только одного читателя.

Оба типа очередей не могут быть пустыми (чтение из пустой очереди приводит к ошибке) и могут иметь только одного отправителя.

Несинхронизированные очереди

У очереди без синхронизации может быть только один отправитель, но любое количество получателей. У очереди есть одна позиция записи, но каждый читатель отслеживает собственную позицию чтения.

Запись в очередь всегда выполняется успешно (не проверяется на переполнение), если размер записываемых данных не превышает настроенную емкость очереди (запись данных, размер которых превышает емкость очереди, сразу же завершается ошибкой). Поскольку у каждого читателя может быть разная позиция чтения, данные удаляются из очереди, когда требуется место для новых записей, а не когда все читатели прочитают все данные.

Читатели должны получать данные до того, как они будут удалены из очереди. При попытке чтения большего объема данных, чем доступно, операция чтения либо немедленно завершается с ошибкой (если она неблокирующая), либо ожидает, пока не будет доступно достаточно данных (если она блокирующая). Попытка чтения большего объема данных, чем вмещает очередь, всегда приводит к немедленному сбою.

Если читатель не успевает за писателем, так что объем записанных и ещё не прочитанных этим читателем данных превышает емкость очереди, следующее чтение не возвращает данные. Вместо этого оно сбрасывает позицию чтения читателя на позицию записи плюс половина емкости, а затем возвращает ошибку. Таким образом, половина буфера остается доступной для чтения, а также резервируется место для новых записей, чтобы избежать переполнения очереди. Если проверить доступные для чтения данные после переполнения, но до следующего чтения, то будет показано, что доступно больше данных, чем вмещает очередь. Это означает, что произошло переполнение. (Если очередь переполняется между проверкой доступных данных и попыткой их прочитать, единственным признаком переполнения будет сбой чтения.)

Синхронизированные очереди

В синхронизированной очереди есть один записывающий и один читающий процесс, а также одна позиция записи и одна позиция чтения. Невозможно записать больше данных, чем позволяет размер очереди, или прочитать больше данных, чем в ней содержится. В зависимости от того, какая функция записи или чтения (блокирующая или неблокирующая) вызывается, попытки превысить доступное пространство или объем данных либо немедленно возвращают ошибку, либо блокируются до тех пор, пока желаемая операция не будет выполнена. Попытки прочитать или записать больше данных, чем вмещает очередь, всегда завершаются неудачей.

Как настроить FMQ

Для очереди сообщений требуется несколько объектов MessageQueue: один для записи и один или несколько для чтения. Нет явной конфигурации того, какой объект используется для записи или чтения. Пользователь должен убедиться, что ни один объект не используется для чтения и записи одновременно, что есть не более одного объекта для записи, а для синхронизированных очередей – не более одного объекта для чтения.

Создайте первый объект MessageQueue

Очередь сообщений создается и настраивается с помощью одного вызова:

#include <fmq/MessageQueue.h>
using android::hardware::kSynchronizedReadWrite;
using android::hardware::kUnsynchronizedWrite;
using android::hardware::MQDescriptorSync;
using android::hardware::MQDescriptorUnsync;
using android::hardware::MessageQueue;
....
// For a synchronized nonblocking FMQ
mFmqSynchronized =
  new (std::nothrow) MessageQueue<uint16_t, kSynchronizedReadWrite>
      (kNumElementsInQueue);
// For an unsynchronized FMQ that supports blocking
mFmqUnsynchronizedBlocking =
  new (std::nothrow) MessageQueue<uint16_t, kUnsynchronizedWrite>
      (kNumElementsInQueue, true /* enable blocking operations */);
  • Инициализатор MessageQueue<T, flavor>(numElements) создает и инициализирует объект, поддерживающий функциональность очереди сообщений.
  • Инициализатор MessageQueue<T, flavor>(numElements, configureEventFlagWord) создает и инициализирует объект, поддерживающий функции очереди сообщений с блокировкой.
  • flavor может быть kSynchronizedReadWrite для синхронизированной очереди или kUnsynchronizedWrite для несинхронизированной очереди.
  • uint16_t (в этом примере) может быть любым типом, определенным в HIDL, который не включает вложенные буферы (не типы string или vec), дескрипторы или интерфейсы.
  • kNumElementsInQueue указывает размер очереди в количестве записей. Он определяет размер буфера общей памяти, выделенного для очереди.

Создайте второй объект MessageQueue

Вторая сторона очереди сообщений создается с помощью объекта MQDescriptor, полученного от первой стороны. Объект MQDescriptor отправляется через вызов HIDL или AIDL RPC в процесс, который содержит второй конец очереди сообщений. MQDescriptor содержит информацию об очереди, в том числе:

  • Информация для сопоставления буфера и указателя записи.
  • Информация для сопоставления указателя чтения (если очередь синхронизирована).
  • Информация для сопоставления слова флага события (если очередь блокируется).
  • Тип объекта (<T, flavor>), который включает тип, определенный в HIDL, элементов очереди и ее тип (синхронизированная или несинхронизированная).

Объект MQDescriptor можно использовать для создания объекта MessageQueue:

MessageQueue<T, flavor>::MessageQueue(const MQDescriptor<T, flavor>& Desc, bool resetPointers)

Параметр resetPointers указывает, нужно ли сбрасывать позиции чтения и записи до 0 при создании объекта MessageQueue. В несинхронизированной очереди позиция чтения (которая является локальной для каждого объекта MessageQueue в несинхронизированных очередях) всегда устанавливается на 0 при создании. Как правило, MQDescriptor инициализируется при создании первого объекта очереди сообщений. Чтобы лучше контролировать общую память, вы можете вручную настроить MQDescriptor (определение MQDescriptor приведено в system/libhidl/base/include/hidl/MQDescriptor.h), а затем создать каждый объект MessageQueue, как описано в этом разделе.

Блокировка очередей и флагов событий

По умолчанию очереди не поддерживают чтение с блокировкой и запись. Существует два типа вызовов чтения с блокировкой и записи:

  • Краткая форма с тремя параметрами (указатель данных, количество элементов, время ожидания) поддерживает блокировку отдельных операций чтения и записи в одной очереди. При использовании этой формы флаг события и битовые маски обрабатываются очередью внутренним образом, а первый объект очереди сообщений должен быть инициализирован с помощью второго параметра true. Пример:
    // For an unsynchronized FMQ that supports blocking
    mFmqUnsynchronizedBlocking =
      new (std::nothrow) MessageQueue<uint16_t, kUnsynchronizedWrite>
          (kNumElementsInQueue, true /* enable blocking operations */);
    
  • Расширенная форма с шестью параметрами (включая флаг события и битовые маски) позволяет использовать общий объект EventFlag в нескольких очередях и указывать битовые маски уведомлений. В этом случае флаг события и битовые маски должны быть предоставлены для каждого вызова чтения и записи.

В длинной форме можно явно указать EventFlag в каждом вызове readBlocking() и writeBlocking(). Одну из очередей можно инициализировать с помощью внутреннего флага события, который затем необходимо извлечь из объектов MessageQueue этой очереди с помощью getEventFlagWord() и использовать для создания объектов EventFlag в каждом процессе для использования с другими очередями FMQ. Кроме того, вы можете инициализировать объекты EventFlag с помощью любой подходящей общей памяти.

Как правило, в каждой очереди должен использоваться только один из следующих типов вызовов: неблокирующий, краткосрочный блокирующий или долгосрочный блокирующий. Их можно использовать вместе, но для получения нужного результата требуется тщательное программирование.

Как сделать воспоминание доступным только для чтения

По умолчанию у общей памяти есть разрешения на чтение и запись. Для несинхронизированных очередей (kUnsynchronizedWrite) записывающий процесс может удалить разрешения на запись для всех считывающих процессов, прежде чем передавать объекты MQDescriptorUnsync. Это гарантирует, что другие процессы не смогут записывать данные в очередь, что рекомендуется для защиты от ошибок или нежелательного поведения в процессах чтения. Если автор хочет, чтобы читатели могли сбрасывать очередь всякий раз, когда они используют MQDescriptorUnsync для создания стороны чтения очереди, то память нельзя пометить как доступную только для чтения. Это режим работы по умолчанию для конструктора MessageQueue. Поэтому, если в очереди есть пользователи, их код нужно изменить, чтобы создать очередь с помощью resetPointer=false.

  • Вызовите ashmem_set_prot_region с дескриптором файла MQDescriptor и регионом, установленным в режим только чтение (PROT_READ):
    int res = ashmem_set_prot_region(mqDesc->handle->data[0], PROT_READ)
  • Читатель: создать очередь сообщений с помощью resetPointer=false (по умолчанию используется true):
    mFmq = new (std::nothrow) MessageQueue(mqDesc, false);

Как использовать MessageQueue

Общедоступный API объекта MessageQueue:

size_t availableToWrite() // Space available (number of elements).
size_t availableToRead() // Number of elements available.
size_t getQuantumSize() // Size of type T in bytes.
size_t getQuantumCount() // Number of items of type T that fit in the FMQ.
bool isValid() // Whether the FMQ is configured correctly.
const MQDescriptor<T, flavor>* getDesc() // Return info to send to other process.

bool write(const T* data)  // Write one T to FMQ; true if successful.
bool write(const T* data, size_t count) // Write count T's; no partial writes.

bool read(T* data); // read one T from FMQ; true if successful.
bool read(T* data, size_t count); // Read count T's; no partial reads.

bool writeBlocking(const T* data, size_t count, int64_t timeOutNanos = 0);
bool readBlocking(T* data, size_t count, int64_t timeOutNanos = 0);

// Allows multiple queues to share a single event flag word
std<::atomic>uint32_t* getEventFlagWord();

bool writeBlocking(const T* data, size_t count, uint32_t readNotification,
uint32_t writeNotification, int64_t timeOutNanos = 0,
android::hardware::EventFlag* evFlag = nullptr); // Blocking write operation for count Ts.

bool readBlocking(T* data, size_t count, uint32_t readNotification,
uint32_t writeNotification, int64_t timeOutNanos = 0,
android::hardware::EventFlag* evFlag = nullptr) // Blocking read operation for count Ts;

// APIs to allow zero copy read/write operations
bool beginWrite(size_t nMessages, MemTransaction* memTx) const;
bool commitWrite(size_t nMessages);
bool beginRead(size_t nMessages, MemTransaction* memTx) const;
bool commitRead(size_t nMessages);

Вы можете использовать availableToWrite() и availableToRead(), чтобы определить, сколько данных можно передать за одну операцию. В очереди без синхронизации:

  • availableToWrite() всегда возвращает емкость очереди.
  • У каждого считывателя есть своя позиция считывания и он выполняет собственный расчет для availableToRead().
  • С точки зрения медленного читателя очередь может переполняться, поэтому функция availableToRead() может возвращать значение, превышающее размер очереди. Первое чтение после переполнения завершается неудачно, и позиция чтения для этого читателя устанавливается равной текущему указателю записи независимо от того, было ли переполнение зарегистрировано с помощью availableToRead().

Методы read() и write() возвращают значение true, если все запрошенные данные были успешно переданы в очередь и из нее. Эти методы не блокируют выполнение программы, а сразу возвращают результат: true в случае успеха или false в случае ошибки.

Методы readBlocking() и writeBlocking() ожидают, пока запрошенная операция не будет выполнена или не истечет время ожидания (значение timeOutNanos, равное 0, означает, что время ожидания не ограничено).

Операции блокировки реализуются с помощью слова флага события. По умолчанию каждая очередь создает и использует собственное ключевое слово для поддержки краткой формы readBlocking() и writeBlocking(). Несколько очередей могут использовать одно слово, чтобы процесс мог ожидать записи или чтения в любую из очередей. Вызвав getEventFlagWord(), вы можете получить указатель на слово флага события очереди и использовать этот указатель (или любой указатель на подходящее место в общей памяти) для создания объекта EventFlag, который можно передать в длинную форму readBlocking() и writeBlocking() для другой очереди. Параметры readNotification и writeNotification указывают, какие биты в флаге события должны использоваться для сигнализации операций чтения и записи в этой очереди. readNotification и writeNotification – это 32-битные битовые маски.

readBlocking() ожидает биты writeNotification; если этот параметр равен 0, вызов всегда завершается неудачно. Если значение readNotification равно 0, вызов не завершится ошибкой, но при успешном чтении не будут установлены биты уведомлений. В синхронизированной очереди это означает, что соответствующий вызов writeBlocking() никогда не будет активирован, если бит не будет задан в другом месте. В несинхронизированной очереди writeBlocking() не ждет (его все равно следует использовать для установки бита уведомления о записи), и для операций чтения не нужно устанавливать биты уведомлений. Аналогично, writeblocking() завершается ошибкой, если readNotification равно 0, а успешная запись устанавливает указанные биты writeNotification.

Чтобы ожидать несколько очередей одновременно, используйте метод wait() объекта EventFlag, который позволяет ожидать битовую маску уведомлений. Метод wait() возвращает слово состояния с битами, которые вызвали пробуждение. Эта информация используется, чтобы проверить, достаточно ли места или данных в очереди для выполнения нужной операции записи или чтения, а также для выполнения неблокирующих операций write() и read(). Чтобы получить уведомление об операции публикации, используйте другой вызов метода wake() объекта EventFlag. Определение абстракции EventFlag приведено в статье system/libfmq/include/fmq/EventFlag.h.

Операции без копирования

Методы read, write, readBlocking и writeBlocking() принимают указатель на буфер ввода-вывода в качестве аргумента и используют вызовы memcpy() для копирования данных между одним и тем же кольцевым буфером FMQ. Чтобы повысить производительность, в Android 8.0 и более поздних версиях есть набор API, которые предоставляют прямой доступ к указателю в кольцевом буфере, устраняя необходимость использовать вызовы memcpy.

Для операций с очередью сообщений FMQ без копирования можно использовать следующие общедоступные API:

bool beginWrite(size_t nMessages, MemTransaction* memTx) const;
bool commitWrite(size_t nMessages);

bool beginRead(size_t nMessages, MemTransaction* memTx) const;
bool commitRead(size_t nMessages);
  • Метод beginWrite предоставляет базовые указатели на кольцевой буфер FMQ. После записи данных зафиксируйте их с помощью commitWrite(). Методы beginRead и commitRead действуют одинаково.
  • Методы beginRead и Write принимают в качестве входных данных количество сообщений, которые нужно прочитать и записать, и возвращают логическое значение, указывающее, возможно ли чтение или запись. Если чтение или запись возможны, структура memTx заполняется базовыми указателями, которые можно использовать для прямого доступа к памяти общего буфера кольца.
  • Структура MemRegion содержит сведения о блоке памяти, в том числе базовый указатель (базовый адрес блока памяти) и длину в T (длина блока памяти в терминах определенного HIDL типа очереди сообщений).
  • Структура MemTransaction содержит две структуры MemRegion, first и second, поскольку для чтения или записи в кольцевой буфер может потребоваться перенос в начало очереди. Это означает, что для чтения и записи данных в кольцевой буфер FMQ необходимы два базовых указателя.

Чтобы получить базовый адрес и длину из структуры MemRegion, выполните следующие действия:

T* getAddress(); // gets the base address
size_t getLength(); // gets the length of the memory region in terms of T
size_t getLengthInBytes(); // gets the length of the memory region in bytes

Чтобы получить ссылки на первую и вторую структуры MemRegion в объекте MemTransaction:

const MemRegion& getFirstRegion(); // get a reference to the first MemRegion
const MemRegion& getSecondRegion(); // get a reference to the second MemRegion

Пример записи в FMQ с помощью API без копирования:

MessageQueueSync::MemTransaction tx;
if (mQueue->beginRead(dataLe&n, tx)) {
    auto first = tx.getFirstRegion();
    auto second = tx.getSecondRegion();

    foo(first.getAddress(), first.getLength()); // method that performs the data write
    foo(second.getAddress(), second.getLength()); // method that performs the data write

    if(commitWrite(dataLen) == false) {
       // report error
    }
} else {
   // report error
}

В MemTransaction также входят следующие вспомогательные методы:

  • T* getSlot(size_t idx); возвращает указатель на рекламное место idx в массиве MemRegions, который является частью объекта MemTransaction. Если объект MemTransaction представляет области памяти для чтения и записи N элементов типа T, то допустимый диапазон idx – от 0 до N-1.
  • bool copyTo(const T* data, size_t startIdx, size_t nMessages = 1); записывает nMessages объектов типа T в области памяти, описанные объектом, начиная с индекса startIdx. Этот метод использует memcpy() и не предназначен для операций без копирования. Если объект MemTransaction представляет память для чтения и записи N элементов типа T, то допустимый диапазон для idx – от 0 до N-1.
  • bool copyFrom(T* data, size_t startIdx, size_t nMessages = 1); – вспомогательный метод для чтения элементов nMessages типа T из областей памяти, описанных объектом, начиная с startIdx. Этот метод использует memcpy() и не предназначен для операций с нулевым копированием.

Отправка очереди через HIDL

При создании контента:

  1. Создайте объект очереди сообщений, как описано выше.
  2. Убедитесь, что объект действителен, с помощью isValid().
  3. Если вы ожидаете несколько очередей, передавая EventFlag в длинную форму readBlocking() или writeBlocking(), вы можете извлечь указатель флага события (с помощью getEventFlagWord()) из объекта MessageQueue, который был инициализирован для создания флага, и использовать этот флаг для создания нужного объекта EventFlag.
  4. Используйте метод getDesc() класса MessageQueue, чтобы получить объект дескриптора.
  5. В файле HAL укажите для метода параметр типа fmq_sync или fmq_unsync, где T – подходящий тип, определенный в HIDL. Используйте его, чтобы отправить объект, возвращенный getDesc(), в процесс получения.

На стороне получателя:

  1. Используйте объект дескриптора, чтобы создать объект MessageQueue. Используйте тот же тип очереди и данных, иначе шаблон не будет скомпилирован.
  2. Если вы извлекли флаг события, извлеките его из соответствующего объекта MessageQueue в принимающем процессе.
  3. Используйте объект MessageQueue для передачи данных.