Browse Source

Refactor: Code cleanup and improvements across multiple modules

- Refactored ll_queue.c/h with improved error handling
- Cleaned up pkt_normalizer.c memory management
- Updated etcp.c with better logging
- Fixed test files compilation warnings
- General code quality improvements
nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
9a5db10939
  1. 46
      lib/ll_queue.c
  2. 33
      lib/ll_queue.h
  3. 141
      src/aa
  4. 38
      src/etcp.c
  5. 751
      src/etcp.c1
  6. 98
      src/pkt_normalizer.c
  7. 197
      src/pkt_normalizer.c1
  8. 2
      src/pkt_normalizer.h
  9. 50
      src/pkt_normalizer.h1
  10. 28
      src/req.txt
  11. BIN
      tests/test_config_debug
  12. BIN
      tests/test_crypto
  13. BIN
      tests/test_debug_categories
  14. BIN
      tests/test_ecc_encrypt
  15. BIN
      tests/test_etcp_100_packets
  16. 8
      tests/test_etcp_100_packets.c
  17. BIN
      tests/test_etcp_crypto
  18. BIN
      tests/test_etcp_minimal
  19. BIN
      tests/test_etcp_simple_traffic
  20. 4
      tests/test_etcp_simple_traffic.c
  21. BIN
      tests/test_etcp_two_instances
  22. BIN
      tests/test_intensive_memory_pool
  23. 6
      tests/test_intensive_memory_pool.c
  24. 6
      tests/test_intensive_memory_pool_new.c
  25. BIN
      tests/test_ll_queue
  26. 28
      tests/test_ll_queue.c
  27. BIN
      tests/test_memory_pool_and_config
  28. 2
      tests/test_memory_pool_and_config.c
  29. BIN
      tests/test_packet_dump
  30. BIN
      tests/test_pkt_normalizer_etcp
  31. 8
      tests/test_pkt_normalizer_etcp.c
  32. BIN
      tests/test_u_async_comprehensive
  33. BIN
      tests/test_u_async_performance

46
lib/ll_queue.c

