Browse Source

1

etcp-inflight-fix
Evgeny 4 months ago
parent
commit
4556f327c2
  1. 41
      lib/ll_queue.c
  2. 27
      lib/ll_queue.h
  3. 5
      src/etcp_router.c
  4. 3
      src/etcp_router.h
  5. 4
      src/pkt_normalizer.c
  6. 4
      src/proxy/tcp_proxy_client.c
  7. 4
      src/proxy/tcp_proxy_server.c
  8. 24
      tests/test_intensive_memory_pool.c
  9. 38
      tests/test_ll_queue.c
  10. 7
      tests/test_memory_pool_and_config.c

41
lib/ll_queue.c

@ -291,13 +291,15 @@ static void check_waiters(struct ll_queue* q) {
struct queue_waiter_handle* h = waiter->handle;
if (h) h->internal = NULL;
queue_threshold_callback_fn cb = waiter->callback;
void* arg = waiter->callback_arg;
u_free(waiter);
if (h) {
if (cb) {
DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "check_waiters: waking head waiter, count=%d<=%d, bytes=%zu<=%zu",
q->count, q->threshold_max_packets, q->total_bytes, q->threshold_max_bytes);
h->callback(q, h->callback_arg);
cb(q, arg);
}
}
@ -554,29 +556,18 @@ void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes)
q->threshold_max_bytes = max_bytes;
}
void queue_waiter_handle_init(struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg) {
if (!h) return;
h->internal = NULL;
h->callback = callback;
h->callback_arg = arg;
}
int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg) {
if (!q || !h) return -1;
if (h->internal) queue_waiter_cancel(q, h);
int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h) {
if (!q || !h) return 1;
if (h->internal) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] queue_waiter_wait: waiter already registered", q->name);
return 1;
}
// Проверить условие немедленно
if (q->count <= q->threshold_max_packets
&& (q->threshold_max_bytes == 0 || q->total_bytes <= q->threshold_max_bytes)) {
h->callback(q, h->callback_arg);
callback(q, arg);
return 1;
}
struct queue_waiter* waiter = u_malloc(sizeof(struct queue_waiter));
if (!waiter) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] queue_waiter_wait: u_malloc failed", q->name);
@ -584,11 +575,13 @@ int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h) {
}
waiter->next = NULL;
waiter->handle = h;
waiter->callback = callback;
waiter->callback_arg = arg;
h->internal = waiter;
if (q->waiter_tail) q->waiter_tail->next = waiter; else q->waiter_head = waiter;
q->waiter_tail = waiter;
return 0;
}

27
lib/ll_queue.h

