You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
1841 lines
84 KiB
1841 lines
84 KiB
// 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 <stdlib.h> |
|
#include <string.h> |
|
#include <sys/time.h> |
|
#include <math.h> // For bandwidth calcs |
|
#include <limits.h> // 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); |
|
etcp_router_conn_destroyed(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 || peer_reset_id == 0) return; /* 0 = неизвестно (STCP-сервер ещё не создал conn) */ |
|
if (peer_reset_id == conn->peer_reset_id) return; // uid пира не изменился |
|
uint64_t old = conn->peer_reset_id; |
|
conn->peer_reset_id = peer_reset_id; |
|
if (old != 0 && 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; i<elm_cnt; i++) { |
|
uint32_t seq=data[-ack_section_len+8+i*8] | (data[-ack_section_len+9+i*8]<<8) | (data[-ack_section_len+10+i*8]<<16) | (data[-ack_section_len+11+i*8]<<24); |
|
uint16_t ts=data[-ack_section_len+12+i*8] | (data[-ack_section_len+13+i*8]<<8); |
|
uint16_t dts=data[-ack_section_len+14+i*8] | (data[-ack_section_len+15+i*8]<<8); |
|
etcp_ack_recv(etcp, seq, ts, dts); |
|
} |
|
while ((int32_t)(etcp->rx_ack_till-till)<0) { etcp->rx_ack_till++; etcp_ack_recv(etcp, etcp->rx_ack_till, -1, -1); }// подтверждаем всё по till |
|
break; |
|
} |
|
case ETCP_SECTION_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; |
|
} |
|
} |
|
} |
|
|
|
|
|
|