Browse Source

fix use-after-free in tests: swap queue_entry_free/queue_dgram_free order, fix dangling timer handles and stack args; add ETCP MTU bounds, assembly overflow check, etcp_update_mtu

congestion
Evgeny 5 months ago
parent
commit
b829ff7001
  1. 38
      src/etcp.c
  2. 3
      src/etcp.h
  3. 13
      src/etcp_connections.c
  4. 2
      src/etcp_connections.h
  5. 12
      tests/test_etcp_api.c
  6. 4
      tests/test_nat_transport.c
  7. 6
      tests/test_remote_proxy.c
  8. 33
      tests/test_u_async_comprehensive.c
  9. 4
      tests/test_u_async_timeouts.c

38
src/etcp.c

@ -72,6 +72,10 @@ static void feed_dgram_to_asm(struct ETCP_CONN* etcp, struct PKTNORM* pn,
uint8_t** asm_buf, uint16_t* asm_len, uint16_t* asm_cap, uint8_t** asm_buf, uint16_t* asm_len, uint16_t* asm_cap,
uint32_t* returned) { uint32_t* returned) {
uint16_t need = *asm_len + dgram_len; uint16_t need = *asm_len + dgram_len;
if (need > ASM_BUF_MAX_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] assembly buffer would exceed max size %d", etcp->log_name, ASM_BUF_MAX_SIZE);
return;
}
if (need > *asm_cap) { if (need > *asm_cap) {
uint16_t new_cap = *asm_cap ? *asm_cap : 256; uint16_t new_cap = *asm_cap ? *asm_cap : 256;
while (new_cap < need) new_cap *= 2; while (new_cap < need) new_cap *= 2;
@ -219,7 +223,7 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n
// etcp->window_size = MAX_INFLIGHT_BYTES; // Not used // etcp->window_size = MAX_INFLIGHT_BYTES; // Not used
etcp->mtu = 1500; // Default MTU etcp->mtu = ETCP_RFC791_MIN_MTU; // Default MTU per RFC 791
etcp->next_tx_id = 1; etcp->next_tx_id = 1;
etcp->rtt_avg_10 = 10; // Initial guess (1ms) etcp->rtt_avg_10 = 10; // Initial guess (1ms)
etcp->rtt_avg_100 = 10; etcp->rtt_avg_100 = 10;
@ -536,8 +540,8 @@ int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
} }
// Check length against maximum packet size // Check length against maximum packet size
if (len > PACKET_DATA_SIZE) { if (len > ETCP_MAX_PAYLOAD_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] packet too large (len=%zu, max=%d)", etcp->log_name, len, PACKET_DATA_SIZE); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] packet too large (len=%zu, max=%d)", etcp->log_name, len, ETCP_MAX_PAYLOAD_SIZE);
return -1; return -1;
} }
@ -1086,7 +1090,9 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
dgram->data[ptr++]=inf_pkt->seq>>16; dgram->data[ptr++]=inf_pkt->seq>>16;
dgram->data[ptr++]=inf_pkt->seq>>24; dgram->data[ptr++]=inf_pkt->seq>>24;
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len; if ((int)(ptr + inf_pkt->ll.len) <= PACKET_DATA_SIZE) {
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len;
} else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload overflow: ptr=%d len=%u max=%d", etcp->log_name, ptr, inf_pkt->ll.len, PACKET_DATA_SIZE);
} }
else { else {
if (link->burst_active) { if (link->burst_active) {
@ -1434,6 +1440,13 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
len = 0; len = 0;
break; break;
} }
if (data + len > pkt->data + pkt->data_len) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload section bounds overflow: data+len=%p > end=%p", etcp->log_name, (void*)(data + len), (void*)(pkt->data + pkt->data_len));
memory_pool_free(etcp->instance->data_pool, payload_data);
queue_entry_free(&rx_pkt->ll);
len = 0;
break;
}
// Copy the actual payload data // Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len); memcpy(payload_data, data + 5, pkt_len);
queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt); queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt);
@ -1546,3 +1559,20 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram
} }
void etcp_update_mtu(struct ETCP_CONN* etcp) {
if (!etcp) return;
int new_mtu = PACKET_DATA_MAX_MTU;
int has_links = 0;
struct ETCP_LINK* link = etcp->links;
while (link) {
has_links = 1;
if (link->mtu > 0 && link->mtu < new_mtu) new_mtu = link->mtu;
link = link->next;
}
if (!has_links) new_mtu = ETCP_RFC791_MIN_MTU;
if (new_mtu != etcp->mtu) {
etcp->mtu = new_mtu;
if (etcp->normalizer) etcp->normalizer->frag_size = etcp->mtu - ACK_REZERV - UDP_HDR_SIZE - UDP_SC_HDR_SIZE;
}
}

