#include "ll_queue.h" #include #include #include "u_async.h" // Предварительное объявление для отложенного возобновления static void queue_resume_timeout_cb(void* arg); // ==================== Управление очередью ==================== 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->size_limit = -1; // По умолчанию без ограничения q->callback = NULL; q->callback_arg = NULL; q->callback_suspended = 0; // Коллбэки разрешены изначально q->resume_timeout_id = 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; } // Отменить отложенное возобновление если запланировано 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++; // Если очередь была пустой до добавления (count был 0) и коллбэки разрешены - вызвать коллбэк // Это запускает автоматическую обработку очереди if (q->count == 1 && !q->callback_suspended && q->callback) { q->callback(q, entry, q->callback_arg); } 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++; // Если очередь была пустой до добавления (count был 0) и коллбэки разрешены - вызвать коллбэк if (q->count == 1 && !q->callback_suspended && q->callback) { q->callback(q, entry, q->callback_arg); } 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--; entry->next = NULL; // Отсоединить от очереди // При извлечении элемента приостанавливаем коллбэки // Это предотвращает рекурсию если во время обработки добавляются новые элементы q->callback_suspended = 1; return entry; } int queue_entry_count(ll_queue_t* q) { if (!q) return 0; return q->count; }