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.
284 lines
8.8 KiB
284 lines
8.8 KiB
#include "ll_queue.h" |
|
#include <stdlib.h> |
|
#include <string.h> |
|
#include <stdio.h> |
|
#include <assert.h> |
|
#include "u_async.h" |
|
|
|
// Предварительное объявление для отложенного возобновления |
|
static void queue_resume_timeout_cb(void* arg); |
|
|
|
// Проверить и запустить ожидающие коллбэки |
|
static void check_waiters(ll_queue_t* q) { |
|
if (!q || !q->waiters) return; |
|
|
|
queue_waiter_t** pprev = &q->waiters; |
|
queue_waiter_t* waiter = q->waiters; |
|
|
|
while (waiter) { |
|
queue_waiter_t* next = waiter->next; |
|
|
|
// Проверить условие: не больше max_packets и не больше max_bytes |
|
if (q->count <= waiter->max_packets && q->total_bytes <= waiter->max_bytes) { |
|
// Условие выполнено - вызвать коллбэк |
|
waiter->callback(q, waiter->callback_arg); |
|
// Удалить waiter из списка |
|
*pprev = next; |
|
free(waiter); |
|
// pprev уже указывает на правильный следующий элемент |
|
} else { |
|
// Условие не выполнено - оставить в списке |
|
pprev = &waiter->next; |
|
} |
|
waiter = next; |
|
} |
|
} |
|
|
|
// ==================== Управление очередью ==================== |
|
|
|
ll_queue_t* queue_new(void) { |
|
ll_queue_t* q = calloc(1, sizeof(ll_queue_t)); |
|
if (!q) return NULL; |
|
|
|
q->head = NULL; |
|
q->tail = NULL; |
|
q->count = 0; |
|
q->total_bytes = 0; |
|
q->size_limit = -1; // По умолчанию без ограничения |
|
q->callback = NULL; |
|
q->callback_arg = NULL; |
|
q->callback_suspended = 0; // Коллбэки разрешены изначально |
|
q->resume_timeout_id = NULL; |
|
q->waiters = NULL; |
|
|
|
return q; |
|
} |
|
|
|
void queue_free(ll_queue_t* q) { |
|
if (!q) return; |
|
|
|
// Освободить все элементы |
|
ll_entry_t* entry = q->head; |
|
while (entry) { |
|
ll_entry_t* next = entry->next; |
|
free(entry); |
|
entry = next; |
|
} |
|
|
|
// Освободить все ожидающие коллбэки |
|
queue_waiter_t* waiter = q->waiters; |
|
while (waiter) { |
|
queue_waiter_t* next = waiter->next; |
|
free(waiter); |
|
waiter = next; |
|
} |
|
|
|
// Отменить отложенное возобновление если запланировано |
|
if (q->resume_timeout_id) { |
|
uasync_cancel_timeout(q->resume_timeout_id); |
|
} |
|
|
|
free(q); |
|
} |
|
|
|
// ==================== Конфигурация очереди ==================== |
|
|
|
void queue_set_callback(ll_queue_t* q, queue_callback_t cbk_fn, void* arg) { |
|
if (!q) return; |
|
q->callback = cbk_fn; |
|
q->callback_arg = arg; |
|
} |
|
|
|
static void queue_resume_timeout_cb(void* arg) { |
|
ll_queue_t* q = (ll_queue_t*)arg; |
|
if (!q || !q->callback) return; |
|
|
|
// Очистить ID таймаута (таймаут сработал) |
|
q->resume_timeout_id = NULL; |
|
|
|
// Разрешить коллбэки |
|
q->callback_suspended = 0; |
|
|
|
// Если в очереди есть элементы, вызвать коллбэк с первым элементом |
|
// Обработчик должен извлечь этот элемент вызовом queue_entry_get() |
|
if (q->head) { |
|
q->callback(q, q->head, q->callback_arg); |
|
} |
|
} |
|
|
|
void queue_resume_callback(ll_queue_t* q) { |
|
if (!q || !q->callback) return; |
|
|
|
// Если уже есть отложенное возобновление, ничего не делать |
|
if (q->resume_timeout_id) { |
|
return; |
|
} |
|
|
|
// Запланировать отложенное возобновление через uasync |
|
q->resume_timeout_id = uasync_set_timeout(0, q, queue_resume_timeout_cb); |
|
} |
|
|
|
void queue_set_size_limit(ll_queue_t* q, int lim) { |
|
if (!q) return; |
|
q->size_limit = lim; |
|
} |
|
|
|
// ==================== Управление элементами ==================== |
|
|
|
ll_entry_t* queue_entry_new(size_t data_size) { |
|
// Выделить память под структуру + область данных |
|
ll_entry_t* entry = malloc(sizeof(ll_entry_t) + data_size); |
|
if (!entry) return NULL; |
|
|
|
entry->next = NULL; |
|
entry->size = data_size; |
|
// Область данных оставить неинициализированной для производительности |
|
|
|
return entry; |
|
} |
|
|
|
void queue_entry_free(ll_entry_t* entry) { |
|
free(entry); |
|
} |
|
|
|
// ==================== Операции с очередью ==================== |
|
|
|
int queue_entry_put(ll_queue_t* q, ll_entry_t* entry) { |
|
if (!q || !entry) return -1; |
|
|
|
// Проверить лимит размера |
|
if (q->size_limit >= 0 && q->count >= q->size_limit) { |
|
queue_entry_free(entry); |
|
return -1; |
|
} |
|
|
|
// Добавить в хвост (FIFO) |
|
entry->next = NULL; |
|
if (q->tail) { |
|
q->tail->next = entry; |
|
} else { |
|
q->head = entry; |
|
} |
|
q->tail = entry; |
|
q->count++; |
|
q->total_bytes += entry->size; |
|
|
|
// Если коллбэки разрешены - вызвать коллбэк |
|
// Это запускает автоматическую обработку очереди |
|
if (!q->callback_suspended && q->callback) { |
|
printf("[LL_QUEUE DEBUG] queue_entry_put: calling callback, count=%d, suspended=%d\n", q->count, q->callback_suspended); |
|
q->callback(q, entry, q->callback_arg); |
|
} |
|
|
|
// Проверить ожидающие коллбэки |
|
check_waiters(q); |
|
|
|
return 0; |
|
} |
|
|
|
int queue_entry_put_first(ll_queue_t* q, ll_entry_t* entry) { |
|
if (!q || !entry) return -1; |
|
|
|
// Проверить лимит размера |
|
if (q->size_limit >= 0 && q->count >= q->size_limit) { |
|
queue_entry_free(entry); |
|
return -1; |
|
} |
|
|
|
// Добавить в голову (LIFO, высокий приоритет) |
|
entry->next = q->head; |
|
q->head = entry; |
|
if (!q->tail) { |
|
q->tail = entry; |
|
} |
|
q->count++; |
|
q->total_bytes += entry->size; |
|
|
|
// Если коллбэки разрешены - вызвать коллбэк |
|
if (!q->callback_suspended && q->callback) { |
|
printf("[LL_QUEUE DEBUG] queue_entry_put_first: calling callback, count=%d, suspended=%d\n", q->count, q->callback_suspended); |
|
q->callback(q, entry, q->callback_arg); |
|
} |
|
|
|
// Проверить ожидающие коллбэки |
|
check_waiters(q); |
|
|
|
return 0; |
|
} |
|
|
|
ll_entry_t* queue_entry_get(ll_queue_t* q) { |
|
if (!q || !q->head) return NULL; |
|
|
|
ll_entry_t* entry = q->head; |
|
q->head = entry->next; |
|
if (!q->head) { |
|
q->tail = NULL; |
|
} |
|
q->count--; |
|
q->total_bytes -= entry->size; |
|
|
|
entry->next = NULL; // Отсоединить от очереди |
|
|
|
// При извлечении элемента приостанавливаем коллбэки |
|
// Это предотвращает рекурсию если во время обработки добавляются новые элементы |
|
q->callback_suspended = 1; |
|
|
|
// Проверить ожидающие коллбэки |
|
check_waiters(q); |
|
|
|
return entry; |
|
} |
|
|
|
int queue_entry_count(ll_queue_t* q) { |
|
if (!q) return 0; |
|
return q->count; |
|
} |
|
|
|
// ==================== Асинхронное ожидание ==================== |
|
|
|
queue_waiter_t* queue_wait_threshold(ll_queue_t* q, int max_packets, size_t max_bytes, |
|
queue_threshold_callback_t callback, void* arg) { |
|
if (!q || !callback) return NULL; |
|
|
|
// Создать новый waiter |
|
queue_waiter_t* waiter = malloc(sizeof(queue_waiter_t)); |
|
if (!waiter) return NULL; |
|
|
|
waiter->max_packets = max_packets; |
|
waiter->max_bytes = max_bytes; |
|
waiter->callback = callback; |
|
waiter->callback_arg = arg; |
|
waiter->next = NULL; |
|
|
|
// Проверить условие немедленно |
|
if (q->count <= max_packets && q->total_bytes <= max_bytes) { |
|
// Условие уже выполнено - вызвать коллбэк и освободить waiter |
|
callback(q, arg); |
|
free(waiter); |
|
return NULL; |
|
} |
|
|
|
// Добавить в список ожидающих |
|
waiter->next = q->waiters; |
|
q->waiters = waiter; |
|
|
|
return waiter; |
|
} |
|
|
|
void queue_cancel_wait(ll_queue_t* q, queue_waiter_t* waiter) { |
|
if (!q || !waiter) return; |
|
|
|
// Найти и удалить waiter из списка |
|
queue_waiter_t** pprev = &q->waiters; |
|
queue_waiter_t* w = q->waiters; |
|
|
|
while (w) { |
|
if (w == waiter) { |
|
*pprev = w->next; |
|
free(w); |
|
return; |
|
} |
|
pprev = &w->next; |
|
w = w->next; |
|
} |
|
}
|
|
|