From add785c8a1e57939d465b47687274fcf6b2b3183 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sun, 1 Feb 2026 22:31:25 +0300 Subject: [PATCH] Test: Added test_pkt_normalizer_etcp for testing pkt_normalizer with ETCP - Created test_pkt_normalizer_etcp.c based on test_etcp_100_packets - Tests bidirectional transfer of 100 packets (10-10000 bytes) via normalizer - Fixed memory management bugs in pkt_normalizer.c: * Fixed double-free in pn_buf_renew() * Added pn_send_to_etcp() to properly create ETCP_FRAGMENT * Fixed memory freeing in pn_unpacker_cb() - Added test to Makefile.am --- src/pkt_normalizer.c | 410 ++++++++++++++----------------- src/pkt_normalizer.h | 38 +-- tests/Makefile.am | 5 + tests/test_pkt_normalizer_etcp.c | 409 ++++++++++++++++++++++++++++++ 4 files changed, 607 insertions(+), 255 deletions(-) create mode 100644 tests/test_pkt_normalizer_etcp.c diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index 9b5bd24a..c15bb685 100644 --- a/src/pkt_normalizer.c +++ b/src/pkt_normalizer.c @@ -6,6 +6,7 @@ #include #include #include // For debugging (can be removed if not needed) +#include "debug_config.h" // Assuming this for DEBUG_ERROR // Internal helper to convert void* data to struct ll_entry* static inline struct ll_entry* data_to_entry(void* data) { @@ -15,8 +16,10 @@ static inline struct ll_entry* data_to_entry(void* data) { // Forward declarations static void packer_cb(struct ll_queue* q, void* arg); -static void flush_cb(void* arg); -static void input_ready_cb(struct ll_queue* q, void* arg); +static void pn_flush_cb(void* arg); +static void etcp_input_ready_cb(struct ll_queue* q, void* arg); +static void pn_unpacker_cb(struct ll_queue* q, void* arg); +static void pn_send_to_etcp(struct PKTNORM* pn, struct ll_entry* entry); // Initialization struct PKTNORM* pn_init(struct ETCP_CONN* etcp) { @@ -27,31 +30,22 @@ struct PKTNORM* pn_init(struct ETCP_CONN* etcp) { pn->etcp = etcp; pn->ua = etcp->instance->ua; - pn->pkt_size = etcp->mtu; // Use MTU as fixed packet size (adjust if headers need subtraction) + pn->frag_size = etcp->mtu - 100; // Use MTU as fixed packet size (adjust if headers need subtraction) + pn->tx_wait_time = 10; pn->input = queue_new(pn->ua, 0); // No hash needed pn->output = queue_new(pn->ua, 0); // No hash needed - pn->pending = queue_new(pn->ua, 0); - pn->send_pending = queue_new(pn->ua, 0); - if (!pn->input || !pn->output || !pn->pending || !pn->send_pending) { + if (!pn->input || !pn->output) { pn_pair_deinit(pn); return NULL; } - pn->in_buf = calloc(1, pn->pkt_size); - pn->cap = pn->pkt_size; // For packer buffer (fixed size) - pn->len = 0; // Current length in packer buffer - - pn->out_buf = calloc(1, pn->pkt_size * 10); // Arbitrary large buffer for assembly - pn->cap = pn->pkt_size * 10; // Initial capacity for unpacker (can reallocate if needed) - pn->len = 0; - pn->total_len = 0; - pn->in_fragment = 0; - - // Set callback for automatic processing when items are added to input queue_set_callback(pn->input, packer_cb, pn); + queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn); + pn->sndpart = NULL; + pn->recvpart = NULL; pn->flush_timer = NULL; return pn; @@ -65,6 +59,10 @@ void pn_pair_deinit(struct PKTNORM* pn) { if (pn->input) { void* data; while ((data = queue_data_get(pn->input)) != NULL) { + struct ll_entry* entry = data_to_entry(data); + if (entry->dgram) { + free(entry->dgram); + } queue_data_free(data); } queue_free(pn->input); @@ -72,284 +70,240 @@ void pn_pair_deinit(struct PKTNORM* pn) { if (pn->output) { void* data; while ((data = queue_data_get(pn->output)) != NULL) { + struct ll_entry* entry = data_to_entry(data); + if (entry->dgram) { + free(entry->dgram); + } queue_data_free(data); } queue_free(pn->output); } - if (pn->pending) { - void* data; - while ((data = queue_data_get(pn->pending)) != NULL) { - queue_data_free(data); - } - queue_free(pn->pending); - } - if (pn->send_pending) { - void* data; - while ((data = queue_data_get(pn->send_pending)) != NULL) { - struct ETCP_FRAGMENT* frag = data_to_entry(data); - if (frag->ll.dgram) memory_pool_free(pn->etcp->instance->data_pool, frag->ll.dgram); - memory_pool_free(pn->etcp->rx_pool, frag); - } - queue_free(pn->send_pending); - } if (pn->flush_timer) { uasync_cancel_timeout(pn->ua, pn->flush_timer); } - free(pn->in_buf); - free(pn->out_buf); + if (pn->sndpart) { + ll_free_dgram(pn->sndpart); + queue_data_free(pn->sndpart); + } + if (pn->recvpart) { + ll_free_dgram(pn->recvpart); + queue_data_free(pn->recvpart); + } + free(pn); } // Reset unpacker state void pn_unpacker_reset_state(struct PKTNORM* pn) { if (!pn) return; - pn->len = 0; - pn->total_len = 0; - pn->in_fragment = 0; + if (pn->recvpart) { + ll_free_dgram(pn->recvpart); + queue_data_free(pn->recvpart); + pn->recvpart = NULL; + } } // Send data to packer (copies and adds to input queue or pending, triggering callback) void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) { if (!pn || !data || len == 0) return; - void* entry_data = queue_data_new(len); - if (!entry_data) return; - - struct ll_entry* entry = data_to_entry(entry_data); - memcpy(entry->data, data, len); + + struct ll_entry* entry = ll_alloc_lldgram(len); + memcpy(entry->dgram, data, len); entry->len = len; - entry->size = len; - - // To minimize delay, add to input only if empty, else to pending and set waiter if not set - if (queue_entry_count(pn->input) == 0 && pn->input->waiter.callback == NULL) { - queue_data_put(pn->input, entry_data, 0); - } else { - queue_data_put(pn->pending, entry_data, 0); - if (pn->input->waiter.callback == NULL) { - queue_wait_threshold(pn->input, 0, 0, input_ready_cb, pn); - } - } -} + entry->dgram_pool = NULL; -// Internal: Callback when input queue becomes empty, move next from pending -static void input_ready_cb(struct ll_queue* q, void* arg) { - struct PKTNORM* pn = (struct PKTNORM*)arg; + queue_data_put(pn->input, entry, 0); - void* data = queue_data_get(pn->pending); - if (data) { - queue_data_put(q, data, 0); // This will trigger packer_cb async if needed - if (queue_entry_count(pn->pending) > 0) { - queue_wait_threshold(q, 0, 0, input_ready_cb, pn); - } + // Cancel flush timer if active + if (pn->flush_timer) { + uasync_cancel_timeout(pn->ua, pn->flush_timer); + pn->flush_timer = NULL; } } -static void etcp_input_ready_cb(struct ll_queue* q, void* arg) { +// Internal: Packer callback +static void packer_cb(struct ll_queue* q, void* arg) { struct PKTNORM* pn = (struct PKTNORM*)arg; + if (!pn) return; - void* data = queue_data_get(pn->send_pending); - if (data) { - queue_data_put(q, data, 0); - if (queue_entry_count(pn->send_pending) > 0) { - queue_wait_threshold(q, 0, 0, etcp_input_ready_cb, pn); - } - } + queue_wait_threshold(pn->etcp->input_queue, 0, 0, etcp_input_ready_cb, pn); } -// Internal: Send a fixed-size chunk to etcp->input_queue (or send_pending) -static void send_chunk(struct PKTNORM* pn, uint8_t* chunk_data, uint16_t chunk_len) { - // Alloc ETCP_FRAGMENT for the chunk (as payload for etcp) +// Helper to send sndpart to ETCP as ETCP_FRAGMENT +static void pn_send_to_etcp(struct PKTNORM* pn, struct ll_entry* entry) { + if (!pn || !entry || entry->len == 0) return; + + // Allocate data + uint8_t* packet_data = malloc(entry->len); + if (!packet_data) return; + + memcpy(packet_data, entry->dgram, entry->len); + + // Allocate ETCP_FRAGMENT from rx_pool struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool); - if (!frag) return; - - frag->ll.dgram = memory_pool_alloc(pn->etcp->instance->data_pool); - if (!frag->ll.dgram) { - memory_pool_free(pn->etcp->rx_pool, frag); + if (!frag) { + free(packet_data); return; } - - memcpy(frag->ll.dgram, chunk_data, chunk_len); - frag->ll.size = chunk_len; - frag->ll.len = chunk_len; - frag->seq = 0; // Not used for input + + frag->seq = 0; frag->timestamp = 0; + frag->ll.dgram = packet_data; + frag->ll.size = entry->len; + frag->ll.len = entry->len; + frag->ll.memlen = entry->len; + frag->ll.dgram_pool = NULL; + + queue_data_put(pn->etcp->input_queue, frag, 0); +} - struct ll_queue* eq = pn->etcp->input_queue; - - // Add to etcp input_queue only if empty, else to send_pending and set waiter - if (queue_entry_count(eq) == 0 && eq->waiter.callback == NULL) { - queue_data_put(eq, frag, 0); - } else { - queue_data_put(pn->send_pending, frag, 0); - if (eq->waiter.callback == NULL) { - queue_wait_threshold(eq, 0, 0, etcp_input_ready_cb, pn); +// Internal: Renew sndpart buffer +static void pn_buf_renew(struct PKTNORM* pn) { + if (pn->sndpart) { + int remain = pn->frag_size - pn->sndpart->len; + if (remain < 3) { + if (pn->sndpart->len > 0) { + // Transfer ownership to ETCP queue + pn_send_to_etcp(pn, pn->sndpart); + } + // Free the ll_entry (but not dgram - it's copied in pn_send_to_etcp) + ll_free_dgram(pn->sndpart); + queue_data_free(pn->sndpart); + pn->sndpart = NULL; + } + } + if (!pn->sndpart) { + pn->sndpart = ll_alloc_lldgram(pn->frag_size); + if (pn->sndpart) { + pn->sndpart->len = 0; } } } -// Internal: Packer callback (aggregates small, fragments large, sends chunks) -static void packer_cb(struct ll_queue* q, void* arg) { +// Internal: Process input when etcp->input_queue is ready (empty) +static void etcp_input_ready_cb(struct ll_queue* q, void* arg) { struct PKTNORM* pn = (struct PKTNORM*)arg; if (!pn) return; - pn->len = 0; // Start with empty buffer for this batch + void* data = queue_data_get(pn->input); + if (!data) { + queue_resume_callback(pn->input); + return; + } - while (q->head) { - void* data = queue_data_get(q); - struct ll_entry* entry = data_to_entry(data); - uint16_t item_len = entry->len; + struct ll_entry* in_dgram = data_to_entry(data); + uint16_t ptr = 0; - if (item_len + 2 > pn->pkt_size) { - // Large item: fragment (must be alone in buffer) - if (pn->len > 0) { - // Flush current aggregated small items first - memset(pn->in_buf + pn->len, 0, pn->pkt_size - pn->len); // Pad - send_chunk(pn, pn->in_buf, pn->pkt_size); - pn->len = 0; - } + while (ptr < in_dgram->len) { + pn_buf_renew(pn); + if (!pn->sndpart) break; // Allocation failed - // Fragment the large item - uint8_t* ptr = entry->data; - uint32_t remaining = item_len; - - // First fragment - *(uint16_t*)(pn->in_buf) = 0xffff; - *(uint32_t*)(pn->in_buf + 2) = remaining; - uint16_t chunk = pn->pkt_size - 6; - memcpy(pn->in_buf + 6, ptr, chunk); - send_chunk(pn, pn->in_buf, pn->pkt_size); - ptr += chunk; - remaining -= chunk; - - // Continuation fragments - while (remaining > 0) { - chunk = (remaining < pn->pkt_size - 2) ? remaining : pn->pkt_size - 2; - *(uint16_t*)(pn->in_buf) = 0x0000; // Continuation marker - memcpy(pn->in_buf + 2, ptr, chunk); - memset(pn->in_buf + 2 + chunk, 0, pn->pkt_size - 2 - chunk); // Pad - send_chunk(pn, pn->in_buf, pn->pkt_size); - ptr += chunk; - remaining -= chunk; - } - } else { - // Small item: try to aggregate - if (pn->len + item_len + 2 > pn->pkt_size) { - // Flush current buffer - memset(pn->in_buf + pn->len, 0, pn->pkt_size - pn->len); // Pad - send_chunk(pn, pn->in_buf, pn->pkt_size); - pn->len = 0; - } + int remain = pn->frag_size - pn->sndpart->len; + if (remain < 3) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_normalizer: part size error, remain=%d", remain); + break; + } - // Add to buffer - *(uint16_t*)(pn->in_buf + pn->len) = item_len; - memcpy(pn->in_buf + pn->len + 2, entry->data, item_len); - pn->len += item_len + 2; + if (ptr == 0) { + pn->sndpart->dgram[pn->sndpart->len++] = in_dgram->len & 0xFF; + pn->sndpart->dgram[pn->sndpart->len++] = (in_dgram->len >> 8) & 0xFF; + remain -= 2; } - queue_data_free(data); // Free the entry + int n = remain; + int rem = in_dgram->len - ptr; + if (n > rem) n = rem; + memcpy(pn->sndpart->dgram + pn->sndpart->len, in_dgram->dgram + ptr, n); + pn->sndpart->len += n; + ptr += n; + + pn_buf_renew(pn); } - // If remaining in buffer, set timeout to flush if queues empty - if (pn->len > 0) { - if (pn->flush_timer) { - uasync_cancel_timeout(pn->ua, pn->flush_timer); - } - pn->flush_timer = uasync_set_timeout(pn->ua, 10, pn, flush_cb); // 1ms = 10 * 0.1ms + ll_free_dgram(in_dgram); + queue_data_free(data); + + // Cancel flush timer if active + if (pn->flush_timer) { + uasync_cancel_timeout(pn->ua, pn->flush_timer); + pn->flush_timer = NULL; } - queue_resume_callback(q); + // Set flush timer if no more input + if (queue_entry_count(pn->input) == 0) { + pn->flush_timer = uasync_set_timeout(pn->ua, pn->tx_wait_time, pn, pn_flush_cb); + } + + queue_resume_callback(pn->input); } // Internal: Flush callback on timeout -static void flush_cb(void* arg) { +static void pn_flush_cb(void* arg) { struct PKTNORM* pn = (struct PKTNORM*)arg; if (!pn) return; pn->flush_timer = NULL; - - if (pn->len > 0 && queue_entry_count(pn->input) == 0 && queue_entry_count(pn->pending) == 0) { - memset(pn->in_buf + pn->len, 0, pn->pkt_size - pn->len); // Pad - send_chunk(pn, pn->in_buf, pn->pkt_size); - pn->len = 0; - } -} - -// Internal: Add assembled data to etcp->output_queue as ETCP_FRAGMENT -static void add_to_output(struct PKTNORM* pn, uint8_t* app_data, uint16_t app_len) { - struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool); - if (!frag) return; - - frag->ll.dgram = memory_pool_alloc(pn->etcp->instance->data_pool); - if (!frag->ll.dgram) { - memory_pool_free(pn->etcp->rx_pool, frag); - return; + if (pn->sndpart && pn->sndpart->len > 0) { + pn_send_to_etcp(pn, pn->sndpart); + // Free the ll_entry after sending + ll_free_dgram(pn->sndpart); + queue_data_free(pn->sndpart); + pn->sndpart = NULL; // Will alloc new when needed } - - memcpy(frag->ll.dgram, app_data, app_len); - frag->ll.size = app_len; - frag->ll.len = app_len; - frag->seq = 0; - frag->timestamp = 0; - - queue_data_put(pn->etcp->output_queue, frag, 0); } -// Internal: Process incoming fixed-size payload chunk from etcp -void pn_unpacker_input(struct PKTNORM* pn, uint8_t* data, uint16_t len) { - if (!pn || len != pn->pkt_size) return; - - uint16_t header = *(uint16_t*)data; +// Internal: Unpacker callback (assembles fragments into original packets) +static void pn_unpacker_cb(struct ll_queue* q, void* arg) { + struct PKTNORM* pn = (struct PKTNORM*)arg; + if (!pn) return; - if (pn->in_fragment) { - // Expect continuation - if (header != 0x0000) { - // Error: invalid continuation - pn_unpacker_reset_state(pn); - return; - } + while (1) { + void* data = queue_data_get(pn->etcp->output_queue); + if (!data) break; + + struct ETCP_FRAGMENT* frag = (struct ETCP_FRAGMENT*)data; // Since data is ll_entry* + uint8_t* payload = frag->ll.dgram; + uint16_t len = frag->ll.len; + uint16_t ptr = 0; + + while (ptr < len) { + if (!pn->recvpart) { + // Need length header for new packet + if (len - ptr < 2) { + // Incomplete header, reset + pn_unpacker_reset_state(pn); + break; + } + uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8); + ptr += 2; + + pn->recvpart = ll_alloc_lldgram(pkt_len); + if (!pn->recvpart) { + break; + } + pn->recvpart->len = 0; + } - uint16_t chunk = pn->pkt_size - 2; - if (pn->len + chunk > pn->cap) { - pn->cap *= 2; - pn->out_buf = realloc(pn->out_buf, pn->cap); - } - memcpy(pn->out_buf + pn->len, data + 2, chunk); - pn->len += chunk; + uint16_t rem = pn->recvpart->memlen - pn->recvpart->len; + uint16_t avail = len - ptr; + uint16_t cp = (rem < avail) ? rem : avail; + memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp); + pn->recvpart->len += cp; + ptr += cp; - if (pn->len >= pn->total_len) { - // Assembly complete: add to output_queue - add_to_output(pn, pn->out_buf, pn->total_len); - pn_unpacker_reset_state(pn); - } - } else { - if (header == 0xffff) { - // Start of fragment - pn->total_len = *(uint32_t*)(data + 2); - uint16_t chunk = pn->pkt_size - 6; - if (pn->total_len <= chunk) { - // Invalid or degenerate case - pn_unpacker_reset_state(pn); - return; - } - memcpy(pn->out_buf, data + 6, chunk); - pn->len = chunk; - pn->in_fragment = 1; - } else { - // Normal aggregated packet: process multiple items (ignore pad) - uint8_t* ptr = data; - size_t processed = 0; - while (processed < pn->pkt_size) { - header = *(uint16_t*)ptr; - if (header == 0 || header == 0xffff || header == 0x0000) break; // End (pad) or error - if (processed + 2 + header > pn->pkt_size) break; - - add_to_output(pn, ptr + 2, header); - - ptr += 2 + header; - processed += 2 + header; + if (pn->recvpart->len == pn->recvpart->memlen) { + queue_data_put(pn->output, pn->recvpart, 0); + pn->recvpart = NULL; } } + + // Free the fragment - dgram was malloc'd in pn_send_to_etcp + free(frag->ll.dgram); + memory_pool_free(pn->etcp->rx_pool, frag); } + + queue_resume_callback(q); } diff --git a/src/pkt_normalizer.h b/src/pkt_normalizer.h index 54a1a9f8..7b5552f8 100644 --- a/src/pkt_normalizer.h +++ b/src/pkt_normalizer.h @@ -8,31 +8,24 @@ // Структура для packer struct PKTNORM { -// public: + // public: struct ll_queue* input; // Входная очередь в packer (через нее отправляем пакеты) struct ll_queue* output; // Выходная очередь из unpacker (через нее принимаем пакеты) + uint16_t tx_wait_time; -// private: + // private: + uasync_t* ua; // uasync instance struct ETCP_CONN* etcp; - uint16_t pkt_size; - -// packer: - uasync_t* ua; // uasync instance - uint8_t* in_buf; // Буфер для упаковки - size_t len; // Текущая длина в буфере - size_t cap; // Емкость буфера (MAX_PACKET_SIZE) + // packer: + uint16_t frag_size; // размер фрагмента на которые разбивать + struct ll_entry* sndpart; // блок ожидающий досборки void* flush_timer; // For timeout flush - struct ll_queue* pending; // For not filling input - -// unpacker: - uint8_t* out_buf; // Буфер для сборки фрагментов -// size_t len; // Текущая накопленная длина - size_t total_len; // Ожидаемая общая длина (для фрагментов) -// size_t cap; // Емкость буфера - int in_fragment; // Флаг: идет сборка фрагмента (1) или нет (0) - struct ll_queue* send_pending; // For not filling etcp->input_queue + struct ll_entry* pending; // Partial processed input entry + uint16_t pending_in_ptr; // Pointer in pending entry + // unpacker: + struct ll_entry* recvpart; // блок ожидающий заполнение }; // Инициализация пары @@ -47,13 +40,4 @@ void pn_unpacker_reset_state(struct PKTNORM* pn); // создаёт malloc data, копирует, помещает в input. void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len); -/* Как работает: -Формат отправки в etcp: 2 байта размер, далее данные (порезанные на куски и отправленные через etcp) -собирает по возможности полные пакеты с размером pkt_size. -+ timeout: неполный пакет отправляется по таймауту 1ms (если очереди пустые) -+ входящая очередь etcp не должне наполняться для минимизации задержки - новый пакет отправляем только когда очередь пустая - -*/ - - #endif // PKT_NORMALIZER_H \ No newline at end of file diff --git a/tests/Makefile.am b/tests/Makefile.am index 3445072b..f94e08f8 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -7,6 +7,7 @@ check_PROGRAMS = test_etcp_crypto$(EXEEXT) \ test_etcp_simple_traffic$(EXEEXT) \ test_etcp_minimal$(EXEEXT) \ test_etcp_100_packets$(EXEEXT) \ + test_pkt_normalizer_etcp$(EXEEXT) \ test_ll_queue$(EXEEXT) \ test_ecc_encrypt$(EXEEXT) \ test_intensive_memory_pool$(EXEEXT) \ @@ -41,6 +42,10 @@ test_etcp_100_packets_SOURCES = test_etcp_100_packets.c test_etcp_100_packets_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source test_etcp_100_packets_LDADD = $(top_builddir)/src/utun-config_parser.o $(top_builddir)/src/utun-config_updater.o $(top_builddir)/src/utun-crc32.o $(top_builddir)/src/utun-etcp.o $(top_builddir)/src/utun-etcp_connections.o $(top_builddir)/src/utun-etcp_loadbalancer.o $(top_builddir)/src/utun-secure_channel.o $(top_builddir)/src/utun-routing.o $(top_builddir)/src/utun-tun_if.o $(top_builddir)/src/utun-utun_instance.o $(top_builddir)/src/utun-pkt_normalizer.o $(top_builddir)/tinycrypt/lib/source/utun-aes_encrypt.o $(top_builddir)/tinycrypt/lib/source/utun-aes_decrypt.o $(top_builddir)/tinycrypt/lib/source/utun-ccm_mode.o $(top_builddir)/tinycrypt/lib/source/utun-cmac_mode.o $(top_builddir)/tinycrypt/lib/source/utun-ctr_mode.o $(top_builddir)/tinycrypt/lib/source/utun-ecc.o $(top_builddir)/tinycrypt/lib/source/utun-ecc_dh.o $(top_builddir)/tinycrypt/lib/source/utun-ecc_dsa.o $(top_builddir)/tinycrypt/lib/source/utun-sha256.o $(top_builddir)/tinycrypt/lib/source/utun-ecc_platform_specific.o $(top_builddir)/tinycrypt/lib/source/utun-utils.o $(top_builddir)/lib/libuasync.a -lpthread -lcrypto +test_pkt_normalizer_etcp_SOURCES = test_pkt_normalizer_etcp.c +test_pkt_normalizer_etcp_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source +test_pkt_normalizer_etcp_LDADD = $(top_builddir)/src/utun-config_parser.o $(top_builddir)/src/utun-config_updater.o $(top_builddir)/src/utun-crc32.o $(top_builddir)/src/utun-etcp.o $(top_builddir)/src/utun-etcp_connections.o $(top_builddir)/src/utun-etcp_loadbalancer.o $(top_builddir)/src/utun-secure_channel.o $(top_builddir)/src/utun-routing.o $(top_builddir)/src/utun-tun_if.o $(top_builddir)/src/utun-utun_instance.o $(top_builddir)/src/utun-pkt_normalizer.o $(top_builddir)/tinycrypt/lib/source/utun-aes_encrypt.o $(top_builddir)/tinycrypt/lib/source/utun-aes_decrypt.o $(top_builddir)/tinycrypt/lib/source/utun-ccm_mode.o $(top_builddir)/tinycrypt/lib/source/utun-cmac_mode.o $(top_builddir)/tinycrypt/lib/source/utun-ctr_mode.o $(top_builddir)/tinycrypt/lib/source/utun-ecc.o $(top_builddir)/tinycrypt/lib/source/utun-ecc_dh.o $(top_builddir)/tinycrypt/lib/source/utun-ecc_dsa.o $(top_builddir)/tinycrypt/lib/source/utun-sha256.o $(top_builddir)/tinycrypt/lib/source/utun-ecc_platform_specific.o $(top_builddir)/tinycrypt/lib/source/utun-utils.o $(top_builddir)/lib/libuasync.a -lpthread -lcrypto + # Basic crypto test test_crypto_SOURCES = test_crypto.c test_crypto_CFLAGS = -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source -I$(top_srcdir)/lib diff --git a/tests/test_pkt_normalizer_etcp.c b/tests/test_pkt_normalizer_etcp.c new file mode 100644 index 00000000..db2e2cdc --- /dev/null +++ b/tests/test_pkt_normalizer_etcp.c @@ -0,0 +1,409 @@ +#include +#include +#include +#include +#include + +#include "../src/etcp.h" +#include "../src/etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/utun_instance.h" +#include "../src/routing.h" +#include "../src/tun_if.h" +#include "../src/secure_channel.h" +#include "../src/pkt_normalizer.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" + +#define TEST_TIMEOUT_MS 30000 // 30 second timeout +#define TOTAL_PACKETS 100 // Total packets to send +#define MAX_QUEUE_SIZE 5 // Max packets in input queue +#define MIN_PACKET_SIZE 10 // Minimum packet size +#define MAX_TEST_PACKET_SIZE 10000 // Maximum packet size (10KB) - renamed to avoid conflict + +static struct UTUN_INSTANCE* server_instance = NULL; +static struct UTUN_INSTANCE* client_instance = NULL; +static struct PKTNORM* server_pn = NULL; +static struct PKTNORM* client_pn = NULL; +static int test_completed = 0; +static void* packet_timeout_id = NULL; + +// Test statistics - forward direction (client -> server) +static int packets_sent_fwd = 0; +static int packets_received_fwd = 0; +static int current_packet_seq_fwd = 0; + +// Test statistics - backward direction (server -> client) +static int packets_sent_back = 0; +static int packets_received_back = 0; +static int current_packet_seq_back = 0; + +// Packet sizes for each packet (random) +static int packet_sizes[TOTAL_PACKETS]; + +// Timing variables +static struct timespec start_time_fwd, end_time_fwd; +static struct timespec start_time_back, end_time_back; +static int phase = 0; // 0 = connecting, 1 = forward transfer, 2 = backward transfer + +// Function to generate packet data with CRC-like pattern +static void generate_packet_data(int seq, uint8_t* buffer, int size) { + buffer[0] = (uint8_t)(seq & 0xFF); + buffer[1] = (uint8_t)((seq >> 8) & 0xFF); + buffer[2] = (uint8_t)(size & 0xFF); + buffer[3] = (uint8_t)((size >> 8) & 0xFF); + + // Fill rest with pattern based on sequence and position + for (int i = 4; i < size; i++) { + buffer[i] = (uint8_t)((seq * 7 + i * 13) % 256); + } +} + +// Verify packet data integrity +static int verify_packet_data(uint8_t* buffer, int size, int expected_seq) { + if (size < 4) return 0; + + int seq = buffer[0] | (buffer[1] << 8); + int pkt_size = buffer[2] | (buffer[3] << 8); + + if (seq != expected_seq || pkt_size != size) { + return 0; + } + + for (int i = 4; i < size; i++) { + if (buffer[i] != (uint8_t)((seq * 7 + i * 13) % 256)) { + return 0; + } + } + return 1; +} + +// Check if connection is established +static int is_connection_established(struct UTUN_INSTANCE* inst) { + if (!inst) return 0; + struct ETCP_CONN* conn = inst->connections; + while (conn) { + struct ETCP_LINK* link = conn->links; + while (link) { + if (link->initialized) return 1; + link = link->next; + } + conn = conn->next; + } + return 0; +} + +// Send packets from client to server (forward direction) via normalizer +static void send_packets_fwd(void) { + if (!client_instance || !client_pn || packets_sent_fwd >= TOTAL_PACKETS) return; + + // Start timing on first packet + if (packets_sent_fwd == 0) { + clock_gettime(CLOCK_MONOTONIC, &start_time_fwd); + phase = 1; + printf("Starting forward transfer (client -> server) via normalizer...\n"); + } + + // Send while we have packets + while (packets_sent_fwd < TOTAL_PACKETS) { + int size = packet_sizes[packets_sent_fwd]; + uint8_t* buffer = malloc(size); + if (!buffer) break; + + generate_packet_data(current_packet_seq_fwd, buffer, size); + + pn_packer_send(client_pn, buffer, size); + + free(buffer); + packets_sent_fwd++; + current_packet_seq_fwd++; + } + + if (packets_sent_fwd >= TOTAL_PACKETS) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All %d forward packets queued to normalizer", TOTAL_PACKETS); + } +} + +// Send packets from server to client (backward direction) via normalizer +static void send_packets_back(void) { + if (!server_instance || !server_pn || packets_sent_back >= TOTAL_PACKETS) return; + + // Start timing on first packet + if (packets_sent_back == 0) { + clock_gettime(CLOCK_MONOTONIC, &start_time_back); + phase = 2; + printf("Starting backward transfer (server -> client) via normalizer...\n"); + } + + // Send while we have packets + while (packets_sent_back < TOTAL_PACKETS) { + int size = packet_sizes[packets_sent_back]; + uint8_t* buffer = malloc(size); + if (!buffer) break; + + generate_packet_data(current_packet_seq_back, buffer, size); + + pn_packer_send(server_pn, buffer, size); + + free(buffer); + packets_sent_back++; + current_packet_seq_back++; + } + + if (packets_sent_back >= TOTAL_PACKETS) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All %d backward packets queued to normalizer", TOTAL_PACKETS); + } +} + +// Check packets received by server (forward direction) via normalizer output +static void check_received_packets_fwd(void) { + if (!server_instance || !server_pn) return; + + // Debug: check output queue count + int output_count = queue_entry_count(server_pn->output); + if (output_count > 0) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Server output queue has %d entries", output_count); + } + + void* data; + while ((data = queue_data_get(server_pn->output)) != NULL) { + struct ll_entry* entry = (struct ll_entry*)data; + + if (entry->len >= 4) { + int seq = entry->dgram[0] | (entry->dgram[1] << 8); + + if (verify_packet_data(entry->dgram, entry->len, seq)) { + packets_received_fwd++; + } else { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Packet verification failed, seq=%d, len=%d", seq, entry->len); + } + } + + ll_free_dgram(entry); + queue_data_free(data); + } +} + +// Check packets received by client (backward direction) via normalizer output +static void check_received_packets_back(void) { + if (!client_instance || !client_pn) return; + + void* data; + while ((data = queue_data_get(client_pn->output)) != NULL) { + struct ll_entry* entry = (struct ll_entry*)data; + + if (entry->len >= 4) { + int seq = entry->dgram[0] | (entry->dgram[1] << 8); + + if (verify_packet_data(entry->dgram, entry->len, seq)) { + packets_received_back++; + } else { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Packet verification failed, seq=%d, len=%d", seq, entry->len); + } + } + + ll_free_dgram(entry); + queue_data_free(data); + } +} + +// Calculate time difference in milliseconds +static double time_diff_ms(struct timespec* start, struct timespec* end) { + double seconds = end->tv_sec - start->tv_sec; + double nanoseconds = end->tv_nsec - start->tv_nsec; + return (seconds * 1000.0) + (nanoseconds / 1000000.0); +} + +// Monitor function +static void monitor_and_send(void* arg) { + (void)arg; + + if (test_completed) { + packet_timeout_id = NULL; + return; + } + + static int connection_checked = 0; + + if (!connection_checked) { + if (is_connection_established(client_instance)) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Connection established, starting transmission via normalizer"); + connection_checked = 1; + + // Initialize normalizers after connection is established + if (!client_pn && client_instance->connections) { + client_pn = pn_init(client_instance->connections); + if (!client_pn) { + printf("Failed to create client normalizer\n"); + test_completed = 2; + return; + } + printf("Client normalizer created (frag_size=%d)\n", client_pn->frag_size); + } + + if (!server_pn && server_instance->connections) { + server_pn = pn_init(server_instance->connections); + if (!server_pn) { + printf("Failed to create server normalizer\n"); + test_completed = 2; + return; + } + printf("Server normalizer created (frag_size=%d)\n\n", server_pn->frag_size); + } + } + } + + if (connection_checked) { + // Phase 1: Forward transfer (client -> server) + if (packets_sent_fwd < TOTAL_PACKETS || packets_received_fwd < TOTAL_PACKETS) { + send_packets_fwd(); + check_received_packets_fwd(); + + // Check if forward phase completed + if (packets_sent_fwd >= TOTAL_PACKETS && packets_received_fwd >= TOTAL_PACKETS) { + if (end_time_fwd.tv_sec == 0) { + clock_gettime(CLOCK_MONOTONIC, &end_time_fwd); + double duration = time_diff_ms(&start_time_fwd, &end_time_fwd); + printf("Forward transfer completed: %d/%d packets in %.2f ms\n", + packets_received_fwd, TOTAL_PACKETS, duration); + } + } + } + // Phase 2: Backward transfer (server -> client) + else if (packets_sent_back < TOTAL_PACKETS || packets_received_back < TOTAL_PACKETS) { + send_packets_back(); + check_received_packets_back(); + } + // Check completion + else { + clock_gettime(CLOCK_MONOTONIC, &end_time_back); + double duration_back = time_diff_ms(&start_time_back, &end_time_back); + double duration_total = time_diff_ms(&start_time_fwd, &end_time_back); + + printf("Backward transfer completed: %d/%d packets in %.2f ms\n", + packets_received_back, TOTAL_PACKETS, duration_back); + + test_completed = 1; + printf("\n=== SUCCESS: Bidirectional transfer via normalizer completed! ===\n"); + printf("Forward (client->server): %d/%d packets in %.2f ms\n", + packets_received_fwd, TOTAL_PACKETS, time_diff_ms(&start_time_fwd, &end_time_fwd)); + printf("Backward (server->client): %d/%d packets in %.2f ms\n", + packets_received_back, TOTAL_PACKETS, duration_back); + printf("Total time: %.2f ms\n", duration_total); + + if (packet_timeout_id) { + uasync_cancel_timeout(server_instance->ua, packet_timeout_id); + packet_timeout_id = NULL; + } + return; + } + } + + if (!test_completed) { + packet_timeout_id = uasync_set_timeout(server_instance->ua, 10, NULL, monitor_and_send); + } +} + +// Timeout handler +static void test_timeout(void* arg) { + (void)arg; + if (!test_completed) { + printf("\n=== TIMEOUT ===\n"); + printf("Forward: Sent: %d/%d, Received: %d/%d\n", + packets_sent_fwd, TOTAL_PACKETS, packets_received_fwd, TOTAL_PACKETS); + printf("Backward: Sent: %d/%d, Received: %d/%d\n", + packets_sent_back, TOTAL_PACKETS, packets_received_back, TOTAL_PACKETS); + test_completed = 2; + if (packet_timeout_id) { + uasync_cancel_timeout(server_instance->ua, packet_timeout_id); + packet_timeout_id = NULL; + } + } +} + +int main() { + printf("=== PKT Normalizer + ETCP Test ===\n"); + printf("Testing with %d packets of random sizes (%d-%d bytes)\n\n", + TOTAL_PACKETS, MIN_PACKET_SIZE, MAX_PACKET_SIZE); + + // Generate random packet sizes + srand((unsigned)time(NULL)); + int total_bytes = 0; + for (int i = 0; i < TOTAL_PACKETS; i++) { + packet_sizes[i] = MIN_PACKET_SIZE + rand() % (MAX_TEST_PACKET_SIZE - MIN_PACKET_SIZE + 1); + total_bytes += packet_sizes[i]; + } + printf("Total data to transfer: %d bytes (%.2f KB average per packet)\n\n", + total_bytes, (float)total_bytes / TOTAL_PACKETS / 1024); + + debug_config_init(); + debug_set_level(DEBUG_LEVEL_DEBUG); + debug_set_categories(DEBUG_CATEGORY_ETCP); + + utun_instance_set_tun_init_enabled(0); + + printf("Creating server...\n"); + struct UASYNC* server_ua = uasync_create(); + server_instance = utun_instance_create(server_ua, "test_server.conf"); + if (!server_instance || init_connections(server_instance) < 0) { + printf("Failed to create server\n"); + return 1; + } + printf("Server created, waiting for connection...\n\n"); + + printf("Creating client...\n"); + struct UASYNC* client_ua = uasync_create(); + client_instance = utun_instance_create(client_ua, "test_client.conf"); + if (!client_instance || init_connections(client_instance) < 0) { + printf("Failed to create client\n"); + return 1; + } + printf("Client created\n\n"); + + printf("Sending %d packets in each direction via normalizer...\n", TOTAL_PACKETS); + packet_timeout_id = uasync_set_timeout(server_ua, 500, NULL, monitor_and_send); + void* global_timeout_id = uasync_set_timeout(server_ua, TEST_TIMEOUT_MS, NULL, test_timeout); + + int elapsed = 0; + while (!test_completed && elapsed < TEST_TIMEOUT_MS + 1000) { + if (server_ua) uasync_poll(server_ua, 5); + if (client_ua) uasync_poll(client_ua, 5); + usleep(5000); + elapsed += 5; + } + + printf("\nCleaning up...\n"); + + if (packet_timeout_id) uasync_cancel_timeout(server_ua, packet_timeout_id); + if (global_timeout_id) uasync_cancel_timeout(server_ua, global_timeout_id); + + if (server_pn) { + pn_pair_deinit(server_pn); + } + if (client_pn) { + pn_pair_deinit(client_pn); + } + + if (server_instance) { + server_instance->running = 0; + utun_instance_destroy(server_instance); + } + if (client_instance) { + client_instance->running = 0; + utun_instance_destroy(client_instance); + } + + if (test_completed == 1) { + printf("\n=== TEST PASSED ===\n"); + printf("All %d packets transmitted in each direction via normalizer\n", TOTAL_PACKETS); + return 0; + } else { + printf("\n=== TEST FAILED ===\n"); + printf("Forward: Sent: %d/%d, Received: %d/%d\n", + packets_sent_fwd, TOTAL_PACKETS, packets_received_fwd, TOTAL_PACKETS); + printf("Backward: Sent: %d/%d, Received: %d/%d\n", + packets_sent_back, TOTAL_PACKETS, packets_received_back, TOTAL_PACKETS); + return 1; + } +}