diff --git a/lib/ll_queue.c b/lib/ll_queue.c index da2426e4..c8d5da8c 100755 --- a/lib/ll_queue.c +++ b/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; } diff --git a/lib/ll_queue.h b/lib/ll_queue.h index 4a6a55ec..257e5741 100644 --- a/lib/ll_queue.h +++ b/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); // Получить текущее количество элементов в очереди diff --git a/src/aa b/src/aa new file mode 100755 index 00000000..e6e74cf1 --- /dev/null +++ b/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 раза в неделю для поддержки diff --git a/src/etcp.c b/src/etcp.c index bd74ca9a..60aba3c7 100644 --- a/src/etcp.c +++ b/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); diff --git a/src/etcp.c1 b/src/etcp.c1 new file mode 100755 index 00000000..35fae7f5 --- /dev/null +++ b/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 +#include +#include +#include // For bandwidth calcs +#include // 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; irx_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 +} + + diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index c15bb685..c3a4b1bd 100644 --- a/src/pkt_normalizer.c +++ b/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) diff --git a/src/pkt_normalizer.c1 b/src/pkt_normalizer.c1 new file mode 100755 index 00000000..0d5aabd4 --- /dev/null +++ b/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 +#include +#include // 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; + +// сборка: выделение нужного кол-ва памяти по заголовку, далее наполнение. когда заполнили - в очередь и наполняем следующий. + +} diff --git a/src/pkt_normalizer.h b/src/pkt_normalizer.h index 7b5552f8..dc266ca0 100644 --- a/src/pkt_normalizer.h +++ b/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 diff --git a/src/pkt_normalizer.h1 b/src/pkt_normalizer.h1 new file mode 100755 index 00000000..cf28512d --- /dev/null +++ b/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 + +// Структура для 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 \ No newline at end of file diff --git a/src/req.txt b/src/req.txt new file mode 100755 index 00000000..1ee6d964 --- /dev/null +++ b/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: + -> diff --git a/tests/test_config_debug b/tests/test_config_debug new file mode 100755 index 00000000..b5eb4cfd Binary files /dev/null and b/tests/test_config_debug differ diff --git a/tests/test_crypto b/tests/test_crypto new file mode 100755 index 00000000..248b02b3 Binary files /dev/null and b/tests/test_crypto differ diff --git a/tests/test_debug_categories b/tests/test_debug_categories new file mode 100755 index 00000000..52f5edc8 Binary files /dev/null and b/tests/test_debug_categories differ diff --git a/tests/test_ecc_encrypt b/tests/test_ecc_encrypt new file mode 100755 index 00000000..ab4724fb Binary files /dev/null and b/tests/test_ecc_encrypt differ diff --git a/tests/test_etcp_100_packets b/tests/test_etcp_100_packets index 4b557ed3..4e0b4a6c 100755 Binary files a/tests/test_etcp_100_packets and b/tests/test_etcp_100_packets differ diff --git a/tests/test_etcp_100_packets.c b/tests/test_etcp_100_packets.c index 5edd3ecd..ef17b942 100644 --- a/tests/test_etcp_100_packets.c +++ b/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); } } diff --git a/tests/test_etcp_crypto b/tests/test_etcp_crypto new file mode 100755 index 00000000..f2aee3a4 Binary files /dev/null and b/tests/test_etcp_crypto differ diff --git a/tests/test_etcp_minimal b/tests/test_etcp_minimal new file mode 100755 index 00000000..6e3df713 Binary files /dev/null and b/tests/test_etcp_minimal differ diff --git a/tests/test_etcp_simple_traffic b/tests/test_etcp_simple_traffic new file mode 100755 index 00000000..a03b9bf0 Binary files /dev/null and b/tests/test_etcp_simple_traffic differ diff --git a/tests/test_etcp_simple_traffic.c b/tests/test_etcp_simple_traffic.c index 0b9847b4..78520d72 100644 --- a/tests/test_etcp_simple_traffic.c +++ b/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"); } diff --git a/tests/test_etcp_two_instances b/tests/test_etcp_two_instances new file mode 100755 index 00000000..fdbd05c5 Binary files /dev/null and b/tests/test_etcp_two_instances differ diff --git a/tests/test_intensive_memory_pool b/tests/test_intensive_memory_pool new file mode 100755 index 00000000..1c6819f7 Binary files /dev/null and b/tests/test_intensive_memory_pool differ diff --git a/tests/test_intensive_memory_pool.c b/tests/test_intensive_memory_pool.c index 4de4032c..895549bc 100644 --- a/tests/test_intensive_memory_pool.c +++ b/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); // Вернется в пул } } diff --git a/tests/test_intensive_memory_pool_new.c b/tests/test_intensive_memory_pool_new.c index 3069b4ad..ec953302 100644 --- a/tests/test_intensive_memory_pool_new.c +++ b/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); // Вернется в пул } } diff --git a/tests/test_ll_queue b/tests/test_ll_queue new file mode 100755 index 00000000..7d682022 Binary files /dev/null and b/tests/test_ll_queue differ diff --git a/tests/test_ll_queue.c b/tests/test_ll_queue.c index 497f5375..0e37b19a 100644 --- a/tests/test_ll_queue.c +++ b/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); diff --git a/tests/test_memory_pool_and_config b/tests/test_memory_pool_and_config new file mode 100755 index 00000000..e4fb62f4 Binary files /dev/null and b/tests/test_memory_pool_and_config differ diff --git a/tests/test_memory_pool_and_config.c b/tests/test_memory_pool_and_config.c index 58dd2c19..c1f94ccf 100644 --- a/tests/test_memory_pool_and_config.c +++ b/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); } } diff --git a/tests/test_packet_dump b/tests/test_packet_dump new file mode 100755 index 00000000..5efad3ef Binary files /dev/null and b/tests/test_packet_dump differ diff --git a/tests/test_pkt_normalizer_etcp b/tests/test_pkt_normalizer_etcp new file mode 100755 index 00000000..e1b90923 Binary files /dev/null and b/tests/test_pkt_normalizer_etcp differ diff --git a/tests/test_pkt_normalizer_etcp.c b/tests/test_pkt_normalizer_etcp.c index db2e2cdc..311c6510 100644 --- a/tests/test_pkt_normalizer_etcp.c +++ b/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); } } diff --git a/tests/test_u_async_comprehensive b/tests/test_u_async_comprehensive new file mode 100755 index 00000000..779d6361 Binary files /dev/null and b/tests/test_u_async_comprehensive differ diff --git a/tests/test_u_async_performance b/tests/test_u_async_performance new file mode 100755 index 00000000..bc2ac05c Binary files /dev/null and b/tests/test_u_async_performance differ