// 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 "../lib/mem.h" #include "../lib/strbuf.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_ETCP, "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_ETCP, "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; random_bytes((uint8_t*)&etcp->reset_id, 8); // начальная случайная эпоха ресета 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; struct strbuf sb = strbuf_new(); int any = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { if (l->link_status) { if (any) strbuf_adds(&sb, ", "); any = 1; strbuf_adds(&sb, 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, any > 0 ? " — up: " : "", strbuf_str(&sb)); strbuf_free(&sb); (void)elapsed_ms; etcp_cbk_fire(etcp, ETCP_CBK_EVENT_UP); etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_UP); } static void etcp_on_down(struct ETCP_CONN* etcp, struct ETCP_LINK* down_link) { struct strbuf sb = strbuf_new(); int any = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { if (l == down_link) continue; if (any) strbuf_adds(&sb, ", "); any = 1; strbuf_adds(&sb, 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, any > 0 ? ", remaining: " : "", strbuf_str(&sb)); else DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — closing, links: %s", etcp->log_name, any ? strbuf_str(&sb) : "(none)"); strbuf_free(&sb); 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_ETCP, "[%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 === etcp->state = 2; // deleted — blocks ref_take and etcp_send immediately if (etcp->links_up != 0) etcp->links_up = 0; etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DELETE); etcp_cbk_fire(etcp, ETCP_CBK_EVENT_DOWN); // 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; } routing_del_conn(etcp); etcp_router_transit_queues_destroy(etcp); if (etcp->normalizer) { pn_deinit((struct PKTNORM*)etcp->normalizer); etcp->normalizer = NULL; etcp->send_input_q = 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; } // === 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) { if (!etcp || etcp->state == 2) { if (etcp) DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] conn_reset on deleted conn", etcp->log_name); return; } // 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) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение if (!etcp || etcp->state == 2) return; // Сбрасываем 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) { etcp_conn_reinit_id(etcp, reason, etcp ? etcp->reset_id : 0); } // Фатальный ресет (normalizer desync / bad fragment size): меняем эпоху reset_id void etcp_conn_fatal_reinit(struct ETCP_CONN* etcp, const char* reason) { uint64_t new_id; random_bytes((uint8_t*)&new_id, 8); etcp_conn_reinit_id(etcp, reason, new_id); } void etcp_conn_reinit_id(struct ETCP_CONN* etcp, const char* reason, uint64_t reset_id) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение if (!etcp) return; if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] conn_reinit on deleted conn", etcp->log_name); return; } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] REINIT: %s (init=%d links=%d reinit=%u tx=%d) rid=%016llx", etcp->log_name, reason, etcp->initialized, etcp->links_up, etcp->reinit_count, etcp->tx_state, (unsigned long long)reset_id); etcp->reset_id = reset_id; 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"); } // Применение чужого reset_id (эпохи ресета пира). Симметрично: // если uid пира изменился — сбрасываем своё состояние (свой uid сохраняем). void etcp_conn_apply_peer_reset_id(struct ETCP_CONN* conn, uint64_t peer_reset_id) { if (!conn) return; if (peer_reset_id == conn->peer_reset_id) return; // uid пира не изменился conn->peer_reset_id = peer_reset_id; if (conn->got_initial_pkt) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] peer reset id changed to %016llx — resetting self", conn->log_name, (unsigned long long)peer_reset_id); etcp_conn_reinit(conn, "peer reset epoch"); } else { DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] peer reset id=%016llx (clean, no reset)", conn->log_name, (unsigned long long)peer_reset_id); } } // внутренняя функция. Вызывается один раз когда первый линк готов. void etcp_conn_ready(struct ETCP_CONN* conn) { if (!conn) return; if (conn->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] conn_ready on deleted conn", conn->log_name); 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); 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* existing = queue_find_data_by_index(conn->instance->connections, (const uint8_t*)&conn->peer_node_id); if (existing) { struct conn_queue_entry* ece = (struct conn_queue_entry*)existing->data; DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] queue_set_ready: DUPLICATE node_id key=0x%016llx already indexed by [%s] conn=%p state=%d (this=%p)", conn->log_name, (unsigned long long)conn->peer_node_id, ece->conn ? ece->conn->log_name : "?", (void*)(ece->conn), ece->conn ? ece->conn->state : -1, (void*)conn); } 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_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->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] int_send on deleted conn", etcp->log_name); 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 || etcp->state == 2) 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++; } 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, ""); if (!etcp || etcp->state == 2) return; 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); } 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) { struct strbuf sb = strbuf_new(); struct ETCP_LINK* link = etcp->links; while (link) { strbuf_addf(&sb, "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, strbuf_str(&sb)); strbuf_free(&sb); 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, ""); if (!etcp) return NULL; if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] request_pkt on deleted conn", etcp->log_name); return NULL; } 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++; if (inf_pkt->last_link->bbr) { 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 } 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; DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] KA-TX: ts_section dt=%llu tcp=%d", link->etcp->log_name, (unsigned long long)dt, dgram->link->is_tcp); 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; if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ack_recv on deleted conn", etcp->log_name); 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 && acked_pkt->last_link->bbr) { 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; DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] KA-RTT: rtt=%u cur=%u ret=%u dlen=%u", etcp->log_name, new_rtt, cur_ts, ret_ts, pkt->data_len); 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); 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) { 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; } 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 || etcp->state == 2) 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) { int max_udp = IP_UDP_OVERHEAD_V4; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) if (l->remote_addr.ss_family == AF_INET6) max_udp = IP_UDP_OVERHEAD_V6; etcp->normalizer->frag_size = etcp->mtu - ACK_REZERV - max_udp - UDP_SC_HDR_SIZE; } } }