@ -99,6 +99,8 @@ typedef void (*queue_threshold_callback_fn)(struct ll_queue* q, void* arg);
struct queue_waiter {
struct queue_waiter* next; ///< Следующий в списке ожидающих
struct queue_waiter_handle* handle; ///< Обратная ссылка на публичную структуру
queue_threshold_callback_fn callback;
void* callback_arg;
};
/**
@ -106,15 +108,12 @@ struct queue_waiter {
* @brief Публичная управляющая структура (встраивается в структуру вызывающей стороны).
*
* Вызывающая сторона:
* 1. Объявляет поле struct queue_waiter_handle в своей структуре
* 2. Вызывает queue_waiter_handle_init() один раз при инициализации
* 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 = не ждём; иначе — внутренний узел
queue_threshold_callback_fn callback;
void* callback_arg;
};
/**
@ -223,22 +222,18 @@ void queue_resume_callback(struct ll_queue* q);
*/
void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes);
/**
* @brief Инициализирует публичную управляющую структуру waiter.
* @param h указатель на handle (встраивается в структуру вызывающей стороны)
* @param callback функция, вызываемая при достижении порога
* @param arg аргумент коллбэка
*/
void queue_waiter_handle_init(struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg);
/**
* @brief Регистрирует ожидание освобождения очереди до общего порога.
* @param q очередь
* @param h указатель на инициализированный handle
* @return 1 — условие уже выполнено (callback вызван немедленно), 0 — зарегистрирован в списке
* @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);
int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg);
/**
* @brief Отменяет ожидание (удаляет из списка, если зарегистрирован).

5
src/etcp_router.c

@ -516,11 +516,12 @@ int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id) {
}
void etcp_router_waiter_register(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
struct queue_waiter_handle* h) {
struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg) {
if (!inst || !inst->bgp || !h) return;
struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, peer_node_id);
if (!conn || !conn->normalizer || !conn->normalizer->input) return;
queue_waiter_wait(conn->normalizer->input, h);
queue_waiter_wait(conn->normalizer->input, h, callback, arg);
}
void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,

3
src/etcp_router.h

@ -93,7 +93,8 @@ int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id);
// Backpressure: зарегистрировать/отменить waiter на normalizer->input очереди
void etcp_router_waiter_register(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
struct queue_waiter_handle* h);
struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg);
void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
struct queue_waiter_handle* h);

4
src/pkt_normalizer.c

@ -69,7 +69,7 @@ struct PKTNORM* pn_init(struct ETCP_CONN* etcp) {
pn->recvpart = NULL;
pn->flush_timer = NULL;
queue_waiter_handle_init(&pn->input_waiter_handle, etcp_input_ready_cb, pn);
pn->input_waiter_handle.internal = NULL;
return pn;
}
@ -215,7 +215,7 @@ static void packer_cb(struct ll_queue* q, void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "input_q->pn: waiting etcp input threshold");
queue_waiter_wait(pn->etcp->input_queue, &pn->input_waiter_handle);
queue_waiter_wait(pn->etcp->input_queue, &pn->input_waiter_handle, etcp_input_ready_cb, pn);
}
// Helper to send block to ETCP as ETCP_FRAGMENT

4
src/proxy/tcp_proxy_client.c

@ -227,7 +227,6 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t
tcp_nagle_disable(newpcb);
pc->to_lwip = queue_new(p->ua, 0, 0, 0, "to_lwip");
queue_waiter_handle_init(&pc->tx_waiter, tcp_proxy_client_tx_waiter_cb, pc);
if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); u_free(pc); return LERR_MEM; }
@ -285,7 +284,8 @@ static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbu
if (ret == 0) { tcp_recved(pcb, len); u_free(data); }
else {
pc->tx_buf = data; pc->tx_len = len;
etcp_router_waiter_register(pc->proxy->inst, pc->proxy->via_node_id, &pc->tx_waiter);
etcp_router_waiter_register(pc->proxy->inst, pc->proxy->via_node_id, &pc->tx_waiter,
tcp_proxy_client_tx_waiter_cb, pc);
}
} else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY recv malloc(%u) failed sid=%08x", len, pc->stream_id);
pbuf_free(p);

4
src/proxy/tcp_proxy_server.c

@ -199,7 +199,8 @@ static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) {
if (rc->pause_buf) { memcpy(rc->pause_buf, buf, n); rc->pause_len = n; }
else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE malloc(%zd) failed sid=%08x", n, rc->stream_id);
if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; }
etcp_router_waiter_register(inst, rc->peer_node_id, &rc->pause_waiter);
etcp_router_waiter_register(inst, rc->peer_node_id, &rc->pause_waiter,
tcp_proxy_server_pause_waiter_cb, rc);
}
} else if (n == 0) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x total=%d cli_closed=%d sock_closed=%d",
@ -298,7 +299,6 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
rc->ctx = ctx; rc->stream_id = stream_id; rc->peer_node_id = src_node_id;
memcpy(rc->dest_ip, dest_ip, 4); rc->dest_port = dest_port;
rc->ua = inst->ua; rc->sock = SOCKET_INVALID;
queue_waiter_handle_init(&rc->pause_waiter, tcp_proxy_server_pause_waiter_cb, rc);
rc->sock = socket(AF_INET, SOCK_STREAM, 0);
if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; }

24
tests/test_intensive_memory_pool.c

@ -37,12 +37,9 @@ static double test_without_pools(int iterations) {
}
// Создать много waiters (общий порог 0, count=10 => все в очередь)
struct queue_waiter_handle handles[32];
for (int i = 0; i < 32; i++) {
queue_waiter_handle_init(&handles[i], intensive_waiter_callback, NULL);
queue_waiter_wait(queue, &handles[i]);
}
struct queue_waiter_handle handles[32] = {0};
for (int i = 0; i < 32; i++) queue_waiter_wait(queue, &handles[i], intensive_waiter_callback, NULL);
// Удалить записи (триггер waiters — по одному за get)
for (int i = 0; i < 10; i++) {
void* retrieved = queue_data_get(queue);
@ -50,9 +47,7 @@ static double test_without_pools(int iterations) {
}
// Отменить оставшиеся waiters
for (int i = 0; i < 32; i++) {
queue_waiter_cancel(queue, &handles[i]);
}
for (int i = 0; i < 32; i++) queue_waiter_cancel(queue, &handles[i]);
}
queue_free(queue);
@ -85,11 +80,8 @@ static double test_with_pools(int iterations) {
}
// Создать много waiters
struct queue_waiter_handle handles[32];
for (int i = 0; i < 32; i++) {
queue_waiter_handle_init(&handles[i], intensive_waiter_callback, NULL);
queue_waiter_wait(queue, &handles[i]);
}
struct queue_waiter_handle handles[32] = {0};
for (int i = 0; i < 32; i++) queue_waiter_wait(queue, &handles[i], intensive_waiter_callback, NULL);
// Удалить записи (триггер waiters)
for (int i = 0; i < 10; i++) {
@ -98,9 +90,7 @@ static double test_with_pools(int iterations) {
}
// Отменить оставшиеся waiters
for (int i = 0; i < 32; i++) {
queue_waiter_cancel(queue, &handles[i]);
}
for (int i = 0; i < 32; i++) queue_waiter_cancel(queue, &handles[i]);
}
queue_free(queue);

38
tests/test_ll_queue.c

@ -177,9 +177,8 @@ static void test_waiter(void) {
queue_set_threshold(q, 2, 0);
int called = 0;
struct queue_waiter_handle h;
queue_waiter_handle_init(&h, waiter_cb, &called);
int ret = queue_waiter_wait(q, &h);
struct queue_waiter_handle h = {0};
int ret = queue_waiter_wait(q, &h, waiter_cb, &called);
ASSERT(ret == 1 && called == 1, "immediate when condition met"); // empty queue
for (int i = 0; i < 5; i++) {
@ -188,7 +187,7 @@ static void test_waiter(void) {
}
called = 0;
ret = queue_waiter_wait(q, &h);
ret = queue_waiter_wait(q, &h, waiter_cb, &called);
ASSERT(ret == 0 && called == 0, "");
for (int i = 0; i < 3; i++) queue_entry_free((struct ll_entry*)queue_data_get(q));
@ -214,11 +213,8 @@ static void test_waiter_multiple(void) {
waiter_order_idx = 0;
int ids[4] = {10, 20, 30, 40};
struct queue_waiter_handle handles[4];
for (int i = 0; i < 4; i++) {
queue_waiter_handle_init(&handles[i], waiter_cb_order, &ids[i]);
ASSERT(queue_waiter_wait(q, &handles[i]) == 0, "should queue");
}
struct queue_waiter_handle handles[4] = {0};
for (int i = 0; i < 4; i++) ASSERT(queue_waiter_wait(q, &handles[i], waiter_cb_order, &ids[i]) == 0, "should queue");
for (int i = 0; i < 5; i++) queue_entry_free((struct ll_entry*)queue_data_get(q));
ASSERT_EQ(waiter_order_idx, 3, "3 waiters fired");
@ -242,11 +238,8 @@ static void test_waiter_cancel(void) {
waiter_order_idx = 0;
int ids[3] = {1, 2, 3};
struct queue_waiter_handle handles[3];
for (int i = 0; i < 3; i++) {
queue_waiter_handle_init(&handles[i], waiter_cb_order, &ids[i]);
queue_waiter_wait(q, &handles[i]); // count=10 > 3, all queued
}
struct queue_waiter_handle handles[3] = {0};
for (int i = 0; i < 3; i++) queue_waiter_wait(q, &handles[i], waiter_cb_order, &ids[i]); // count=10 > 3, all queued
queue_waiter_cancel(q, &handles[1]); // cancel middle (ID 2)
@ -270,9 +263,8 @@ static void test_waiter_queue_free(void) {
}
int called = 0;
struct queue_waiter_handle h;
queue_waiter_handle_init(&h, waiter_cb, &called);
queue_waiter_wait(q, &h);
struct queue_waiter_handle h = {0};
queue_waiter_wait(q, &h, waiter_cb, &called);
ASSERT(h.internal != NULL, "waiter is queued");
queue_free(q);
@ -293,9 +285,8 @@ static void test_waiter_threshold(void) {
}
int called = 0;
struct queue_waiter_handle h;
queue_waiter_handle_init(&h, waiter_cb, &called);
int ret = queue_waiter_wait(q, &h);
struct queue_waiter_handle h = {0};
int ret = queue_waiter_wait(q, &h, waiter_cb, &called);
ASSERT(ret == 0 && called == 0, "queued (10 > 0)");
queue_set_threshold(q, 10, 0); // raise threshold
@ -313,10 +304,9 @@ static void test_waiter_reuse(void) {
struct ll_queue *q = queue_new(ua, 0, 0, 0, "wr");
int called = 0;
struct queue_waiter_handle h;
queue_waiter_handle_init(&h, waiter_cb, &called);
struct queue_waiter_handle h = {0};
int ret = queue_waiter_wait(q, &h);
int ret = queue_waiter_wait(q, &h, waiter_cb, &called);
ASSERT(ret == 1 && called == 1, "immediate (empty queue)");
for (int i = 0; i < 3; i++) {
@ -325,7 +315,7 @@ static void test_waiter_reuse(void) {
}
called = 0;
ret = queue_waiter_wait(q, &h);
ret = queue_waiter_wait(q, &h, waiter_cb, &called);
ASSERT(ret == 0 && called == 0, "queued");
for (int i = 0; i < 3; i++) queue_entry_free((struct ll_entry*)queue_data_get(q));

7
tests/test_memory_pool_and_config.c

@ -45,12 +45,11 @@ int main() {
return 1;
}
queue_set_threshold(queue, 2, 0);
// Test multiple waiters with new linked-list backpressure
struct queue_waiter_handle handles[10];
struct queue_waiter_handle handles[10] = {0};
for (int i = 0; i < 10; i++) {
queue_waiter_handle_init(&handles[i], test_waiter_callback, NULL);
queue_waiter_wait(queue, &handles[i]); // queue empty => immediate callback
queue_waiter_wait(queue, &handles[i], test_waiter_callback, NULL); // queue empty => immediate callback
}
// Add some entries — waiters already fired, none queued

Loading…
Cancel
Save