You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

11 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_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). Только для отладки.