3
src/etcp.h

@ -59,6 +59,7 @@ uint16_t get_current_timestamp(void);
#define INFLIGHT_INITIAL_HASH_SIZE 1024 #define INFLIGHT_INITIAL_HASH_SIZE 1024
#define MAX_INFLIGHT_SIZE 16384 // максимальное число элементов в inflight приёмной очереди (для предотвращения атак) #define MAX_INFLIGHT_SIZE 16384 // максимальное число элементов в inflight приёмной очереди (для предотвращения атак)
#define ASM_BUF_MAX_SIZE (64 * 1024) // максимальный размер буфера сборки фрагментов
// в этот список пакет добавляется когда перемещается из input_queue в input_send_q, при этом к пакету добавляется struct INFLIGHT_PACKET из inflight_pool. // в этот список пакет добавляется когда перемещается из input_queue в input_send_q, при этом к пакету добавляется struct INFLIGHT_PACKET из inflight_pool.
// пакет полностью удаляется когда приходит ACK (либо conn_reset/close) // пакет полностью удаляется когда приходит ACK (либо conn_reset/close)
@ -234,6 +235,8 @@ void etcp_conn_ready(struct ETCP_CONN* conn);
void etcp_on_link_down(struct ETCP_CONN* etcp); void etcp_on_link_down(struct ETCP_CONN* etcp);
void etcp_update_mtu(struct ETCP_CONN* etcp);
#ifdef __cplusplus #ifdef __cplusplus
} }
#endif #endif

13
src/etcp_connections.c

