Browse Source

ll_queue: отложенный первый запуск автозабора (callback_defer) + тест

queue_set_callback_defer(q, 1) переносит первый запуск callback из
queue_data_put* в uasync_call_soon (очистка стека производителя).
Повторные запуски (queue_resume_callback) и так всегда deferred.
proxy
evgeny 2 weeks ago
parent
commit
8aaac45020
  1. 34
      lib/ll_queue.c
  2. 12
      lib/ll_queue.h
  3. 1
      lib/ll_queue_doc.md
  4. 30
      tests/test_ll_queue.c

34
lib/ll_queue.c

@ -194,6 +194,19 @@ void queue_resume_callback(struct ll_queue* q) {
}
}
// Общий запуск коллбэка автозабора: immediate (default) или deferred через uasync_call_soon.
// Deferred-режим очищает стек производителя перед обработкой первого элемента.
static void queue_trigger_callback(struct ll_queue* q) {
if (!q->head || q->callback_suspended || !q->callback) return;
if (q->callback_defer) {
if (!q->resume_timeout_id)
q->resume_timeout_id = uasync_call_soon(q->ua, q, queue_resume_timeout_cb);
} else {
q->callback(q, q->callback_arg);
}
}
void queue_set_size_limit(struct ll_queue* q, int lim) {
if (!q) return;
q->size_limit = lim;
@ -384,9 +397,7 @@ int queue_data_put(struct ll_queue* q, struct ll_entry* entry) {
// НЕ вызываем add_to_hash
if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg);
}
if (q->count == 1) queue_trigger_callback(q);
#ifdef QUEUE_DEBUG
queue_check_consistency(q);
@ -433,9 +444,7 @@ int queue_data_put_with_index(struct ll_queue* q, struct ll_entry* entry) {
add_to_hash(q, entry);
if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg);
}
if (q->count == 1) queue_trigger_callback(q);
#ifdef QUEUE_DEBUG
queue_check_consistency(q);
@ -471,9 +480,7 @@ int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry) {
entry->int_len = entry->len;
q->total_bytes += entry->int_len;
if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg);
}
if (q->count == 1) queue_trigger_callback(q);
#ifdef QUEUE_DEBUG
queue_check_consistency(q);
@ -520,9 +527,7 @@ int queue_data_put_first_with_index(struct ll_queue* q, struct ll_entry* entry)
add_to_hash(q, entry);
if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg);
}
if (q->count == 1) queue_trigger_callback(q);
#ifdef QUEUE_DEBUG
queue_check_consistency(q);
@ -647,6 +652,11 @@ void queue_set_waiter_defer(struct ll_queue* q, int enable) {
q->waiter_defer = enable;
}
void queue_set_callback_defer(struct ll_queue* q, int enable) {
if (!q) return;
q->callback_defer = enable;
}
void queue_set_empty_callback(struct ll_queue* q, queue_callback_fn cbk_fn, void* arg) {
if (!q) return;
void* old = q->empty_call_soon_id;

12
lib/ll_queue.h

@ -149,6 +149,7 @@ struct ll_queue {
int threshold_max_packets; // Общий порог: макс. кол-во элементов (по умолчанию 0)
size_t threshold_max_bytes; // Общий порог: макс. объём данных (0 = не проверять)
int waiter_defer; // 0=immediate callback, 1=defer via uasync_call_soon
int callback_defer; // 0=immediate callback из put (default), 1=defer via uasync_call_soon (очистка стека)
queue_callback_fn empty_callback; // одноразовый коллбэк при опустошении очереди
void* empty_callback_arg;
@ -241,6 +242,17 @@ void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes);
void queue_set_waiter_defer(struct ll_queue* q, int enable);
/**
* @brief Управляет способом вызова коллбэка автозабора при добавлении элемента.
* @param q очередь
* @param enable 0 — первый запуск коллбэка из queue_data_put* вызывается синхронно (по умолчанию);
* 1 — первый запуск откладывается через uasync_call_soon (очищает стек вызывающего).
*
* @note Влияет только на первый запуск из put. Повторные запуски (queue_resume_callback)
* всегда идут через uasync_call_soon.
*/
void queue_set_callback_defer(struct ll_queue* q, int enable);
/**
* @brief Устанавливает одноразовый коллбэк при опустошении очереди (count==0).
* @param q очередь

1
lib/ll_queue_doc.md

@ -195,6 +195,7 @@ queue_set_empty_callback(q, on_empty, my_data);
| `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_set_callback_defer(q, enable)` | `enable=1` — первый запуск callback из `put` откладывается через uasync_call_soon (очистка стека). По умолчанию `0` (синхронно). Повторные запуски всегда deferred. |
### Запись

30
tests/test_ll_queue.c

@ -170,6 +170,35 @@ static void test_callback(void) {
PASS();
}
static void test_callback_defer(void) {
TEST("callback defer (call_soon) vs immediate");
struct UASYNC *ua = uasync_create();
struct ll_queue *q = queue_new(ua, 0, 0, 0, "cd");
int cnt = 0;
queue_set_callback(q, queue_cb, &cnt);
queue_set_callback_defer(q, 1);
test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d));
d->id = 1; d->value = 1; d->checksum = checksum(d);
queue_data_put(q, (struct ll_entry*)d);
ASSERT_EQ(cnt, 0, "deferred: no sync callback on first put");
for (int i = 0; i < 10 && cnt < 1; i++) uasync_poll(ua, 5);
ASSERT_EQ(cnt, 1, "deferred: callback fired via call_soon");
ASSERT_EQ(queue_entry_count(q), 0, "");
queue_set_callback_defer(q, 0);
cnt = 0;
d = (test_data_t*)queue_entry_new(sizeof(*d));
d->id = 2; d->value = 2; d->checksum = checksum(d);
queue_data_put(q, (struct ll_entry*)d);
ASSERT_EQ(cnt, 1, "immediate: sync callback on first put");
queue_free(q); uasync_destroy(ua, 0);
PASS();
}
static void test_waiter(void) {
TEST("wait_threshold + cancel");
struct UASYNC *ua = uasync_create();
@ -421,6 +450,7 @@ int main(void) {
test_fifo();
test_lifo_priority();
test_callback();
test_callback_defer();
test_waiter();
test_waiter_multiple();
test_waiter_cancel();

Loading…
Cancel
Save