@ -47,6 +47,21 @@ struct ll_queue* queue_new(struct UASYNC* ua, size_t hash_size) {
return q;
}
struct ll_entry* ll_alloc_lldgram(uint16_t len) {
struct ll_entry* entry = queue_data_new(0);
if (!entry) return NULL;
entry->len=0;
entry->memlen=len;
entry->dgram = malloc(len);
if (!entry->dgram) {
queue_entry_free(entry);
return NULL;
}
return entry;
}
void queue_free(struct ll_queue* q) {
if (!q) return;
@ -122,7 +137,7 @@ void* queue_data_new(size_t data_size) {
return (void*)(entry + xxx);
}
void* queue_data_new_from_pool(struct memory_pool* pool) {
void* queue_entry_new_from_pool(struct memory_pool* pool) {
if (!pool) return NULL;
struct ll_entry* entry = memory_pool_alloc(pool);
@ -133,15 +148,30 @@ void* queue_data_new_from_pool(struct memory_pool* pool) {
entry->len = 0;
entry->pool = pool; // Выделено из пула
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_new_from_pool: created entry %p from pool %p", entry, pool);
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_entry_new_from_pool: created entry %p from pool %p", entry, pool);
return (void*)(entry + xxx);
}
void queue_data_free(void* data) {
if (!data) return;
struct ll_entry* entry = data_to_entry(data);
//void ll_free_dgram(struct ll_entry* entry) {
void queue_dgram_free(struct ll_entry* entry) {
if (!entry) return;
if (entry->dgram) {
if (entry->dgram_free_fn) {
entry->dgram_free_fn(entry->dgram, NULL); // arg=NULL, если не задан
} else if (entry->dgram_pool) {
memory_pool_free(entry->dgram_pool, entry->dgram);
} else {
free(entry->dgram);
}
entry->dgram = NULL;
entry->len = 0; // Опционально сброс len
}
}
void queue_entry_free(struct ll_entry* entry) {
if (!entry) return;
if (entry->pool) {
memory_pool_free(entry->pool, entry);
@ -206,7 +236,7 @@ int queue_data_put(struct ll_queue* q, void* data, uint32_t id) {
// Проверить лимит размера
if (q->size_limit >= 0 && q->count >= q->size_limit) {
queue_data_free(data); // Освободить элемент если превышен лимит
queue_entry_free(data); // Освободить элемент если превышен лимит
return -1;
}
@ -249,7 +279,7 @@ int queue_data_put_first(struct ll_queue* q, void* data, uint32_t id) {
// Проверить лимит размера
if (q->size_limit >= 0 && q->count >= q->size_limit) {
queue_data_free(data); // Освободить элемент если превышен лимит
queue_entry_free(data); // Освободить элемент если превышен лимит
return -1;
}

33
lib/ll_queue.h

@ -24,8 +24,9 @@ struct ll_entry {
struct ll_entry* prev; // Указатель на предыдущий элемент в очереди
uint16_t size; // Размер доступной памяти после блока ll_entry - т.е. data[size]. используется для добавления доп. параметров
uint16_t len; // размер пакета
uint8_t* dgram; // данные пакета
uint16_t len; // размер пакета (dgram)
uint16_t memlen; // размер выделенной памяти (dgram)
uint8_t* dgram; // данные пакета
void (*dgram_free_fn)(uint8_t* data, void* arg); // функция освобождения блока
struct memory_pool* dgram_pool; // Пул, из которого выделен этот элемент (NULL, если выделен через malloc)
@ -80,6 +81,17 @@ struct ll_queue* queue_new(struct UASYNC* ua, size_t hash_size);
// Элементы должны быть предварительно извлечены через queue_data_get() и освобождены через queue_data_free()
void queue_free(struct ll_queue* q);
// выделить ll_entry и память под кодограмму
struct ll_entry* ll_alloc_lldgram(uint16_t len);
// Функция для освобождения только dgram в ll_entry
// Принимает void* data (как возвращается из queue_data_new или queue_data_get)
// Освобождает dgram с использованием dgram_pool (если указан) или free (если нет)
// Если есть dgram_free_fn, использует её; иначе - pool или free
// Устанавливает dgram в NULL после освобождения
// Не освобождает саму структуру ll_entry
//void ll_free_dgram(struct ll_entry*);
// ==================== Конфигурация очереди ====================
// Установить функцию и аргумент коллбэка для автозабора из очереди
@ -102,36 +114,35 @@ void queue_set_size_limit(struct ll_queue* q, int lim);
// Создать новый элемент с областью данных указанного размера
// Память выделяется одним блоком: [struct ll_entry][область данных data_size байт]
// Возвращает: указатель на структуру элемента (struct ll_entry*) или NULL при ошибке выделения памяти
// ПРИМЕЧАНИЕ: ref_count = 1, элемент не в очереди
void* queue_data_new(size_t data_size);
// Создать новый элемент из пула (размер был определен при создании пула)
// Возвращает: указатель на структуру элемента (struct ll_entry*) или NULL при ошибке выделения памяти
// ПРИМЕЧАНИЕ: ref_count = 1, элемент не в очереди
void* queue_data_new_from_pool(struct memory_pool* pool);
void* queue_entry_new_from_pool(struct memory_pool* pool);
// Освободить элемент (не влияет на очереди - должен быть предварительно взят из всех очередей)
// Уменьшает ref_count, если становится 0 - освобождает память
void queue_data_free(void* data);
// Освободить только entry (не влияет на очереди, dgram не освобождает)
void queue_entry_free(struct ll_entry* entry);
// Освободить только entry->dgram (не влияет на очереди)
void queue_dgram_free(struct ll_entry* entry);
//void queue_data_free(void* data);
// ==================== Операции с очередью ====================
// Добавить элемент в конец очереди (FIFO)
// Если очередь была пустой и коллбэки разрешены - вызывает коллбэк
// Возвращает: 0 при успехе, -1 если превышен лимит размера (элемент освобожден)
// ПРИМЕЧАНИЕ: НЕ изменяет ref_count элемента
int queue_data_put(struct ll_queue* q, void* data, uint32_t id);
// Добавить элемент в начало очереди (LIFO, высокий приоритет)
// Если очередь была пустой и коллбэки разрешены - вызывает коллбэк
// Возвращает: 0 при успехе, -1 если превышен лимит размера (элемент освобожден)
// ПРИМЕЧАНИЕ: НЕ изменяет ref_count элемента
int queue_data_put_first(struct ll_queue* q, void* data, uint32_t id);
// Извлечь элемент из начала очереди
// При извлечении приостанавливает коллбэки (callback_suspended = 1) чтобы предотвратить рекурсию
// Возвращает: указатель на структуру элемента (struct ll_entry*) или NULL если очередь пуста
// ПРИМЕЧАНИЕ: НЕ изменяет ref_count элемента, просто удаляет из очереди
// ПРИМЕЧАНИЕ: не освобождает память элемента
void* queue_data_get(struct ll_queue* q);
// Получить текущее количество элементов в очереди

141
src/aa

@ -0,0 +1,141 @@
ЭТАП 1. СНЯТИЕ СПАЗМА (5–7 минут)
Подходит даже в обострение.
1️⃣ Диафрагмальное дыхание лёжа
Исходное: лёжа на спине, ноги согнуты, стопы на полу, одна рука на животе
Как делать:
вдох носом → живот поднимается
выдох ртом → живот мягко опускается
Важно: поясница расслаблена
🕐 2 минуты
🎯 снимает спазм поясницы и таза
2️⃣ Поясничные перекаты
Исходное: лёжа, ноги согнуты
Как:
на выдохе слегка прижать поясницу к полу
на вдохе — отпустить (не прогибаться активно)
🕐 10–12 раз
🎯 возвращает подвижность, уменьшает «слиплось»
3️⃣ Колени к груди (по одному)
Исходное: лёжа
Как:
подтяни одно колено к груди
держи 10–15 сек
поменяй ногу
🕐 2–3 раза на каждую
🎯 разгружает крестец
🔹 ЭТАП 2. ВКЛЮЧЕНИЕ ГЛУБОКИХ МЫШЦ (8–10 минут)
Когда боль ≤ 3/10.
4️⃣ Ягодичный мост (медленный)
Исходное: лёжа, ноги согнуты
Как:
выдох → напряги ягодицы
медленно подними таз
держи 5 сек
опусти
🕐 8–10 раз
❗ не прогибаться в пояснице
🎯 стабилизирует крестец
5️⃣ “Пятка скользит”
Исходное: лёжа
Как:
медленно выпрямляй одну ногу, скользя пяткой по полу
возвращай обратно
🕐 6–8 раз на каждую
🎯 мягкая активация без нагрузки
6️⃣ Кошка (очень мягко)
Исходное: на четвереньках
Как:
выдох → округли спину
вдох → вернись в нейтраль (НЕ прогиб!)
🕐 8–10 раз
🎯 улучшает кровоток
🔹 ЭТАП 3. СТАБИЛИЗАЦИЯ (5–8 минут)
Добавляй через 7–10 дней.
7️⃣ Подъём руки + противоположной ноги (укороченный)
Исходное: на четвереньках
Как:
подними руку
если комфортно — добавь противоположную ногу (не высоко)
🕐 5–6 раз на сторону
🎯 стабилизация без оси
8️⃣ Растяжка ягодиц
Исходное: лёжа
Как:
щиколотку на колено
мягко подтяни к себе
🕐 15–20 сек × 2
🎯 снимает давление с крестца
⏱ Время
минимум: 10–15 минут
оптимально: 20 минут
❌ Чего НЕ делать сейчас
скручивания
«лодочку»
глубокие наклоны
резкие растяжки
вис на турнике
📅 Как долго
первые улучшения: 7–14 дней
стойкий эффект: 6–8 недель
дальше — 2–3 раза в неделю для поддержки

38
src/etcp.c

@ -177,7 +177,7 @@ void etcp_connection_close(struct ETCP_CONN* etcp) {
if (etcp->ack_q) {
struct ACK_PACKET* pkt;
while ((pkt = queue_data_get(etcp->ack_q)) != NULL) {
queue_data_free(pkt);
queue_entry_free(pkt);
}
queue_free(etcp->ack_q);
etcp->ack_q = NULL;
@ -272,7 +272,7 @@ int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
// Create queue entry - this allocates ll_entry + data pointer
struct ETCP_FRAGMENT* pkt = queue_data_new_from_pool(etcp->rx_pool);
struct ETCP_FRAGMENT* pkt = queue_entry_new_from_pool(etcp->rx_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate queue entry");
memory_pool_free(etcp->instance->data_pool, packet_data);
@ -282,7 +282,7 @@ int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->ll.dgram = packet_data; // Point to data_pool allocation
pkt->ll.size = len; // размер packet_data
pkt->ll.len = len; // размер packet_data
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_send: created PACKET %p with data %p (len=%zu)", pkt, packet_data, len);
@ -326,7 +326,7 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: processing ETCP_FRAGMENT %p (seq=%u, len=%u)", in_pkt, in_pkt->seq, in_pkt->ll.size);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: processing ETCP_FRAGMENT %p (seq=%u, len=%u)", in_pkt, in_pkt->seq, in_pkt->ll.len);
memory_pool_free(etcp->rx_pool, in_pkt);// перемещаем из rx_pool в inflight_pool
@ -334,7 +334,7 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->inflight_pool, etcp);
queue_data_free(in_pkt); // Free the ETCP_FRAGMENT
queue_entry_free(in_pkt); // Free the ETCP_FRAGMENT
queue_resume_callback(q);
return;
}
@ -345,12 +345,12 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
p->state = INFLIGHT_STATE_WAIT_SEND;
p->last_timestamp = 0;
p->ll.dgram = in_pkt->ll.dgram;
p->ll.size = in_pkt->ll.size;
p->ll.len = in_pkt->ll.len;
// Add to send queue
if (queue_data_put(etcp->input_send_q, p, p->seq) != 0) {
memory_pool_free(etcp->inflight_pool, p);
queue_data_free(in_pkt);
queue_entry_free(in_pkt);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_cb: EXIT (queue put failed)");
return;
}
@ -445,7 +445,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: getting packet from input_send_q");
struct INFLIGHT_PACKET* inf_pkt = queue_data_get(etcp->input_send_q);
if (inf_pkt) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: prepare udp dgram for send packet %p (seq=%u, len=%u)", inf_pkt, inf_pkt->seq, inf_pkt->ll.size);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: prepare udp dgram for send packet %p (seq=%u, len=%u)", inf_pkt, inf_pkt->seq, inf_pkt->ll.len);
inf_pkt->last_timestamp=get_current_time_units();
inf_pkt->send_count++;
@ -497,9 +497,9 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
dgram->data[ptr++]=dly;
dgram->data[ptr++]=dly>>8;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: add ACK N%d dTS=%d", ack_pkt->seq, dly);
queue_data_free(ack_pkt);
queue_entry_free(ack_pkt);
if (inf_pkt && inf_pkt->ll.size+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки)
if (inf_pkt && inf_pkt->ll.len+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки)
if (ptr>500) break;
}
@ -507,14 +507,14 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
if (inf_pkt) {
// фрейм data (0) обязательно в конец
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: packet with payload (seq=%u, len=%u), ack_size=%d", inf_pkt->seq, inf_pkt->ll.size, dgram->data[1]);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: packet with payload (seq=%u, len=%u), ack_size=%d", inf_pkt->seq, inf_pkt->ll.len, dgram->data[1]);
dgram->data[ptr++]=0;// payload
dgram->data[ptr++]=inf_pkt->seq;
dgram->data[ptr++]=inf_pkt->seq>>8;
dgram->data[ptr++]=inf_pkt->seq>>16;
dgram->data[ptr++]=inf_pkt->seq>>24;
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.size); ptr+=inf_pkt->ll.size;
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len;
}
else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: built packet with %d bytes total", dgram->data_len);
@ -572,7 +572,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: assembling packet id=%u (len=%u)",
rx_pkt->seq, rx_pkt->ll.size);
rx_pkt->seq, rx_pkt->ll.len);
// Simply move ETCP_FRAGMENT from recv_q to output_queue - no data copying needed
// Remove from recv_q first
@ -580,7 +580,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// Add to output_queue using the same ETCP_FRAGMENT structure
if (queue_data_put(etcp->output_queue, rx_pkt, next_expected_id) == 0) {
delivered_bytes += rx_pkt->ll.size;
delivered_bytes += rx_pkt->ll.len;
delivered_count++;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: moved packet id=%u to output_queue",
next_expected_id);
@ -652,8 +652,8 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
}
// Update connection statistics
etcp->unacked_bytes -= acked_pkt->ll.size;
etcp->bytes_sent_total += acked_pkt->ll.size;
etcp->unacked_bytes -= acked_pkt->ll.len;
etcp->bytes_sent_total += acked_pkt->ll.len;
etcp->ack_packets_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: removed packet seq=%u from wait_ack, unacked_bytes now %u", seq, etcp->unacked_bytes);
@ -706,7 +706,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
if (len>=5) {
// формируем ACK
struct ACK_PACKET* p = queue_data_new_from_pool(etcp->instance->ack_pool);
struct ACK_PACKET* p = queue_entry_new_from_pool(etcp->instance->ack_pool);
uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24);
p->seq=seq;
p->pkt_timestamp=pkt->timestamp;
@ -721,11 +721,11 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: adding packet seq=%u to recv_q (last_delivered_id=%u)", seq, etcp->last_delivered_id);
// отправляем пакет в очередь на сборку
uint8_t* payload_data = memory_pool_alloc(etcp->instance->data_pool);
struct ETCP_FRAGMENT* rx_pkt = queue_data_new_from_pool(etcp->rx_pool);
struct ETCP_FRAGMENT* rx_pkt = queue_entry_new_from_pool(etcp->rx_pool);
rx_pkt->seq=seq;
rx_pkt->timestamp=pkt->timestamp;
rx_pkt->ll.dgram=payload_data;
rx_pkt->ll.size=pkt_len;
rx_pkt->ll.len=pkt_len;
// Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len);
queue_data_put(etcp->recv_q, rx_pkt, seq);

751
src/etcp.c1

@ -0,0 +1,751 @@
// etcp.c - ETCP Protocol Implementation (refactored and expanded based on etcp_protocol.txt)
#include "etcp.h"
#include "etcp_loadbalancer.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
#include "../lib/debug_config.h"
#include "crc32.h" // For potential hashing, though not used yet.
#include <stdlib.h>
#include <string.h>
#include <sys/time.h>
#include <math.h> // For bandwidth calcs
#include <limits.h> // For UINT16_MAX
// Enable comprehensive debug output for ETCP module
#define DEBUG_CATEGORY_ETCP_DETAILED 1
// Constants from spec (adjusted for completeness)
#define MAX_INFLIGHT_BYTES 65536 // Initial window
#define RETRANS_K1 2.0f // RTT multiplier for retrans timeout
#define RETRANS_K2 1.5f // Jitter multiplier
#define ACK_DELAY_TB 20 // ACK timer delay (2ms in 0.1ms units)
#define BURST_DELAY_FACTOR 4 // Delay before burst
#define BURST_SIZE 5 // Packets in burst (1 delayed + 4 burst)
#define RTT_HISTORY_SIZE 10 // For jitter calc
#define MAX_PENDING 32 // For ACKs/retrans (arbitrary; adjust)
#define SECTION_HEADER_SIZE 3 // type(1) + len(2)
// Container-of macro for getting struct from data pointer
//#define CONTAINER_OF(ptr, type, member) ((type *)((char *)(ptr) - offsetof(type, member)))
// Forward declarations
static void input_queue_cb(struct ll_queue* q, void* arg);
static void etcp_link_ready_callback(struct ETCP_CONN* etcp);
static void input_send_q_cb(struct ll_queue* q, void* arg);
static void wait_ack_cb(struct ll_queue* q, void* arg);
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp);
// Get current time in 0.1ms units
uint64_t get_current_time_units() {
struct timeval tv;
gettimeofday(&tv, NULL);
uint64_t time_units = ((uint64_t)tv.tv_sec * 10000ULL) + (tv.tv_usec / 100);
// DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "get_current_time_units: tv_sec=%ld, tv_usec=%ld, result=%llu",
// tv.tv_sec, tv.tv_usec, (unsigned long long)time_units);
return time_units;
}
uint16_t get_current_timestamp() {
uint16_t timestamp = (uint16_t)(get_current_time_units() & 0xFFFF);
// DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "get_current_timestamp: result=%u", timestamp);
return timestamp;
}
// Timestamp diff (with wrap-around)
static uint16_t timestamp_diff(uint16_t t1, uint16_t t2) {
// DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "timestamp_diff: t1=%u, t2=%u", t1, t2);
if (t1 >= t2) {
return t1 - t2;
}
return (UINT16_MAX - t2) + t1 + 1;
}
// Create new ETCP connection
struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance) {
if (!instance) return NULL;
struct ETCP_CONN* etcp = calloc(1, sizeof(struct ETCP_CONN));
if (!etcp) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "etcp_connection_create: creating connection failed for instance %p", instance);
return NULL;
}
etcp->instance = instance;
etcp->input_queue = queue_new(instance->ua, 0); // No hash for input_queue
etcp->output_queue = queue_new(instance->ua, 0); // No hash for output_queue
etcp->input_send_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE); // Hash for send_q
etcp->input_wait_ack = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE); // Hash for wait_ack
etcp->recv_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE); // Hash for send_q
etcp->ack_q = queue_new(instance->ua, 0);
etcp->inflight_pool = memory_pool_init(sizeof(struct INFLIGHT_PACKET));
etcp->rx_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT));
if (!etcp->input_queue || !etcp->output_queue || !etcp->input_send_q || !etcp->recv_q || !etcp->ack_q ||
!etcp->input_wait_ack || !etcp->inflight_pool || !etcp->rx_pool) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "etcp_connection_create: error - closing - input:%p output:%p send:%p wait:%x pool:%p",
etcp->input_queue, etcp->output_queue, etcp->input_send_q, etcp->input_wait_ack, etcp->inflight_pool);
etcp_connection_close(etcp);
return NULL;
}
// etcp->normalizer = pn_pair_init(instance->ua, etcp->mtu);
// if (!etcp->normalizer) {
// etcp_connection_close(etcp);
// return NULL;
// }
etcp->mtu = 1500; // Default MTU
etcp->window_size = MAX_INFLIGHT_BYTES;
etcp->next_tx_id = 1;
etcp->rtt_avg_10 = 10; // Initial guess (1ms)
etcp->rtt_avg_100 = 10;
etcp->rtt_history_idx = 0;
memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history));
// Set input queue callback
queue_set_callback(etcp->input_queue, input_queue_cb, etcp);
queue_set_callback(etcp->input_send_q, input_send_q_cb, etcp);
queue_set_callback(etcp->input_wait_ack, wait_ack_cb, etcp);
etcp->link_ready_for_send_fn = etcp_link_ready_callback;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connection_create: connection initialized. ETCP=%p mtu=%d, window_size=%u, next_tx_id=%u",
etcp, etcp->mtu, etcp->window_size, etcp->next_tx_id);
return etcp;
}
// Close connection with NULL pointer safety (prevents double free)
void etcp_connection_close(struct ETCP_CONN* etcp) {
if (!etcp) return;
// Drain and free input_queue (contains ETCP_FRAGMENT with pkt_data from data_pool)
if (etcp->input_queue) {
struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(etcp->input_queue)) != NULL) {
if (pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram);
}
memory_pool_free(etcp->rx_pool, pkt);
}
queue_free(etcp->input_queue);
etcp->input_queue = NULL;
}
// Drain and free output_queue (contains ETCP_FRAGMENT with pkt_data from data_pool)
if (etcp->output_queue) {
struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(etcp->output_queue)) != NULL) {
if (pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram);
}
memory_pool_free(etcp->rx_pool, pkt);
}
queue_free(etcp->output_queue);
etcp->output_queue = NULL;
}
// Drain and free input_send_q (contains INFLIGHT_PACKET with pkt_data from data_pool)
if (etcp->input_send_q) {
struct INFLIGHT_PACKET* pkt;
while ((pkt = queue_data_get(etcp->input_send_q)) != NULL) {
if (pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram);
}
memory_pool_free(etcp->inflight_pool, pkt);
}
queue_free(etcp->input_send_q);
etcp->input_send_q = NULL;
}
// Drain and free input_wait_ack (contains INFLIGHT_PACKET with pkt_data from data_pool)
if (etcp->input_wait_ack) {
struct INFLIGHT_PACKET* pkt;
while ((pkt = queue_data_get(etcp->input_wait_ack)) != NULL) {
if (pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram);
}
memory_pool_free(etcp->inflight_pool, pkt);
}
queue_free(etcp->input_wait_ack);
etcp->input_wait_ack = NULL;
}
// Drain and free ack_q (contains ACK_PACKET from ack_pool)
if (etcp->ack_q) {
struct ACK_PACKET* pkt;
while ((pkt = queue_data_get(etcp->ack_q)) != NULL) {
queue_entry_free(pkt);
}
queue_free(etcp->ack_q);
etcp->ack_q = NULL;
}
// Drain and free recv_q (contains ETCP_FRAGMENT with pkt_data from data_pool)
if (etcp->recv_q) {
struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(etcp->recv_q)) != NULL) {
if (pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram);
}
memory_pool_free(etcp->rx_pool, pkt);
}
queue_free(etcp->recv_q);
etcp->recv_q = NULL;
}
// Free memory pools after all elements are returned
if (etcp->inflight_pool) {
memory_pool_destroy(etcp->inflight_pool);
etcp->inflight_pool = NULL;
}
if (etcp->rx_pool) {
memory_pool_destroy(etcp->rx_pool);
etcp->rx_pool = NULL;
}
// Clear links list safely
if (etcp->links) {
struct ETCP_LINK* link = etcp->links;
while (link) {
struct ETCP_LINK* next = link->next;
etcp_link_close(link);
link = next;
}
etcp->links = NULL;
}
// Clear next pointer to prevent dangling references
etcp->next = NULL;
// TODO: Free rx_list, etc.
free(etcp);
}
// Reset connection (stub)
void etcp_conn_reset(struct ETCP_CONN* etcp) {
// Reset IDs, queues, etc. as per protocol.txt
etcp->next_tx_id = 1;
etcp->last_rx_id = 0;
etcp->last_delivered_id = 0;
// Clear inflight, rx_list, etc.
}
// ====================================================================== Отправка данных
// Send data through ETCP connection
// Allocates memory from data_pool and places in input queue
// Returns: 0 on success, -1 on failure
int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
// DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_send: ENTER etcp=%p, data=%p, len=%zu", etcp, data, len);
if (!etcp || !data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: invalid parameters (etcp=%p, data=%p, len=%zu)", etcp, data, len);
return -1;
}
if (!etcp->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: input_queue is NULL for etcp=%p", etcp);
return -1;
}
// Check length against maximum packet size
if (len > PACKET_DATA_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: packet too large (len=%zu, max=%d)", len, PACKET_DATA_SIZE);
return -1;
}
// Allocate packet data from data_pool (following ETCP reception pattern)
uint8_t* packet_data = memory_pool_alloc(etcp->instance->data_pool);
if (!packet_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate packet data from data_pool");
return -1;
}
// Copy user data to packet buffer
memcpy(packet_data, data, len);
// Create queue entry - this allocates ll_entry + data pointer
struct ETCP_FRAGMENT* pkt = queue_data_new_from_pool(etcp->rx_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate queue entry");
memory_pool_free(etcp->instance->data_pool, packet_data);
return -1;
}
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->ll.dgram = packet_data; // Point to data_pool allocation
pkt->ll.size = len; // размер packet_data
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_send: created PACKET %p with data %p (len=%zu)", pkt, packet_data, len);
// Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, pkt, 0) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to add to input queue");
memory_pool_free(etcp->instance->data_pool, packet_data);
memory_pool_free(etcp->rx_pool, pkt);
return -1;
}
return 0;
}
static void input_queue_try_resume(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_try_resume: ENTER etcp=%p", etcp);
// если размер input_wait_ack+input_send_q в байтах < optimal_inflight то resume сейчас.
size_t wait_ack_bytes = queue_total_bytes(etcp->input_wait_ack);
size_t send_q_bytes = queue_total_bytes(etcp->input_send_q);
size_t total_bytes = wait_ack_bytes + send_q_bytes;
if (total_bytes < etcp->optimal_inflight) {
queue_resume_callback(etcp->input_send_q);// вызвать лишний раз resume не страшно.
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_try_resume: resumed input_send_q callback");
}
}
// Input callback for input_queue (добавление новых кодограмм в стек)
// input_queue -> input_send_q
static void input_queue_cb(struct ll_queue* q, void* arg) {
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
struct ETCP_FRAGMENT* in_pkt = queue_data_get(q);
if (!in_pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot get element (pool=%p etcp=%p)", etcp->inflight_pool, etcp);
queue_resume_callback(q);
return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: processing ETCP_FRAGMENT %p (seq=%u, len=%u)", in_pkt, in_pkt->seq, in_pkt->ll.size);
memory_pool_free(etcp->rx_pool, in_pkt);// перемещаем из rx_pool в inflight_pool
// Create INFLIGHT_PACKET
struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->inflight_pool, etcp);
queue_entry_free(in_pkt); // Free the ETCP_FRAGMENT
queue_resume_callback(q);
return;
}
// Setup inflight packet (based on protocol.txt)
memset(p, 0, sizeof(*p));
p->seq = etcp->next_tx_id++; // Assign seq
p->state = INFLIGHT_STATE_WAIT_SEND;
p->last_timestamp = 0;
p->ll.dgram = in_pkt->ll.dgram;
p->ll.size = in_pkt->ll.size;
// Add to send queue
if (queue_data_put(etcp->input_send_q, p, p->seq) != 0) {
memory_pool_free(etcp->inflight_pool, p);
queue_entry_free(in_pkt);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_cb: EXIT (queue put failed)");
return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: successfully moved from input_queue to input_send_q");
input_queue_try_resume(etcp);
// Resume input_queue callback to process next packet if any
queue_resume_callback(q);
}
static void input_send_q_cb(struct ll_queue* q, void* arg) {// etcp->input_send_q processing
struct ETCP_CONN* etcp=(struct ETCP_CONN*)arg;
size_t send_q_bytes = queue_total_bytes(etcp->input_send_q);
size_t send_q_pkts = queue_entry_count(etcp->input_send_q);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_send_q_cb: input_send_q status: %d pkt %d bytes", send_q_pkts, send_q_bytes);
etcp_conn_process_send_queue(etcp);
}
static void ack_timeout_check(void* arg) {
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
uint64_t now = get_current_time_units();
uint64_t timeout = 1000;//(uint64_t)(etcp->rtt_avg_10 * RETRANS_K1) + (uint64_t)(etcp->jitter * RETRANS_K2);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "ack_timeout_check: starting check, now=%llu, timeout=%llu, rtt_avg_10=%u, jitter=%u",
(unsigned long long)now, (unsigned long long)timeout, etcp->rtt_avg_10, etcp->jitter);
struct ll_entry* current = etcp->input_wait_ack->head;
while (current) {
struct ll_entry* next = current->next; // Save next as we may remove current
struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current;
uint64_t elapsed = now - pkt->last_timestamp;
if (elapsed > timeout) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "ack_timeout_check: timeout for seq=%u, elapsed=%llu, timeout=%llu, send_count=%u",
pkt->seq, (unsigned long long)elapsed, (unsigned long long)timeout, pkt->send_count);
// Increment counters
pkt->send_count++;
pkt->retrans_req_count++; // Optional, if used for retrans request logic
pkt->last_timestamp = now;
pkt->last_link = NULL; // Reset last link for re-selection
// Remove from wait_ack
queue_remove_data(etcp->input_wait_ack, pkt);
// Change state and add to send_q for retransmission
pkt->state = INFLIGHT_STATE_WAIT_SEND;
queue_data_put(etcp->input_send_q, pkt, pkt->seq);
// Update stats
etcp->retransmissions_count++;
}
else {// не надо до конца сканировать - они уже сортированы по таймстемпу т.к. очередь fifo, а timestamp = время добавления в очередь = время отправки
// shedule timer
int64_t next_timeout=timeout - elapsed;
etcp->retrans_timer = uasync_set_timeout(etcp->instance->ua, next_timeout+10, etcp, ack_timeout_check);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "ack_timeout_check: rescheduled timer for %llu units", next_timeout);
break;
}
current = next;
}
}
static void wait_ack_cb(struct ll_queue* q, void* arg) {
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
ack_timeout_check(etcp);
}
// Подготовить и отправить кодограмму
// вызывается линком когда освобождается или очередью если появляются данные на передачу
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp);
if (!link) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no link available");
return NULL;// если линков нет - ждём появления свободного
}
size_t send_q_size = queue_entry_count(etcp->input_send_q);
if (send_q_size == 0) {// сгребаем из input_queue
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: input_send_q empty, check if avail input_queue -> inflight");
input_queue_try_resume(etcp);
// return NULL;
}
// First, check if there's a packet in input_send_q (retrans or new)
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: getting packet from input_send_q");
struct INFLIGHT_PACKET* inf_pkt = queue_data_get(etcp->input_send_q);
if (inf_pkt) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: prepare udp dgram for send packet %p (seq=%u, len=%u)", inf_pkt, inf_pkt->seq, inf_pkt->ll.size);
inf_pkt->last_timestamp=get_current_time_units();
inf_pkt->send_count++;
inf_pkt->state=INFLIGHT_STATE_WAIT_ACK;
queue_data_put(etcp->input_wait_ack, inf_pkt, inf_pkt->seq);// move dgram to wait_ack queue
}
size_t ack_q_size = queue_entry_count(etcp->ack_q);
if (!inf_pkt && ack_q_size == 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no data/ack to send");
return NULL;
}
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: failed to allocate ETCP_DGRAM");
return NULL;
}
dgram->link = link;
dgram->noencrypt_len=0;
dgram->timestamp=get_current_timestamp();
// формат ack: [01] [elements count] [4 байта last_delivered_id] и <[4 байта seq][2 байта recv_ts][2 байта txrx delay ts]> x count
dgram->data[0]=1;// ack
int ptr=2;
dgram->data[ptr++]=etcp->last_delivered_id;
dgram->data[ptr++]=etcp->last_delivered_id>>8;
dgram->data[ptr++]=etcp->last_delivered_id>>16;
dgram->data[ptr++]=etcp->last_delivered_id>>24;
// тут (потом) добавим опциональные заголовки
struct ACK_PACKET* ack_pkt;
while (ack_pkt = queue_data_get(etcp->ack_q)) {
// seq 4 байта
dgram->data[ptr++]=ack_pkt->seq;
dgram->data[ptr++]=ack_pkt->seq>>8;
dgram->data[ptr++]=ack_pkt->seq>>16;
dgram->data[ptr++]=ack_pkt->seq>>24;
// ts приема 2 байта
dgram->data[ptr++]=ack_pkt->recv_timestamp;
dgram->data[ptr++]=ack_pkt->recv_timestamp>>8;
// время задержки 2 байта между recv и ack
uint16_t dly=get_current_timestamp()-ack_pkt->recv_timestamp;
dgram->data[ptr++]=dly;
dgram->data[ptr++]=dly>>8;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: add ACK N%d dTS=%d", ack_pkt->seq, dly);
queue_entry_free(ack_pkt);
if (inf_pkt && inf_pkt->ll.size+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки)
if (ptr>500) break;
}
dgram->data[1]=ptr/8;
if (inf_pkt) {
// фрейм data (0) обязательно в конец
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: packet with payload (seq=%u, len=%u), ack_size=%d", inf_pkt->seq, inf_pkt->ll.size, dgram->data[1]);
dgram->data[ptr++]=0;// payload
dgram->data[ptr++]=inf_pkt->seq;
dgram->data[ptr++]=inf_pkt->seq>>8;
dgram->data[ptr++]=inf_pkt->seq>>16;
dgram->data[ptr++]=inf_pkt->seq>>24;
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.size); ptr+=inf_pkt->ll.size;
}
else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: built packet with %d bytes total", dgram->data_len);
}
dgram->data_len=ptr;
return dgram;
}
// Callback for when a link is ready to send data
static void etcp_link_ready_callback(struct ETCP_CONN* etcp) {
if (!etcp) return;
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_link_ready_callback: processing send queue for etcp=%p", etcp);
etcp_conn_process_send_queue(etcp);
}
// Process packets in send queue and transmit them
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {
struct ETCP_DGRAM* dgram;
while(dgram = etcp_request_pkt(etcp)) {
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: sending packet");
etcp_loadbalancer_send(dgram);
}
}
static void ack_response_timer_cb(void* arg) {// проверяем неотправленные ack response и отправляем если надо.
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
etcp_conn_process_send_queue(etcp);// проталкиваем (она же должна отправлять только ack если больше ничего нет)
// если ack все еще заняты - обновляем таймаут
if (etcp->ack_q->count) etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb);
else etcp->ack_resp_timer=NULL;
}
// ====================================================================== Прием данных
void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// пробуем собрать выходную очередь из фрагментов
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: etcp=%p, last_delivered_id=%u, recv_q_count=%d",
etcp, etcp->last_delivered_id, queue_entry_count(etcp->recv_q));
uint32_t next_expected_id = etcp->last_delivered_id + 1;
int delivered_count = 0;
uint32_t delivered_bytes = 0;
// Look for contiguous packets starting from next_expected_id
while (1) {
struct ETCP_FRAGMENT* rx_pkt = queue_find_data_by_id(etcp->recv_q, next_expected_id);
if (!rx_pkt) {
// No more contiguous packets found
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: no packet found for id=%u, stopping", next_expected_id);
break;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: assembling packet id=%u (len=%u)",
rx_pkt->seq, rx_pkt->ll.size);
// Simply move ETCP_FRAGMENT from recv_q to output_queue - no data copying needed
// Remove from recv_q first
queue_remove_data(etcp->recv_q, rx_pkt);
// Add to output_queue using the same ETCP_FRAGMENT structure
if (queue_data_put(etcp->output_queue, rx_pkt, next_expected_id) == 0) {
delivered_bytes += rx_pkt->ll.size;
delivered_count++;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: moved packet id=%u to output_queue",
next_expected_id);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: failed to add packet id=%u to output_queue",
next_expected_id);
// Put it back in recv_q if we can't add to output_queue
queue_data_put(etcp->recv_q, rx_pkt, next_expected_id);
break;
}
// Update state for next iteration
etcp->last_delivered_id = next_expected_id;
next_expected_id++;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: delivered %u contiguous packets (%u bytes), last_delivered_id=%u, output_queue_count=%d",
delivered_count, delivered_bytes, etcp->last_delivered_id, queue_entry_count(etcp->output_queue));
}
// Process ACK receipt - remove acknowledged packet from inflight queues
void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t dts) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: processing ACK for seq=%u, ts=%u, dts=%u", seq, ts, dts);
// Find the acknowledged packet in the wait_ack queue
struct INFLIGHT_PACKET* acked_pkt = queue_find_data_by_id(etcp->input_wait_ack, seq);
if (acked_pkt) queue_remove_data(etcp->input_wait_ack, acked_pkt);
else { acked_pkt = queue_find_data_by_id(etcp->input_send_q, seq);
queue_remove_data(etcp->input_send_q, acked_pkt);
}
if (!acked_pkt) {
// Packet might be already acknowledged or not found
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: packet seq=%u not found in wait_ack queue", seq);
return;
}
// Calculate RTT if timestamps are valid
if (ts != (uint16_t)-1 && dts != (uint16_t)-1) {
uint16_t rtt = timestamp_diff(ts, dts);
etcp->rtt_last = rtt;
// Update RTT averages (exponential smoothing)
if (etcp->rtt_avg_10 == 0) {
etcp->rtt_avg_10 = rtt;
etcp->rtt_avg_100 = rtt;
} else {
// RTT average over 10 packets
etcp->rtt_avg_10 = (etcp->rtt_avg_10 * 9 + rtt) / 10;
// RTT average over 100 packets
etcp->rtt_avg_100 = (etcp->rtt_avg_100 * 99 + rtt) / 100;
}
// Update jitter calculation (max - min of last 10 RTT samples)
etcp->rtt_history[etcp->rtt_history_idx] = rtt;
etcp->rtt_history_idx = (etcp->rtt_history_idx + 1) % 10;
uint16_t rtt_min = UINT16_MAX, rtt_max = 0;
for (int i = 0; i < 10; i++) {
if (etcp->rtt_history[i] < rtt_min) rtt_min = etcp->rtt_history[i];
if (etcp->rtt_history[i] > rtt_max) rtt_max = etcp->rtt_history[i];
}
etcp->jitter = rtt_max - rtt_min;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: RTT updated - last=%u, avg_10=%u, avg_100=%u, jitter=%u",
rtt, etcp->rtt_avg_10, etcp->rtt_avg_100, etcp->jitter);
}
// Update connection statistics
etcp->unacked_bytes -= acked_pkt->ll.size;
etcp->bytes_sent_total += acked_pkt->ll.size;
etcp->ack_packets_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: removed packet seq=%u from wait_ack, unacked_bytes now %u", seq, etcp->unacked_bytes);
if (acked_pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, acked_pkt->ll.dgram);
}
memory_pool_free(etcp->inflight_pool, acked_pkt);
// Try to resume sending more packets if window space opened up
input_queue_try_resume(etcp);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: completed for seq=%u", seq);
}
// Process incoming decrypted packet
void etcp_conn_input(struct ETCP_DGRAM* pkt) {
if (!pkt || !pkt->data_len) return;
struct ETCP_CONN* etcp = pkt->link->etcp;
uint8_t* data = pkt->data;
uint16_t len = pkt->data_len;
uint16_t ts = pkt->timestamp; // Received timestamp
// Note: Assume packet starts with sections after timestamp (but timestamp is already extracted in connections?).
// Protocol.txt: timestamp is first 2B, then sections.
// But in conn_input, pkt->data is after timestamp? Assume data starts with first section.
while (len >= 1) {
uint8_t type = data[0];
// Process sections as per protocol.txt
switch (type) {
case ETCP_SECTION_ACK: {
int elm_cnt=data[1];
uint32_t till=data[2] | (data[3]<<8) | (data[4]<<16) | (data[5]<<24);
data+=6;
for (int i=0; i<elm_cnt; i++) {
uint32_t seq=data[0] | (data[1]<<8) | (data[2]<<16) | (data[3]<<24);
uint16_t ts=data[4] | (data[5]<<8);
uint16_t dts=data[6] | (data[7]<<8);
etcp_ack_recv(etcp, seq, ts, dts);
data+=8;
}
while (etcp->rx_ack_till-till<0) { etcp->rx_ack_till++; etcp_ack_recv(etcp, etcp->rx_ack_till, -1, -1); }// подтверждаем всё по till
break;
}
case ETCP_SECTION_PAYLOAD: {
if (len>=5) {
// формируем ACK
struct ACK_PACKET* p = queue_data_new_from_pool(etcp->instance->ack_pool);
uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24);
p->seq=seq;
p->pkt_timestamp=pkt->timestamp;
p->recv_timestamp=get_current_timestamp();
queue_data_put(etcp->ack_q, p, p->seq);
if (etcp->ack_resp_timer == NULL) {
etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: set ack_timer for delayed ACK send");
}
if ((int32_t)(etcp->last_delivered_id-seq)<0) if (queue_find_data_by_id(etcp->recv_q, seq)==NULL) {// проверяем есть ли пакет с этим seq
uint32_t pkt_len=len-5;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: adding packet seq=%u to recv_q (last_delivered_id=%u)", seq, etcp->last_delivered_id);
// отправляем пакет в очередь на сборку
uint8_t* payload_data = memory_pool_alloc(etcp->instance->data_pool);
struct ETCP_FRAGMENT* rx_pkt = queue_data_new_from_pool(etcp->rx_pool);
rx_pkt->seq=seq;
rx_pkt->timestamp=pkt->timestamp;
rx_pkt->ll.dgram=payload_data;
rx_pkt->ll.size=pkt_len;
// Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len);
queue_data_put(etcp->recv_q, rx_pkt, seq);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: packet seq=%u added to recv_q, calling assembly (last_delivered_id=%u)", seq, etcp->last_delivered_id);
if (etcp->last_delivered_id+1==seq) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов
}
}
len=0;
break;
}
default:
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_conn_input: unknown section type=0x%02x", type);
len=0;
break;
}
}
memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram
}