@ -775,6 +775,7 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
if (mtu > PACKET_DATA_MAX_MTU) mtu = PACKET_DATA_MAX_MTU; if (mtu > PACKET_DATA_MAX_MTU) mtu = PACKET_DATA_MAX_MTU;
link->mtu_local = mtu; link->mtu_local = mtu;
link->mtu = mtu; link->mtu = mtu;
etcp_update_mtu(etcp);
link->initialized = 0; link->initialized = 0;
link->init_timer = NULL; link->init_timer = NULL;
link->init_timeout = 0; link->init_timeout = 0;
@ -1075,11 +1076,11 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
// Mark that packet was sent (for keepalive logic) // Mark that packet was sent (for keepalive logic)
dgram->link->pkt_sent_since_keepalive = 1; dgram->link->pkt_sent_since_keepalive = 1;
dgram->flag_up=dgram->link->recv_keepalive; dgram->flag_up=dgram->link->recv_keepalive;
// 28 байт = udp headers. MTU=UDP payload+28 (1472 bytes max) // Размер зашифрованной части не должен превышать UDP payload = MTU - заголовки
int errcode=0; int errcode=0;
sc_context_t* sc = &dgram->link->etcp->crypto_ctx; sc_context_t* sc = &dgram->link->etcp->crypto_ctx;
int len=dgram->data_len-dgram->noencrypt_len;// не забываем добавить timestamp (2 bytes) int len=dgram->data_len-dgram->noencrypt_len;// не забываем добавить timestamp (2 bytes)
if (len<0 || len>1472) { dgram->link->send_errors++; errcode=1; goto es_err; } if (len<0 || len>dgram->link->mtu - UDP_HDR_SIZE) { dgram->link->send_errors++; errcode=1; goto es_err; }
uint8_t enc_buf[1600]; uint8_t enc_buf[1600];
size_t enc_buf_len=0; size_t enc_buf_len=0;
dgram->timestamp=get_current_timestamp(); dgram->timestamp=get_current_timestamp();
@ -1090,7 +1091,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
if (enc_buf_len == 0) { if (enc_buf_len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "eencryption failed for node %016llx", (unsigned long long)dgram->link->etcp->instance->node_id); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "eencryption failed for node %016llx", (unsigned long long)dgram->link->etcp->instance->node_id);
dgram->link->send_errors++; errcode=2; goto es_err; } dgram->link->send_errors++; errcode=2; goto es_err; }
if (enc_buf_len + dgram->noencrypt_len > 1472) { if (enc_buf_len + dgram->noencrypt_len > (size_t)(dgram->link->mtu - UDP_HDR_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "packet too long len=%d ne_len=%d", enc_buf_len, dgram->noencrypt_len); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "packet too long len=%d ne_len=%d", enc_buf_len, dgram->noencrypt_len);
dgram->link->send_errors++; errcode=3; goto es_err; } dgram->link->send_errors++; errcode=3; goto es_err; }
memcpy(enc_buf+enc_buf_len, dgram->data+len, dgram->noencrypt_len); memcpy(enc_buf+enc_buf_len, dgram->data+len, dgram->noencrypt_len);
@ -1126,7 +1127,7 @@ static int etcp_send_ping_raw(struct ETCP_DGRAM* dgram, socket_t fd, sc_context_
return -1; return -1;
} }
int len = dgram->data_len - dgram->noencrypt_len; int len = dgram->data_len - dgram->noencrypt_len;
if (len < 0 || len > 1472) { if (len < 0 || len > PACKET_DATA_SIZE - UDP_HDR_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping packet data invalid len=%d", len); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping packet data invalid len=%d", len);
return -1; return -1;
} }
@ -1139,7 +1140,7 @@ static int etcp_send_ping_raw(struct ETCP_DGRAM* dgram, socket_t fd, sc_context_
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "encryption failed for ping"); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "encryption failed for ping");
return -1; return -1;
} }
if (enc_buf_len + dgram->noencrypt_len > 1472) { if (enc_buf_len + dgram->noencrypt_len > (size_t)(PACKET_DATA_SIZE - UDP_HDR_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping packet too long enc=%zu ne=%d", enc_buf_len, dgram->noencrypt_len); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping packet too long enc=%zu ne=%d", enc_buf_len, dgram->noencrypt_len);
return -1; return -1;
} }
@ -1632,6 +1633,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
link->mtu_remote = (ack_hdr->mtu[0] << 8) | ack_hdr->mtu[1]; link->mtu_remote = (ack_hdr->mtu[0] << 8) | ack_hdr->mtu[1];
if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU; if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU;
link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote; link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote;
etcp_update_mtu(link->etcp);
struct { struct {
uint8_t code; uint8_t code;
@ -1780,6 +1782,7 @@ process_decrypted:
offset += 2; offset += 2;
if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU; if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU;
link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote; link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote;
etcp_update_mtu(link->etcp);
link->remote_link_id = pkt->data[offset++]; link->remote_link_id = pkt->data[offset++];
link->remote_socket_id = pkt->data[offset++]; link->remote_socket_id = pkt->data[offset++];
link->remote_only_local = pkt->data[offset++]; link->remote_only_local = pkt->data[offset++];

2
src/etcp_connections.h

@ -16,6 +16,8 @@
#define PACKET_DATA_SIZE 1600//1536 #define PACKET_DATA_SIZE 1600//1536
#define PACKET_DATA_MAX_MTU 1600 #define PACKET_DATA_MAX_MTU 1600
#define ETCP_MAX_PAYLOAD_SIZE (PACKET_DATA_SIZE - ETCP_ACK_BASE_SIZE - 5) /* 1587 */
#define ETCP_RFC791_MIN_MTU 576
// Типы кодограмм протокола // Типы кодограмм протокола
#define ETCP_INIT_REQUEST 0x02 #define ETCP_INIT_REQUEST 0x02

12
tests/test_etcp_api.c

@ -234,8 +234,8 @@ static void server_recv_callback(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!entry || !entry->dgram || entry->len < PACKET_HEADER_SIZE + 1) { if (!entry || !entry->dgram || entry->len < PACKET_HEADER_SIZE + 1) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Server received invalid packet"); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Server received invalid packet");
if (entry) { if (entry) {
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
} }
return; return;
} }
@ -248,8 +248,8 @@ static void server_recv_callback(struct ETCP_CONN* conn, struct ll_entry* entry)
seq, packets_received_fwd); seq, packets_received_fwd);
} }
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
} }
// Callback для получения пакетов на клиенте (backward direction) // Callback для получения пакетов на клиенте (backward direction)
@ -259,8 +259,8 @@ static void client_recv_callback(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!entry || !entry->dgram || entry->len < PACKET_HEADER_SIZE + 1) { if (!entry || !entry->dgram || entry->len < PACKET_HEADER_SIZE + 1) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Client received invalid packet"); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Client received invalid packet");
if (entry) { if (entry) {
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
} }
return; return;
} }
@ -273,8 +273,8 @@ static void client_recv_callback(struct ETCP_CONN* conn, struct ll_entry* entry)
seq, packets_received_back); seq, packets_received_back);
} }
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
} }
// Check if connection is established and crypto session key is ready // Check if connection is established and crypto session key is ready
@ -347,8 +347,8 @@ static void send_packets_fwd(void) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send failed for packet %d", DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send failed for packet %d",
packets_sent_fwd); packets_sent_fwd);
// При ошибке освобождаем entry сами // При ошибке освобождаем entry сами
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
break; break;
} }
@ -397,8 +397,8 @@ static void send_packets_back(void) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send failed for packet %d", DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send failed for packet %d",
packets_sent_back); packets_sent_back);
// При ошибке освобождаем entry сами // При ошибке освобождаем entry сами
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
break; break;
} }

