Browse Source

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
nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
add785c8a1
  1. 396
      src/pkt_normalizer.c
  2. 30
      src/pkt_normalizer.h
  3. 5
      tests/Makefile.am
  4. 409
      tests/test_pkt_normalizer_etcp.c

396
src/pkt_normalizer.c

@ -6,6 +6,7 @@
#include <stdlib.h>
#include <string.h>
#include <stdio.h> // 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) {
queue_data_free(data);
}
queue_free(pn->output);
struct ll_entry* entry = data_to_entry(data);
if (entry->dgram) {
free(entry->dgram);
}
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);
queue_free(pn->output);
}
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;
entry->dgram_pool = NULL;
// 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);
}
queue_data_put(pn->input, entry, 0);
// Cancel flush timer if active
if (pn->flush_timer) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
pn->flush_timer = NULL;
}
}
// Internal: Callback when input queue becomes empty, move next from pending
static void 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->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);
}
}
queue_wait_threshold(pn->etcp->input_queue, 0, 0, etcp_input_ready_cb, pn);
}
static void etcp_input_ready_cb(struct ll_queue* q, void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
// 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;
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);
}
}
}
// Allocate data
uint8_t* packet_data = malloc(entry->len);
if (!packet_data) return;
// 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)
struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool);
if (!frag) return;
memcpy(packet_data, entry->dgram, entry->len);
frag->ll.dgram = memory_pool_alloc(pn->etcp->instance->data_pool);
if (!frag->ll.dgram) {
memory_pool_free(pn->etcp->rx_pool, frag);
// Allocate ETCP_FRAGMENT from rx_pool
struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool);
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;
struct ll_queue* eq = pn->etcp->input_queue;
queue_data_put(pn->etcp->input_queue, frag, 0);
}
// 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;
}
struct ll_entry* in_dgram = data_to_entry(data);
uint16_t ptr = 0;
while (q->head) {
void* data = queue_data_get(q);
struct ll_entry* entry = data_to_entry(data);
uint16_t item_len = entry->len;
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;
}
// 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;
}
// 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;
}
queue_data_free(data); // Free the entry
}
// If remaining in buffer, set timeout to flush if queues empty
if (pn->len > 0) {
while (ptr < in_dgram->len) {
pn_buf_renew(pn);
if (!pn->sndpart) break; // Allocation failed
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;
}
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;
}
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);
}
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;
}
pn->flush_timer = uasync_set_timeout(pn->ua, 10, pn, flush_cb); // 1ms = 10 * 0.1ms
// 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(q);
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;
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
}
}
// 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;
}
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: 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;
// 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;
while (1) {
void* data = queue_data_get(pn->etcp->output_queue);
if (!data) break;
uint16_t header = *(uint16_t*)data;
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;
if (pn->in_fragment) {
// Expect continuation
if (header != 0x0000) {
// Error: invalid continuation
while (ptr < len) {
if (!pn->recvpart) {
// Need length header for new packet
if (len - ptr < 2) {
// Incomplete header, reset
pn_unpacker_reset_state(pn);
return;
break;
}
uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8);
ptr += 2;
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);
pn->recvpart = ll_alloc_lldgram(pkt_len);
if (!pn->recvpart) {
break;
}
memcpy(pn->out_buf + pn->len, data + 2, chunk);
pn->len += chunk;
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;
pn->recvpart->len = 0;
}
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);
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;
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);
}

30
src/pkt_normalizer.h

@ -11,28 +11,21 @@ struct PKTNORM {
// public:
struct ll_queue* input; // Входная очередь в packer (через нее отправляем пакеты)
struct ll_queue* output; // Выходная очередь из unpacker (через нее принимаем пакеты)
uint16_t tx_wait_time;
// 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)
uint16_t frag_size; // размер фрагмента на которые разбивать
struct ll_entry* sndpart; // блок ожидающий досборки
void* flush_timer; // For timeout flush
struct ll_queue* pending; // For not filling input
struct ll_entry* pending; // Partial processed input entry
uint16_t pending_in_ptr; // Pointer in pending entry
// 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* 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

5
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

409
tests/test_pkt_normalizer_etcp.c

@ -0,0 +1,409 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <time.h>
#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;
}
}
Loading…
Cancel
Save