From a3783b5f557e2d51154205e45a32ccdba7918cb5 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 6 Aug 2026 23:14:57 +0300 Subject: [PATCH] fix: re-index conn_queue_entry after setting peer_node_id tcp_server_on_link sets conn->peer_node_id after etcp_connection_create which adds the entry with peer_node_id=0. instance_find_conn couldnt find TCP connections until the queue entry was re-indexed. --- src/transport_layer/etcp.c | 2012 ------------------------ src/transport_layer/etcp_connections.c | 1 + src/transport_layer/stcp_link.c | 286 ---- tests/test_etcp_stcp.c | 153 -- tests/test_stcp_traffic.c | 234 --- 5 files changed, 1 insertion(+), 2685 deletions(-) delete mode 100644 src/transport_layer/etcp.c delete mode 100644 src/transport_layer/stcp_link.c delete mode 100644 tests/test_etcp_stcp.c delete mode 100644 tests/test_stcp_traffic.c diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c deleted file mode 100644 index d0c1faa5..00000000 --- a/src/transport_layer/etcp.c +++ /dev/null @@ -1,2012 +0,0 @@ -// etcp.c - ETCP Protocol Implementation (refactored and expanded based on etcp_protocol.txt) - -#include "etcp.h" -#include "etcp_connect.h" -#include "etcp_debug.h" -#include "etcp_loadbalancer.h" -#include "etcp_router.h" -#include "routing.h" -#include "topo_group.h" -#include "route_ping.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 -#include // For strftime in metrics snapshot -#include "../lib/mem.h" -#include "../lib/memory_pool.h" - -// 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 32 // RTT multiplier for retrans timeout -#define RETRANS_K2 32 // 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 send_ack_req_cb(struct ll_queue* q, void* arg); -static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp); -static void etcp_connection_free_deferred(void* arg); -void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status); -struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp); - -static void clear_queue(struct ll_queue* q) { - if (!q) return; - struct ll_entry* pkt; - while ((pkt = queue_data_get(q)) != NULL) { - queue_dgram_free(pkt); - queue_entry_free(pkt); - } -} - -static void drain_and_free_queue(struct ll_queue** q) { - if (!*q) return; - clear_queue(*q); - queue_free(*q); - *q = NULL; -} - -static int inflight_seq_cmp(const void* a, const void* b) { - const struct INFLIGHT_PACKET* pa = *(const struct INFLIGHT_PACKET**)a; - const struct INFLIGHT_PACKET* pb = *(const struct INFLIGHT_PACKET**)b; - if (pa->seq < pb->seq) return -1; - if (pa->seq > pb->seq) return 1; - return 0; -} - -static void feed_dgram_to_asm(struct ETCP_CONN* etcp, struct PKTNORM* pn, - uint8_t* dgram, uint16_t dgram_len, - uint8_t** asm_buf, uint16_t* asm_len, uint32_t* asm_cap, - uint32_t* returned) { - uint32_t need = (uint32_t)*asm_len + dgram_len; - if (need >= ASM_BUF_MAX_SIZE) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] assembly buffer would exceed max size %d", etcp->log_name, ASM_BUF_MAX_SIZE); - return; - } - if (need > *asm_cap) { - uint32_t new_cap = *asm_cap ? *asm_cap : 256; - while (new_cap < need) new_cap *= 2; - if (new_cap > 0xFFFF) new_cap = 0xFFFF; - uint8_t* new_buf = u_realloc(*asm_buf, (uint16_t)new_cap); - if (!new_buf) return; - *asm_buf = new_buf; - *asm_cap = new_cap; - } - memcpy(*asm_buf + *asm_len, dgram, dgram_len); - *asm_len = need; - - while (*asm_len >= 2) { - uint16_t pkt_len = (*asm_buf)[0] | ((*asm_buf)[1] << 8); - if (pkt_len == 0 || pkt_len > 16384 || *asm_len < 2 + pkt_len) break; - - const uint8_t* pkt_data = *asm_buf + 2; - if (pkt_data[0] != ETCP_ID_TOPO_ENTRY) { - struct ll_entry* e = ll_alloc_lldgram(pkt_len); - if (e) { - memcpy(e->dgram, pkt_data, pkt_len); - e->len = pkt_len; - queue_data_put(pn->input, e); - (*returned)++; - } - } - - uint16_t consumed = 2 + pkt_len; - *asm_len -= consumed; - if (*asm_len > 0) memmove(*asm_buf, *asm_buf + consumed, *asm_len); - } -} - -static void etcp_return_inflight_to_normalizer(struct ETCP_CONN* etcp) { - struct PKTNORM* pn = etcp->normalizer; - if (!pn) return; - - struct INFLIGHT_PACKET** inflight = NULL; - int inflight_count = 0; - struct ll_entry** input_frags = NULL; - int input_count = 0; - struct ll_entry* entry; - - while ((entry = queue_data_get(etcp->input_wait_ack))) { - inflight = u_realloc(inflight, (inflight_count + 1) * sizeof(*inflight)); - inflight[inflight_count++] = (struct INFLIGHT_PACKET*)entry; - } - while ((entry = queue_data_get(etcp->input_send_q))) { - inflight = u_realloc(inflight, (inflight_count + 1) * sizeof(*inflight)); - inflight[inflight_count++] = (struct INFLIGHT_PACKET*)entry; - } - while ((entry = queue_data_get(etcp->input_queue))) { - input_frags = u_realloc(input_frags, (input_count + 1) * sizeof(*input_frags)); - input_frags[input_count++] = entry; - } - - if (inflight_count + input_count == 0 && (!pn->data || pn->data_ptr == 0)) return; - - if (inflight_count > 0) qsort(inflight, inflight_count, sizeof(*inflight), inflight_seq_cmp); - - uint8_t* asm_buf = NULL; - uint16_t asm_len = 0; - uint32_t asm_cap = 0; - uint32_t returned_packets = 0; - - for (int i = 0; i < inflight_count; i++) { - feed_dgram_to_asm(etcp, pn, inflight[i]->ll.dgram, inflight[i]->ll.len, - &asm_buf, &asm_len, &asm_cap, &returned_packets); - memory_pool_free(etcp->instance->data_pool, inflight[i]->ll.dgram); - memory_pool_free(etcp->inflight_pool, inflight[i]); - } - u_free(inflight); - - for (int i = 0; i < input_count; i++) { - feed_dgram_to_asm(etcp, pn, input_frags[i]->dgram, input_frags[i]->len, - &asm_buf, &asm_len, &asm_cap, &returned_packets); - memory_pool_free(etcp->instance->data_pool, input_frags[i]->dgram); - memory_pool_free(etcp->io_pool, input_frags[i]); - } - u_free(input_frags); - - if (pn->data && pn->data_ptr > 0) { - feed_dgram_to_asm(etcp, pn, pn->data, pn->data_ptr, - &asm_buf, &asm_len, &asm_cap, &returned_packets); - } - - if (asm_buf) { u_free(asm_buf); } - - if (returned_packets > 0) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] returned %u packets to normalizer input during reinit", - etcp->log_name, returned_packets); - } -} - -uint16_t get_current_timestamp() { - return (uint16_t)get_time_tb(); -} - -// Timestamp diff (with wrap-around) -static uint16_t timestamp_diff(uint16_t t1, uint16_t t2) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (t1 >= t2) { - return t1 - t2; - } - return (UINT16_MAX - t2) + t1 + 1; -} - -static const char EMPTY_NAME[] = ""; - -// Create new ETCP connection -struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (!instance) return NULL; - - - struct ETCP_CONN* etcp = u_calloc(1, sizeof(struct ETCP_CONN)); - if (!etcp) { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "creating connection failed for instance %p", instance); - return NULL; - } - - etcp->instance = instance; - etcp->input_queue = queue_new(instance->ua, 0, 0, 0, "ETCP input"); // No hash for input_queue - queue_set_threshold(etcp->input_queue, 0, 0); // Backpressure: ждать полного освобождения - etcp->output_queue = queue_new(instance->ua, 0, 0, 0, "ETCP output"); // No hash for output_queue - etcp->transit_queues = queue_new(instance->ua, TRANSIT_QUEUE_HASH, 0, 16, "transit_q_reg"); // Hash for transit queues - etcp->input_send_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "input_send_q"); // Hash for send_q - etcp->input_wait_ack = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "input_wait_ack"); // Hash for wait_ack - etcp->recv_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "recv_q"); // Hash for recv_q - etcp->ack_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "ack_q"); - etcp->inflight_pool = memory_pool_init(sizeof(struct INFLIGHT_PACKET), "inflight_pool"); - etcp->io_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT), "io_pool"); - etcp->max_inflight = (uint32_t)instance->config->global.bbr_max_cwnd; - etcp->optimal_inflight=100000; - etcp->initialized=0; - etcp->links_up=0; - etcp->setup_start_tb = get_time_tb(); - etcp->reset_done=0; - etcp->callbacks_running=0; - etcp->ref_count=0; - etcp->last_rr_link=NULL; - etcp->name = u_strdup(name); - // Initialize log_name with local node_id (peer will be updated later when known) - snprintf(etcp->log_name, sizeof(etcp->log_name), "%04X->???? [%s]", (uint16_t)instance->node_id, etcp->name); - - if (!etcp->input_queue || !etcp->output_queue || !etcp->transit_queues || - !etcp->input_send_q || !etcp->recv_q || !etcp->ack_q || - !etcp->input_wait_ack || !etcp->inflight_pool || !etcp->io_pool) { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "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->window_size = MAX_INFLIGHT_BYTES; // Not used - etcp->mtu = ETCP_RFC791_MIN_MTU; // Default MTU per RFC 791 - etcp->next_tx_id = 1; - etcp->rtt_avg_10 = 10; // Initial guess (1ms) - etcp->rtt_history_idx = 0; - memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history)); - - - etcp->normalizer = pn_init(etcp); - if (!etcp->normalizer) { - etcp_connection_close(etcp); - return NULL; - } - etcp->send_input_q = etcp->normalizer->input; - - // 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); -// queue_set_callback(etcp->ack_q, send_ack_req_cb, etcp);// отправляем через таймаут - - etcp->link_ready_for_send_fn = etcp_link_ready_callback; - - // Add to instance's connections queue (peer_node_id=0 for pending, reindexed later) - { - struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry)); - if (!qe) { etcp_connection_close(etcp); return NULL; } - struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; - ce->peer_node_id = 0; - ce->conn = etcp; - etcp->conn_queue_entry = qe; - etcp->conn_queue = instance->connections; - queue_data_put_with_index(instance->connections, qe); - } - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] connection initialized. ETCP=%p mtu=%d, next_tx_id=%u state=pending", - etcp->log_name, etcp, etcp->mtu, etcp->next_tx_id); - - // Вызываем callback для нового соединения если установлен - if (instance) { - struct etcp_inst_cbk_entry* cbe = instance->new_conn_cbks; - while (cbe) { struct etcp_inst_cbk_entry* n = cbe->next; cbe->fn(etcp, cbe->arg); cbe = n; } - } - etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_NEW); - - return etcp; -} - - -void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status) { - if (!etcp || !etcp->instance) return; - static const char* names[] = { "NEW", "UP", "DOWN", "DELETE" }; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[%s] Connection status changed to %s", etcp->log_name, - (status >= 0 && status < (int)(sizeof(names)/sizeof(names[0]))) ? names[status] : "?"); - struct etcp_status_cbk_entry* cbe = etcp->instance->conn_status_cbks; - while (cbe) { struct etcp_status_cbk_entry* n = cbe->next; cbe->fn(etcp, status, cbe->arg); cbe = n; } -} - -void etcp_cbk_fire(struct ETCP_CONN* conn, int event) { - if (!conn) return; - static const char* names[] = { "INIT", "REINIT", "UP", "DOWN", "NODE_CHANGED" }; - int idx = 0, e = event; while (e >>= 1) idx++; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] callback event: %s", conn->log_name, (idx >= 0 && idx < (int)(sizeof(names)/sizeof(names[0]))) ? names[idx] : "?"); - conn->callbacks_running = 1; - struct etcp_cbk_entry* cbe = conn->cbks; - while (cbe) { struct etcp_cbk_entry* n = cbe->next; if (cbe->event_mask & event) cbe->fn(conn, event, cbe->arg); cbe = n; } - conn->callbacks_running = 0; -} - -static void etcp_on_up(struct ETCP_CONN* etcp) { - int total_links = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) total_links++; - uint64_t elapsed_ms = (get_time_tb() - etcp->setup_start_tb) / 10; - char links_str[256] = {0}; int pp = 0; - for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { - if (l->link_status) { - if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", "); - pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str); - } - } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection UP (%d/%d links, mtu=%d, reinit=%u)%s%s", - etcp->log_name, etcp->links_up, total_links, etcp->mtu, etcp->reinit_count, - pp > 0 ? " — up: " : "", links_str); - (void)elapsed_ms; - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_ON_UP_BEFORE: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x", - etcp->log_name, - etcp->crypto_ctx.session_key[0], etcp->crypto_ctx.session_key[1], - etcp->crypto_ctx.session_key[2], etcp->crypto_ctx.session_key[3], - etcp->crypto_ctx.peer_public_key[0], etcp->crypto_ctx.peer_public_key[1], - etcp->crypto_ctx.peer_public_key[2], etcp->crypto_ctx.peer_public_key[3]); - etcp_cbk_fire(etcp, ETCP_CBK_EVENT_UP); - etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_UP); - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_ON_UP_AFTER: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x", - etcp->log_name, - etcp->crypto_ctx.session_key[0], etcp->crypto_ctx.session_key[1], - etcp->crypto_ctx.session_key[2], etcp->crypto_ctx.session_key[3], - etcp->crypto_ctx.peer_public_key[0], etcp->crypto_ctx.peer_public_key[1], - etcp->crypto_ctx.peer_public_key[2], etcp->crypto_ctx.peer_public_key[3]); -} - -static void etcp_on_down(struct ETCP_CONN* etcp, struct ETCP_LINK* down_link) { - char links_str[256] = {0}; int pp = 0; - for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { - if (l == down_link) continue; - if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", "); - pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str); - } - if (down_link) - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — link %s lost%s%s", - etcp->log_name, sockaddr_storage_to_str(&down_link->remote_addr).str, - pp > 0 ? ", remaining: " : "", links_str); - else - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — closing, links: %s", etcp->log_name, pp ? links_str : "(none)"); - etcp_cbk_fire(etcp, ETCP_CBK_EVENT_DOWN); - etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DOWN); -} - - -// Phase 2 resources cleanup (callable both sync and async via call_soon) -static void etcp_connection_free_resources(struct ETCP_CONN* etcp) { - if (!etcp) return; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] freeing resources phase 2", etcp->log_name); - - // Close links - 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; - } - - // Drain and free all queues - drain_and_free_queue(&etcp->input_queue); - drain_and_free_queue(&etcp->output_queue); - drain_and_free_queue(&etcp->input_send_q); - drain_and_free_queue(&etcp->input_wait_ack); - drain_and_free_queue(&etcp->recv_q); - drain_and_free_queue(&etcp->ack_q); - - // Free callback chain - { struct etcp_cbk_entry* cbe = etcp->cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->cbks = NULL; } - - if (etcp->inflight_pool) { memory_pool_destroy(etcp->inflight_pool); etcp->inflight_pool = NULL; } - if (etcp->io_pool) { memory_pool_destroy(etcp->io_pool); etcp->io_pool = NULL; } - - etcp_connect_cancel_for_conn(etcp->instance, etcp); - - u_free(etcp->name); - u_free(etcp); -} - -static void etcp_connection_free_deferred(void* arg) { - etcp_connection_free_resources((struct ETCP_CONN*)arg); -} - -// Close connection: phase 1 detach + deferred phase 2 cleanup -void etcp_connection_close(struct ETCP_CONN* etcp) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (!etcp) return; - if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] already deleted", etcp->log_name); return; } - - if (etcp->callbacks_running) { - DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "[%s] FATAL: etcp_connection_close called from inside callback chain — SEGFAULTING to show backtrace", - etcp->log_name); - *(volatile int*)0 = 0; - } - - // === PHASE 1: detach from external world === - - if (etcp->links_up != 0) { etcp->links_up = 0; etcp_on_down(etcp, NULL); } - - // Cancel active timers - if (etcp->retrans_timer) { uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer); etcp->retrans_timer = NULL; } - if (etcp->ack_resp_timer) { uasync_cancel_timeout(etcp->instance->ua, etcp->ack_resp_timer); etcp->ack_resp_timer = NULL; } - etcp_metrics_stop_timer(etcp); - - routing_del_conn(etcp); - etcp_router_transit_queues_destroy(etcp); - - if (etcp->normalizer) { pn_deinit((struct PKTNORM*)etcp->normalizer); etcp->normalizer = NULL; } - - // Close links (detach from socket scan queues) - 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; - } - - // Remove from instance queue - if (etcp->conn_queue && etcp->conn_queue_entry) { - queue_remove_data(etcp->conn_queue, etcp->conn_queue_entry); - queue_entry_free(etcp->conn_queue_entry); - etcp->conn_queue_entry = NULL; - etcp->conn_queue = NULL; - } - - etcp->state = 2; // deleted - - etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DELETE); - - // === PHASE 2: deferred resource cleanup (only if no outstanding refs) === - if (etcp->ref_count == 0) - uasync_call_soon(etcp->instance->ua, etcp, etcp_connection_free_deferred); - // else: cleanup deferred until last etcp_conn_ref_free() -} - -// Take a reference on connection. Returns -1 if already deleted (state==2). -int etcp_conn_ref_take(struct ETCP_CONN* conn) { - if (!conn) return -1; - if (conn->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ref_take on deleted conn", conn->log_name); return -1; } - conn->ref_count++; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ref_take ref_count=%d", conn->log_name, conn->ref_count); - return 0; -} - -// Release a reference. If last ref and conn is deleted, schedule deferred free. -void etcp_conn_ref_free(struct ETCP_CONN* conn) { - if (!conn) return; - if (conn->ref_count <= 0) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ref_free underflow ref_count=%d", conn->log_name, conn->ref_count); return; } - conn->ref_count--; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ref_free ref_count=%d state=%d", conn->log_name, conn->ref_count, conn->state); - if (conn->ref_count == 0 && conn->state == 2) - uasync_call_soon(conn->instance->ua, conn, etcp_connection_free_deferred); -} - -// Reset connection -void etcp_conn_reset(struct ETCP_CONN* etcp) { - // Reset IDs - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "Resetting ETCP instance [%s]", etcp->log_name); - etcp->next_tx_id = 1; - etcp->last_rx_id = 0; - etcp->last_delivered_id = 0; - etcp->rx_ack_till = 0; - - // Устанавливаем флаг ожидания первого пакета - etcp->got_initial_pkt = 0; - - // Reset metrics - etcp->unacked_bytes = 0; - etcp->rtt_last = 0; - etcp->rtt_avg_10 = 0; - etcp->jitter = 0; - etcp->last_rtt_cb_time = 0; - etcp->bytes_sent_total = 0; - etcp->retransmissions_count = 0; - etcp->ack_packets_count = 0; - etcp->rx_dup_count = 0; - etcp->tx_dup_count = 0; - - // Reset RTT history - memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history)); - etcp->rtt_history_idx = 0; - etcp->last_rr_link = NULL; - - // Return unconfirmed inflight data to normalizer input (до нормалайзера) - if (etcp->normalizer) { - etcp->normalizer->input->callback_suspended = 1; - etcp_return_inflight_to_normalizer(etcp); - } - - // Clear queues (keep queue structures) - clear_queue(etcp->input_queue); - clear_queue(etcp->output_queue); - clear_queue(etcp->input_send_q); - clear_queue(etcp->input_wait_ack); - clear_queue(etcp->recv_q); - clear_queue(etcp->ack_q); - - // clear_queue leaves callback_suspended=1, resume to prevent deadlock - queue_resume_callback(etcp->input_queue); - queue_resume_callback(etcp->input_send_q); - queue_resume_callback(etcp->input_wait_ack); - queue_resume_callback(etcp->output_queue); - queue_resume_callback(etcp->ack_q); - if (etcp->normalizer && etcp->normalizer->input) { - clear_queue(etcp->normalizer->input); - queue_resume_callback(etcp->normalizer->input); - } - -// В etcp_conn_reset(), после очистки очередей добавьте: - struct ETCP_LINK* l = etcp->links; - while (l) { - l->inflight_bytes = 0; - l->inflight_packets = 0; - l->delivered_bytes = 0; - l->acked_bytes = 0; - l->acked_packets = 0; - l->last_ack_time_tb = get_time_tb(); - l = l->next; - } - - // Cancel active timers to prevent memory leaks - if (etcp->retrans_timer) { - uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer); - etcp->retrans_timer = NULL; - } - if (etcp->ack_resp_timer) { - uasync_cancel_timeout(etcp->instance->ua, etcp->ack_resp_timer); - etcp->ack_resp_timer = NULL; - } - - etcp->reset_count++; - - // Reset normalizer (packer/unpacker state for reconnection) - if (etcp->normalizer) { - pn_reset(etcp->normalizer); - } - - queue_resume_callback(etcp->input_queue); - queue_resume_callback(etcp->output_queue); - queue_resume_callback(etcp->input_send_q); - queue_resume_callback(etcp->input_wait_ack); - queue_resume_callback(etcp->recv_q); - queue_resume_callback(etcp->ack_q); - - etcp->tx_state=ETCP_TX_STATE_DATA_WAIT; - - - int up=0; - struct ETCP_LINK* link = etcp->links; - while (link) { - if (link->link_status) up=1; - link = link->next; - } - - if (up) { - if (etcp->links_up==0) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "reset conn: set link up"); - etcp->links_up=1; - etcp_on_up(etcp); - } else { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "reset conn: set link down/up"); - etcp_on_down(etcp, NULL); - etcp_on_up(etcp); - } - } else { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "links not up - skip on_down/on_up"); - if (etcp->links_up) { - etcp->links_up=0; - etcp->setup_start_tb = get_time_tb(); - etcp_on_down(etcp, NULL); - } - } - - clear_queue(etcp->input_send_q); - clear_queue(etcp->input_wait_ack); - queue_resume_callback(etcp->input_send_q); - queue_resume_callback(etcp->input_wait_ack); - - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "end"); -} - -void etcp_links_reset(struct ETCP_CONN* etcp) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение - // Сбрасываем initialized во всех линках - struct ETCP_LINK* link = etcp->links; - while (link) { - link->initialized = 0; - link = link->next; - } -} - -void etcp_conn_reinit(struct ETCP_CONN* etcp, const char* reason) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] REINIT: %s (init=%d links=%d reinit=%u tx=%d)", - etcp->log_name, reason, etcp->initialized, etcp->links_up, etcp->reinit_count, etcp->tx_state); - - etcp->setup_start_tb = get_time_tb(); - etcp->reinit_count++; - etcp->reinit_pending = 1; - etcp->reset_done = 0; - etcp->initialized = 0;// еще раз придёт conn_ready_callback - - // Отменяем висящие NAT-ping'и для этого соединения - if (etcp->instance && etcp->instance->nat_det) { - nat_detection_cancel_for_conn(etcp->instance->nat_det, etcp); - } - - // Сбрасываем ретрансмиты роутера ДО очистки ETCP-очередей (избегаем гонки с data_pool) - if (etcp->peer_node_id) etcp_router_pause_retrans_for_node(etcp->instance, etcp->peer_node_id); - - // Сбрасываем NAT-статус у всех линков этого соединения - struct ETCP_LINK* l = etcp->links; - while (l) { - l->nat_type = NAT_TYPE_UNKNOWN; - l->nat_check_status = NAT_CHECK_NONE; - l = l->next; - } - - // Вызываем etcp_conn_reset для сброса состояния - etcp_conn_reset(etcp); - - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "end"); -} - -// внутренняя функция. Вызывается один раз когда первый линк готов. -void etcp_conn_ready(struct ETCP_CONN* conn) { - if (!conn) return; - if (conn->initialized) return; // already ready - - conn->initialized = 1; - conn->reset_done = 1; - if (conn->tx_state == 0) { conn->tx_state = ETCP_TX_STATE_DATA_WAIT; } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Connection ready", conn->log_name); - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_CONN_READY: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x my_pub=%02x%02x%02x%02x links=%p", - conn->log_name, - conn->crypto_ctx.session_key[0], conn->crypto_ctx.session_key[1], - conn->crypto_ctx.session_key[2], conn->crypto_ctx.session_key[3], - conn->crypto_ctx.peer_public_key[0], conn->crypto_ctx.peer_public_key[1], - conn->crypto_ctx.peer_public_key[2], conn->crypto_ctx.peer_public_key[3], - conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[0] : 0, - conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[1] : 0, - conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[2] : 0, - conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[3] : 0, - (void*)conn->links); - - etcp_conn_queue_set_ready(conn); -} - -// Move from pending to indexed connections, fire ready callbacks -void etcp_conn_queue_set_ready(struct ETCP_CONN* conn) { - if (!conn || conn->state != 0) return; - - // Reindex: remove old entry (key=0), add new entry with real peer_node_id - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] queue_set_ready: removing old entry=%p key=0 -> new key=0x%llx", - conn->log_name, (void*)conn->conn_queue_entry, (unsigned long long)conn->peer_node_id); - if (conn->conn_queue && conn->conn_queue_entry) { - queue_remove_data(conn->conn_queue, conn->conn_queue_entry); - queue_entry_free(conn->conn_queue_entry); - } - - struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry)); - if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to alloc queue entry for ready", conn->log_name); return; } - struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; - ce->peer_node_id = conn->peer_node_id; - ce->conn = conn; - conn->conn_queue_entry = qe; - conn->conn_queue = conn->instance->connections; - conn->state = 1; - queue_data_put_with_index(conn->instance->connections, qe); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] queue_set_ready: new entry=%p put in queue, count=%d", - conn->log_name, (void*)qe, queue_entry_count(conn->instance->connections)); - - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] moved to ready queue, state=%d peer_node_id=0x%llx", - conn->log_name, conn->state, (unsigned long long)conn->peer_node_id); - - etcp_metrics_start_timer(conn); - - etcp_cbk_fire(conn, ETCP_CBK_EVENT_INIT); - if (conn->reinit_pending) { conn->reinit_pending = 0; etcp_cbk_fire(conn, ETCP_CBK_EVENT_REINIT); } - if (conn->links_up) etcp_on_up(conn); -} - - - - -// Update log_name when peer_node_id becomes known -void etcp_update_log_name(struct ETCP_CONN* etcp) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (!etcp || !etcp->instance) return; - uint16_t local_id = etcp->instance->node_id; - uint16_t peer_id = etcp->peer_node_id; - const char* name = etcp->name ? etcp->name : EMPTY_NAME; - if (etcp->instance->topo_groups && peer_id) { - struct TOPO_NODE* ni = topo_node_registry_find(etcp->instance->topo_groups, etcp->peer_node_id); - if (ni && ni->node_name && ni->node_name[0]) name = ni->node_name; - } - snprintf(etcp->log_name, sizeof(etcp->log_name), "%04X->%04X [%s]", local_id, peer_id, name); -} - - -// ====================================================================== Отправка данных в etcp после нормализации - -// Send data through ETCP connection -// Allocates memory from data_pool and places in input queue -// Returns: 0 on success, -1 on failure -int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - - if (!etcp || !data || len == 0) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] invalid parameters (etcp=%p, data=%p, len=%zu)", etcp->log_name, etcp, data, len); - return -1; - } - - if (!etcp->input_queue) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] input_queue is NULL for etcp=%p", etcp->log_name, etcp); - return -1; - } - - // Check length against maximum packet size - if (len > ETCP_MAX_PAYLOAD_SIZE) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] packet too large (len=%zu, max=%d)", etcp->log_name, len, ETCP_MAX_PAYLOAD_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, "[%s] failed to allocate packet data from data_pool", etcp->log_name); - return -1; - } - - // Copy user data to packet buffer - memcpy(packet_data, data, len); - - struct ETCP_FRAGMENT* pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool); - if (!pkt) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate queue entry", etcp->log_name); - 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.len = len; // размер packet_data - pkt->ll.dgram_pool = etcp->instance->data_pool; - pkt->ll.memlen = etcp->instance->data_pool->object_size; - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] created PACKET %p with data %p (len=%zu)", etcp->log_name, pkt, packet_data, len); - - // Add to input queue - input_queue_cb will process it - if (queue_data_put(etcp->input_queue, (struct ll_entry*)pkt) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name); - queue_dgram_free(&pkt->ll); - queue_entry_free(&pkt->ll); - return -1; - } - - return 0; -} - - - - - -static void input_queue_try_push(struct ETCP_CONN* etcp) {// пробуем протолкнуть при отправке - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - - // когда очередь отправки пуста - пробуем взять новый пакет на обработку - size_t wait_ack_bytes = queue_total_bytes(etcp->input_wait_ack); - if (wait_ack_bytes <= etcp->optimal_inflight) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] resume input queue: inflight_bytes=%d, input_len=%d", etcp->log_name, wait_ack_bytes, etcp->input_queue->total_bytes); - if (etcp->input_queue->count > 0) input_queue_cb(etcp->input_queue, etcp); - else queue_resume_callback(etcp->input_queue);// и только когда больше нечего отправлять - забираем новый пакет - } -} - -static void input_queue_try_resume(struct ETCP_CONN* etcp) {// при ACK - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - - // Сперва отправим всё из очереди отправки - size_t send_q_bytes = queue_total_bytes(etcp->input_send_q); -// queue_resume_callback(etcp->input_send_q);// вызвать лишний раз resume не страшно. - if (send_q_bytes>0) return; - - // когда очередь отправки пуста - пробуем взять новый пакет на обработку - size_t wait_ack_bytes = queue_total_bytes(etcp->input_wait_ack); - if (wait_ack_bytes <= etcp->optimal_inflight) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] resume input queue: inflight_bytes=%d, input_len=%d", etcp->log_name, wait_ack_bytes, etcp->input_queue->total_bytes); - queue_resume_callback(etcp->input_queue);// и только когда больше нечего отправлять - забираем новый пакет - } -} - -// Called from link-level etcp_link_update_inflight_lim() after inflight_lim_bytes changes. -// Recalculates connection-level optimal_inflight and resumes input_queue if room opened up. -void etcp_conn_on_inflight_lim_changed(struct ETCP_CONN* etcp) { - if (!etcp) return; - uint32_t sum = 0; - for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes; - etcp->optimal_inflight = sum; -// input_queue_try_resume(etcp); -} - -void etcp_stats(struct ETCP_CONN* etcp) { - if (!etcp) return; - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] stats for conn=%p:", etcp->log_name, etcp); - - // Queue statistics - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Queues:", etcp->log_name); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_queue: %zu pkts, %zu bytes", etcp->log_name, - queue_entry_count(etcp->input_queue), queue_total_bytes(etcp->input_queue)); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_send_q: %zu pkts, %zu bytes", etcp->log_name, - queue_entry_count(etcp->input_send_q), queue_total_bytes(etcp->input_send_q)); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_wait_ack: %zu pkts, %zu bytes", etcp->log_name, - queue_entry_count(etcp->input_wait_ack), queue_total_bytes(etcp->input_wait_ack)); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_q: %zu pkts", etcp->log_name, - queue_entry_count(etcp->ack_q)); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] recv_q: %zu pkts", etcp->log_name, - queue_entry_count(etcp->recv_q)); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] output_queue: %zu pkts, %zu bytes", etcp->log_name, - queue_entry_count(etcp->output_queue), queue_total_bytes(etcp->output_queue)); - - // RTT metrics - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RTT metrics:", etcp->log_name); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_last: %u (0.1ms)", etcp->log_name, etcp->rtt_last); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_avg_10: %u (0.1ms)", etcp->log_name, etcp->rtt_avg_10); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] jitter: %u (0.1ms)", etcp->log_name, etcp->jitter); - - // Counters - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Counters:", etcp->log_name); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] bytes_sent_total: %u", etcp->log_name, etcp->bytes_sent_total); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] retransmissions_count: %u", etcp->log_name, etcp->retransmissions_count); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_packets_count: %u", etcp->log_name, etcp->ack_packets_count); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] unacked_bytes: %u", etcp->log_name, etcp->unacked_bytes); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rx_dup_count: %u", etcp->log_name, etcp->rx_dup_count); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] tx_dup_count: %u", etcp->log_name, etcp->tx_dup_count); - - // IDs - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] IDs:", etcp->log_name); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] next_tx_id: %u", etcp->log_name, etcp->next_tx_id); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] last_rx_id: %u", etcp->log_name, etcp->last_rx_id); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] last_delivered_id:%u", etcp->log_name, etcp->last_delivered_id); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rx_ack_till: %u", etcp->log_name, etcp->rx_ack_till); -} - -// Input callback for input_queue (добавление новых кодограмм в стек) -// input_queue -> input_send_q -static void input_queue_cb(struct ll_queue* q, void* arg) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; - struct ETCP_FRAGMENT* in_pkt = (struct ETCP_FRAGMENT*)queue_data_get(q); - - if (!in_pkt) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot get element (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp); - queue_resume_callback(q); - return; - } - - - - // Create INFLIGHT_PACKET - - struct INFLIGHT_PACKET* p = (struct INFLIGHT_PACKET*)queue_entry_new_from_pool(etcp->inflight_pool); - if (!p) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp); - queue_dgram_free(&in_pkt->ll); - queue_entry_free((struct ll_entry*)in_pkt); - 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.dgram_pool = in_pkt->ll.dgram_pool; - p->ll.len = in_pkt->ll.len; - -// memory_pool_free(etcp->io_pool, in_pkt);// перемещаем из io_pool в inflight_pool - queue_entry_free(&in_pkt->ll); - - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] TX input -> inflight (seq=%u, len=%u); Qlen: in=%d snd=%d wait=%d", etcp->log_name, p->seq, p->ll.len, etcp->input_queue->count, etcp->input_send_q->count, etcp->input_wait_ack->count); - int len=p->ll.len;// сохраним len - - // Add to send queue - if (queue_data_put_with_index(etcp->input_send_q, &p->ll) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet seq=%u to input_send_q", etcp->log_name, p->seq); - queue_dgram_free(&p->ll); - queue_entry_free(&p->ll); - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] EXIT (queue put failed)", etcp->log_name); - return; - } - etcp->unacked_bytes += len; - -// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "successfully moved from input_queue to input_send_q"); - -// etcp_conn_process_send_queue(etcp);// сразу обработаем этот пакет - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] nextloop, input_queue size=%d ", etcp->log_name, q->count); -} - -static void ack_timeout_cb(void* arg); -static void ack_timeout_check(struct ETCP_CONN* etcp) { - uint64_t now = get_time_tb(); - int32_t timeout = (etcp->rtt_avg_10 * RETRANS_K1 + etcp->jitter * RETRANS_K2) / 16; - if (timeout<50) timeout=50; - if (timeout>10000) timeout=10000; - - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "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; - current = etcp->input_wait_ack->head; - while (current) { - struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current; - - int64_t elapsed = now - pkt->last_timestamp; - if (elapsed > timeout) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] timeout for seq=%u, elapsed=%lld, now=%llu, timeout=%llu, send_count=%u. Moving: wait_ack -> send_q", - etcp->log_name, pkt->seq, (unsigned long long)elapsed, (unsigned long long)now, (unsigned long long)timeout, pkt->send_count); - - // Remove from wait_ack - pkt=(struct INFLIGHT_PACKET*)queue_data_get(etcp->input_wait_ack); - if (!pkt) break; - // Increment counters - pkt->retrans_req_count++; // Optional, if used for retrans request logic - pkt->last_timestamp = now; - - // Change state and add to send_q for retransmission - pkt->state = INFLIGHT_STATE_WAIT_SEND; - queue_data_put_with_index(etcp->input_send_q, (struct ll_entry*)pkt); - - // Update stats - etcp->retransmissions_count++; - etcp_metrics_add_loss(etcp, 1); - } - else {// не надо до конца сканировать - они уже сортированы по таймстемпу т.к. очередь fifo, а timestamp = время добавления в очередь = время отправки - // shedule timer - int64_t next_timeout=timeout - elapsed; - if (next_timeout<0) next_timeout=0; - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] retransmission timer set for %llu units", etcp->log_name, next_timeout); - etcp->retrans_timer = uasync_set_timeout(etcp->instance->ua, next_timeout+10, etcp, ack_timeout_cb, "etcp_retrans"); - return; - } - current = etcp->input_wait_ack->head; - } - queue_resume_callback(etcp->input_wait_ack); -} - -static void ack_timeout_cb(void* arg) {// сработал таймер переотправки - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; - etcp->retrans_timer=NULL; - ack_timeout_check(etcp);// он установит новый таймер или выгребет всё и установит ожидание на очередь -} - -static void wait_ack_cb(struct ll_queue* q, void* arg) {// добавили пакет в ожидание ACK - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; - if (!etcp->retrans_timer) ack_timeout_check(etcp);// если таймер ретрансмиссий взведен - дожидаемся таймера. -} - -static void input_send_q_cb(struct ll_queue* q, void* arg) {// etcp->input_send_q data ready - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_CONN* etcp=(struct ETCP_CONN*)arg; - etcp_conn_process_send_queue(etcp); - if (etcp->tx_state==ETCP_TX_STATE_DATA_WAIT) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "resume input_send_q"); - queue_resume_callback(etcp->input_send_q); - } else DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "no resume - link busy"); - } -/* -static void send_ack_req_cb(struct ll_queue* q, void* arg) {// etcp->ack_q data ready - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_CONN* etcp=(struct ETCP_CONN*)arg; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] processing", etcp->log_name); - - etcp_conn_process_send_queue(etcp); - if (etcp->tx_state==ETCP_TX_STATE_DATA_WAIT) queue_resume_callback(etcp->ack_q); - } -*/ -void etcp_on_link_down(struct ETCP_CONN* etcp, struct ETCP_LINK* down_link) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - int up=0; - struct ETCP_LINK* link = etcp->links; - while (link) { - if (link->link_status) up=1; - link = link->next; - } - int was_up = etcp->links_up; - etcp->links_up = up; - if (up == 0 && was_up != 0) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All links fall down"); - etcp_on_down(etcp, down_link); - } -} - -static void ack_response_timer_cb(void* arg) {// проверяем неотправленные ack response и отправляем если надо. - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; - etcp->ack_resp_timer=NULL; - if (etcp->ack_q->count==0) return;// нечего отправлять - 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, "etcp_ack_resp"); -// else etcp->ack_resp_timer=NULL; -} - -static void etcp_link_ready_callback(struct ETCP_CONN* etcp) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (!etcp) return; - - if (etcp->tx_state==0) { - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "ETCP not ready, skip link state change"); - return;// not initialized - } - - if (etcp->links_up==0) { - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] etcp_link_ready_callback: links_up 0→1, calling etcp_on_up (initialized=%d tx_state=%d)", etcp->log_name, etcp->initialized, etcp->tx_state); - etcp->links_up=1; - etcp_on_up(etcp); - } else DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[%s] etcp_link_ready_callback: links_up=%d already up", etcp->log_name, etcp->links_up); - - - if (etcp->tx_state!=ETCP_TX_STATE_LINK_WAIT) return; - etcp->tx_state = ETCP_TX_STATE_DATA_WAIT; - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "resume input_send_q+ack_q; link_wait->data_wait"); - queue_resume_callback(etcp->input_send_q); -// queue_resume_callback(etcp->ack_q); - if (etcp->ack_q->count && etcp->ack_resp_timer == NULL) { - ack_response_timer_cb(etcp); - } -} - -// Process packets in send queue and transmit them -static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {// вызываем когда есть элемент в send_q или надо отправить ack - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_DGRAM* dgram; - if (etcp->tx_state!=ETCP_TX_STATE_DATA_WAIT) { - char l_status[256]={0}; - struct ETCP_LINK* link = etcp->links; - while (link) { - snprintf (l_status+strlen(l_status), 256-strlen(l_status), "L%d%d%s ", link->recv_keepalive, link->remote_keepalive, (link->shaper_timer)?"wait":"rdy"); - link = link->next; - } - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX state: %d (link not ready, skip send) %s", etcp->log_name, etcp->tx_state, l_status); - return; - } - dgram = etcp_request_pkt(etcp); - while(dgram) { - etcp_loadbalancer_send(dgram); - dgram = etcp_request_pkt(etcp); - } - if (etcp->tx_state==ETCP_TX_STATE_DATA_WAIT) { - queue_resume_callback(etcp->input_send_q); - } -} - -// Подготовить и отправить кодограмму -// вызывается линком когда освобождается или очередью если появляются данные на передачу -struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_LINK* link = NULL; -// если есть активный burst — используем этот линк напрямую - for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { - if (l->burst_active) { link = l; break; } - } - if (!link) link = etcp_loadbalancer_select_link(etcp); - if (!link) { - etcp->tx_state=ETCP_TX_STATE_LINK_WAIT; - etcp->cnt_link_wait++; - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no link available", etcp->log_name); - return NULL;// если линков нет - ждём появления свободного - } - etcp->tx_state=ETCP_TX_STATE_DATA_WAIT; - - size_t send_q_size = queue_entry_count(etcp->input_send_q); - - if (send_q_size == 0) {// сгребаем из input_queue -// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_send_q empty, check if avail input_queue -> inflight"); - input_queue_try_push(etcp); - } - - - // First, check if there's a packet in input_send_q (retrans or new) -// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "getting packet from input_send_q"); - struct INFLIGHT_PACKET* inf_pkt = (struct INFLIGHT_PACKET*)queue_data_get(etcp->input_send_q); - if (inf_pkt) { - uint64_t now=get_time_tb(); - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] send_q->wait_ack seq=%d TS=%llu", etcp->log_name, inf_pkt->seq, now); - -// === NEW: per-link inflight + send_hist logic === - // Compute 1-based link number in the linked list (as requested) - uint8_t link_num = 0; - struct ETCP_LINK* tmp = etcp->links; - while (tmp) { - link_num++; - if (tmp == link) break; - tmp = tmp->next; - } - if (link_num == 0) link_num = 255; // safety (should never happen) - - // Fill history (send_count is the index before this attempt) - inf_pkt->send_hist[inf_pkt->send_count % 8] = link_num; - - // Always subtract from previous last_link (if any) – this handles retransmission - if (inf_pkt->last_link) { - inf_pkt->last_link->total_retransmissions++; - inf_pkt->last_link->bbr_loss_since_ack += inf_pkt->ll.len; - bbr_note_loss(inf_pkt->last_link->bbr); - inf_pkt->last_link->inflight_bytes -= inf_pkt->ll.len; - inf_pkt->last_link->inflight_packets--; - } - - // Always add to the CURRENT link (first send or retransmission) - link->inflight_bytes += inf_pkt->ll.len; - link->inflight_packets++; - - // BBR: сохраняем снэпшот на момент отправки - inf_pkt->delivered_at_send = link->delivered_bytes; - inf_pkt->inflight_at_send = link->inflight_bytes; - inf_pkt->is_app_limited = (etcp->input_queue->count == 0) ? 1 : 0; - - // Update last_link for future ACK/retrans - inf_pkt->last_link = link; - - inf_pkt->last_timestamp=now; - inf_pkt->send_count++; - inf_pkt->state=INFLIGHT_STATE_WAIT_ACK; - - queue_data_put_with_index(etcp->input_wait_ack, &inf_pkt->ll);// move dgram to wait_ack queue - etcp_metrics_add_sent(etcp, inf_pkt->ll.len); - } - size_t ack_q_size = queue_entry_count(etcp->ack_q); - - if (!inf_pkt && ack_q_size == 0 && !link->burst_active) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no data/ack to send", etcp->log_name); - return NULL; - } - - struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool); - if (!dgram) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ETCP_DGRAM", etcp->log_name); - return NULL; - } - - dgram->link = link; - dgram->noencrypt_len=0; - dgram->timestamp = get_current_timestamp(); - -// формат ack: [01] [elements count] [4 байта last_delivered_id] [2 байта rx_dup_count] и <[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; - - dgram->data[ptr++]=etcp->rx_dup_count; - dgram->data[ptr++]=etcp->rx_dup_count>>8; - - int data_len=0; - if (inf_pkt) data_len=inf_pkt->ll.len; - int remain_len = link->mtu - - 28 /*udp headers*/ - - SC_NONCE_SIZE - SC_TAG_SIZE - SC_CRC32_SIZE - - (ETCP_ENCRYPTED_HDR_SIZE + ETCP_ACK_BASE_SIZE) - - (1 + sizeof(inf_pkt->seq)) /*payload hdr*/ - - data_len; -// DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "remain_len= %d pl=%d", remain_len, data_len); - -// добавим опциональные заголовки - struct ACK_PACKET* ack_pkt; - while (remain_len>=8) { - ack_pkt = (struct ACK_PACKET*)queue_data_get(etcp->ack_q); - if (!ack_pkt) break; - remain_len-=8; - // 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, "[%s] RX Send ACK seq=%d need=%d", etcp->log_name, ack_pkt->seq, etcp->last_delivered_id+1); - queue_entry_free((struct ll_entry*)ack_pkt); - - if (inf_pkt && inf_pkt->ll.len+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки) - if (ptr>500) break; - } - - dgram->data[1]=(ptr - 8) / 8; - - int max_enc = link->mtu - 28 - SC_NONCE_SIZE - SC_TAG_SIZE - SC_CRC32_SIZE - ETCP_ENCRYPTED_HDR_SIZE; - if (max_enc > PACKET_DATA_SIZE) max_enc = PACKET_DATA_SIZE; - if (max_enc < 0) max_enc = 0; - uint8_t burst_flags = 0; int pkt_sz_pos = 0; - -// === MEAS_RESP piggyback (если есть готовый ответ и не активен burst) === - if (link->burst_resp_pending && !link->burst_active && max_enc - ptr >= MEAS_RESP_SECTION_SIZE) { - dgram->data[ptr++] = ETCP_SECTION_MEAS_RESP; - dgram->data[ptr++] = link->burst_recv_id & 0xFF; dgram->data[ptr++] = (link->burst_recv_id >> 8) & 0xFF; - dgram->data[ptr++] = link->burst_resp_valid; - uint32_t v = link->burst_resp_gap_avg; - dgram->data[ptr++] = v & 0xFF; dgram->data[ptr++] = (v >> 8) & 0xFF; dgram->data[ptr++] = (v >> 16) & 0xFF; dgram->data[ptr++] = (v >> 24) & 0xFF; - v = link->burst_resp_gap_min; - dgram->data[ptr++] = v & 0xFF; dgram->data[ptr++] = (v >> 8) & 0xFF; dgram->data[ptr++] = (v >> 16) & 0xFF; dgram->data[ptr++] = (v >> 24) & 0xFF; - dgram->data[ptr++] = link->burst_resp_pkt_count; - link->burst_resp_pending = 0; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX burst resp id=%u valid=%u gap_min=%u", etcp->log_name, link->burst_recv_id, link->burst_resp_valid, link->burst_resp_gap_min); - } - -// === MEAS_TS section (burst) === - if (link->burst_active && max_enc - ptr >= MEAS_TS_SECTION_SIZE) { - burst_flags = 0; - if (!inf_pkt) burst_flags |= MEAS_FLAG_IS_FILLER; - if (link->burst_seq == 0) burst_flags |= MEAS_FLAG_IS_FIRST; - if (link->burst_seq + 1 >= link->burst_count) burst_flags |= MEAS_FLAG_IS_LAST; - dgram->data[ptr++] = ETCP_SECTION_MEAS_TS; - dgram->data[ptr++] = link->burst_id & 0xFF; dgram->data[ptr++] = (link->burst_id >> 8) & 0xFF; - dgram->data[ptr++] = burst_flags; - dgram->data[ptr++] = link->burst_seq; - uint16_t ts_us = (uint16_t)get_time_us(); - dgram->data[ptr++] = ts_us & 0xFF; dgram->data[ptr++] = (ts_us >> 8) & 0xFF; - pkt_sz_pos = ptr; - dgram->data[ptr++] = 0; dgram->data[ptr++] = 0; - } - - - if (link->last_recv_updated && remain_len>=5) {// если есть данные - добавим channel_timestamp - uint64_t now=get_time_tb(); - uint64_t dt=now - link->last_recv_local_time; - link->last_recv_updated=0; - if (dt<1000000) { - dgram->data[ptr++]=ETCP_SECTION_TIMESTAMP; - - uint16_t t=link->last_recv_timestamp + dt; - dgram->data[ptr++]=t; - dgram->data[ptr++]=t>>8; - - t=link->last_recv_local_time - link->last_recv_timestamp; - dgram->data[ptr++]=t; - dgram->data[ptr++]=t>>8; - remain_len-=5; - } - } - -// DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "remain_len(2)= %d", remain_len); - - if (inf_pkt) { -// === FILLER section before PAYLOAD (burst padding) === - if (link->burst_active) { - int payload_overhead = 1 + 4 + data_len; // type + seq + data - int buf_avail = PACKET_DATA_SIZE - ptr - payload_overhead; - int filler_data = max_enc - ptr - FILLER_HDR_SIZE - payload_overhead; - if (filler_data > buf_avail) filler_data = buf_avail; - if (filler_data > 0) { - dgram->data[ptr++] = ETCP_SECTION_FILLER; - dgram->data[ptr++] = filler_data & 0xFF; - dgram->data[ptr++] = (filler_data >> 8) & 0xFF; - memset(&dgram->data[ptr], 0, filler_data); ptr += filler_data; - } - } - // фрейм data (0) обязательно в конец - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX DATA: seq=%u len=%u retry=%d; Qlen: in=%d snd=%d wait=%d", etcp->log_name, inf_pkt->seq, inf_pkt->ll.len, inf_pkt->send_count, etcp->input_queue->count, etcp->input_send_q->count, etcp->input_wait_ack->count); - - 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; - - if ((int)(ptr + inf_pkt->ll.len) <= PACKET_DATA_SIZE) { - memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len; - } else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload overflow: ptr=%d len=%u max=%d", etcp->log_name, ptr, inf_pkt->ll.len, PACKET_DATA_SIZE); - } - else { - if (link->burst_active) { - int buf_avail = PACKET_DATA_SIZE - ptr; - int filler_data = max_enc - ptr - FILLER_HDR_SIZE; - if (filler_data > buf_avail - FILLER_HDR_SIZE) filler_data = buf_avail - FILLER_HDR_SIZE; - if (filler_data > 0) { - dgram->data[ptr++] = ETCP_SECTION_FILLER; - dgram->data[ptr++] = filler_data & 0xFF; - dgram->data[ptr++] = (filler_data >> 8) & 0xFF; - memset(&dgram->data[ptr], 0, filler_data); ptr += filler_data; - } - } - int chk=queue_check_consistency(etcp->ack_q); - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] only ACK (size=%d) packet with %d bytes total (chk=%d) rem=%d", etcp->log_name, ack_q_size, ptr, chk, remain_len); - } - if (ptr>=PACKET_DATA_SIZE-50) DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] SIZE ERROR!!! %d", etcp->log_name, ptr); - - dgram->data_len=ptr; - -// === Fill pkt_sz in MEAS_TS and advance burst state === - if (link->burst_active && pkt_sz_pos > 0) { - uint16_t pkt_sz = (uint16_t)dgram->data_len; - dgram->data[pkt_sz_pos] = pkt_sz & 0xFF; - dgram->data[pkt_sz_pos + 1] = (pkt_sz >> 8) & 0xFF; - link->burst_pkt_size = pkt_sz; - link->burst_seq++; - if (link->burst_seq >= link->burst_count) etcp_link_burst_finish(link); - } - - etcp_dump_pkt_sections(dgram, link, 1); - - return dgram; -} - -// ====================================================================== Прием данных - - -void etcp_output_try_assembly(struct ETCP_CONN* etcp) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - // пробуем собрать выходную очередь из фрагментов -// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] etcp=%p, last_delivered_id=%u, recv_q_count=%d", -// etcp->log_name, 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 = (struct ETCP_FRAGMENT*)queue_find_data_by_index(etcp->recv_q, &next_expected_id); - if (!rx_pkt) { - // No more contiguous packets found - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no packet found for id=%u, stopping", etcp->log_name, next_expected_id); - break; - } - -// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] assembling packet id=%u (len=%u)", etcp->log_name, -// 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 - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "Remove from assembly queue"); - queue_remove_data(etcp->recv_q, (struct ll_entry*)rx_pkt); - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "move: ETCP -> PN"); - - // Add to output_queue using the same ETCP_FRAGMENT structure - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] moving packet id=%u to output_queue (qlen=%d)", etcp->log_name, - next_expected_id, etcp->output_queue->count); - uint16_t pkt_len = rx_pkt->ll.len; - if (queue_data_put(etcp->output_queue, (struct ll_entry*)rx_pkt) == 0) { - delivered_bytes += pkt_len; - delivered_count++; - } else { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet id=%u to output_queue", etcp->log_name, - next_expected_id); - // Put it back in recv_q if we can't add to output_queue - queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt); - break; - } - - // Update state for next iteration - etcp->last_delivered_id = next_expected_id; - next_expected_id++; - } - - if (delivered_count>0) DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] delivered %u contiguous packets (%u bytes), last_delivered_id=%u, asm_queue=%d, output_queue=%d", - etcp->log_name, delivered_count, delivered_bytes, etcp->last_delivered_id, queue_entry_count(etcp->recv_q), 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) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (!etcp) return; - -// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] processing ACK for seq=%u, ts=%u, dts=%u", etcp->log_name, seq, ts, dts); - - // Find the acknowledged packet in the wait_ack queue - struct INFLIGHT_PACKET* acked_pkt; - acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_index(etcp->input_wait_ack, &seq); - if (acked_pkt) { - etcp->cnt_ack_hit_inf++; - queue_remove_data(etcp->input_wait_ack, (struct ll_entry*)acked_pkt); - } - else { - acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_index(etcp->input_send_q, &seq); - if (acked_pkt) { - etcp->cnt_ack_hit_sndq++; - queue_remove_data(etcp->input_send_q, (struct ll_entry*)acked_pkt); - } - else etcp->cnt_ack_miss++; - } - - if (!acked_pkt) { - // Packet might be already acknowledged or not found - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "packet seq=%u not found in wait_ack queue", seq); - return; - } - - // === subtract inflight from the LAST link the packet was sent on === - if (acked_pkt->last_link) { - struct ETCP_LINK* link = acked_pkt->last_link; - link->inflight_bytes -= acked_pkt->ll.len; - link->inflight_packets--; - - // BBR: собираем rate_sample и вызываем bbr_main - uint64_t now_tb = get_time_tb(); - uint32_t interval_us = link->last_ack_time_tb ? (uint32_t)((now_tb - link->last_ack_time_tb) * 100) : 0; - struct bbr_rate_sample rs = { - .delivered = acked_pkt->ll.len, - .interval_us = interval_us, - .rtt_us = (uint32_t)link->rtt_last * 100, - .acked_sacked = acked_pkt->ll.len, - .prior_delivered = (uint32_t)acked_pkt->delivered_at_send, - .tx_in_flight = acked_pkt->inflight_at_send, - .lost = (int)link->bbr_loss_since_ack, - .is_app_limited = acked_pkt->is_app_limited, - }; - link->delivered_bytes += rs.delivered; - link->last_ack_time_tb = now_tb; - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[L%u] BBR in: seq=%u del=%u iv=%uus rtt=%uus pri=%u tx_infl=%u cur=%u/%u lost=%d app=%d", - link->local_link_id, seq, rs.delivered, rs.interval_us, rs.rtt_us, - rs.prior_delivered, rs.tx_in_flight, link->inflight_bytes, link->inflight_lim_bytes, - rs.lost, rs.is_app_limited); - - uint32_t old_cwnd = link->inflight_lim_bytes; - uint32_t old_pacing = link->bbr_pacing_rate; - uint32_t new_cwnd = old_cwnd; - uint32_t new_pacing = old_pacing; - bbr_main(link->bbr, &rs, &new_cwnd, &new_pacing, link->mtu, link->inflight_bytes, - (link->inflight_bytes >= link->inflight_lim_bytes)); - etcp_link_update_inflight_lim(link, new_cwnd); - link->bbr_pacing_rate = new_pacing; - link->bandwidth = (uint32_t)((uint64_t)new_pacing * 8 / 1000); - link->bbr_loss_since_ack = 0; - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[L%u] BBR out: cwnd %u→%u pace %u→%u bw=%uK mode=%d cyc=%d infl=%u/%u", - link->local_link_id, old_cwnd, new_cwnd, old_pacing, new_pacing, - link->bandwidth, link->bbr->mode, link->bbr->cycle_idx, link->inflight_bytes, new_cwnd); - - link->acked_bytes += acked_pkt->ll.len; - link->acked_packets++; - } - - - // Update connection statistics - if (etcp->unacked_bytes >= acked_pkt->ll.len) etcp->unacked_bytes -= acked_pkt->ll.len; - else { - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] unacked_bytes underflow prevented: %u < %u", etcp->log_name, etcp->unacked_bytes, acked_pkt->ll.len); - etcp->unacked_bytes = 0; - } - etcp->bytes_sent_total += acked_pkt->ll.len; - etcp->ack_packets_count++; - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX removed packet seq=%u from wait_ack, unacked_bytes now %u total acked=%u", etcp->log_name, seq, etcp->unacked_bytes, etcp->ack_packets_count); - - 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, "[%s] completed for seq=%u", etcp->log_name, seq); -} - - -// Process incoming decrypted packet -void etcp_conn_input(struct ETCP_DGRAM* pkt) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - if (!pkt) return; - if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt); - volatile int _halt = 1; while (_halt) {} - } - if (!pkt->data_len) { - memory_pool_free(pkt->link->etcp->instance->pkt_pool, pkt); - return; - } - - etcp_dump_pkt_sections(pkt, pkt->link, 0); - - pkt->link->last_recv_local_time = get_time_tb(); - - 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 - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX pkt dlen=%d", etcp->log_name, len); - - while (len >= 1) { - uint8_t type = data[0]; - - // Process sections as per protocol.txt - switch (type) { - case ETCP_SECTION_ACK: { - if (len < 2) { len = 0; break; } - int elm_cnt=data[1]; - uint32_t till=data[2] | (data[3]<<8) | (data[4]<<16) | (data[5]<<24); - uint16_t rx_dup_count_16 = data[6] | (data[7]<<8); - uint16_t old_tx_dup_16 = etcp->tx_dup_count & 0xFFFF; - int16_t diff = (int16_t)(rx_dup_count_16 - old_tx_dup_16); - etcp->tx_dup_count += diff; - int ack_section_len = 8 + elm_cnt * 8; - if (ack_section_len > len) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] ACK section too large: elm_cnt=%d len=%d", etcp->log_name, elm_cnt, len); - len = 0; - break; - } - data+=ack_section_len; - len-=ack_section_len; - 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_TIMESTAMP: { - if (len < 5) { len = 0; break; } - uint16_t cur_ts=get_current_timestamp(); - uint16_t ret_ts=data[1] | (data[2]<<8);// cur_ts=ret_ts = RTT - uint16_t new_rtt=cur_ts-ret_ts; - pkt->link->rtt_last=new_rtt; - etcp_metrics_add_rtt(etcp, new_rtt); - - int recv_dt_tx1=data[3] | (data[4]<<8);// localtime удаленной стороны момента принятия пакета - timestamp этого пакета (на стороне отправителя, т.е. у нас) - - int recv_dt_rx=cur_ts - pkt->link->rtt_last/2 - 1000 - ts; - int recv_dt_tx=recv_dt_tx1 - pkt->link->rtt_last/2 - 1000; - if (pkt->link->recv_dt_avg_rx==0) pkt->link->recv_dt_avg_rx=recv_dt_rx*65536; - if (pkt->link->recv_dt_avg_tx==0) pkt->link->recv_dt_avg_tx=recv_dt_tx*65536; - pkt->link->recv_dt_avg_rx +=((int32_t)(recv_dt_rx*65536 - pkt->link->recv_dt_avg_rx))/32; - pkt->link->recv_dt_avg_tx +=((int32_t)(recv_dt_tx*65536 - pkt->link->recv_dt_avg_tx))/32; - pkt->link->rt_last = cur_ts - ts - pkt->link->recv_dt_avg_rx/65536; - pkt->link->tt_last = recv_dt_tx1 - pkt->link->recv_dt_avg_tx/65536; - //tts_correction += ((NOW - RTT/2 - TTS) - tts_correction)/16 (инициализируем сразу по 1 пакету) - data+=5; len-=5; - - int prev_rtt = pkt->link->rtt_last; - int d_rtt = new_rtt - prev_rtt; - if (d_rtt<0) d_rtt=-d_rtt; - pkt->link->jitter +=((int32_t)(d_rtt*65536 - pkt->link->jitter))/32; - - struct ETCP_LINK* c=etcp->links; - int rtt_sum=0; - int rtt_max=0; - int tt_sum=0; - int j_sum=0; - int cnt=0; - while (c) { - rtt_sum += c->rtt_last; - if (rtt_max < c->rtt_last) rtt_max = c->rtt_last; - tt_sum += c->tt_last; - j_sum += c->jitter/32768; - cnt++; - c=c->next; - } - etcp->rtt_last=rtt_sum/cnt; - etcp->rtt_avg_10=rtt_max; - etcp->tt_last=tt_sum/cnt; - etcp->jitter=j_sum/cnt; - if (etcp->rtt_last) { - uint64_t now_tb = get_time_tb(); - if (now_tb - etcp->last_rtt_cb_time > RTT_CB_PERIOD_TB) { - topo_node_ping_update_rtt(etcp->instance->topo_groups, etcp->peer_node_id, etcp->rtt_last); - etcp->last_rtt_cb_time = now_tb; - } - } - break; - } - case ETCP_SECTION_PAYLOAD: { - - if (len>=5) { - // формируем ACK - uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24); - if (etcp->got_initial_pkt == 0) { - if (seq==1) { etcp->got_initial_pkt = 1; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Initial packet seq=1 received", etcp->log_name); } - else { - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Out-of-order pkt before init seq=%d, waiting seq=1", - etcp->log_name, seq); - len=0; - break; - } - } - else { - int32_t d=seq - etcp->last_delivered_id; - if (d>MAX_INFLIGHT_SIZE || d<-MAX_INFLIGHT_SIZE) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] Received packet out of inflight bounds: seq=%d last delivered=%d", etcp->log_name, seq, etcp->last_delivered_id); - len=0; - break; - } - } - if (queue_find_data_by_index(etcp->ack_q, &seq) == NULL) { - struct ACK_PACKET* p = (struct ACK_PACKET*)queue_entry_new_from_pool(etcp->instance->ack_pool); - if (!p) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ACK_PACKET", etcp->log_name); - len = 0; - break; - } - p->seq=seq; - p->pkt_timestamp=pkt->timestamp; - p->recv_timestamp=get_current_timestamp(); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX add to ack_q seq=%d", etcp->log_name, seq); - - queue_data_put_with_index(etcp->ack_q, (struct ll_entry*)p); - if (etcp->ack_resp_timer == NULL) { - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] set ack_timer for delayed ACK send", etcp->log_name); - etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb, "etcp_ack_resp"); - } - } - if (((int32_t)(etcp->last_delivered_id-seq)<0) && (queue_find_data_by_index(etcp->recv_q, &seq)==NULL)) {// проверяем есть ли пакет с этим seq - uint32_t pkt_len=len-5; - DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] adding packet seq=%u to recv_q (last_delivered_id=%u)", etcp->log_name, seq, etcp->last_delivered_id); - // отправляем пакет в очередь на сборку - uint8_t* payload_data = memory_pool_alloc(etcp->instance->data_pool); - if (!payload_data) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate payload_data from data_pool", etcp->log_name); - len=0; - break; - } - struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool); - if (!rx_pkt) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate rx_pkt from io_pool", etcp->log_name); - memory_pool_free(etcp->instance->data_pool, payload_data); - len=0; - break; - } - rx_pkt->seq = seq; - rx_pkt->timestamp = pkt->timestamp; - rx_pkt->ll.dgram = payload_data; - rx_pkt->ll.len = pkt_len; - rx_pkt->ll.dgram_pool = etcp->instance->data_pool; - rx_pkt->ll.memlen = etcp->instance->data_pool->object_size; - if (pkt_len > (int)etcp->instance->data_pool->object_size) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload too large for data_pool: %u > %zu", etcp->log_name, pkt_len, etcp->instance->data_pool->object_size); - memory_pool_free(etcp->instance->data_pool, payload_data); - queue_entry_free(&rx_pkt->ll); - len = 0; - break; - } - if (data + len > pkt->data + pkt->data_len) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload section bounds overflow: data+len=%p > end=%p", etcp->log_name, (void*)(data + len), (void*)(pkt->data + pkt->data_len)); - memory_pool_free(etcp->instance->data_pool, payload_data); - queue_entry_free(&rx_pkt->ll); - len = 0; - break; - } - // Copy the actual payload data - memcpy(payload_data, data + 5, pkt_len); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX seq=%u need=%u asm_len=%d", etcp->log_name, seq, etcp->last_delivered_id+1, etcp->recv_q->count); - queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt); - etcp_metrics_add_rcvd(etcp, pkt_len); - if ((int32_t)(seq - etcp->last_delivered_id) == 1) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов - } else { - etcp->rx_dup_count++; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX dup: seq=%u need=%u asm_len=%d", etcp->log_name, seq, etcp->last_delivered_id+1, etcp->recv_q->count); - } - } else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "payload len %d < 5", len); - len=0; - break; - } - - case ETCP_SECTION_MEAS_TS: { - if (len < MEAS_TS_SECTION_SIZE) { len = 0; break; } - uint16_t b_id = data[1] | (data[2] << 8); - uint8_t flags = data[3]; - uint8_t seq = data[4]; - uint16_t pkt_sz = data[7] | (data[8] << 8); - struct ETCP_LINK* link = pkt->link; - if (b_id != link->burst_recv_id) { - link->burst_recv_id = b_id; - link->burst_recv_count = BURST_PACKET_COUNT; - link->burst_recv_next_seq = 0; - link->burst_recv_valid = 1; - link->burst_recv_received = 0; - } - if (link->burst_recv_valid) { - if (seq != link->burst_recv_next_seq) { - if (seq > link->burst_recv_next_seq) etcp_metrics_add_loss(etcp, seq - link->burst_recv_next_seq); - link->burst_recv_valid = 0; - } - if (seq < 16) { - link->burst_recv_times[seq] = get_time_us(); - link->burst_recv_next_seq = seq + 1; - link->burst_recv_received++; - link->burst_recv_pkt_size = pkt_sz; - } - } - if ((flags & MEAS_FLAG_IS_LAST) || link->burst_recv_received >= link->burst_recv_count) { - uint8_t total = (flags & MEAS_FLAG_IS_LAST) ? link->burst_recv_received : link->burst_recv_count; - if (link->burst_recv_valid && total >= BURST_SKIP_COUNT + 2) { - uint32_t sum_gap = 0, min_gap = UINT32_MAX; - uint8_t measured = 0; - for (uint8_t i = BURST_SKIP_COUNT; i + 1 < total; i++) { - uint64_t dt = link->burst_recv_times[i + 1] - link->burst_recv_times[i]; - uint32_t gap = (uint32_t)(dt > 0 ? dt : 1); - sum_gap += gap; measured++; - if (gap < min_gap) min_gap = gap; - } - if (measured > 0) { - link->burst_resp_gap_avg = sum_gap / measured; - link->burst_resp_gap_min = min_gap; - link->burst_resp_pkt_count = measured; - link->burst_resp_valid = 1; - } - } - link->burst_resp_pending = 1; - } - data += MEAS_TS_SECTION_SIZE; len -= MEAS_TS_SECTION_SIZE; - break; - } - - case ETCP_SECTION_MEAS_RESP: { - if (len < MEAS_RESP_SECTION_SIZE) { len = 0; break; } - uint16_t b_id = data[1] | (data[2] << 8); - uint8_t valid = data[3]; - uint32_t gap_avg = data[4] | (data[5] << 8) | ((uint32_t)data[6] << 16) | ((uint32_t)data[7] << 24); - uint32_t gap_min = data[8] | (data[9] << 8) | ((uint32_t)data[10] << 16) | ((uint32_t)data[11] << 24); - uint8_t pkt_cnt = data[12]; - struct ETCP_LINK* link = pkt->link; - if (b_id == link->burst_id - 1 || b_id == link->burst_id) { - if (link->burst_resp_timer) { uasync_cancel_timeout(link->etcp->instance->ua, link->burst_resp_timer); link->burst_resp_timer = NULL; } - if (valid == MEAS_RESP_VALID && gap_min > 0 && link->burst_pkt_size > 0) { - uint64_t bw_kbps = (uint64_t)link->burst_pkt_size * 8000ULL / gap_min; - if (bw_kbps > 10000000) bw_kbps = 10000000; - link->bandwidth = (uint32_t)bw_kbps; - float rtt_sec = (float)link->bbr->min_rtt_us / 1000000.f; - if (rtt_sec < 0.005f) rtt_sec = 0.005f; - link->burst_target_bdp = (uint32_t)((float)bw_kbps * 1000.f / 8.f * rtt_sec); - // BBR: инициализируем верхнюю границу из burst - if (link->bbr->inflight_hi == ~0U) link->bbr->inflight_hi = link->burst_target_bdp; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] burst resp: BW=%u Kbps, BDP=%u gap_min=%u us pkt_cnt=%u", - link->etcp->log_name, (uint32_t)bw_kbps, link->burst_target_bdp, gap_min, pkt_cnt); - } - } - data += MEAS_RESP_SECTION_SIZE; len -= MEAS_RESP_SECTION_SIZE; - break; - } - - case ETCP_SECTION_FILLER: { - if (len < FILLER_HDR_SIZE) { len = 0; break; } - uint16_t fill_len = data[1] | (data[2] << 8); - if ((uint16_t)(fill_len + FILLER_HDR_SIZE) > len) { len = 0; break; } - data += FILLER_HDR_SIZE + fill_len; len -= FILLER_HDR_SIZE + fill_len; - break; - } - - case ETCP_SECTION_METRICS: { - if (len < 2 + SC_SIGN_SIZE) { len = 0; break; } - uint16_t csv_len = data[1] | (data[2] << 8); - if ((uint32_t)(2 + SC_SIGN_SIZE + csv_len) > len) { len = 0; break; } - uint8_t* sig = data + 2; - uint8_t* csv = data + 2 + SC_SIGN_SIZE; - uint8_t* peer_pubkey = pkt->link->remote_ed25519_pubkey; - struct sc_stream_sign_state verify_state; - int sig_ok = (sc_stream_sign_verify_init(&verify_state, peer_pubkey) == SC_OK - && sc_stream_sign_update(&verify_state, csv, csv_len) == SC_OK - && sc_stream_sign_verify(&verify_state, sig, SC_SIGN_SIZE) == SC_OK); - if (!sig_ok && verify_state.initialized) /* update failed, then verify didn't run */ - sc_stream_sign_cleanup(&verify_state); - if (sig_ok) { - const char* stats_dir = etcp->instance->stats_dir; - if (stats_dir[0] && etcp->peer_node_id != 0) { - char path[768]; - snprintf(path, sizeof(path), "%s/from_%016llx.dat", stats_dir, (unsigned long long)etcp->peer_node_id); - FILE* f = fopen(path, "a"); - if (f) { - fprintf(f, "%.*s,", csv_len, csv); - for (int i = 0; i < SC_SIGN_SIZE; i++) fprintf(f, "%02x", sig[i]); - fprintf(f, "\n"); - fclose(f); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Received metrics from peer (%u bytes)", etcp->log_name, csv_len); - if (etcp->instance->api_bindings.on_metrics_rcvd) - etcp->instance->api_bindings.on_metrics_rcvd(etcp->instance->api_bindings.metrics_user_ptr, etcp, csv, csv_len, sig); - } else { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] Cannot open from_metrics file: %s", etcp->log_name, path); - } - } - } else { - DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "[%s] Metrics signature verification FAILED", etcp->log_name); - } - data += 2 + SC_SIGN_SIZE + csv_len; len -= 2 + SC_SIGN_SIZE + csv_len; - break; - } - - default: - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "unknown section type=0x%02x", type); - len=0; - break; - } - - } - - if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt); - volatile int _halt = 1; while (_halt) {} - } - - - memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram -} - -void etcp_update_mtu(struct ETCP_CONN* etcp) { - if (!etcp) return; - int new_mtu = PACKET_DATA_MAX_MTU; - int has_links = 0; - struct ETCP_LINK* link = etcp->links; - while (link) { - has_links = 1; - if (link->mtu > 0 && link->mtu < new_mtu) new_mtu = link->mtu; - link = link->next; - } - if (!has_links) new_mtu = ETCP_RFC791_MIN_MTU; - if (new_mtu != etcp->mtu) { - etcp->mtu = new_mtu; - if (etcp->normalizer) etcp->normalizer->frag_size = etcp->mtu - ACK_REZERV - UDP_HDR_SIZE - UDP_SC_HDR_SIZE; - } -} - -// ====================================================================== Метрики: гистограммы RTT/потерь + тотальные счётчики - -static const uint16_t etcp_metrics_rtt_bounds[] = { 10, 20, 30, 50, 80, 130, 210, 340, 550, 900, 1500 }; -#define ETCP_METRICS_RTT_BOUNDS_COUNT (sizeof(etcp_metrics_rtt_bounds) / sizeof(etcp_metrics_rtt_bounds[0])) - -void etcp_metrics_init(struct etcp_metrics* m) { - if (!m) return; - memset(m, 0, sizeof(*m)); -} - -void etcp_metrics_add_rtt(struct ETCP_CONN* etcp, uint16_t rtt_tb) { - if (!etcp) return; - struct etcp_metrics* m = &etcp->metrics; - int i; - for (i = 0; i < (int)ETCP_METRICS_RTT_BOUNDS_COUNT; i++) { - if (rtt_tb < etcp_metrics_rtt_bounds[i]) break; - } - if (i >= ETCP_METRICS_RTT_BUCKETS) i = ETCP_METRICS_RTT_BUCKETS - 1; - m->rtt_hist[i]++; - m->work_samples++; -} - -void etcp_metrics_add_sent(struct ETCP_CONN* etcp, uint32_t len) { - if (!etcp) return; - struct etcp_metrics* m = &etcp->metrics; - m->work_sent++; - m->work_bytes_sent += len; - m->total_sent++; - m->total_bytes_sent += len; -} - -void etcp_metrics_add_rcvd(struct ETCP_CONN* etcp, uint32_t len) { - if (!etcp) return; - struct etcp_metrics* m = &etcp->metrics; - m->work_rcvd++; - m->work_bytes_rcvd += len; - m->total_rcvd++; - m->total_bytes_rcvd += len; -} - -void etcp_metrics_add_loss(struct ETCP_CONN* etcp, uint32_t count) { - if (!etcp) return; - struct etcp_metrics* m = &etcp->metrics; - m->work_lost += count; - m->total_lost += count; -} - -static void etcp_metrics_do_snapshot(struct ETCP_CONN* etcp) { - if (!etcp || !etcp->instance) return; - struct etcp_metrics* m = &etcp->metrics; - - // Копируем working → final - memcpy(m->rtt_hist_final, m->rtt_hist, sizeof(m->rtt_hist_final)); - memcpy(m->loss_hist_final, m->loss_hist, sizeof(m->loss_hist_final)); - m->samples_final = m->work_samples; - m->sent_final = m->work_sent; - m->rcvd_final = m->work_rcvd; - m->bytes_sent_final = m->work_bytes_sent; - m->bytes_rcvd_final = m->work_bytes_rcvd; - m->lost_final = m->work_lost; - - // Рассчитываем loss-гистограмму - uint32_t total = m->work_sent + m->work_lost; - if (total > 0) { - uint32_t loss_pct = (m->work_lost * 100 + total / 2) / total; - int bucket = (int)loss_pct; - if (bucket >= ETCP_METRICS_LOSS_BUCKETS) bucket = ETCP_METRICS_LOSS_BUCKETS - 1; - m->loss_hist[bucket]++; - } - - // Обнуляем рабочую копию - memset(m->rtt_hist, 0, sizeof(m->rtt_hist)); - memset(m->loss_hist, 0, sizeof(m->loss_hist)); - m->work_samples = 0; - m->work_sent = 0; m->work_rcvd = 0; - m->work_bytes_sent = 0; m->work_bytes_rcvd = 0; - m->work_lost = 0; - - m->last_snapshot_tb = get_time_tb(); - - // Строим CSV строку - char csv_buf[4096]; - int pos = 0; - time_t now = time(NULL); - struct tm* tm = localtime(&now); - char date_str[16], time_str[16]; - if (tm) { strftime(date_str, sizeof(date_str), "%Y-%m-%d", tm); strftime(time_str, sizeof(time_str), "%H:%M:%S", tm); } - else { strcpy(date_str, "----"); strcpy(time_str, "----"); } - pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, "%s,%s", date_str, time_str); - for (int i = 0; i < ETCP_METRICS_RTT_BUCKETS; i++) pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%u", m->rtt_hist_final[i]); - for (int i = 0; i < ETCP_METRICS_LOSS_BUCKETS; i++) pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%u", m->loss_hist_final[i]); - pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%u,%u,%u,%u,%u,%u", m->samples_final, m->sent_final, m->rcvd_final, m->bytes_sent_final, m->bytes_rcvd_final, m->lost_final); - pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%llu,%llu,%llu,%llu,%llu", (unsigned long long)m->total_sent, (unsigned long long)m->total_rcvd, (unsigned long long)m->total_bytes_sent, (unsigned long long)m->total_bytes_rcvd, (unsigned long long)m->total_lost); - uint32_t csv_len = (uint32_t)pos; - - // Запись локального файла - const char* stats_dir = etcp->instance->stats_dir; - if (stats_dir[0] && etcp->peer_node_id != 0) { - char path[768]; - snprintf(path, sizeof(path), "%s/%016llx.dat", stats_dir, (unsigned long long)etcp->peer_node_id); - FILE* f = fopen(path, "a"); - if (f) { fprintf(f, "%s\n", csv_buf); fclose(f); } - else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] Cannot open metrics file: %s", etcp->log_name, path); - } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Metrics snapshot: samples=%u sent=%u rcvd=%u lost=%u", etcp->log_name, m->samples_final, m->sent_final, m->rcvd_final, m->lost_final); - - // Подпись Ed25519 и отправка пиру - if (etcp->crypto_ctx.initialized && etcp->initialized) { - struct sc_stream_sign_state sign_state; - if (sc_stream_sign_init(&etcp->crypto_ctx, &sign_state) == SC_OK) { - uint8_t sig[SC_SIGN_SIZE]; - size_t sig_len = sizeof(sig); - if (sc_stream_sign_update(&sign_state, (uint8_t*)csv_buf, csv_len) == SC_OK - && sc_stream_sign_final(&sign_state, sig, &sig_len) == SC_OK) { - // Построить и отправить дграмму с секцией METRICS - struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp); - if (link) { - struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool); - if (dgram) { - dgram->link = link; - dgram->noencrypt_len = 0; - dgram->timestamp = get_current_timestamp(); - int ptr = 0; - // Минимальная ACK-секция - dgram->data[ptr++] = 1; dgram->data[ptr++] = 0; - 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; - dgram->data[ptr++] = etcp->rx_dup_count; dgram->data[ptr++] = etcp->rx_dup_count>>8; - // METRICS секция - dgram->data[ptr++] = ETCP_SECTION_METRICS; - uint16_t csv16 = (uint16_t)csv_len; - dgram->data[ptr++] = csv16 & 0xFF; dgram->data[ptr++] = (csv16 >> 8) & 0xFF; - memcpy(dgram->data + ptr, sig, SC_SIGN_SIZE); ptr += SC_SIGN_SIZE; - memcpy(dgram->data + ptr, csv_buf, csv_len); ptr += csv_len; - dgram->data_len = ptr; - etcp_encrypt_send(dgram); - memory_pool_free(etcp->instance->pkt_pool, dgram); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Metrics sent to peer (%u bytes)", etcp->log_name, csv_len); - } - } - } else { - sc_stream_sign_cleanup(&sign_state); - } - } - } -} - -static void metrics_snapshot_timer_cb(void* arg) { - struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; - if (!etcp) return; - etcp->metrics.timer = NULL; - - if (etcp->metrics.work_sent + etcp->metrics.work_rcvd >= ETCP_METRICS_MIN_PACKETS) - etcp_metrics_do_snapshot(etcp); - - // Перезапускаем таймер - etcp->metrics.timer = uasync_set_timeout(etcp->instance->ua, ETCP_METRICS_INTERVAL_TB, etcp, metrics_snapshot_timer_cb, "etcp_metrics"); -} - -void etcp_metrics_start_timer(struct ETCP_CONN* etcp) { - if (!etcp || etcp->metrics.timer) return; - etcp_metrics_init(&etcp->metrics); - etcp->metrics.timer = uasync_set_timeout(etcp->instance->ua, ETCP_METRICS_INTERVAL_TB, etcp, metrics_snapshot_timer_cb, "etcp_metrics"); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Metrics timer started (interval=%ums)", etcp->log_name, ETCP_METRICS_INTERVAL_TB / 10); -} - -void etcp_metrics_stop_timer(struct ETCP_CONN* etcp) { - if (!etcp) return; - if (etcp->metrics.timer && etcp->instance) { - uasync_cancel_timeout(etcp->instance->ua, etcp->metrics.timer); - etcp->metrics.timer = NULL; - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Metrics timer stopped", etcp->log_name); - } -} - diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 64a913cd..790b5bce 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -60,6 +60,7 @@ static void tcp_server_on_link(struct stcp_link *link, void *arg) { conn->peer_node_id = node_id; sc_set_peer_public_key(&conn->crypto_ctx, pubkey, SC_PEER_PUBKEY_BIN); etcp_update_log_name(conn); + { struct conn_queue_entry* ce = (struct conn_queue_entry*)conn->conn_queue_entry->data; ce->peer_node_id = node_id; queue_remove_data(conn->conn_queue, conn->conn_queue_entry); queue_data_put_with_index(conn->conn_queue, conn->conn_queue_entry); } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: new ETCP_CONN peer=0x%016llx", (unsigned long long)node_id); } diff --git a/src/transport_layer/stcp_link.c b/src/transport_layer/stcp_link.c deleted file mode 100644 index 41c8a255..00000000 --- a/src/transport_layer/stcp_link.c +++ /dev/null @@ -1,286 +0,0 @@ -// stcp_link.c — STCP link management implementation -#include "stcp_link.h" -#include "stcp.h" -#include "stcp_server.h" -#include "stcp_client.h" -#include "secure_channel.h" -#include "etcp.h" -#include "etcp_connections.h" -#include "etcp_router.h" -#include "utun_instance.h" -#include "../lib/debug_config.h" -#include "../lib/ll_queue.h" -#include "../lib/mem.h" -#include "../lib/memory_pool.h" -#include -#include - -struct stcp_server { - struct stcp_server *next; // linked list in UTUN_INSTANCE - struct stcp_server *srv; // stcp_server from stcp_server.h - struct stcp_link_config cfg; - stcp_server_on_link_cb on_link; - void *on_link_arg; -}; - -struct stcp_link { - struct stcp_link_config cfg; - struct stcp_client *cli; - struct stcp_conn *conn; - - struct ETCP_CONN *etcp_conn; // parent ETCP_CONN (for rx dispatch / etcp_send compat) - struct ETCP_LINK *etcp_link; // owning ETCP_LINK (for etcp_conn_input pkt->link) - - uint8_t ready; - uint8_t closing; - - stcp_link_cb on_ready_cb; - void *ready_arg; - void (*on_close_cb)(struct stcp_link *link, int err, void *arg); - void *close_arg; - - struct ll_queue *tx_queue; // owned by this link -}; - -// ====== rx dispatch ====== - -static void link_rx_cb(struct ll_queue *q, void *arg) { - struct stcp_link *link = (struct stcp_link *)arg; - struct ll_entry *e = queue_data_get(q); - if (!e) { queue_resume_callback(q); return; } - - if (link->etcp_conn && link->etcp_link && link->cfg.inst) { - if (e->len >= 1 && e->dgram[0] == ETCP_KEEPALIVE) { - struct ETCP_LINK *l = link->etcp_link; - if (e->len >= 3) { uint16_t pp = e->dgram[1] | ((uint16_t)e->dgram[2] << 8); l->keepalive_timeout = (uint32_t)pp * KA_TIMEOUT_MULT; } - l->recv_keepalive = 1; l->remote_keepalive = 1; l->link_status = 1; l->keepalive_recv_count++; - l->last_recv_local_time = get_time_tb(); - queue_dgram_free(e); queue_entry_free(e); - } else { - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "stcp_link rx → etcp_conn_input len=%zu link=%p", e->len, (void*)link->etcp_link); - struct ETCP_DGRAM *pkt = memory_pool_alloc(link->cfg.inst->pkt_pool); - if (pkt) { - pkt->link = link->etcp_link; pkt->data_len = (uint16_t)e->len; pkt->noencrypt_len = 0; - if (e->len > 0) memcpy(pkt->data, e->dgram, e->len); - etcp_conn_input(pkt); - } - queue_dgram_free(e); queue_entry_free(e); - } - } else { - struct UTUN_INSTANCE *inst = link->cfg.inst; - uint8_t id = (e->dgram && e->len > 0) ? e->dgram[0] : 0; - if (inst && inst->api_bindings.callbacks[id]) - inst->api_bindings.callbacks[id](link->etcp_conn ? link->etcp_conn : NULL, e); - else if (inst && inst->api_bindings.callbacks[0]) - inst->api_bindings.callbacks[0](link->etcp_conn ? link->etcp_conn : NULL, e); - else { queue_dgram_free(e); queue_entry_free(e); } - } - queue_resume_callback(q); -} - -// ====== server accept → link ====== - -static void server_accept_cb(struct stcp_conn *conn, void *arg) { - struct stcp_server *ss = (struct stcp_server *)arg; - struct stcp_link *link = u_calloc(1, sizeof(struct stcp_link)); - if (!link) { stcp_conn_free(conn); return; } - link->cfg = ss->cfg; - link->ready = 1; - link->conn = conn; - - struct ll_queue *rx = queue_new(conn->ua, 0, 0, 0, "srx"); - queue_set_callback(rx, link_rx_cb, link); - stcp_conn_set_rx_queue(conn, rx); - link->tx_queue = queue_new(conn->ua, 0, 0, 0, "stx"); - queue_set_threshold(link->tx_queue, 0, 0); - queue_set_waiter_defer(link->tx_queue, 1); - stcp_conn_set_tx_queue(conn, link->tx_queue); - - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_link: server accepted connection"); - if (ss->on_link) ss->on_link(link, ss->on_link_arg); -} - -// ====== client connect → link ====== - -static void client_ready_cb(struct stcp_conn *conn, void *arg) { - struct stcp_link *link = (struct stcp_link *)arg; - if (!conn) return; - link->ready = 1; - link->conn = conn; - - struct ll_queue *rx = queue_new(conn->ua, 0, 0, 0, "crx"); - queue_set_callback(rx, link_rx_cb, link); - stcp_conn_set_rx_queue(conn, rx); - link->tx_queue = queue_new(conn->ua, 0, 0, 0, "ctx"); - queue_set_threshold(link->tx_queue, 0, 0); - queue_set_waiter_defer(link->tx_queue, 1); - stcp_conn_set_tx_queue(conn, link->tx_queue); - - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_link: client handshake OK, link=%p etcp_link=%p", - (void*)link, (void*)link->etcp_link); - if (link->etcp_link) etcp_link_enter_ready_tcp(link->etcp_link); - if (link->on_ready_cb) link->on_ready_cb(link, link->ready_arg); -} - -// ====== API ====== - -struct stcp_server *stcp_server_listen(struct stcp_link_config *cfg, uint16_t port, - stcp_server_on_link_cb on_link, void *arg) { - if (!cfg || !cfg->ua || !cfg->my_keys) return NULL; - struct stcp_server *ss = u_calloc(1, sizeof(struct stcp_server)); - if (!ss) return NULL; - ss->cfg = *cfg; - ss->on_link = on_link; - ss->on_link_arg = arg; - ss->srv = stcp_server_create(cfg->ua, port, cfg->my_keys, server_accept_cb, ss, NULL, NULL, cfg->listen_family); - if (!ss->srv) { u_free(ss); return NULL; } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "port=%u", port); - return ss; -} - -void stcp_link_server_destroy(struct stcp_server *ss) { - if (!ss) return; - if (ss->srv) stcp_server_destroy(ss->srv); - u_free(ss); -} - -struct stcp_link *stcp_link_connect(struct stcp_link_config *cfg) { - if (!cfg || !cfg->ua || !cfg->my_keys || !cfg->peer_pubkey || !cfg->remote_addr) - return NULL; - - int family = cfg->remote_addr->ss_family; - char addr_str[64]; - uint16_t port; - if (family == AF_INET) { - struct sockaddr_in *sa = (struct sockaddr_in *)cfg->remote_addr; - inet_ntop(AF_INET, &sa->sin_addr, addr_str, sizeof(addr_str)); - port = cfg->remote_port ? cfg->remote_port : ntohs(sa->sin_port); - } else if (family == AF_INET6) { - struct sockaddr_in6 *sa6 = (struct sockaddr_in6 *)cfg->remote_addr; - getnameinfo((struct sockaddr*)sa6, sizeof(*sa6), addr_str, sizeof(addr_str), NULL, 0, NI_NUMERICHOST); - port = cfg->remote_port ? cfg->remote_port : ntohs(sa6->sin6_port); - } else { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "stcp_link: unsupported address family %d", family); - return NULL; - } - - struct stcp_link *link = u_calloc(1, sizeof(struct stcp_link)); - if (!link) return NULL; - link->cfg = *cfg; - - uint8_t pubkey[SC_PUBKEY_SIZE]; - if (cfg->peer_pubkey_mode) { - struct secure_channel sc_tmp; - sc_init_ctx(&sc_tmp, cfg->my_keys); - if (sc_set_peer_public_key(&sc_tmp, cfg->peer_pubkey, 1) != SC_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid peer pubkey hex"); - u_free(link); return NULL; - } - memcpy(pubkey, sc_tmp.peer_public_key, SC_PUBKEY_SIZE); - } else { - memcpy(pubkey, cfg->peer_pubkey, SC_PUBKEY_SIZE); - } - - link->cli = stcp_client_connect(cfg->ua, addr_str, port, cfg->my_keys, pubkey, - client_ready_cb, link, NULL, NULL); - if (!link->cli) { u_free(link); return NULL; } - - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_link: connecting to %s:%u pubkey=%016llx", - addr_str, port, (unsigned long long)*(const uint64_t*)pubkey); - return link; -} - -static void stcp_link_close_impl(void *arg) { - struct stcp_link *link = (struct stcp_link *)arg; - if (link->on_close_cb) link->on_close_cb(link, 0, link->close_arg); - if (link->conn) { - if (link->conn->rx_queue) queue_free(link->conn->rx_queue); - if (link->tx_queue) queue_free(link->tx_queue); - stcp_conn_free(link->conn); - } - if (link->cli) stcp_client_destroy(link->cli); - u_free(link); -} - -void stcp_link_close(struct stcp_link *link) { - if (!link) return; - if (link->closing) return; - link->closing = 1; - uasync_call_soon(link->cfg.ua, link, stcp_link_close_impl); -} - -int stcp_link_send(struct stcp_link *link, const uint8_t *data, size_t len) { - if (!link || !link->ready) return -1; - struct ll_entry *e = queue_entry_new(0); - if (!e) return -1; - e->dgram = u_malloc(len ? len : 1); - if (!e->dgram) { queue_entry_free(e); return -1; } - if (len) memcpy(e->dgram, data, len); - e->len = (uint16_t)len; - queue_data_put(link->conn->tx_queue, e); - return 0; -} - -int stcp_link_is_ready(struct stcp_link *link) { - return link ? link->ready : 0; -} - -struct ETCP_CONN *stcp_link_get_etcp_conn(struct stcp_link *link) { - return link ? link->etcp_conn : NULL; -} - -void stcp_link_set_etcp_conn(struct stcp_link *link, struct ETCP_CONN *conn) { - if (!link) return; - link->etcp_conn = conn; -} - -void stcp_link_set_etcp_link(struct stcp_link *link, struct ETCP_LINK *elink) { - if (!link) return; - link->etcp_link = elink; -} - -const struct sockaddr_storage *stcp_link_get_remote_addr(struct stcp_link *link) { - return link && link->cfg.remote_addr ? link->cfg.remote_addr : NULL; -} - -void stcp_link_set_on_ready(struct stcp_link *link, stcp_link_cb cb, void *arg) { - if (!link) return; - link->on_ready_cb = cb; - link->ready_arg = arg; -} - -void stcp_link_set_on_close(struct stcp_link *link, void (*cb)(struct stcp_link *link, int err, void *arg), void *arg) { - if (!link) return; - link->on_close_cb = cb; - link->close_arg = arg; -} - -const uint8_t *stcp_link_get_peer_pubkey(struct stcp_link *link) { - return link && link->conn ? link->conn->peer_pubkey : NULL; -} - -void stcp_server_list_add(struct UTUN_INSTANCE *inst, struct stcp_server *srv) { - if (!inst || !srv) return; - srv->next = inst->stcp_servers; - inst->stcp_servers = srv; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server_list_add: %p total=%d", (void*)srv, stcp_server_list_count(inst)); -} - -void stcp_server_list_destroy_all(struct UTUN_INSTANCE *inst) { - if (!inst) return; - int count = 0; - while (inst->stcp_servers) { - struct stcp_server *next = inst->stcp_servers->next; - stcp_link_server_destroy(inst->stcp_servers); - inst->stcp_servers = next; - count++; - } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server_list_destroy_all: destroyed %d servers", count); -} - -int stcp_server_list_count(struct UTUN_INSTANCE *inst) { - if (!inst) return 0; - int n = 0; - for (struct stcp_server *s = inst->stcp_servers; s; s = s->next) n++; - return n; -} diff --git a/tests/test_etcp_stcp.c b/tests/test_etcp_stcp.c deleted file mode 100644 index cafb0286..00000000 --- a/tests/test_etcp_stcp.c +++ /dev/null @@ -1,153 +0,0 @@ -// test_etcp_stcp.c — integration: etcp_send + etcp_bind over STCP link -#include "stcp_link.h" -#include "etcp_api.h" -#include "etcp.h" -#include "etcp_connections.h" -#include "secure_channel.h" -#include "topo_group.h" -#include "topo_node.h" -#include "../src/utun_instance.h" -#include "../lib/u_async.h" -#include "../lib/ll_queue.h" -#include "../lib/debug_config.h" -#include "../lib/mem.h" -#include -#include -#include -#include - -#define TEST_PORT 25678 - -static int test_failed = 0; -static struct SC_MYKEYS s_keys, c_keys; -static struct UTUN_INSTANCE *srv_inst, *cli_inst; -static struct ETCP_CONN *srv_conn, *cli_conn; -static int g_srv_recv_count = 0; -static uint8_t g_srv_recv_data[256]; - -#define TASSERT(cond) do { \ - if (!(cond)) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, " FAIL: %s", #cond); test_failed = 1; return test_failed; } \ -} while(0) - -static void etcp_recv_cb(struct ETCP_CONN *conn, struct ll_entry *entry) { - (void)conn; - g_srv_recv_count++; - if (entry->dgram && entry->len) { - size_t n = entry->len < 256 ? entry->len : 255; - memcpy(g_srv_recv_data, entry->dgram, n); - } - queue_dgram_free(entry); - queue_entry_free(entry); -} - -static void connect_cb(void *arg, struct ETCP_CONN *conn, int type) { - (void)arg; - if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "etcp_connect failed"); test_failed = 1; return; } - if (type & ETCP_CONNECT_EARLY) { - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "etcp_connect: EARLY ready"); - cli_conn = conn; - } -} - -static char *build_server_config(int port) { - static char buf[512]; - snprintf(buf, sizeof(buf), - "[global]\n" - "my_node_id=0xAAAAAAAAAAAAAAAA\n" - "my_private_key=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68\n" - "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" - "\n" - "[server: stcp]\n" - "addr=127.0.0.1:%d\n" - "type=public\n" - "transport=tcp\n", - port); - return buf; -} - -static char *build_client_config(int port) { - static char buf[512]; - snprintf(buf, sizeof(buf), - "[global]\n" - "my_node_id=0xBBBBBBBBBBBBBBBB\n" - "my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" - "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" - "\n" - "[server: stcp]\n" - "addr=127.0.0.1:%d\n" - "type=public\n" - "transport=tcp\n", - port); - return buf; -} - -int main(void) { - debug_config_init(); - debug_set_level(DEBUG_LEVEL_INFO); - debug_set_categories(DEBUG_CATEGORY_GENERAL | DEBUG_CATEGORY_SOCKET | DEBUG_CATEGORY_CRYPTO); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "=== etcp_send/bind over STCP link ==="); - - srand((unsigned)time(NULL)); - struct UASYNC *ua = uasync_create(); TASSERT(ua); - utun_instance_set_tun_init_enabled(0); - - int port = TEST_PORT + rand() % 1000; - - srv_inst = utun_instance_create_from_str(ua, build_server_config(port)); - TASSERT(srv_inst); - TASSERT(utun_instance_init(srv_inst) >= 0); - - cli_inst = utun_instance_create_from_str(ua, build_client_config(port + 1)); - TASSERT(cli_inst); - TASSERT(utun_instance_init(cli_inst) >= 0); - - TASSERT(etcp_bind(srv_inst, 0, etcp_recv_cb) >= 0); - - struct TOPO_GROUP_NODE *srv_node = topo_groups_get_default(srv_inst->topo_groups) ? topo_groups_get_default(srv_inst->topo_groups)->local_node : NULL; - TASSERT(srv_node); - { struct TOPO_NODE* srv_ni = topo_node_registry_find(srv_inst->topo_groups, srv_node->node_id); TASSERT(srv_ni); - struct TOPO_NODE* cli_ni = u_calloc(1, sizeof(struct TOPO_NODE)); - TASSERT(cli_ni); memcpy(cli_ni, srv_ni, sizeof(struct TOPO_NODE)); cli_ni->group_ref_count = 0; - cli_ni->v4_sock_meta = NULL; cli_ni->v4_addrs = NULL; cli_ni->v6_sock_meta = NULL; cli_ni->v6_addrs = NULL; - struct TOPO_ADDR4* a = memory_pool_alloc(cli_inst->topo_groups->v4_addr_pool); - TASSERT(a); a->addr[0]=127; a->addr[1]=0; a->addr[2]=0; a->addr[3]=1; a->port = port; - a->type = TOPO_ADDR_INTERFACE; a->socket_id = 0; a->protocol = TOPO_PROTO_TCP; - a->next = cli_ni->v4_addrs; cli_ni->v4_addrs = a; - topo_node_registry_store(cli_inst->topo_groups, cli_ni); } - cli_inst->etcp_connect_timeout_tb = 100000; - TASSERT(etcp_connect(cli_inst, srv_node, connect_cb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE) == 0); - - int ticks = 0; - while (!cli_conn && ticks < 5000) { uasync_poll(ua, 10); ticks++; } - TASSERT(cli_conn); - - { struct ll_entry* e = srv_inst->connections->head; - while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; - struct ETCP_CONN *c = ce->conn; - struct ETCP_LINK *l = c->links; while (l) { if (l->is_tcp) { srv_conn = c; break; } l = l->next; } - if (srv_conn) break; - e = e->next; } - } - TASSERT(srv_conn); - - const char *msg = "hello via etcp_send!"; - struct ll_entry *e = queue_entry_new(0); - e->dgram = u_malloc(strlen(msg) + 1); - strcpy((char *)e->dgram, msg); - e->len = (uint16_t)strlen(msg); - TASSERT(etcp_send(cli_conn, e) == 0); - - ticks = 0; - while (g_srv_recv_count < 1 && ticks < 2000) { uasync_poll(ua, 10); ticks++; } - TASSERT(g_srv_recv_count == 1); - TASSERT(memcmp(g_srv_recv_data, msg, strlen(msg)) == 0); - - srv_inst->running = 0; cli_inst->running = 0; - utun_instance_destroy(srv_inst); srv_inst = NULL; - utun_instance_destroy(cli_inst); cli_inst = NULL; - uasync_destroy(ua, 0); - - if (test_failed) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "=== FAILED ==="); return 1; } - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "=== PASSED ==="); - return 0; -} diff --git a/tests/test_stcp_traffic.c b/tests/test_stcp_traffic.c deleted file mode 100644 index 922aeaad..00000000 --- a/tests/test_stcp_traffic.c +++ /dev/null @@ -1,234 +0,0 @@ -// test_stcp_traffic.c — STCP transport test: 2 nodes, 1MB bidirectional, random packets, content verify -#include "stcp_link.h" -#include "etcp_api.h" -#include "etcp.h" -#include "etcp_connections.h" -#include "secure_channel.h" -#include "topo_group.h" -#include "topo_node.h" -#include "../src/utun_instance.h" -#include "../lib/u_async.h" -#include "../lib/ll_queue.h" -#include "../lib/debug_config.h" -#include "../lib/mem.h" -#include -#include -#include -#include - -#define TEST_PORT 25680 -#define TRAFFIC_MB (1024*1024) -#define PKT_MIN 10 -#define PKT_MAX 1800 - -static int test_failed = 0; - -#define TASSERT(cond) do { \ - if (!(cond)) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, " FAIL: %s", #cond); test_failed = 1; return 1; } \ -} while(0) - -struct rx_ctx { - uint8_t *buf; - size_t len, cap; - int pkts; - int done; - size_t expected_total; -}; - -static struct rx_ctx srv_rx, cli_rx; -static struct UTUN_INSTANCE *srv_inst, *cli_inst; -static struct ETCP_CONN *srv_conn, *cli_conn; -static int srv_send_ready, cli_send_ready; - -static void srv_recv_cb(struct ETCP_CONN *conn, struct ll_entry *entry) { - (void)conn; - if (!entry->dgram || entry->len < 5) { queue_dgram_free(entry); queue_entry_free(entry); return; } - srv_rx.pkts++; - size_t pay = entry->len - 5; - size_t need = srv_rx.len + pay; - if (need > srv_rx.cap) { srv_rx.cap = need + 65536; srv_rx.buf = u_realloc(srv_rx.buf, srv_rx.cap); } - memcpy(srv_rx.buf + srv_rx.len, entry->dgram + 5, pay); - srv_rx.len += pay; - queue_dgram_free(entry); - queue_entry_free(entry); - if (srv_rx.len >= srv_rx.expected_total) srv_rx.done = 1; -} - -static void cli_recv_cb(struct ETCP_CONN *conn, struct ll_entry *entry) { - (void)conn; - if (!entry->dgram || entry->len < 5) { queue_dgram_free(entry); queue_entry_free(entry); return; } - cli_rx.pkts++; - size_t pay = entry->len - 5; - size_t need = cli_rx.len + pay; - if (need > cli_rx.cap) { cli_rx.cap = need + 65536; cli_rx.buf = u_realloc(cli_rx.buf, cli_rx.cap); } - memcpy(cli_rx.buf + cli_rx.len, entry->dgram + 5, pay); - cli_rx.len += pay; - queue_dgram_free(entry); - queue_entry_free(entry); - if (cli_rx.len >= cli_rx.expected_total) cli_rx.done = 1; -} - -static void connect_cb(void *arg, struct ETCP_CONN *conn, int type) { - (void)arg; - if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "etcp_connect failed"); test_failed = 1; return; } - if (type & ETCP_CONNECT_EARLY) { - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "etcp_connect: EARLY ready"); - cli_conn = conn; - cli_send_ready = 1; - } -} - -static size_t gen_packet(uint8_t *buf, int side, int seq, size_t paylen) { - buf[0] = (uint8_t)side; - buf[1] = (uint8_t)(seq >> 0); - buf[2] = (uint8_t)(seq >> 8); - buf[3] = (uint8_t)(seq >> 16); - buf[4] = (uint8_t)(seq >> 24); - size_t i; - for (i = 5; i < paylen + 5; i++) buf[i] = (uint8_t)((i * 7 + seq * 13 + side) & 0xFF); - return paylen + 5; -} - -static int send_etcp_packet(struct ETCP_CONN *conn, int side, int seq, size_t paylen) { - uint8_t buf[PKT_MAX + 5]; - size_t tot = gen_packet(buf, side, seq, paylen); - struct ll_entry *e = queue_entry_new(0); - if (!e) return -1; - e->dgram = u_malloc(tot); - if (!e->dgram) { queue_entry_free(e); return -1; } - memcpy(e->dgram, buf, tot); - e->len = (uint16_t)tot; - return etcp_send(conn, e); -} - -static char *build_server_config(int port) { - static char buf[512]; - snprintf(buf, sizeof(buf), - "[global]\n" - "my_node_id=0xAAAAAAAAAAAAAAAA\n" - "my_private_key=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68\n" - "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" - "\n" - "[server: stcp]\n" - "addr=127.0.0.1:%d\n" - "type=public\n" - "transport=tcp\n", - port); - return buf; -} - -static char *build_client_config(int port) { - static char buf[512]; - snprintf(buf, sizeof(buf), - "[global]\n" - "my_node_id=0xBBBBBBBBBBBBBBBB\n" - "my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" - "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" - "\n" - "[server: stcp]\n" - "addr=127.0.0.1:%d\n" - "type=public\n" - "transport=tcp\n", - port); - return buf; -} - -int main(void) { - debug_config_init(); - debug_set_level(DEBUG_LEVEL_INFO); - debug_set_categories(DEBUG_CATEGORY_GENERAL | DEBUG_CATEGORY_SOCKET | DEBUG_CATEGORY_CRYPTO | DEBUG_CATEGORY_ETCP); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "=== STCP Traffic Test ==="); - - srand((unsigned)time(NULL)); - struct UASYNC *ua = uasync_create(); TASSERT(ua); - utun_instance_set_tun_init_enabled(0); - - int port = TEST_PORT + rand() % 1000; - - srv_inst = utun_instance_create_from_str(ua, build_server_config(port)); - TASSERT(srv_inst); - TASSERT(utun_instance_init(srv_inst) >= 0); - - cli_inst = utun_instance_create_from_str(ua, build_client_config(port + 1)); - TASSERT(cli_inst); - TASSERT(utun_instance_init(cli_inst) >= 0); - - memset(&srv_rx, 0, sizeof(srv_rx)); srv_rx.expected_total = TRAFFIC_MB; - memset(&cli_rx, 0, sizeof(cli_rx)); cli_rx.expected_total = TRAFFIC_MB; - TASSERT(etcp_bind(srv_inst, 1, srv_recv_cb) >= 0); - TASSERT(etcp_bind(cli_inst, 2, cli_recv_cb) >= 0); - - struct TOPO_GROUP_NODE *srv_node = topo_groups_get_default(srv_inst->topo_groups) ? topo_groups_get_default(srv_inst->topo_groups)->local_node : NULL; - TASSERT(srv_node); - { struct TOPO_NODE* srv_ni = topo_node_registry_find(srv_inst->topo_groups, srv_node->node_id); TASSERT(srv_ni); - struct TOPO_NODE* cli_ni = u_calloc(1, sizeof(struct TOPO_NODE)); - TASSERT(cli_ni); memcpy(cli_ni, srv_ni, sizeof(struct TOPO_NODE)); cli_ni->group_ref_count = 0; - cli_ni->v4_sock_meta = NULL; cli_ni->v4_addrs = NULL; cli_ni->v6_sock_meta = NULL; cli_ni->v6_addrs = NULL; - struct TOPO_ADDR4* a = memory_pool_alloc(cli_inst->topo_groups->v4_addr_pool); - TASSERT(a); a->addr[0]=127; a->addr[1]=0; a->addr[2]=0; a->addr[3]=1; a->port = port; - a->type = TOPO_ADDR_INTERFACE; a->socket_id = 0; a->protocol = TOPO_PROTO_TCP; - a->next = cli_ni->v4_addrs; cli_ni->v4_addrs = a; - topo_node_registry_store(cli_inst->topo_groups, cli_ni); } - cli_inst->etcp_connect_timeout_tb = 100000; - TASSERT(etcp_connect(cli_inst, srv_node, connect_cb, NULL, ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE) == 0); - - int ticks = 0; - while (!cli_send_ready && ticks < 5000) { uasync_poll(ua, 10); ticks++; } - TASSERT(cli_send_ready); - TASSERT(cli_conn); - - { struct ll_entry* e = srv_inst->connections->head; - while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; - struct ETCP_CONN *c = ce->conn; - struct ETCP_LINK *l = c->links; while (l) { if (l->is_tcp) { srv_conn = c; break; } l = l->next; } - if (srv_conn) break; - e = e->next; } - } - TASSERT(srv_conn); - - size_t remaining_s = TRAFFIC_MB, remaining_c = TRAFFIC_MB; - int seq_s = 0, seq_c = 0; - int srv_done = 0, cli_done = 0; - - ticks = 0; - while (!srv_done || !cli_done) { - uasync_poll(ua, 1); - - if (!srv_done && remaining_s > 0) { - size_t sz = PKT_MIN + (rand() % (PKT_MAX - PKT_MIN + 1)); - if (sz > remaining_s) sz = remaining_s; - if (send_etcp_packet(srv_conn, 2, seq_s++, sz) == 0) remaining_s -= sz; - } - if (remaining_s == 0 && seq_s > 0) srv_done = 1; - - if (!cli_done && remaining_c > 0) { - size_t sz = PKT_MIN + (rand() % (PKT_MAX - PKT_MIN + 1)); - if (sz > remaining_c) sz = remaining_c; - if (send_etcp_packet(cli_conn, 1, seq_c++, sz) == 0) remaining_c -= sz; - } - if (remaining_c == 0 && seq_c > 0) cli_done = 1; - - if (++ticks > 50000) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "send TIMEOUT"); test_failed = 1; break; } - } - - ticks = 0; - while ((!srv_rx.done || !cli_rx.done) && ticks < 20000) { uasync_poll(ua, 10); ticks++; } - - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "srv_rx: %zu bytes, %d pkts, done=%d", srv_rx.len, srv_rx.pkts, srv_rx.done); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "cli_rx: %zu bytes, %d pkts, done=%d", cli_rx.len, cli_rx.pkts, cli_rx.done); - - TASSERT(srv_rx.len >= TRAFFIC_MB); - TASSERT(cli_rx.len >= TRAFFIC_MB); - - if (srv_rx.buf) u_free(srv_rx.buf); - if (cli_rx.buf) u_free(cli_rx.buf); - - srv_inst->running = 0; cli_inst->running = 0; - utun_instance_destroy(srv_inst); srv_inst = NULL; - utun_instance_destroy(cli_inst); cli_inst = NULL; - uasync_destroy(ua, 0); - - if (test_failed) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "=== FAILED ==="); return 1; } - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "=== PASSED ==="); - return 0; -}