4
tests/test_nat_transport.c

@ -302,8 +302,8 @@ static int test_provider_egress(void) {
int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry); int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry);
if (ret != 0) { if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NAT, "etcp_route_send failed"); DEBUG_ERROR(DEBUG_CATEGORY_NAT, "etcp_route_send failed");
queue_entry_free(entry);
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry);
return 0; return 0;
} }
DEBUG_INFO(DEBUG_CATEGORY_NAT, "ETCP_ID_NAT sent from client to provider via etcp_router"); DEBUG_INFO(DEBUG_CATEGORY_NAT, "ETCP_ID_NAT sent from client to provider via etcp_router");
@ -469,7 +469,7 @@ static int test_full_roundtrip(void) {
entry->dgram = dgram; entry->dgram = dgram;
entry->len = total; entry->len = total;
int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry); int ret = etcp_route_send(inst_client, NODE_ID_PROVIDER, entry);
if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_NAT, "etcp_route_send failed"); queue_entry_free(entry); queue_dgram_free(entry); return 0; } if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_NAT, "etcp_route_send failed"); queue_dgram_free(entry); queue_entry_free(entry); return 0; }
// Poll for provider to process egress // Poll for provider to process egress
int cycles = 0; int cycles = 0;

6
tests/test_remote_proxy.c

