#ifndef LL_QUEUE_H #define LL_QUEUE_H #include // для size_t // Предварительные объявления typedef struct ll_queue ll_queue_t; typedef struct uasync_s uasync_t; typedef struct ll_entry ll_entry_t; // Тип коллбэка: вызывается при добавлении элемента в пустую очередь или для продолжения обработки // Параметры: указатель на очередь, указатель на элемент (первый в очереди), пользовательский аргумент typedef void (*queue_callback_t)(ll_queue_t* q, ll_entry_t* entry, void* arg); // Структура элемента - переменный размер, данные расположены сразу после структуры struct ll_entry { struct ll_entry* next; // Указатель на следующий элемент в очереди size_t size; // Размер данных элемента (байт) int ref_count; // Счетчик ссылок для предотвращения double-free }; // Структура пула памяти для оптимизации аллокаций #define LL_QUEUE_POOL_SIZE 64 // Размер пула для объектов одинакового размера typedef struct memory_pool { void* free_list[LL_QUEUE_POOL_SIZE]; // Список свободных блоков int free_count; // Количество свободных блоков size_t object_size; // Размер объектов в пуле size_t allocations; // Статистика: всего аллокаций size_t reuse_count; // Статистика: повторное использование } memory_pool_t; // Структура условия ожидания (waiter) struct queue_waiter { int max_packets; // Максимальное количество пакетов size_t max_bytes; // Максимальное количество байт void (*callback)(ll_queue_t* q, void* arg); // Коллбэк для вызова void* callback_arg; // Аргумент коллбэка struct queue_waiter* next; // Следующий ожидающий в списке }; typedef struct queue_waiter queue_waiter_t; typedef void (*queue_threshold_callback_t)(ll_queue_t* q, void* arg); // Структура очереди struct ll_queue { ll_entry_t* head; // Первый элемент (извлекается отсюда) ll_entry_t* tail; // Последний элемент (добавляется сюда) int count; // Текущее количество элементов size_t total_bytes; // Общий размер данных всех элементов (байт) int size_limit; // Максимальное количество (-1 = без ограничения) queue_callback_t callback; // Функция коллбэка void* callback_arg; // Пользовательский аргумент для коллбэка int callback_suspended; // 1 если коллбэки приостановлены (во время обработки) void* resume_timeout_id; // ID таймаута uasync для отложенного возобновления uasync_t* ua; // Экземпляр uasync для таймеров queue_waiter_t* waiters; // Список ожидающих коллбэков // Пулы памяти для оптимизации аллокаций memory_pool_t waiter_pool; // Пул для структур queue_waiter_t memory_pool_t entry_pool; // Пул для структур ll_entry_t (фиксированный размер) int use_pools; // Флаг использования пулов (включается при частых аллокациях) }; // ==================== Управление пулами памяти ==================== // Инициализировать пул памяти void memory_pool_init(memory_pool_t* pool, size_t object_size); // Выделить объект из пула или из malloc void* memory_pool_alloc(memory_pool_t* pool); // Освободить объект в пул или в free void memory_pool_free(memory_pool_t* pool, void* obj); // Получить статистику пула void memory_pool_get_stats(memory_pool_t* pool, size_t* allocations, size_t* reuse_count); // Очистить пул памяти void memory_pool_destroy(memory_pool_t* pool); // ==================== Управление очередью ==================== // Создать новую пустую очередь // ua - экземпляр uasync для таймеров (обязательный параметр) // use_pools - использовать пулы памяти для оптимизации (1) или нет (0) // Возвращает: указатель на очередь или NULL при ошибке выделения памяти ll_queue_t* queue_new_with_pools(uasync_t* ua, int use_pools); // Создать новую пустую очередь (обратная совместимость) ll_queue_t* queue_new(uasync_t* ua); // Освободить очередь и все её элементы // Также отменяет отложенное возобновление если оно запланировано void queue_free(ll_queue_t* q); // ==================== Конфигурация очереди ==================== // Установить функцию и аргумент коллбэка для очереди // Коллбэк вызывается при добавлении элемента в пустую очередь (разрешенные коллбэки) // обработчик должен обработать этот пакет и когда будет готов к приёму следующего - вызывает resume_callback. обработка строго по одному пакету. void queue_set_callback(ll_queue_t* q, queue_callback_t cbk_fn, void* arg); // Возобновить коллбэки после обработки элемента переданного в коллбэке (тянуть дополнительные элементы из очереди не предусмотернные api нельзя). // эта функция должна вызываться всегда после того как cbk_fn обработала пакет (можно с ожиданием через async), иначе очередь застрянет. // Если в очереди остались элементы, запланирует вызов коллбэка через uasync_set_timeout(0) // Это предотвращает накопление рекурсии в стеке вызовов void queue_resume_callback(ll_queue_t* q); // Установить максимальное количество элементов в очереди // При превышении лимита новый элемент автоматически освобождается void queue_set_size_limit(ll_queue_t* q, int lim); // ==================== Управление элементами ==================== // Создать новый элемент с областью данных указанного размера // Память выделяется одним блоком: [ll_entry_t][область данных data_size байт] // Возвращает: указатель на элемент или NULL при ошибке выделения памяти ll_entry_t* queue_entry_new(size_t data_size); // Освободить элемент (не влияет на связи в очереди) void queue_entry_free(ll_entry_t* entry); // ==================== Операции с очередью ==================== // Добавить элемент в конец очереди (FIFO) // Если очередь была пустой и коллбэки разрешены - вызывает коллбэк // Возвращает: 0 при успехе, -1 если превышен лимит размера (элемент освобожден) int queue_entry_put(ll_queue_t* q, ll_entry_t* entry); // Добавить элемент в начало очереди (LIFO, высокий приоритет) // Если очередь была пустой и коллбэки разрешены - вызывает коллбэк // Возвращает: 0 при успехе, -1 если превышен лимит размера (элемент освобожден) int queue_entry_put_first(ll_queue_t* q, ll_entry_t* entry); // Извлечь элемент из начала очереди // При извлечении приостанавливает коллбэки (callback_suspended = 1) чтобы предотвратить рекурсию // Возвращает: указатель на элемент или NULL если очередь пуста ll_entry_t* queue_entry_get(ll_queue_t* q); // Получить текущее количество элементов в очереди int queue_entry_count(ll_queue_t* q); // ==================== Вспомогательные функции ==================== // Получить указатель на область данных элемента // Данные расположены сразу после структуры ll_entry_t static inline void* ll_entry_data(ll_entry_t* entry) { return (void*)(entry + 1); } // Получить размер данных элемента static inline size_t ll_entry_size(ll_entry_t* entry) { return entry->size; } // ==================== Асинхронное ожидание ==================== // Зарегистрировать одноразовый коллбэк, который будет вызван когда очередь будет иметь // не более max_packets пакетов и не более max_bytes байт. // Если условие уже выполнено, коллбэк вызывается немедленно. // Можно зарегистрировать несколько ожиданий на одной очереди. // Возвращает указатель на waiter для возможной отмены через queue_cancel_wait queue_waiter_t* queue_wait_threshold(ll_queue_t* q, int max_packets, size_t max_bytes, queue_threshold_callback_t callback, void* arg); // Отменить ожидание (удалить waiter из списка) void queue_cancel_wait(ll_queue_t* q, queue_waiter_t* waiter); // Получить общий размер данных в очереди (байт) static inline size_t queue_total_bytes(ll_queue_t* q) { if (!q) return 0; return q->total_bytes; } // Получить статистику использования пулов памяти void queue_get_pool_stats(ll_queue_t* q, size_t* waiter_allocations, size_t* waiter_reuse); #endif // LL_QUEUE_H