Browse Source

Bound ETCP receive window and fix timestamp and allocator UB

proxy
evgeny 4 days ago
parent
commit
03f71fd94e
  1. 64
      doc/etcp_full_suite_findings.md
  2. 6
      doc/etcp_protocol.txt
  3. 12
      doc/tasks.md
  4. 10
      lib/mem.c
  5. 20
      lib/memory_pool.c
  6. 34
      src/transport_layer/etcp.c
  7. 11
      tests/test_auto_socket_dynamic.c
  8. 27
      tests/test_etcp_connect.c
  9. 49
      tests/test_etcp_session.c
  10. 23
      tests/test_memory_pool_and_config.c

64
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 не наблюдаются.
## Попутная проблема диагностики

6
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 по каждому линку адаптивно подстраиваем:

12
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 на индексы.

10
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);

20
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) {}
}

34
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; i<elm_cnt; i++) {
uint32_t seq=data[-ack_section_len+8+i*8] | (data[-ack_section_len+9+i*8]<<8) | (data[-ack_section_len+10+i*8]<<16) | (data[-ack_section_len+11+i*8]<<24);
uint32_t seq=data[-ack_section_len+8+i*8] | (data[-ack_section_len+9+i*8]<<8) | (data[-ack_section_len+10+i*8]<<16) | ((uint32_t)data[-ack_section_len+11+i*8]<<24);
uint16_t ts=data[-ack_section_len+12+i*8] | (data[-ack_section_len+13+i*8]<<8);
uint16_t dts=data[-ack_section_len+14+i*8] | (data[-ack_section_len+15+i*8]<<8);
etcp_ack_recv(etcp, seq, ts, dts);
@ -1619,6 +1626,8 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
while ((int32_t)(etcp->rx_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) {
}
}
}

11
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;
}

27
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;
}

49
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);

23
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...");

Loading…
Cancel
Save