@ -70,18 +70,18 @@ done: close(cli); close(srv);
static void test_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { static void test_handler(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* i = conn ? conn->instance : inst; struct UTUN_INSTANCE* i = conn ? conn->instance : inst;
if (!i || !entry || entry->dgram == NULL || entry->len < TCP_PROXY_HDR_SIZE) { if (!i || !entry || entry->dgram == NULL || entry->len < TCP_PROXY_HDR_SIZE) {
if (entry) { queue_entry_free(entry); queue_dgram_free(entry); } return; if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return;
} }
uint8_t subcmd = entry->dgram[1]; uint8_t subcmd = entry->dgram[1];
uint64_t sid = 0; memcpy(&sid, entry->dgram + 2, 8); uint64_t sid = 0; memcpy(&sid, entry->dgram + 2, 8);
uint64_t src = conn ? conn->peer_node_id : i->node_id; uint64_t src = conn ? conn->peer_node_id : i->node_id;
if (subcmd == TCP_PROXY_SUBCMD_CONNECT && sid == stream_id) { remote_proxy_handle_connect(i, entry, sid, src); return; } if (subcmd == TCP_PROXY_SUBCMD_CONNECT && sid == stream_id) { remote_proxy_handle_connect(i, entry, sid, src); return; }
if (sid != stream_id) { queue_entry_free(entry); queue_dgram_free(entry); return; } if (sid != stream_id) { queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_CONNECTED) { if (subcmd == TCP_PROXY_SUBCMD_CONNECTED) {
uint8_t status = (entry->len >= TCP_PROXY_CONNECTED_HDR_SIZE) ? entry->dgram[TCP_PROXY_HDR_SIZE + 2] : 1; uint8_t status = (entry->len >= TCP_PROXY_CONNECTED_HDR_SIZE) ? entry->dgram[TCP_PROXY_HDR_SIZE + 2] : 1;
connected_ok = (status == TCP_PROXY_CONNECTED_OK) ? 1 : -1; connected_ok = (status == TCP_PROXY_CONNECTED_OK) ? 1 : -1;
} }
queue_entry_free(entry); queue_dgram_free(entry); queue_dgram_free(entry); queue_entry_free(entry);
} }
static void monitor(void* arg) { static void monitor(void* arg) {

33
tests/test_u_async_comprehensive.c

@ -306,10 +306,14 @@ static void test_memory_leak_detection(void) {
ASSERT_NOT_NULL(ua, "Failed to create uasync instance"); ASSERT_NOT_NULL(ua, "Failed to create uasync instance");
/* Create and destroy multiple timers without proper cleanup */ /* Create and destroy multiple timers without proper cleanup */
test_context_t* timer_ctxs[10] = {NULL};
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
test_context_t timer_ctx = {0}; test_context_t* timer_ctx = u_malloc(sizeof(test_context_t));
timer_ctx.timeout_ms = 100; // Set proper timeout to avoid immediate timeout ASSERT_NOT_NULL(timer_ctx, "Failed to allocate timer context");
void* timer = uasync_set_timeout(ua, 100, &timer_ctx, test_timer_callback, "test_cancel"); memset(timer_ctx, 0, sizeof(*timer_ctx));
timer_ctx->timeout_ms = 100;
timer_ctxs[i] = timer_ctx;
void* timer = uasync_set_timeout(ua, 100, timer_ctx, test_timer_callback, "test_cancel");
ASSERT_NOT_NULL(timer, "Failed to set timer"); ASSERT_NOT_NULL(timer, "Failed to set timer");
/* Cancel some, leave others to timeout */ /* Cancel some, leave others to timeout */
@ -323,6 +327,11 @@ static void test_memory_leak_detection(void) {
uasync_poll(ua, 10); uasync_poll(ua, 10);
} }
/* Free allocated contexts */
for (int i = 0; i < 10; i++) {
if (timer_ctxs[i]) u_free(timer_ctxs[i]);
}
/* Destroy should detect any leaks and abort if found */ /* Destroy should detect any leaks and abort if found */
/* This test passes if we don't abort */ /* This test passes if we don't abort */
uasync_destroy(ua, 0); uasync_destroy(ua, 0);
@ -462,15 +471,16 @@ static void test_concurrent_operations(void) {
/* Create multiple timers with different timeouts */ /* Create multiple timers with different timeouts */
void* timers[10]; void* timers[10];
test_context_t* contexts[10] = {NULL};
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
// Use unique contexts to avoid all timers having timeout_ms = 0
test_context_t* individual_ctx = u_malloc(sizeof(test_context_t)); test_context_t* individual_ctx = u_malloc(sizeof(test_context_t));
ASSERT_NOT_NULL(individual_ctx, "Failed to allocate timer context"); ASSERT_NOT_NULL(individual_ctx, "Failed to allocate timer context");
individual_ctx->callback_count = 0; individual_ctx->callback_count = 0;
individual_ctx->expected_count = 0; individual_ctx->expected_count = 0;
individual_ctx->timeout_ms = (i + 1) * 5; // Different timeouts: 5, 10, 15, ..., 50ms individual_ctx->timeout_ms = (i + 1) * 5;
individual_ctx->callback_arg = 0; individual_ctx->callback_arg = 0;
contexts[i] = individual_ctx;
timers[i] = uasync_set_timeout(ua, (i + 1) * 5, individual_ctx, test_timer_callback, "test_stress"); timers[i] = uasync_set_timeout(ua, (i + 1) * 5, individual_ctx, test_timer_callback, "test_stress");
ASSERT_NOT_NULL(timers[i], "Failed to set timer"); ASSERT_NOT_NULL(timers[i], "Failed to set timer");
} }
@ -490,7 +500,7 @@ static void test_concurrent_operations(void) {
if (cycle % 3 == 0) { if (cycle % 3 == 0) {
char data = 'x'; char data = 'x';
ssize_t wret = write(sockets[0], &data, 1); ssize_t wret = write(sockets[0], &data, 1);
(void)wret; /* Suppress warning - best effort write for testing */ (void)wret;
} }
/* Poll */ /* Poll */
@ -510,18 +520,13 @@ static void test_concurrent_operations(void) {
} }
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
if (timers[i]) { if (contexts[i] && contexts[i]->callback_count == 0 && timers[i]) {
// Note: We cannot safely access the context here because the timer might be expired
// The context will be freed by the timeout_heap free callback during uasync_destroy
uasync_cancel_timeout(ua, timers[i]); uasync_cancel_timeout(ua, timers[i]);
timers[i] = NULL;
} }
timers[i] = NULL;
if (contexts[i]) { u_free(contexts[i]); contexts[i] = NULL; }
} }
// Free any remaining allocated contexts (for timers that were cancelled)
// This is a simplified approach - in production code, you'd maintain a list of allocated contexts
// For this test, we'll let uasync_destroy handle the cleanup through the timeout heap's free callback
uasync_destroy(ua, 0); uasync_destroy(ua, 0);
TEST_PASS(); TEST_PASS();
} }

4
tests/test_u_async_timeouts.c

@ -188,10 +188,10 @@ int main(void) {
DEBUG_INFO(DEBUG_CATEGORY_TEST, "Total fires during test: %d", total_fired); DEBUG_INFO(DEBUG_CATEGORY_TEST, "Total fires during test: %d", total_fired);
/* Отменяем стоп-таймер если он ещё не сработал */ /* Отменяем стоп-таймер если он ещё не сработал */
if (stop_timer) { if (!stop_flag && stop_timer) {
uasync_cancel_timeout(ua, stop_timer); uasync_cancel_timeout(ua, stop_timer);
u_free(stop_ctx);
} }
u_free(stop_ctx);
uasync_destroy(ua, 0); uasync_destroy(ua, 0);

Loading…
Cancel
Save