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.
382 lines
19 KiB
382 lines
19 KiB
#ifndef LL_QUEUE_H |
|
#define LL_QUEUE_H |
|
|
|
#include <stddef.h> // для size_t |
|
#include <stdint.h> // для uint64_t |
|
#include "memory_pool.h" // для struct memory_pool |
|
|
|
#ifdef _WIN32 |
|
#include <winsock2.h> |
|
#include <windows.h> |
|
#else |
|
#include <pthread.h> |
|
#endif |
|
|
|
#define QUEUE_DEBUG 1 |
|
//#define QUEUE_THREAD_CHECK 1 // 0 to disable |
|
|
|
/** |
|
* @file ll_queue.h |
|
* @brief Упрощённая двусвязная очередь с поддержкой автозабора элементов через callback, |
|
* порогового ожидания освобождения места, поиска по ID и гибкого управления памятью. |
|
* |
|
* Основные особенности: |
|
* - FIFO-очередь (с возможностью добавления в начало для приоритета) |
|
* - Автоматический вызов callback при появлении элементов (один элемент за раз) |
|
* - Обязательный вызов queue_resume_callback после обработки элемента |
|
* - Поддержка пулов памяти для entry и отдельно для dgram |
|
* - Хеш-таблица для быстрого поиска по ID (опционально) |
|
* - Связный список waiters для backpressure c общим порогом (round-robin) |
|
* - Использует uasync для отложенного возобновления callback'ов (без рекурсии в стеке) |
|
* |
|
* @note Важные правила использования автозабора (queue_set_callback): |
|
* 1. Коллбэк получает управление только если в очереди есть элементы И !callback_suspended |
|
* 2. Внутри коллбэка ОБЯЗАТЕЛЬНО: queue_data_get() → обработать элемент → queue_resume_callback() |
|
* 3. Без вызова queue_resume_callback очередь навсегда застрянет |
|
* 4. Не извлекайте элементы вручную вне коллбэка — это нарушит логику |
|
* 5. queue_data_get() автоматически suspend'ит коллбэки для предотвращения рекурсии |
|
* |
|
* @note Память: |
|
* - queue_free() освобождает структуру очереди, хеш-таблицу и список waiters |
|
* (handle->internal каждого ожидающего сбрасывается в NULL) |
|
* - Все элементы должны быть извлечены через queue_data_get() и освобождены через queue_entry_free() |
|
* - dgram освобождается отдельно через queue_dgram_free() или кастомную dgram_free_fn |
|
*/ |
|
|
|
// Forward declarations |
|
struct ll_queue; |
|
struct queue_waiter_handle; |
|
|
|
/** |
|
* @struct ll_entry |
|
* @brief Элемент очереди (переменного размера). |
|
* |
|
* Память выделяется одним блоком: [struct ll_entry + data[size]]. |
|
* Поле dgram — отдельный блок (malloc / пул / внешняя память). |
|
*/ |
|
struct ll_entry { |
|
char* name; |
|
struct ll_entry* next; // Следующий элемент в очереди |
|
struct ll_entry* prev; // Предыдущий элемент в очереди |
|
uint16_t size; // Размер пользовательского буфера data[] (байт) |
|
uint16_t len; // Актуальная длина данных в dgram |
|
uint16_t memlen; // Выделенный размер под dgram |
|
uint16_t int_len; // Внутреннее (не использовать) |
|
uint8_t* dgram; // Указатель на данные пакета |
|
void (*dgram_free_fn)(uint8_t* data); // Кастомная функция освобождения dgram |
|
struct memory_pool* dgram_pool; // Пул для dgram (если используется) |
|
uint32_t index_hash; // Хеш по первым 4 байтам индекса |
|
struct ll_entry* hash_next; // Следующий в хеш-цепочке |
|
struct memory_pool* pool; // Пул, из которого выделен сам entry (NULL = malloc) |
|
uint8_t data[0]; // Гибкий массив пользовательских данных (размер = size) |
|
}; |
|
|
|
/** |
|
* @typedef queue_callback_fn |
|
* @brief Коллбэк автозабора элементов из очереди. |
|
* @param q указатель на очередь |
|
* @param arg пользовательский аргумент (из queue_set_callback) |
|
* |
|
* @note Внутри функции: |
|
* - Вызвать queue_data_get(q) для получения элемента |
|
* - Обработать элемент (можно асинхронно) |
|
* - После завершения обработки вызвать queue_resume_callback(q) |
|
*/ |
|
typedef void (*queue_callback_fn)(struct ll_queue* q, void* arg); |
|
|
|
/** |
|
* @typedef queue_threshold_callback_fn |
|
* @brief Одноразовый коллбэк при достижении порога очереди. |
|
* @param q указатель на очередь |
|
* @param arg пользовательский аргумент |
|
*/ |
|
typedef void (*queue_threshold_callback_fn)(struct ll_queue* q, void* arg); |
|
|
|
/** |
|
* @struct queue_waiter |
|
* @brief Внутренний узел в связном списке ожидающих (очередь владеет, аллоцирует/освобождает сама). |
|
*/ |
|
struct queue_waiter { |
|
struct queue_waiter* next; ///< Следующий в списке ожидающих |
|
struct queue_waiter_handle* handle; ///< Обратная ссылка на публичную структуру |
|
queue_threshold_callback_fn callback; |
|
void* callback_arg; |
|
}; |
|
|
|
/** |
|
* @struct queue_waiter_handle |
|
* @brief Публичная управляющая структура (встраивается в структуру вызывающей стороны). |
|
* |
|
* Вызывающая сторона: |
|
* 1. Объявляет поле struct queue_waiter_handle в своей структуре (= {0} или u_calloc) |
|
* 3. Вызывает queue_waiter_wait() когда хочет дождаться освобождения очереди |
|
* 4. Вызывает queue_waiter_cancel() при деинициализации для отмены ожидания |
|
*/ |
|
struct queue_waiter_handle { |
|
struct queue_waiter* internal; ///< NULL = не ждём; иначе — внутренний узел |
|
}; |
|
|
|
/** |
|
* @struct ll_queue |
|
* @brief Структура очереди. |
|
*/ |
|
struct ll_queue { |
|
char* name; |
|
struct ll_entry* head; // Голова очереди (отсюда извлекаем) |
|
struct ll_entry* tail; // Хвост очереди (сюда добавляем) |
|
int count; // Текущее количество элементов |
|
size_t total_bytes; // Суммарный объём данных (сумма int_len) |
|
int size_limit; // Максимальное количество элементов (-1 = без лимита) |
|
|
|
queue_callback_fn callback; // Коллбэк автозабора |
|
void* callback_arg; // Аргумент коллбэка |
|
int callback_suspended; // 1 = коллбэки временно приостановлены |
|
|
|
void* resume_timeout_id; // ID таймера uasync для отложенного resume |
|
struct UASYNC* ua; // Экземпляр uasync (обязателен для таймеров) |
|
|
|
struct queue_waiter* waiter_head; // Голова списка ожидающих |
|
struct queue_waiter* waiter_tail; // Хвост списка ожидающих |
|
int threshold_max_packets; // Общий порог: макс. кол-во элементов (по умолчанию 0) |
|
size_t threshold_max_bytes; // Общий порог: макс. объём данных (0 = не проверять) |
|
|
|
struct ll_entry** hash_table; // Хеш-таблица для поиска по id (если hash_size > 0) |
|
size_t hash_size; // Размер хеш-таблицы |
|
uint16_t index_offset; // Смещение индекса в data[] (для всех entry очереди) |
|
uint16_t index_size; // Размер индекса (0 = без индекса) |
|
|
|
#ifdef QUEUE_THREAD_CHECK |
|
#ifdef _WIN32 |
|
DWORD owner_thread; |
|
#else |
|
pthread_t owner_thread; |
|
#endif |
|
#endif |
|
}; |
|
|
|
/* ==================== Создание / уничтожение ==================== */ |
|
|
|
/** |
|
* @brief Создаёт новую пустую очередь. |
|
* @param ua экземпляр uasync (обязателен для отложенного resume коллбэков) |
|
* @param hash_size размер хеш-таблицы (0 = поиск по id отключён) |
|
* @param index_offset смещение ключа индекса в data[] каждого entry (0 если hash_size=0) |
|
* @param index_size размер ключа индекса (0 если hash_size=0) |
|
* @param name имя очереди для отладки |
|
* @return указатель на очередь или NULL при ошибке выделения памяти |
|
*/ |
|
struct ll_queue* queue_new(struct UASYNC* ua, size_t hash_size, uint16_t index_offset, uint16_t index_size, char* name); |
|
|
|
/** |
|
* @brief Освобождает структуру очереди. |
|
* @param q очередь |
|
* |
|
* @warning НЕ ОСВОБОЖДАЕТ элементы в очереди! |
|
* Их нужно извлечь через queue_data_get() и освободить через queue_entry_free(). |
|
*/ |
|
void queue_free(struct ll_queue* q); |
|
|
|
/* ==================== Конфигурация ==================== */ |
|
|
|
/** |
|
* @brief Устанавливает максимальное количество элементов в очереди. |
|
* @param q очередь |
|
* @param lim лимит (-1 = без ограничения) |
|
* |
|
* При превышении лимита новые элементы автоматически освобождаются. |
|
*/ |
|
void queue_set_size_limit(struct ll_queue* q, int lim); |
|
|
|
/* ==================== Автозабор элементов ==================== */ |
|
|
|
/** |
|
* @brief Устанавливает коллбэк для автоматического извлечения элементов. |
|
* @param q очередь |
|
* @param cbk_fn функция-коллбэк |
|
* @param arg пользовательский аргумент |
|
* |
|
* Коллбэк вызывается автоматически при наличии элементов и разрешённом состоянии. |
|
*/ |
|
void queue_set_callback(struct ll_queue* q, queue_callback_fn cbk_fn, void* arg); |
|
|
|
/** |
|
* @brief Возобновляет работу коллбэка после обработки элемента в следующем event loop. |
|
* @param q очередь |
|
* |
|
* @warning ОБЯЗАТЕЛЬНО вызывать после завершения обработки элемента в коллбэке. |
|
* Иначе автозабор элементов навсегда остановится. |
|
*/ |
|
void queue_resume_callback(struct ll_queue* q); |
|
|
|
/* ==================== Пороговое ожидание (backpressure) ==================== */ |
|
|
|
/** |
|
* @brief Устанавливает общий порог освобождения очереди. |
|
* @param q очередь |
|
* @param max_packets максимальное количество элементов (условие: count <= max_packets) |
|
* @param max_bytes максимальный объём данных (0 = не проверять) |
|
* |
|
* Порог — общий на всю очередь. Все waiters в связном списке используют этот порог. |
|
* При достижении порога вызывается следующий waiter из списка (FIFO / round-robin). |
|
* По умолчанию: max_packets=0, max_bytes=0 (пустая очередь). |
|
*/ |
|
void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes); |
|
|
|
/** |
|
* @brief Регистрирует ожидание освобождения очереди до общего порога. |
|
* @param q очередь |
|
* @param h указатель на handle (создаётся как `struct queue_waiter_handle h = {0}` или `h.internal = NULL`) |
|
* @param callback функция, вызываемая при достижении порога |
|
* @param arg аргумент коллбэка |
|
* @return 1 — условие уже выполнено (callback вызван немедленно), 0 — зарегистрирован в списке, -1 — ошибка |
|
* |
|
* Если handle уже ожидает (h->internal != NULL) — старый waiter отменяется и создаётся новый. |
|
*/ |
|
int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h, |
|
queue_threshold_callback_fn callback, void* arg); |
|
|
|
/** |
|
* @brief Отменяет ожидание (удаляет из списка, если зарегистрирован). |
|
* @param q очередь |
|
* @param h указатель на handle |
|
* |
|
* Безопасно вызывать в любом состоянии: если waiter не зарегистрирован — ничего не делает. |
|
*/ |
|
void queue_waiter_cancel(struct ll_queue* q, struct queue_waiter_handle* h); |
|
|
|
/* ==================== Работа с данными ==================== */ |
|
|
|
/** |
|
* @brief Добавляет элемент в конец очереди (FIFO). |
|
* @param q очередь |
|
* @param entry элемент |
|
* @param id идентификатор для поиска |
|
* @return 0 — успех, -1 — превышен лимит (элемент освобождён) |
|
*/ |
|
int queue_data_put(struct ll_queue* q, struct ll_entry* entry); |
|
|
|
/** |
|
* @brief Добавляет элемент в начало очереди (LIFO, высокий приоритет). |
|
* @param q очередь |
|
* @param entry элемент |
|
* @param id идентификатор |
|
* @return 0 — успех, -1 — превышен лимит (элемент освобождён) |
|
*/ |
|
int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry); |
|
|
|
/** |
|
* @brief Добавляет элемент в конец очереди с хеш-индексом (использует index_offset/index_size из очереди). |
|
* @param q очередь |
|
* @param entry элемент |
|
* @return 0 — успех, -1 — ошибка |
|
*/ |
|
int queue_data_put_with_index(struct ll_queue* q, struct ll_entry* entry); |
|
|
|
/** |
|
* @brief Добавляет элемент в начало очереди с хеш-индексом (использует index_offset/index_size из очереди). |
|
* @param q очередь |
|
* @param entry элемент |
|
* @return 0 — успех, -1 — ошибка |
|
*/ |
|
int queue_data_put_first_with_index(struct ll_queue* q, struct ll_entry* entry); |
|
|
|
|
|
/** |
|
* @brief Извлекает элемент из начала очереди. |
|
* @param q очередь |
|
* @return элемент или NULL, если очередь пуста |
|
* |
|
* @note Автоматически приостанавливает коллбэки (callback_suspended = 1) |
|
* После обработки элемента необходимо вызвать queue_resume_callback(). |
|
*/ |
|
struct ll_entry* queue_data_get(struct ll_queue* q); |
|
|
|
/** |
|
* @brief Возвращает текущее количество элементов в очереди. |
|
* @param q очередь |
|
* @return количество элементов |
|
*/ |
|
int queue_entry_count(struct ll_queue* q); |
|
|
|
/* ==================== Управление памятью элементов ==================== */ |
|
|
|
// проверить правильность очереди (корректность цепочки, счетчики элементов и байтов) |
|
int queue_check_consistency(struct ll_queue* q); |
|
|
|
/** |
|
* @brief Выделяет элемент с отдельным буфером dgram (malloc). |
|
* @param len требуемый размер dgram |
|
* @return элемент или NULL при ошибке |
|
*/ |
|
struct ll_entry* ll_alloc_lldgram(uint16_t len); |
|
|
|
/** |
|
* @brief Создаёт элемент с пользовательским буфером data[] (malloc одним блоком). |
|
* @param data_size размер data[] |
|
* @return элемент или NULL при ошибке |
|
*/ |
|
struct ll_entry* queue_entry_new(size_t data_size); |
|
|
|
/** |
|
* @brief Создаёт элемент из пула памяти. |
|
* @param pool пул |
|
* @return элемент или NULL, если пул пуст |
|
*/ |
|
struct ll_entry* queue_entry_new_from_pool(struct memory_pool* pool); |
|
|
|
/** |
|
* @brief Освобождает структуру элемента (возвращает в пул или free). |
|
* @param entry элемент |
|
* |
|
* @note НЕ освобождает dgram! Используйте queue_dgram_free() отдельно. |
|
*/ |
|
void queue_entry_free(struct ll_entry* entry); |
|
|
|
/** |
|
* @brief Освобождает буфер dgram элемента. |
|
* @param entry элемент |
|
*/ |
|
void queue_dgram_free(struct ll_entry* entry); |
|
|
|
/* ==================== Поиск и удаление ==================== */ |
|
|
|
/** |
|
* @brief Находит элемент по ключу индекса (требуется hash_size > 0 при создании). |
|
* @param q очередь (размер ключа берётся из q->index_size) |
|
* @param index_key искомый ключ |
|
* @return элемент или NULL |
|
*/ |
|
struct ll_entry* queue_find_data_by_index(struct ll_queue* q, |
|
const void* index_key); |
|
|
|
/** |
|
* @brief Находит следующий элемент с тем же ключом (продолжение поиска). |
|
* @param q очередь (размер ключа берётся из q->index_size) |
|
* @param index_key искомый ключ |
|
* @param prev_entry предыдущий найденный элемент (от queue_find_data_by_index или предыдущего queue_find_next_by_index) |
|
* @return следующий элемент с тем же ключом или NULL |
|
*/ |
|
struct ll_entry* queue_find_next_by_index(struct ll_queue* q, |
|
const void* index_key, |
|
struct ll_entry* prev_entry); |
|
|
|
/** |
|
* @brief Удаляет элемент из очереди (не освобождает память). |
|
* @param q очередь |
|
* @param entry элемент |
|
* @return 0 — успех, -1 — элемент не найден в очереди |
|
*/ |
|
int queue_remove_data(struct ll_queue* q, struct ll_entry* entry); |
|
|
|
/* ==================== Утилиты ==================== */ |
|
|
|
/** |
|
* @brief Возвращает суммарный объём данных в очереди. |
|
* @param q очередь |
|
* @return количество байт |
|
*/ |
|
static inline size_t queue_total_bytes(struct ll_queue* q) { |
|
return q->total_bytes; |
|
} |
|
|
|
#endif // LL_QUEUE_H
|
|
|