From 03f71fd94e21149fbdf9ee58ce256de25d6e485b Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 27 Sep 2026 12:16:15 +0300 Subject: [PATCH] Bound ETCP receive window and fix timestamp and allocator UB --- doc/etcp_full_suite_findings.md | 64 ++++++++++++++++++++++++----- doc/etcp_protocol.txt | 6 +++ doc/tasks.md | 12 +++--- lib/mem.c | 10 +++-- lib/memory_pool.c | 20 +++++---- src/transport_layer/etcp.c | 34 +++++++++------ tests/test_auto_socket_dynamic.c | 11 ++++- tests/test_etcp_connect.c | 27 ++++++++++-- tests/test_etcp_session.c | 49 ++++++++++++++++++++++ tests/test_memory_pool_and_config.c | 25 ++++++++++- 10 files changed, 211 insertions(+), 47 deletions(-) diff --git a/doc/etcp_full_suite_findings.md b/doc/etcp_full_suite_findings.md index 22323bb4..f84188a3 100644 --- a/doc/etcp_full_suite_findings.md +++ b/doc/etcp_full_suite_findings.md @@ -6,10 +6,10 @@ Промежуточный прогон: 96 passed, 1 failed, 1 skipped; оставался `test_chat_join_e2e`. После KEY_REGISTER_ACK и исправления невыровненных чтений повторный полный `make check -j4`: **97 passed, 0 failed, 1 skipped (98)**. -Лог: `/tmp/utun-ack-final-check.log`. Пропущенный `test_auto_socket_dynamic` -запущен отдельно с root в изолированной сети и упал — подробности ниже. -Внешние `check-proxy/check-burst/check-load` не выполнялись: цепочка root-проверок -остановилась на этом тесте. +Лог: `/tmp/utun-ack-final-check.log`. Исходное падение отдельного root-теста и +последующая доработка описаны ниже. После исправления autosockets общий набор дал +97 PASS / 1 SKIP, а `check-proxy`, `check-burst`, `check-load` прошли: +`/tmp/root-fix-full-check.log`. `test_etcp_lifecycle`, `test_etcp_link_stress`, `test_conn_mgr_phases` дополнительно прошли ASan/UBSan/LeakSanitizer. Инструментированы тесты и изменённые транспортные/ @@ -110,15 +110,59 @@ GUI получает request_id в событии и игнорирует уст В конфиге теста включён auto_sockets, но не отключён стандартный фильтр `auto_socket_skip_no_default_route=1`. Dummy-интерфейсы теста не имеют default route, а `auto_socket.c:reconcile_iface()` исключает такие интерфейсы. Лог показывает -`no links created`, затем TIMEOUT. Для локального сценария тест должен явно задать -`auto_socket_skip_no_default_route=no` и проверить создание сокетов до трафика. +`no links created`, затем TIMEOUT. Исправленный тест создаёт собственный network +namespace, поднимает loopback и добавляет default route только клиентским dummy. +Фильтр остаётся включённым; отдельно проверяется исключение серверного dummy, +иначе его сокет скрывал бы полный обрыв клиентской сети. Кроме того, тест сохраняет handle общего дедлайна, traffic_send и traffic_monitor в одну переменную timeout_handle; фазовые таймеры вообще не сохраняются. При выходе -не отменяются все таймеры; зафиксирован остаток 9 таймеров (665/656). В обработчиках -трафика освобождается entry без queue_dgram_free. Нужна отдельная доработка fixture: -владение таймерами и payload, backpressure отправителя, завершение deferred cleanup. -Производственный auto_socket в рамках исправления приглашений не менялся. +не отменялись все таймеры; зафиксирован остаток 9 таймеров (665/656). В обработчиках +трафика освобождался entry без queue_dgram_free. Теперь у каждого таймера свой slot, +waiters отменяются до teardown, payload освобождается, deferred cleanup завершается +до проверки баланса таймеров, сокетов и `u_get_allocated_count()`. + +Новый генератор шлёт 1024-байтовые пакеты независимо в обе стороны, по одному через +threshold waiter, с бюджетом 64 пакета на tick. Проверяются содержимое, строгий seq, +ограниченность TX/RX очередей, реальное срабатывание backpressure, остановка +генераторов при полном обрыве и полное опустошение после восстановления. NCD handle +удерживает логическую сессию: смена интерфейса не является разрешением терять её данные. +Остатки после действительно удалённого соединения — отдельный допустимый случай. + +Прогон нагрузки обнаружил ошибки, скрытые прежними проверками: + +- `tcp_socket_add` не сохранял `netif_index`: монитор принимал interface-bound сокет + за сокет default route и менял его адрес при событиях другого интерфейса. +- Удаление TCP модели искало её по имени уже удалённого интерфейса и порту из БД. + Теперь запись интерфейса хранит имя и прямую ссылку на модель; один путь закрывает + listener, линки и модель. Портовые записи удаляются с сохранённым именем. +- `tcp_socket_remove` не закрывал линки, созданные вне NCD. Теперь закрывает все + связанные линки до освобождения модели; root-тест создаёт такой линк отдельно. +- Поиск UDP модели не исключал TCP. Добавлена проверка типа сокета. +- При selective ACK sender мог выдавать новые seq за границей `MAX_INFLIGHT_SIZE` + относительно cumulative ACK. Теперь генерация seq останавливается на границе, + cumulative ACK возобновляет очередь даже при отсутствии уже selectively-acked записей. + Тесты проверяют блокировку, возобновление и uint32 wrap. ACK не может подтвердить + ещё не назначенный seq; декодирование старшего байта seq/ACK не делает signed shift. +- Timestamp Q16 умножался как signed int и переполнялся на циклических/отрицательных + смещениях. Перевод выполняется в uint32. Jitter раньше всегда получал нулевую разницу: + предыдущий RTT читался после перезаписи; теперь берётся до неё, EWMA использует int64. +- UBSan выявил невыровненные canary/указатели в `mem.c` и `memory_pool.c`. Метаданные + читаются/пишутся через memcpy; покрыты размеры блоков 1..17 и повторное использование пула. +- Повторный общий прогон поймал `test_etcp_connect`: PID-порты пересекались с занятым + UDP-портом (`EADDRINUSE`). Теперь сервер bind-ится на порт 0, клиент получает фактический + порт через getsockname; собственного алгоритма поиска свободного порта нет. + +Root ASan/UBSan/LSan: PASS, 153335 / 175912 пакетов, reinit=0, все очереди пусты, +баланс памяти/таймеров/сокетов совпал (`/tmp/root-fix-asan.log`). Инструментированы +тест, ETCP/session/router/NCD/autosocket/STCP/normalizer, очереди, event loop и аллокаторы; +прочие объекты архивов обычные. Unit-логи: `/tmp/root-session-asan.log`, +`/tmp/root-memory-asan.log`. LSan запускается вне sandbox: под ptrace он не работает. + +Подробные socket/connection INFO включаются только через `UTUN_TEST_DEBUG=1`. +`ENETUNREACH` непосредственно при принудительном удалении IP/интерфейса ожидаем до +обработки netlink; тест требует последующей доставки всех данных. Ошибки выхода +за receive window недопустимы и после исправления sender gate не наблюдаются. ## Попутная проблема диагностики diff --git a/doc/etcp_protocol.txt b/doc/etcp_protocol.txt index dcaca3c3..e17a458b 100644 --- a/doc/etcp_protocol.txt +++ b/doc/etcp_protocol.txt @@ -41,6 +41,12 @@ inflight - две очереди (ll_queue): - send_hist[8] - история через какие линки отправлялся (индекс линка) - inflight_bytes/inflight_packets per-link отслеживается в ETCP_LINK Размер буфера inflight лимитируем как сумму по всем активным линкам optimal_inflight. +Отдельно ограничиваем размах sequence относительно кумулятивного ACK: новый seq выдаётся +только если (uint32_t)(next_tx_id - rx_ack_till) <= MAX_INFLIGHT_SIZE (16384). +Selective ACK освобождает отдельную inflight-запись, но не продвигает эту границу: +приёмник пока не может освободить пакеты за дыркой. Кумулятивный ACK возобновляет +input_queue даже если подтверждённые записи уже удалены по selective ACK. +ACK ограничивается последним назначенным seq (next_tx_id - 1); сравнения учитывают uint32 wrap. optimal_inflight расчитывается из RTT и bandwidth линка. bandwidth по каждому линку адаптивно подстраиваем: diff --git a/doc/tasks.md b/doc/tasks.md index c6204ca7..c015280d 100644 --- a/doc/tasks.md +++ b/doc/tasks.md @@ -3,11 +3,13 @@ Сюда пишется список задач с короткой аннотацией. Если задача большая и имеет ТЗ - то ТЗ оформляется отдельным файлом, а сюда помещается аннотация и ссылка на ТЗ. ## Текущая задача -[ ] **test_auto_socket_dynamic: восстановить root-регрессию** — тест оставляет включённым - фильтр default route для dummy-интерфейсов без default route и получает TIMEOUT. - Исправить конфигурацию, владение таймерами (общий handle перезаписывается), - освобождение payload, backpressure и очистку fixture. Root-прогон 2026-09-27: - TIMEOUT, 9 оставшихся таймеров. [Разбор](etcp_full_suite_findings.md). +[+] **test_auto_socket_dynamic: root-регрессия и нагрузка** — изолированный namespace, + корректные маршруты, двусторонний поток через backpressure, строгие seq/payload, + пределы TX/RX, остановка при обрыве и drain после восстановления, баланс ресурсов. + Исправлены привязка TCP к интерфейсу, удаление TCP моделей/линков/портов из БД, + sender receive window при selective ACK, timestamp overflow и нулевой jitter. + ASan/UBSan/LSan: PASS; дополнительно исправлены невыровненные метки аллокаторов + и PID-порты в test_etcp_connect. [Разбор и логи](etcp_full_suite_findings.md). [ ] **debug_parse_config: формат и индексы категорий** — parser ожидает `:`, документация описывает `=`; преобразование битовой маски осталось после перехода enum на индексы. diff --git a/lib/mem.c b/lib/mem.c index a100e97a..e149cd50 100644 --- a/lib/mem.c +++ b/lib/mem.c @@ -109,9 +109,10 @@ void u_check(void* ptr, const char* text, const char* location) { error = true; } } - uint32_t* post_canary = (uint32_t*)(user + size); - if (*post_canary != CANARY) { - DEBUG_ERROR(DEBUG_CATEGORY_SYS, "%s: %s: overflow canary corrupted at %p (expected 0x%X, got 0x%X)", location, text, ptr, CANARY, *post_canary); + uint32_t post_canary; + memcpy(&post_canary, user + size, sizeof(post_canary)); + if (post_canary != CANARY) { + DEBUG_ERROR(DEBUG_CATEGORY_SYS, "%s: %s: overflow canary corrupted at %p (expected 0x%X, got 0x%X)", location, text, ptr, CANARY, post_canary); error = true; } if (error) { @@ -141,7 +142,8 @@ void* u_malloc_impl(uint32_t size, const char* location) { uint8_t* prefix_start = base + METADATA_SIZE; *(uint32_t*)(prefix_start + BOUNDARY_CHECK_SIZE) = CANARY; *(uint32_t*)(prefix_start + BOUNDARY_CHECK_SIZE + 4) = size; - *(uint32_t*)(prefix_start + BOUNDARY_CHECK_SIZE + 8 + size) = CANARY; + const uint32_t post_canary = CANARY; /* размер пользователя не обязан быть кратен выравниванию uint32_t */ + memcpy(prefix_start + BOUNDARY_CHECK_SIZE + 8 + size, &post_canary, sizeof(post_canary)); // Add to list (thread-safe) LOCK(); void** prev_ptr = (void**)(base + PREV_OFFSET); diff --git a/lib/memory_pool.c b/lib/memory_pool.c index 75d87a79..298b565f 100644 --- a/lib/memory_pool.c +++ b/lib/memory_pool.c @@ -14,25 +14,27 @@ static void pool_init_tags(struct memory_pool* pool, void* obj, const char* location) { - uint32_t* canary = (uint32_t*)((uint8_t*)obj + POOL_CANARY_OFF); - const char** loc = (const char**)((uint8_t*)obj + POOL_LOC_OFF); + const uint32_t canary = POOL_CANARY_VAL; uint8_t* counter = (uint8_t*)obj + POOL_COUNTER_OFF; - *canary = POOL_CANARY_VAL; - *loc = location; + /* Метаданные следуют за payload произвольного размера и могут быть невыровненными. */ + memcpy((uint8_t*)obj + POOL_CANARY_OFF, &canary, sizeof(canary)); + memcpy((uint8_t*)obj + POOL_LOC_OFF, &location, sizeof(location)); *counter = 1; } static void pool_check_and_clear_tags(struct memory_pool* pool, void* obj, const char* location) { - uint32_t* canary = (uint32_t*)((uint8_t*)obj + POOL_CANARY_OFF); - const char** loc = (const char**)((uint8_t*)obj + POOL_LOC_OFF); + uint32_t canary; + const char* loc; + memcpy(&canary, (uint8_t*)obj + POOL_CANARY_OFF, sizeof(canary)); + memcpy(&loc, (uint8_t*)obj + POOL_LOC_OFF, sizeof(loc)); uint8_t* counter = (uint8_t*)obj + POOL_COUNTER_OFF; - if (*canary != POOL_CANARY_VAL) { + if (canary != POOL_CANARY_VAL) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "pool_free BUFFER OVERFLOW pool=%p name=%s sz=%zu allocs=%zu reuse=%zu free=%d", pool, pool->name ? pool->name : "?", pool->object_size, pool->allocations, pool->reuse_count, pool->free_count); DEBUG_ERROR(DEBUG_CATEGORY_SYS, " obj=%p alloc=%s free=%s canary=0x%08x expected=0x%08x", - obj, *loc ? *loc : "(null)", location, *canary, POOL_CANARY_VAL); + obj, loc ? loc : "(null)", location, canary, POOL_CANARY_VAL); if (pool->object_size) log_dump(DEBUG_LEVEL_ERROR, DEBUG_CATEGORY_SYS, " user_data", obj, pool->object_size); volatile int _halt = 1; @@ -41,7 +43,7 @@ static void pool_check_and_clear_tags(struct memory_pool* pool, void* obj, const if (*counter == 0) { DEBUG_ERROR(DEBUG_CATEGORY_SYS, "pool_free DOUBLE FREE pool=%p name=%s sz=%zu allocs=%zu reuse=%zu obj=%p alloc=%s free=%s", pool, pool->name ? pool->name : "?", pool->object_size, pool->allocations, pool->reuse_count, - obj, *loc ? *loc : "(null)", location); + obj, loc ? loc : "(null)", location); volatile int _halt = 1; while (_halt) {} } diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c index fc35654e..b564c822 100644 --- a/src/transport_layer/etcp.c +++ b/src/transport_layer/etcp.c @@ -927,6 +927,12 @@ void etcp_stats(struct ETCP_CONN* etcp) { static void input_queue_cb(struct ll_queue* q, void* arg) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; + /* Selective ACK освобождает inflight, но не место за дыркой в приёмном окне пира. */ + if ((uint32_t)(etcp->next_tx_id - etcp->rx_ack_till) > MAX_INFLIGHT_SIZE) { + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] peer receive window full: next=%u cumulative_ack=%u queued=%d", + etcp->log_name, etcp->next_tx_id, etcp->rx_ack_till, q->count); + return; + } struct ETCP_FRAGMENT* in_pkt = (struct ETCP_FRAGMENT*)queue_data_get(q); if (!in_pkt) { @@ -1595,8 +1601,9 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { case ETCP_SECTION_ACK: { if (len < ETCP_ACK_BASE_SIZE) { len = 0; break; } int elm_cnt=data[1]; - uint32_t till=data[2] | (data[3]<<8) | (data[4]<<16) | (data[5]<<24); - if ((int32_t)(till - etcp->next_tx_id) > 0) till = etcp->next_tx_id; // wrap-aware: пир не мог подтвердить больше, чем мы отправили + uint32_t till=data[2] | (data[3]<<8) | (data[4]<<16) | ((uint32_t)data[5]<<24); + uint32_t last_tx = etcp->next_tx_id - 1; + if ((int32_t)(till - last_tx) > 0) till = last_tx; // нельзя подтвердить ещё не назначенный seq uint16_t rx_dup_count_16 = data[6] | (data[7]<<8); uint16_t old_tx_dup_16 = etcp->tx_dup_count & 0xFFFF; int16_t diff = (int16_t)(rx_dup_count_16 - old_tx_dup_16); @@ -1610,7 +1617,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { data+=ack_section_len; len-=ack_section_len; for (int i=0; irx_ack_till-till)<0 && ack_gap < ETCP_ACK_GAP_MAX_PACKETS) { etcp->rx_ack_till++; etcp_ack_recv(etcp, etcp->rx_ack_till, -1, -1); ack_gap++; }// подтверждаем всё по till + /* Все записи могли уже уйти по selective ACK; тогда ack_recv не вызовет resume. */ + if (ack_gap) input_queue_try_resume(etcp); break; } case ETCP_SECTION_TIMESTAMP: { @@ -1626,6 +1635,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { uint16_t cur_ts=get_current_timestamp(); uint16_t ret_ts=data[1] | (data[2]<<8);// cur_ts=ret_ts = RTT uint16_t new_rtt=cur_ts-ret_ts; + int prev_rtt = pkt->link->rtt_last; pkt->link->rtt_last=new_rtt; DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] KA-RTT: rtt=%u cur=%u ret=%u dlen=%u", etcp->log_name, new_rtt, cur_ts, ret_ts, pkt->data_len); @@ -1634,19 +1644,21 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { int recv_dt_rx=cur_ts - pkt->link->rtt_last/2 - 1000 - ts; int recv_dt_tx=recv_dt_tx1 - pkt->link->rtt_last/2 - 1000; - if (pkt->link->recv_dt_avg_rx==0) pkt->link->recv_dt_avg_rx=recv_dt_rx*65536; - if (pkt->link->recv_dt_avg_tx==0) pkt->link->recv_dt_avg_tx=recv_dt_tx*65536; - pkt->link->recv_dt_avg_rx +=((int32_t)(recv_dt_rx*65536 - pkt->link->recv_dt_avg_rx))/32; - pkt->link->recv_dt_avg_tx +=((int32_t)(recv_dt_tx*65536 - pkt->link->recv_dt_avg_tx))/32; + /* Timestamp циклический uint16: Q16 также вычисляем по модулю 2^32, без signed overflow. */ + uint32_t sample_rx = (uint32_t)recv_dt_rx * 65536u; + uint32_t sample_tx = (uint32_t)recv_dt_tx * 65536u; + if (pkt->link->recv_dt_avg_rx==0) pkt->link->recv_dt_avg_rx=sample_rx; + if (pkt->link->recv_dt_avg_tx==0) pkt->link->recv_dt_avg_tx=sample_tx; + pkt->link->recv_dt_avg_rx +=((int32_t)(sample_rx - pkt->link->recv_dt_avg_rx))/32; + pkt->link->recv_dt_avg_tx +=((int32_t)(sample_tx - pkt->link->recv_dt_avg_tx))/32; pkt->link->rt_last = cur_ts - ts - pkt->link->recv_dt_avg_rx/65536; pkt->link->tt_last = recv_dt_tx1 - pkt->link->recv_dt_avg_tx/65536; //tts_correction += ((NOW - RTT/2 - TTS) - tts_correction)/16 (инициализируем сразу по 1 пакету) data+=5; len-=5; - int prev_rtt = pkt->link->rtt_last; int d_rtt = new_rtt - prev_rtt; if (d_rtt<0) d_rtt=-d_rtt; - pkt->link->jitter +=((int32_t)(d_rtt*65536 - pkt->link->jitter))/32; + pkt->link->jitter += ((int64_t)d_rtt * 65536 - pkt->link->jitter) / 32; struct ETCP_LINK* c=etcp->links; int rtt_sum=0; @@ -1679,7 +1691,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { if (len>=5) { // формируем ACK - uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24); + uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | ((uint32_t)data[4]<<24); if (etcp->got_initial_pkt == 0) { if (seq==1) { etcp->got_initial_pkt = 1; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Initial packet seq=1 received", etcp->log_name); } else { @@ -1889,5 +1901,3 @@ void etcp_update_mtu(struct ETCP_CONN* etcp) { } } } - - diff --git a/tests/test_auto_socket_dynamic.c b/tests/test_auto_socket_dynamic.c index 1c135f21..046669d6 100644 --- a/tests/test_auto_socket_dynamic.c +++ b/tests/test_auto_socket_dynamic.c @@ -37,7 +37,7 @@ #define ETCP_RT_ID_TEST 0xF0 #define PROBE_SIZE 1024 #define SEND_BUDGET 64 -#define RX_BUDGET_BYTES (16 * 1024 * 1024) /* бюджет памяти сценария при переупорядочении */ +#define RX_BUDGET_BYTES (MAX_INFLIGHT_SIZE * PACKET_DATA_SIZE) /* лимит окна приёма в худшем случае */ typedef void (*timeout_cb)(void*); enum timer_id { TIMER_DEADLINE, TIMER_PHASE, TIMER_SEND, TIMER_MONITOR, TIMER_COUNT }; @@ -57,6 +57,7 @@ static struct test_ctx { struct UASYNC* ua; struct NODE_CONN_DIRECT* handle; struct ETCP_CONN* tcp_probe; + size_t allocations_before; int round; /* 0..7 */ int step; int ip_changes_on_last; @@ -222,9 +223,11 @@ static void check_queues(struct pending_send* p) { if (bytes > p->peak_bytes) p->peak_bytes = bytes; /* Один входной пакет может дать несколько фрагментов плюс flush хвоста. */ struct PKTNORM* pn = c->normalizer; + if ((uint32_t)(c->next_tx_id - c->rx_ack_till) > MAX_INFLIGHT_SIZE + 1u) + fail("sender exceeded peer receive window"); size_t rx_bytes = queue_total_bytes(c->recv_q) + queue_total_bytes(c->output_queue) + queue_total_bytes(pn->output); if (rx_bytes > p->peak_rx_bytes) p->peak_rx_bytes = rx_bytes; - if (rx_bytes > RX_BUDGET_BYTES) { + if (rx_bytes > RX_BUDGET_BYTES || c->recv_q->count > MAX_INFLIGHT_SIZE) { fprintf(stderr, "receive queue overflow dir=%ld bytes=%zu budget=%u\n", (long)(p - pending), rx_bytes, RX_BUDGET_BYTES); fail("receive queues exceeded scenario memory budget"); } @@ -778,6 +781,7 @@ int main(void) { debug_set_category_level(DEBUG_CATEGORY_CONNECTION, DEBUG_LEVEL_INFO); } utun_instance_set_tun_init_enabled(0); + ctx.allocations_before = u_get_allocated_count(); setup(); if (ctx.result) return 1; @@ -821,6 +825,9 @@ done: uasync_destroy(ctx.ua, 0); ctx.ua = NULL; } cleanup(); + if (u_get_allocated_count() != ctx.allocations_before) { + u_report_unfreed_blocks(); fail("allocations leaked after fixture cleanup"); + } fprintf(stderr, "=== %s === timers/sockets cleaned\n", ctx.result == 2 ? "PASS" : "FAIL"); return (ctx.result == 2) ? 0 : 1; } diff --git a/tests/test_etcp_connect.c b/tests/test_etcp_connect.c index a9b34752..79e2b3a4 100644 --- a/tests/test_etcp_connect.c +++ b/tests/test_etcp_connect.c @@ -205,7 +205,7 @@ static void test10(void* arg) { static void setup(void) { test_mkdtemp(tdir); - int base = 48000 + (getpid() % 10000); pa = base; pb = base + 1; + pa = pb = 0; /* ОС атомарно выбирает свободные порты при bind. */ snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir); wf(ca, "[global]\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pa); wf(cb, "[global]\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pb); @@ -214,7 +214,6 @@ static void setup(void) { nid_a = cfa->global.my_node_id; nid_b = cfb->global.my_node_id; free_config(cfa); free_config(cfb); } char *p0 = gv(ca,"pub"), *r0 = gv(ca,"priv"), *p1 = gv(cb,"pub"), *r1 = gv(cb,"priv"); - wf(ca, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", r0, p0, pa, p1, pb); wf(cb, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", r1, p1, pb); u_free(p0); u_free(r0); u_free(p1); u_free(r1); } @@ -224,7 +223,22 @@ int main(void) { debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); utun_instance_set_tun_init_enabled(0); setup(); ua = uasync_create(); - g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb); + g_b = utun_instance_create(ua, cb); + if (!g_b || !g_b->etcp_sockets) goto done; + struct sockaddr_in bound = {0}; socklen_t bound_len = sizeof(bound); + if (getsockname(g_b->etcp_sockets->fd, (struct sockaddr*)&bound, &bound_len) != 0) { + fail("cannot read assigned server port"); goto done; + } + pb = ntohs(bound.sin_port); + if (!pb) { fail("server did not acquire an ephemeral port"); goto done; } + char *p0 = gv(ca, "pub"), *r0 = gv(ca, "priv"), *p1 = gv(cb, "pub"); + if (!p0 || !r0 || !p1) { u_free(p0); u_free(r0); u_free(p1); fail("missing fixture keys"); goto done; } + int rc = wf(ca, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n" + "[server: s1]\naddr=127.0.0.1:0\ntype=public\n[client: to_b]\nkeepalive=1\n" + "peer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", r0, p0, p1, pb); + u_free(p0); u_free(r0); u_free(p1); + if (rc) { fail("cannot write client config"); goto done; } + g_a = utun_instance_create(ua, ca); if (!g_a || !g_b) goto done; utun_instance_init(g_a); utun_instance_init(g_b); uasync_call_soon(ua, NULL, test1); @@ -236,7 +250,12 @@ done: if (ttimer) uasync_cancel_timeout(ua, ttimer); if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); } if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); } - if (ua) { uasync_destroy(ua, 0); ua = NULL; } + if (ua) { + uasync_poll(ua, 0); // завершить deferred cleanup после уничтожения обоих инстансов + if (ua->timer_alloc_count != ua->timer_free_count || ua->socket_alloc_count != ua->socket_free_count) + fail("fixture leaked timers or sockets"); + uasync_destroy(ua, 0); ua = NULL; + } cleanup(); return (result == 2) ? 0 : 1; } diff --git a/tests/test_etcp_session.c b/tests/test_etcp_session.c index da5d3e38..73ee22f9 100644 --- a/tests/test_etcp_session.c +++ b/tests/test_etcp_session.c @@ -110,11 +110,60 @@ static struct frame stream_packet(int data) { CHECK(etcp_encrypt_send(p)>0); memory_pool_free(ep[0].inst.pkt_pool,p); CHECK(tail == head+1); struct frame f=frames[head++ % 4096]; return f; } +static void timestamp_boundaries(void) { + setup(0); + struct ETCP_LINK* link = &ep[0].links[0]; + link->rtt_last = 1000; + for (int i = 0; i < 2; i++) { + struct ETCP_DGRAM* p = memory_pool_alloc(ep[0].inst.pkt_pool); CHECK(p); + p->link = link; p->timestamp = UINT16_MAX; p->data_len = 5; p->data[0] = ETCP_SECTION_TIMESTAMP; + uint16_t ret = get_current_timestamp() - (i ? 30000 : 10); + p->data[1] = ret; p->data[2] = ret >> 8; + p->data[3] = p->data[4] = i ? 0 : 255; + ep[0].c->last_rtt_cb_time = get_time_tb(); // fixture не содержит topology + etcp_conn_input(p); + CHECK(link->jitter > 0); // предыдущий RTT нельзя перезаписывать до вычисления разницы + CHECK(link->jitter <= (uint64_t)UINT16_MAX * 65536); + CHECK(link->recv_dt_avg_tx != 0 && link->recv_dt_avg_rx != 0); + } + teardown(); + puts("PASS timestamp wrap, negative offsets and changing RTT"); +} +static void receive_window_case(int wrap) { + setup(0); + struct ETCP_CONN* c = ep[0].c; + ep[0].inst.data_pool = memory_pool_init(PACKET_DATA_SIZE, "window payload"); CHECK(ep[0].inst.data_pool); + queue_set_callback(c->input_send_q, NULL, NULL); // выделенные seq наблюдаем до физической отправки + c->next_tx_id = wrap ? 5 : MAX_INFLIGHT_SIZE + 1; + c->rx_ack_till = c->next_tx_id - MAX_INFLIGHT_SIZE - 1; + uint32_t next = c->next_tx_id, till = c->rx_ack_till; + struct ETCP_FRAGMENT* f = (struct ETCP_FRAGMENT*)queue_entry_new(sizeof(*f) - sizeof(struct ll_entry)); CHECK(f); + f->ll.dgram = memory_pool_alloc(ep[0].inst.data_pool); CHECK(f->ll.dgram); + f->ll.dgram_pool = ep[0].inst.data_pool; f->ll.len = 1; + CHECK(queue_data_put(c->input_queue, &f->ll) == 0); + CHECK(c->input_queue->count == 1 && c->next_tx_id == next); + for (int cumulative = 0; cumulative < 2; cumulative++) { + struct ETCP_DGRAM* p = memory_pool_alloc(ep[0].inst.pkt_pool); CHECK(p); + p->link = &ep[0].links[0]; p->data_len = cumulative ? 8 : 16; + memset(p->data, 0, p->data_len); p->data[0] = ETCP_SECTION_ACK; p->data[1] = !cumulative; + uint32_t ack = till + cumulative, selective = next - 1; + for (int i = 0; i < 4; i++) { p->data[2+i] = ack >> (8*i); p->data[8+i] = selective >> (8*i); } + etcp_conn_input(p); uasync_poll(ua, 0); + if (!cumulative) CHECK(c->input_queue->count == 1 && c->next_tx_id == next); + else CHECK(c->input_queue->count == 0 && c->next_tx_id == next + 1 && c->input_send_q->count == 1); + } + teardown(); + memory_pool_destroy(ep[0].inst.data_pool); ep[0].inst.data_pool = NULL; + printf("PASS receive-window backpressure and cumulative ACK resume (wrap=%d)\n", wrap); +} int main(void) { debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); const char* log=getenv("SESSION_TEST_LOG"); if (log) { debug_enable_file_output(log,1); debug_set_category_level(DEBUG_CATEGORY_ETCP,DEBUG_LEVEL_DEBUG); } size_t baseline=u_get_allocated_count(); + timestamp_boundaries(); CHECK(u_get_allocated_count() == baseline); + receive_window_case(0); CHECK(u_get_allocated_count() == baseline); + receive_window_case(1); CHECK(u_get_allocated_count() == baseline); for (int scenario=0;scenario<9;scenario++) { setup(scenario == 5); saved[2]=stream_packet(0); saved[3]=stream_packet(1); diff --git a/tests/test_memory_pool_and_config.c b/tests/test_memory_pool_and_config.c index 4668634b..e90cfecb 100644 --- a/tests/test_memory_pool_and_config.c +++ b/tests/test_memory_pool_and_config.c @@ -13,6 +13,7 @@ #include "../lib/memory_pool.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" +#include "../lib/mem.h" static int test_callback_count = 0; @@ -27,6 +28,28 @@ int main() { debug_set_categories(DEBUG_CATEGORY_ALL); DEBUG_INFO(DEBUG_CATEGORY_SYS, "=== Memory Pool and Config File Test ==="); + /* Любой размер блока, включая невыровненный хвост, допустим для malloc/realloc/free. */ + for (uint32_t size = 1; size <= 17; size++) { + unsigned char* p = u_malloc(size); + assert(p); + memset(p, 0x5a, size); + u_check(p, "odd-size allocation", __FILE__); + p = u_realloc(p, size + 3); + assert(p); + for (uint32_t i = 0; i < size; i++) assert(p[i] == 0x5a); + u_free(p); + struct memory_pool* pool = memory_pool_init(sizeof(void*) + size, "odd-size pool"); + assert(pool); + void* obj = memory_pool_alloc(pool); + assert(obj); + memset(obj, 0x5a, pool->object_size); + memory_pool_free(pool, obj); + obj = memory_pool_alloc(pool); + assert(obj); + memory_pool_free(pool, obj); + memory_pool_destroy(pool); + } + DEBUG_INFO(DEBUG_CATEGORY_SYS, "Unaligned tail canary/realloc: PASS"); // Test 1: Memory pool optimization DEBUG_INFO(DEBUG_CATEGORY_SYS, "Test 1: Memory pool optimization..."); @@ -120,4 +143,4 @@ int main() { DEBUG_INFO(DEBUG_CATEGORY_SYS, "=== All tests completed successfully ==="); return 0; -} \ No newline at end of file +}