98
src/pkt_normalizer.c

@ -19,7 +19,7 @@ static void packer_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);
static void pn_send_to_etcp(struct PKTNORM* pn);
// Initialization
struct PKTNORM* pn_init(struct ETCP_CONN* etcp) {
@ -44,7 +44,8 @@ struct PKTNORM* pn_init(struct ETCP_CONN* etcp) {
queue_set_callback(pn->input, packer_cb, pn);
queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn);
pn->sndpart = NULL;
pn->sndpart.dgram = NULL;
pn->sndpart.len = 0;
pn->recvpart = NULL;
pn->flush_timer = NULL;
@ -63,7 +64,7 @@ void pn_pair_deinit(struct PKTNORM* pn) {
if (entry->dgram) {
free(entry->dgram);
}
queue_data_free(data);
queue_entry_free(data);
}
queue_free(pn->input);
}
@ -74,7 +75,7 @@ void pn_pair_deinit(struct PKTNORM* pn) {
if (entry->dgram) {
free(entry->dgram);
}
queue_data_free(data);
queue_entry_free(data);
}
queue_free(pn->output);
}
@ -83,13 +84,12 @@ void pn_pair_deinit(struct PKTNORM* pn) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
}
if (pn->sndpart) {
ll_free_dgram(pn->sndpart);
queue_data_free(pn->sndpart);
if (pn->sndpart.dgram) {
queue_entry_free(&pn->sndpart);
}
if (pn->recvpart) {
ll_free_dgram(pn->recvpart);
queue_data_free(pn->recvpart);
queue_dgram_free(pn->recvpart);
queue_entry_free(pn->recvpart);
}
free(pn);
@ -99,8 +99,8 @@ void pn_pair_deinit(struct PKTNORM* pn) {
void pn_unpacker_reset_state(struct PKTNORM* pn) {
if (!pn) return;
if (pn->recvpart) {
ll_free_dgram(pn->recvpart);
queue_data_free(pn->recvpart);
queue_dgram_free(pn->recvpart);
queue_entry_free(pn->recvpart);
pn->recvpart = NULL;
}
}
@ -133,53 +133,37 @@ static void packer_cb(struct ll_queue* q, void* 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;
// Allocate data
uint8_t* packet_data = malloc(entry->len);
if (!packet_data) return;
memcpy(packet_data, entry->dgram, entry->len);
static void pn_send_to_etcp(struct PKTNORM* pn) {
if (!pn || !pn->sndpart.data || !pn->sndpart.len==0) return;
// Allocate ETCP_FRAGMENT from rx_pool
struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool);
if (!frag) {
free(packet_data);
struct ETCP_FRAGMENT* frag = queue_entry_new_from_pool(pn->etcp->rx_pool);
if (!frag) {// drop data
pn->sndpart.len = 0;
return;
}
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;
frag->ll.dgram = pn->sndpart.dgram;
frag->ll.len = pn->sndpart.len;
frag->ll.memlen = pn->sndpart.memlen;
queue_data_put(pn->etcp->input_queue, frag, 0);
queue_entry_free(&pn->sndpart);
}
// 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.dgram) {
int remain = pn->frag_size - pn->sndpart.len;
if (remain < 3) pn_send_to_etcp(pn);
}
if (!pn->sndpart) {
pn->sndpart = ll_alloc_lldgram(pn->frag_size);
if (pn->sndpart) {
pn->sndpart->len = 0;
}
if (!pn->sndpart.dgram) {
pn->sndpart.len=0;
pn->sndpart.dgram_pool = pn->etcp->instance->data_pool;
pn->sndpart.memlen=pn->etcp->instance->data_pool->object_size;//pn->frag_size;
pn->sndpart.dgram = memory_pool_alloc(pn->etcp->instance->data_pool);
}
}
@ -199,32 +183,32 @@ static void etcp_input_ready_cb(struct ll_queue* q, void* arg) {
while (ptr < in_dgram->len) {
pn_buf_renew(pn);
if (!pn->sndpart) break; // Allocation failed
if (!pn->sndpart.dgram) break; // Allocation failed
int remain = pn->frag_size - pn->sndpart->len;
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;
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;
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);
queue_dgram_free(in_dgram);
queue_entry_free(data);
// Cancel flush timer if active
if (pn->flush_timer) {
@ -246,13 +230,7 @@ static void pn_flush_cb(void* arg) {
if (!pn) return;
pn->flush_timer = NULL;
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
}
pn_send_to_etcp(pn);
}
// Internal: Unpacker callback (assembles fragments into original packets)

197
src/pkt_normalizer.c1

@ -0,0 +1,197 @@
// pkt_normalizer.c - Implementation of packet normalizer for ETCP
#include "pkt_normalizer.h"
#include "etcp.h" // For ETCP_CONN and related structures
#include "ll_queue.h" // For queue operations
#include "u_async.h" // For UASYNC
#include <stdlib.h>
#include <string.h>
#include <stdio.h> // For debugging (can be removed if not needed)
// Internal helper to convert void* data to struct ll_entry*
static inline struct ll_entry* data_to_entry(void* data) {
if (!data) return NULL;
return (struct ll_entry*)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);
// доработать init/deinit/reset под новые структуры
// Initialization
struct PKTNORM* pn_init(struct ETCP_CONN* etcp) {
if (!etcp) return NULL;
struct PKTNORM* pn = calloc(1, sizeof(struct PKTNORM));
if (!pn) return NULL;
pn->etcp = etcp;
pn->ua = etcp->instance->ua;
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
if (!pn->input || !pn->output || !pn->pending || !pn->send_pending) {
pn_pair_deinit(pn);
return NULL;
}
return pn;
}
// Deinitialization
void pn_pair_deinit(struct PKTNORM* pn) {
if (!pn) return;
// Drain and free queues
if (pn->input) {
void* data;
while ((data = queue_data_get(pn->input)) != NULL) {
struct ll_entry* entry = data_to_entry(data);
if (entry->dgram) {
memory_pool_free(entry->dgram_pool, entry->dgram);
}
queue_data_free(data);
}
queue_free(pn->input);
}
if (pn->output) {
void* data;
while ((data = queue_data_get(pn->output)) != NULL) {
struct ll_entry* entry = data_to_entry(data);
if (entry->dgram) {
memory_pool_free(entry->dgram_pool, entry->dgram);
}
queue_data_free(data);
}
queue_free(pn->output);
}
if (pn->flush_timer) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
}
free(pn);
}
// Reset unpacker state
void pn_unpacker_reset_state(struct PKTNORM* pn) {
if (!pn) return;
free recvpart
}
// 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;
struct ll_entry* entry = queue_data_new(0);
if (!entry) return;
entry->dgram = malloc(len);
if (!entry->dgram) {
queue_data_free(entry);
return;
}
memcpy(entry->dgram, data, len);
entry->len = len;
entry->dgram_pool = NULL;
queue_data_put(pn->input, entry, 0);
}
static void pn_buf_renew(struct PKTNORM* pn) {
if (!pn->sndpart) { pn->sndpart=ll_alloc_lldgram(pn); return; }
int remain=pn->sndpart->len - pn->frag_size;
if (remain<3) {
queue_data_put(pn->etcp->input_queue, pn->sndpart, 0);
pn->sndpart=ll_alloc_lldgram(pn);
}
}
// Internal: Flush callback on timeout
static void pn_flush_cb(void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
pn->flush_timer = NULL;
if (pn->sndpart && pn->sndpart->len>0) {
queue_data_put(pn->etcp->input_queue, pn->sndpart, 0);
pn->sndpart=ll_alloc_lldgram(pn);
}
}
// это основа. под нее надо подстроить остальное.
static void etcp_input_ready_cb(struct ll_queue* q, void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
// собираем и по мере отправляем
struct ll_queue* in_dgram = queue_data_get(pn->inpupt);
if (in_dgram) {
uint16_t ptr=0;//in_dgram->len;
while (ptr < in_dgram->len) {
pn_buf_renew(pn);
if (pn->sndpart) {
int remain=pn->sndpart->len-pn->frag_size;// свободного места в пакете
if (remain<3) {DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_normalizer: part size error, remain=%d", remain); return;}
if (len==0) {
pn->sndpart->dgram[ptr++]=in_dgram->len;
pn->sndpart->dgram[ptr++]=in_dgram->len>>8;
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);
}
}
if (pn->flush_timer) uasync_cancel_timeout(pn->flush_timer);
if (in->input->count==0) pn->flush_timer=uasync_set_timeout(pn->ua, pn->tx_wait_time, pn, pn_flush_cb);
}
queue_resume_callback(pn->inpupt);
}
// Internal: Packer callback (aggregates small, fragments large, sends chunks)
static void packer_cb(struct ll_queue* q, void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
queue_wait_threshold(pn->etcp->input_queue, 0, 0, etcp_input_ready_cb, pn);// ждём освобождение входной очереди etcp
}
// 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: Process incoming fixed-size payload chunk from etcp
// надо доработать под новый формат
void pn_unpacker_cb(struct ll_queue* q, void* arg) {// получение из выходной очереди etcp - надо подвязать в init
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
// сборка: выделение нужного кол-ва памяти по заголовку, далее наполнение. когда заполнили - в очередь и наполняем следующий.
}

2
src/pkt_normalizer.h

@ -19,7 +19,7 @@ struct PKTNORM {
// packer:
uint16_t frag_size; // размер фрагмента на которые разбивать
struct ll_entry* sndpart; // блок ожидающий досборки
struct ll_entry sndpart; // блок ожидающий досборки
void* flush_timer; // For timeout flush
struct ll_entry* pending; // Partial processed input entry
uint16_t pending_in_ptr; // Pointer in pending entry

50
src/pkt_normalizer.h1

@ -0,0 +1,50 @@
// pkt_normalizer.h (упрощенная версия)
#ifndef PKT_NORMALIZER_H
#define PKT_NORMALIZER_H
#include "../lib/ll_queue.h"
#include "../lib/u_async.h"
#include <stdint.h>
// Структура для packer
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;
// packer:
uint16_t frag_size; // размер фрагмента на которые разбивать
struct ll_entry* sndpart; // блок ожидающий досборки
void* flush_timer; // For timeout flush
// unpacker:
struct ll_entry* recvpart; // блок ожидающий заполнение
};
// Инициализация пары
struct PKTNORM* pn_init(struct ETCP_CONN* etcp);// все что нужно (в т.ч. mtu и ua) берет из etcp
// Деинициализация пары
void pn_pair_deinit(struct PKTNORM* pn);
// Сброс состояния (для unpacker)
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)
собирает по возможности полные пакеты с размером frag_size.
+ timeout: неполный пакет отправляется по таймауту 1ms (если очереди пустые)
+ входящая очередь etcp не должне наполняться для минимизации задержки - новый пакет отправляем только когда очередь пустая
*/
#endif // PKT_NORMALIZER_H

28
src/req.txt

@ -0,0 +1,28 @@
Возможности utun:
- роутинг между нодами - напрямую (предпочтительный вариант), или через наилучшего кандидата.
Наилучший кандидат - узел с минимальным overall score.
overall score = правило выбираемок пользователем
% loss multiplier:
<0.3% - x1
<1% - x0.5
<5% - x0.2
<10% - x0.1
> - x0.03
- либо minimum delay
Формат данных:
tun -> (src/dst ip) -> routing table -> next hop -> transmit to -> etcp -> routing -> tun
routing: выбирает по dst ip next hop.
next hop = {type + conrol struct}. type= {1-tun, 2-etcp}
control dgram = connect request - если нет прямого соединения
при установке соединения сервер отправляет ip:port клиенту
routing dgram:
<routing c.s.> -> <type, dgram>

BIN
tests/test_config_debug

Binary file not shown.

BIN
tests/test_crypto

Binary file not shown.

BIN
tests/test_debug_categories

Binary file not shown.

BIN
tests/test_ecc_encrypt

Binary file not shown.

BIN
tests/test_etcp_100_packets

Binary file not shown.

8
tests/test_etcp_100_packets.c

@ -150,7 +150,7 @@ static void check_received_packets_fwd(void) {
struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(conn->output_queue)) != NULL) {
if (pkt->ll.size >= PACKET_SIZE) {
if (pkt->ll.len >= PACKET_SIZE) {
int seq = pkt->ll.dgram[0];
uint8_t expected[PACKET_SIZE];
@ -167,7 +167,7 @@ static void check_received_packets_fwd(void) {
if (pkt->ll.dgram) {
memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
}
queue_data_free(pkt);
queue_entry_free(pkt);
}
}
@ -180,7 +180,7 @@ static void check_received_packets_back(void) {
struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(conn->output_queue)) != NULL) {
if (pkt->ll.size >= PACKET_SIZE) {
if (pkt->ll.len >= PACKET_SIZE) {
int seq = pkt->ll.dgram[0];
uint8_t expected[PACKET_SIZE];
@ -197,7 +197,7 @@ static void check_received_packets_back(void) {
if (pkt->ll.dgram) {
memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
}
queue_data_free(pkt);
queue_entry_free(pkt);
}
}

BIN
tests/test_etcp_crypto

Binary file not shown.

BIN
tests/test_etcp_minimal

Binary file not shown.

BIN
tests/test_etcp_simple_traffic

Binary file not shown.

4
tests/test_etcp_simple_traffic.c

@ -160,7 +160,7 @@ static void check_packet_received(void) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Found packet in output queue");
// ETCP_FRAGMENT содержит ll_entry в начале, данные в pkt_data
size_t size = pkt->ll.size;
size_t size = pkt->ll.len;
uint8_t* data = pkt->ll.dgram;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Packet size=%zu, expected=%d", size, PACKET_SIZE);
@ -179,7 +179,7 @@ static void check_packet_received(void) {
if (pkt->ll.dgram) {
memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
}
queue_data_free(pkt);
queue_entry_free(pkt);
} else {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "check_packet_received: No packet found in output queue");
}

BIN
tests/test_etcp_two_instances

Binary file not shown.

BIN
tests/test_intensive_memory_pool

Binary file not shown.

6
tests/test_intensive_memory_pool.c

@ -46,7 +46,7 @@ static double test_without_pools(int iterations) {
for (int i = 0; i < 10; i++) {
void* retrieved = queue_data_get(queue);
if (retrieved) {
queue_data_free(retrieved);
queue_entry_free(retrieved);
}
}
@ -89,7 +89,7 @@ static double test_with_pools(int iterations) {
// Добавить записи из пула
for (int i = 0; i < 10; i++) {
void* data = queue_data_new_from_pool(pool);
void* data = queue_entry_new_from_pool(pool);
if (data) {
queue_data_put(queue, data, i); // Используем ID = i
}
@ -99,7 +99,7 @@ static double test_with_pools(int iterations) {
for (int i = 0; i < 10; i++) {
void* retrieved = queue_data_get(queue);
if (retrieved) {
queue_data_free(retrieved); // Вернется в пул
queue_entry_free(retrieved); // Вернется в пул
}
}

6
tests/test_intensive_memory_pool_new.c

@ -46,7 +46,7 @@ static double test_without_pools(int iterations) {
for (int i = 0; i < 10; i++) {
void* retrieved = queue_data_get(queue);
if (retrieved) {
queue_data_free(retrieved);
queue_entry_free(retrieved);
}
}
@ -89,7 +89,7 @@ static double test_with_pools(int iterations) {
// Добавить записи из пула
for (int i = 0; i < 10; i++) {
void* data = queue_data_new_from_pool(pool);
void* data = queue_entry_new_from_pool(pool);
if (data) {
queue_data_put(queue, data, i); // Используем ID = i
}
@ -99,7 +99,7 @@ static double test_with_pools(int iterations) {
for (int i = 0; i < 10; i++) {
void* retrieved = queue_data_get(queue);
if (retrieved) {
queue_data_free(retrieved); // Вернется в пул
queue_entry_free(retrieved); // Вернется в пул
}
}

BIN
tests/test_ll_queue

Binary file not shown.

28
tests/test_ll_queue.c

@ -62,7 +62,7 @@ static void queue_cb(struct ll_queue *q, void *arg) {
test_data_t *taken = queue_data_get(q);
ASSERT(checksum(taken) == taken->checksum, "corruption in cb");
queue_data_free(taken);
queue_entry_free(taken);
queue_resume_callback(q); // next item when ready
}
@ -100,7 +100,7 @@ static void test_fifo(void) {
for (int i = 0; i < 10; i++) {
test_data_t *d = queue_data_get(q);
ASSERT(d && d->id == i, "FIFO violation");
queue_data_free(d);
queue_entry_free(d);
}
ASSERT_EQ(queue_entry_count(q), 0, "");
queue_free(q); uasync_destroy(ua, 0);
@ -126,12 +126,12 @@ static void test_lifo_priority(void) {
test_data_t *first = queue_data_get(q);
ASSERT(first && first->id == 999, "priority first");
queue_data_free(first);
queue_entry_free(first);
for (int i = 0; i < 3; i++) {
test_data_t *d = queue_data_get(q);
ASSERT(d && d->id == i, "remaining FIFO");
queue_data_free(d);
queue_entry_free(d);
}
queue_free(q); uasync_destroy(ua, 0);
PASS();
@ -177,7 +177,7 @@ static void test_waiter(void) {
w = queue_wait_threshold(q, 2, 0, waiter_cb, &called);
ASSERT(w != NULL && called == 0, "");
for (int i = 0; i < 3; i++) queue_data_free(queue_data_get(q));
for (int i = 0; i < 3; i++) queue_entry_free(queue_data_get(q));
for (int i = 0; i < 15; i++) uasync_poll(ua, 1);
ASSERT_EQ(called, 1, "");
@ -205,7 +205,7 @@ static void test_limits_hash(void) {
queue_remove_data(q, found);
ASSERT(queue_find_data_by_id(q, 21) == NULL, "removed");
while (queue_entry_count(q)) queue_data_free(queue_data_get(q));
while (queue_entry_count(q)) queue_entry_free(queue_data_get(q));
queue_free(q); uasync_destroy(ua, 0);
PASS();
}
@ -219,23 +219,23 @@ static void test_pool(void) {
size_t alloc1 = 0, reuse1 = 0;
memory_pool_get_stats(pool, &alloc1, &reuse1);
test_data_t *d1 = queue_data_new_from_pool(pool);
test_data_t *d1 = queue_entry_new_from_pool(pool);
d1->id = 1; d1->checksum = checksum(d1);
queue_data_put(q, d1, 1);
test_data_t *d2 = queue_data_new_from_pool(pool);
test_data_t *d2 = queue_entry_new_from_pool(pool);
d2->id = 2; d2->checksum = checksum(d2);
queue_data_put(q, d2, 2);
queue_data_free(queue_data_get(q)); // free d1 back to pool
queue_data_free(queue_data_get(q)); // free d2 back to pool
queue_entry_free(queue_data_get(q)); // free d1 back to pool
queue_entry_free(queue_data_get(q)); // free d2 back to pool
// Now allocate again to trigger reuse
test_data_t *d3 = queue_data_new_from_pool(pool);
test_data_t *d3 = queue_entry_new_from_pool(pool);
ASSERT(d3 != NULL, "alloc after free failed");
d3->id = 3; d3->checksum = checksum(d3);
queue_data_put(q, d3, 3);
queue_data_free(queue_data_get(q)); // free d3 back
queue_entry_free(queue_data_get(q)); // free d3 back
size_t alloc2 = 0, reuse2 = 0;
memory_pool_get_stats(pool, &alloc2, &reuse2);
@ -253,7 +253,7 @@ static void test_stress(void) {
for (int i = 0; i < 10000; i++) {
if (queue_entry_count(q) > 80) {
queue_data_free(queue_data_get(q));
queue_entry_free(queue_data_get(q));
}
test_data_t *d = queue_data_new(sizeof(*d));
d->id = rand() % 10000;
@ -261,7 +261,7 @@ static void test_stress(void) {
queue_data_put(q, d, d->id);
stats.ops++;
}
while (queue_entry_count(q)) queue_data_free(queue_data_get(q));
while (queue_entry_count(q)) queue_entry_free(queue_data_get(q));
stats.time_ms += now_ms() - start;
queue_free(q); uasync_destroy(ua, 0);

BIN
tests/test_memory_pool_and_config

Binary file not shown.

2
tests/test_memory_pool_and_config.c

@ -65,7 +65,7 @@ int main() {
for (int i = 0; i < 5; i++) {
void* retrieved = queue_data_get(queue);
if (retrieved) {
queue_data_free(retrieved);
queue_entry_free(retrieved);
}
}

BIN
tests/test_packet_dump

Binary file not shown.

BIN
tests/test_pkt_normalizer_etcp

Binary file not shown.

8
tests/test_pkt_normalizer_etcp.c

@ -180,8 +180,8 @@ static void check_received_packets_fwd(void) {
}
}
ll_free_dgram(entry);
queue_data_free(data);
queue_dgram_free(entry);
queue_entry_free(data);
}
}
@ -203,8 +203,8 @@ static void check_received_packets_back(void) {
}
}
ll_free_dgram(entry);
queue_data_free(data);
queue_dgram_free(entry);
queue_entry_free(data);
}
}

BIN
tests/test_u_async_comprehensive

Binary file not shown.

BIN
tests/test_u_async_performance

Binary file not shown.
Loading…
Cancel
Save