12 KiB
ll_queue — Двусвязная FIFO-очередь с автозабором и backpressure
1. Назначение
Двусвязная очередь элементов (FIFO) с автоматическим вызовом callback при появлении данных, пороговым ожиданием освобождения (backpressure), хеш-индексом для поиска и поддержкой пулов памяти. Используется как основной механизм передачи данных между компонентами внутри одного потока uasyc: TUN ↔ ETCP ↔ routing ↔ normalizer и т.д.
Работает строго в одном потоке (один uasyc). Для многопоточного доступа не предназначена.
2. Как пользоваться
2.1. Создание очереди
// Без хеша (простая очередь)
struct ll_queue* q = queue_new(ua, 0, 0, 0, "my_queue");
// С хешем — для быстрого поиска по ключу внутри data[]
// hash_size=256, ключ лежит в data[] по смещению 0, размер ключа 4 байта
struct ll_queue* q = queue_new(ua, 256, 0, 4, "my_hash_queue");
2.2. Создание элемента
Три способа, по убыванию производительности:
// 1. Из memory pool (самый быстрый)
struct ll_entry* e = queue_entry_new_from_pool(instance->pkt_pool);
// 2. Через malloc (data[] фиксированного размера одним блоком)
struct ll_entry* e = queue_entry_new(sizeof(struct my_data));
// 3. Элемент + отдельный буфер dgram (для переменного размера данных)
struct ll_entry* e = ll_alloc_lldgram(dgram_len);
memcpy(e->dgram, src_data, dgram_len);
e->len = dgram_len;
2.3. Запись в очередь (backpressure)
Главное правило: не забиваем очередь! Добавляем следующий элемент только когда очередь опустела до заданного порога. Для этого используется Пороговое ожидание:
// 1. Настроить порог (например: ждать пока count<=0, т.е. очередь совсем пуста)
queue_set_threshold(q, 0, 0);
// 2. В структуре-производителе завести handle:
struct queue_waiter_handle my_waiter = {0};
// 3. При готовности отправить — зарегистрировать ожидание:
static void my_send_cb(struct ll_queue* q, void* arg) {
struct my_producer* p = (struct my_producer*)arg;
struct ll_entry* e = ... // создать элемент
e->len = data_len;
queue_data_put_with_index(q, e);
// Снова зарегистрироваться на следующую порцию
queue_waiter_wait(q, &p->waiter, my_send_cb, p);
}
// Первичная регистрация:
queue_waiter_wait(q, &my_waiter, my_send_cb, producer);
// 4. При деинициализации — отменить ожидание:
queue_waiter_cancel(q, &my_waiter);
2.4. Чтение из очереди (callback)
// Установить callback — будет вызываться при появлении элементов
queue_set_callback(q, my_callback, my_data);
static void my_callback(struct ll_queue* q, void* arg) {
// 1. ИЗВЛЕЧЬ элемент
struct ll_entry* e = queue_data_get(q);
if (!e) { queue_resume_callback(q); return; }
// 2. ОБРАБОТАТЬ
if (e->dgram && e->len > 0) {
do_something(e->dgram, e->len);
queue_dgram_free(e); // сначала освободить dgram
}
queue_entry_free(e); // потом освободить сам entry
// 3. ОБЯЗАТЕЛЬНО возобновить callback
queue_resume_callback(q);
}
Критически важно: queue_data_get() автоматически приостанавливает callback-и (callback_suspended = 1). Без вызова queue_resume_callback() очередь навсегда застрянет — callback больше никогда не вызовется.
Во всех ветках (включая ранние return по ошибке) должен быть queue_resume_callback(q):
// ПРАВИЛЬНО:
if (!e) { queue_resume_callback(q); return; }
if (!e->dgram) { queue_dgram_free(e); queue_entry_free(e); queue_resume_callback(q); return; }
// НЕПРАВИЛЬНО — очередь зависнет:
if (!e) return;
2.5. Поиск по индексу
// Создать очередь с хешем (ключ 4 байта по смещению 0):
struct ll_queue* q = queue_new(ua, 256, 0, 4, "indexed_q");
// Добавлять элементы через _with_index (хеш вычисляется автоматически):
queue_data_put_with_index(q, entry);
// Искать:
uint32_t key = 42;
struct ll_entry* found = queue_find_data_by_index(q, &key);
// Если несколько элементов с одним ключом — продолжать поиск:
struct ll_entry* next = queue_find_next_by_index(q, &key, found);
2.6. Освобождение ресурсов
// 1. Снять callback
queue_set_callback(q, NULL, NULL);
// 2. Отменить всех waiter-ов
queue_waiter_cancel(q, &my_waiter);
// 3. Извлечь и освободить ВСЕ оставшиеся элементы
struct ll_entry* e;
while ((e = queue_data_get(q)) != NULL) {
queue_dgram_free(e);
queue_entry_free(e);
}
// После извлечения всех элементов resume не нужен — callback уже снят.
// 4. Освободить очередь
queue_free(q);
Важно: queue_free() не освобождает элементы. Их нужно извлечь через queue_data_get() и освободить вручную.
2.7. Приоритетная вставка (LIFO)
// Вставить в начало очереди (высокий приоритет):
queue_data_put_first(q, entry);
queue_data_put_first_with_index(q, entry);
2.8. Удаление из середины
queue_remove_data(q, entry); // удаляет из очереди, НЕ освобождает память
queue_dgram_free(entry); // освободить dgram отдельно
queue_entry_free(entry); // освободить entry отдельно
После queue_remove_data() элемент больше не принадлежит очереди — вызывающий сам отвечает за его освобождение.
2.9. Callback при опустошении
static void on_empty(struct ll_queue* q, void* arg) {
// Очередь стала пустой (count==0)
// Вызывается однократно, после вызова сбрасывается
}
queue_set_empty_callback(q, on_empty, my_data);
3. API
Создание / уничтожение
| Функция | Назначение |
|---|---|
queue_new(ua, hash_size, index_offset, index_size, name) |
Создать очередь. hash_size=0 — без хеша. ua обязателен (для uasync_call_soon). |
queue_free(q) |
Освободить очередь (НЕ элементы — их надо извлечь и освободить отдельно). |
Создание / освобождение элементов
| Функция | Назначение |
|---|---|
queue_entry_new(data_size) |
Выделить entry с data[] через malloc. |
queue_entry_new_from_pool(pool) |
Выделить entry из memory pool. |
ll_alloc_lldgram(len) |
Выделить entry + отдельный буфер dgram (malloc). |
queue_entry_free(entry) |
Освободить entry (возврат в pool или free). НЕ освобождает dgram! |
queue_dgram_free(entry) |
Освободить dgram (через dgram_free_fn, pool или free). |
Чтение (callback)
| Функция | Назначение |
|---|---|
queue_set_callback(q, cbk_fn, arg) |
Установить callback автозабора. Вызывается при наличии элементов. |
queue_data_get(q) |
Извлечь элемент из головы. Приостанавливает callback до queue_resume_callback(). |
queue_resume_callback(q) |
Обязательно вызвать после обработки элемента. Планирует отложенный вызов callback через uasync_call_soon. |
queue_set_callback_defer(q, enable) |
enable=1 — первый запуск callback из put откладывается через uasync_call_soon (очистка стека). По умолчанию 0 (синхронно). Повторные запуски всегда deferred. |
Запись
| Функция | Назначение |
|---|---|
queue_data_put(q, entry) |
Добавить в конец (FIFO). Для очередей БЕЗ хеша. |
queue_data_put_with_index(q, entry) |
Добавить в конец с хеш-индексом. |
queue_data_put_first(q, entry) |
Добавить в начало (LIFO, высокий приоритет). |
queue_data_put_first_with_index(q, entry) |
Добавить в начало с хеш-индексом. |
Backpressure
| Функция | Назначение |
|---|---|
queue_set_threshold(q, max_packets, max_bytes) |
Задать общий порог. Когда count <= max_packets И total_bytes <= max_bytes, waiter-ы пробуждаются (FIFO). |
queue_waiter_wait(q, h, callback, arg) |
Зарегистрировать ожидание. Возвращает 1 (уже выполнено), 0 (в очереди), -1 (ошибка). |
queue_waiter_cancel(q, h) |
Отменить ожидание. Безопасно вызывать в любом состоянии. |
queue_set_waiter_defer(q, enable) |
Отложенный вызов waiter-ов (через uasync_call_soon). |
queue_set_empty_callback(q, cbk_fn, arg) |
Одноразовый callback при count==0. |
queue_set_on_get(q, cbk_fn, arg) |
Callback после каждого queue_data_get (deferred). |
Поиск / удаление
| Функция | Назначение |
|---|---|
queue_find_data_by_index(q, index_key) |
Найти по ключу (требуется hash_size>0). Размер ключа = q->index_size. |
queue_find_next_by_index(q, index_key, prev_entry) |
Следующий элемент с тем же ключом. |
queue_remove_data(q, entry) |
Удалить из очереди (память НЕ освобождает). |
Утилиты
| Функция | Назначение |
|---|---|
queue_entry_count(q) |
Текущее количество элементов. |
queue_total_bytes(q) |
Суммарный объём данных (сумма len). |
queue_set_size_limit(q, lim) |
Ограничение на количество элементов (при превышении новые удаляются). |
queue_check_consistency(q) |
Проверка целостности (счётчики, циклы, prev/next). Только для отладки. |