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.
 
 
 
 
 
 

2012 lines
94 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 <time.h> // For strftime in metrics snapshot
#include "../lib/mem.h"
#include "../lib/memory_pool.h"
// Enable comprehensive debug output for ETCP module
#define DEBUG_CATEGORY_ETCP_DETAILED 1
// Constants from spec (adjusted for completeness)
#define MAX_INFLIGHT_BYTES 65536 // Initial window
#define RETRANS_K1 32 // RTT multiplier for retrans timeout
#define RETRANS_K2 32 // Jitter multiplier
#define ACK_DELAY_TB 20 // ACK timer delay (2ms in 0.1ms units)
#define BURST_DELAY_FACTOR 4 // Delay before burst
#define BURST_SIZE 5 // Packets in burst (1 delayed + 4 burst)
#define RTT_HISTORY_SIZE 10 // For jitter calc
#define MAX_PENDING 32 // For ACKs/retrans (arbitrary; adjust)
#define SECTION_HEADER_SIZE 3 // type(1) + len(2)
// Container-of macro for getting struct from data pointer
//#define CONTAINER_OF(ptr, type, member) ((type *)((char *)(ptr) - offsetof(type, member)))
// Forward declarations
static void input_queue_cb(struct ll_queue* q, void* arg);
static void etcp_link_ready_callback(struct ETCP_CONN* etcp);
static void input_send_q_cb(struct ll_queue* q, void* arg);
static void wait_ack_cb(struct ll_queue* q, void* arg);
static void send_ack_req_cb(struct ll_queue* q, void* arg);
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp);
static void etcp_connection_free_deferred(void* arg);
void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status);
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp);
static void clear_queue(struct ll_queue* q) {
if (!q) return;
struct ll_entry* pkt;
while ((pkt = queue_data_get(q)) != NULL) {
queue_dgram_free(pkt);
queue_entry_free(pkt);
}
}
static void drain_and_free_queue(struct ll_queue** q) {
if (!*q) return;
clear_queue(*q);
queue_free(*q);
*q = NULL;
}
static int inflight_seq_cmp(const void* a, const void* b) {
const struct INFLIGHT_PACKET* pa = *(const struct INFLIGHT_PACKET**)a;
const struct INFLIGHT_PACKET* pb = *(const struct INFLIGHT_PACKET**)b;
if (pa->seq < pb->seq) return -1;
if (pa->seq > pb->seq) return 1;
return 0;
}
static void feed_dgram_to_asm(struct ETCP_CONN* etcp, struct PKTNORM* pn,
uint8_t* dgram, uint16_t dgram_len,
uint8_t** asm_buf, uint16_t* asm_len, uint32_t* asm_cap,
uint32_t* returned) {
uint32_t need = (uint32_t)*asm_len + dgram_len;
if (need >= ASM_BUF_MAX_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] assembly buffer would exceed max size %d", etcp->log_name, ASM_BUF_MAX_SIZE);
return;
}
if (need > *asm_cap) {
uint32_t new_cap = *asm_cap ? *asm_cap : 256;
while (new_cap < need) new_cap *= 2;
if (new_cap > 0xFFFF) new_cap = 0xFFFF;
uint8_t* new_buf = u_realloc(*asm_buf, (uint16_t)new_cap);
if (!new_buf) return;
*asm_buf = new_buf;
*asm_cap = new_cap;
}
memcpy(*asm_buf + *asm_len, dgram, dgram_len);
*asm_len = need;
while (*asm_len >= 2) {
uint16_t pkt_len = (*asm_buf)[0] | ((*asm_buf)[1] << 8);
if (pkt_len == 0 || pkt_len > 16384 || *asm_len < 2 + pkt_len) break;
const uint8_t* pkt_data = *asm_buf + 2;
if (pkt_data[0] != ETCP_ID_TOPO_ENTRY) {
struct ll_entry* e = ll_alloc_lldgram(pkt_len);
if (e) {
memcpy(e->dgram, pkt_data, pkt_len);
e->len = pkt_len;
queue_data_put(pn->input, e);
(*returned)++;
}
}
uint16_t consumed = 2 + pkt_len;
*asm_len -= consumed;
if (*asm_len > 0) memmove(*asm_buf, *asm_buf + consumed, *asm_len);
}
}
static void etcp_return_inflight_to_normalizer(struct ETCP_CONN* etcp) {
struct PKTNORM* pn = etcp->normalizer;
if (!pn) return;
struct INFLIGHT_PACKET** inflight = NULL;
int inflight_count = 0;
struct ll_entry** input_frags = NULL;
int input_count = 0;
struct ll_entry* entry;
while ((entry = queue_data_get(etcp->input_wait_ack))) {
inflight = u_realloc(inflight, (inflight_count + 1) * sizeof(*inflight));
inflight[inflight_count++] = (struct INFLIGHT_PACKET*)entry;
}
while ((entry = queue_data_get(etcp->input_send_q))) {
inflight = u_realloc(inflight, (inflight_count + 1) * sizeof(*inflight));
inflight[inflight_count++] = (struct INFLIGHT_PACKET*)entry;
}
while ((entry = queue_data_get(etcp->input_queue))) {
input_frags = u_realloc(input_frags, (input_count + 1) * sizeof(*input_frags));
input_frags[input_count++] = entry;
}
if (inflight_count + input_count == 0 && (!pn->data || pn->data_ptr == 0)) return;
if (inflight_count > 0) qsort(inflight, inflight_count, sizeof(*inflight), inflight_seq_cmp);
uint8_t* asm_buf = NULL;
uint16_t asm_len = 0;
uint32_t asm_cap = 0;
uint32_t returned_packets = 0;
for (int i = 0; i < inflight_count; i++) {
feed_dgram_to_asm(etcp, pn, inflight[i]->ll.dgram, inflight[i]->ll.len,
&asm_buf, &asm_len, &asm_cap, &returned_packets);
memory_pool_free(etcp->instance->data_pool, inflight[i]->ll.dgram);
memory_pool_free(etcp->inflight_pool, inflight[i]);
}
u_free(inflight);
for (int i = 0; i < input_count; i++) {
feed_dgram_to_asm(etcp, pn, input_frags[i]->dgram, input_frags[i]->len,
&asm_buf, &asm_len, &asm_cap, &returned_packets);
memory_pool_free(etcp->instance->data_pool, input_frags[i]->dgram);
memory_pool_free(etcp->io_pool, input_frags[i]);
}
u_free(input_frags);
if (pn->data && pn->data_ptr > 0) {
feed_dgram_to_asm(etcp, pn, pn->data, pn->data_ptr,
&asm_buf, &asm_len, &asm_cap, &returned_packets);
}
if (asm_buf) { u_free(asm_buf); }
if (returned_packets > 0) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] returned %u packets to normalizer input during reinit",
etcp->log_name, returned_packets);
}
}
uint16_t get_current_timestamp() {
return (uint16_t)get_time_tb();
}
// Timestamp diff (with wrap-around)
static uint16_t timestamp_diff(uint16_t t1, uint16_t t2) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (t1 >= t2) {
return t1 - t2;
}
return (UINT16_MAX - t2) + t1 + 1;
}
static const char EMPTY_NAME[] = "";
// Create new ETCP connection
struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!instance) return NULL;
struct ETCP_CONN* etcp = u_calloc(1, sizeof(struct ETCP_CONN));
if (!etcp) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "creating connection failed for instance %p", instance);
return NULL;
}
etcp->instance = instance;
etcp->input_queue = queue_new(instance->ua, 0, 0, 0, "ETCP input"); // No hash for input_queue
queue_set_threshold(etcp->input_queue, 0, 0); // Backpressure: ждать полного освобождения
etcp->output_queue = queue_new(instance->ua, 0, 0, 0, "ETCP output"); // No hash for output_queue
etcp->transit_queues = queue_new(instance->ua, TRANSIT_QUEUE_HASH, 0, 16, "transit_q_reg"); // Hash for transit queues
etcp->input_send_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "input_send_q"); // Hash for send_q
etcp->input_wait_ack = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "input_wait_ack"); // Hash for wait_ack
etcp->recv_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "recv_q"); // Hash for recv_q
etcp->ack_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "ack_q");
etcp->inflight_pool = memory_pool_init(sizeof(struct INFLIGHT_PACKET), "inflight_pool");
etcp->io_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT), "io_pool");
etcp->max_inflight = (uint32_t)instance->config->global.bbr_max_cwnd;
etcp->optimal_inflight=100000;
etcp->initialized=0;
etcp->links_up=0;
etcp->setup_start_tb = get_time_tb();
etcp->reset_done=0;
etcp->callbacks_running=0;
etcp->ref_count=0;
etcp->last_rr_link=NULL;
etcp->name = u_strdup(name);
// Initialize log_name with local node_id (peer will be updated later when known)
snprintf(etcp->log_name, sizeof(etcp->log_name), "%04X->???? [%s]", (uint16_t)instance->node_id, etcp->name);
if (!etcp->input_queue || !etcp->output_queue || !etcp->transit_queues ||
!etcp->input_send_q || !etcp->recv_q || !etcp->ack_q ||
!etcp->input_wait_ack || !etcp->inflight_pool || !etcp->io_pool) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "error - closing - input:%p output:%p send:%p wait:%x pool:%p",
etcp->input_queue, etcp->output_queue, etcp->input_send_q, etcp->input_wait_ack, etcp->inflight_pool);
etcp_connection_close(etcp);
return NULL;
}
// etcp->window_size = MAX_INFLIGHT_BYTES; // Not used
etcp->mtu = ETCP_RFC791_MIN_MTU; // Default MTU per RFC 791
etcp->next_tx_id = 1;
etcp->rtt_avg_10 = 10; // Initial guess (1ms)
etcp->rtt_history_idx = 0;
memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history));
etcp->normalizer = pn_init(etcp);
if (!etcp->normalizer) {
etcp_connection_close(etcp);
return NULL;
}
etcp->send_input_q = etcp->normalizer->input;
// Set input queue callback
queue_set_callback(etcp->input_queue, input_queue_cb, etcp);
queue_set_callback(etcp->input_send_q, input_send_q_cb, etcp);
queue_set_callback(etcp->input_wait_ack, wait_ack_cb, etcp);
// queue_set_callback(etcp->ack_q, send_ack_req_cb, etcp);// отправляем через таймаут
etcp->link_ready_for_send_fn = etcp_link_ready_callback;
// Add to instance's connections queue (peer_node_id=0 for pending, reindexed later)
{
struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry));
if (!qe) { etcp_connection_close(etcp); return NULL; }
struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
ce->peer_node_id = 0;
ce->conn = etcp;
etcp->conn_queue_entry = qe;
etcp->conn_queue = instance->connections;
queue_data_put_with_index(instance->connections, qe);
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] connection initialized. ETCP=%p mtu=%d, next_tx_id=%u state=pending",
etcp->log_name, etcp, etcp->mtu, etcp->next_tx_id);
// Вызываем callback для нового соединения если установлен
if (instance) {
struct etcp_inst_cbk_entry* cbe = instance->new_conn_cbks;
while (cbe) { struct etcp_inst_cbk_entry* n = cbe->next; cbe->fn(etcp, cbe->arg); cbe = n; }
}
etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_NEW);
return etcp;
}
void etcp_fire_conn_status(struct ETCP_CONN* etcp, int status) {
if (!etcp || !etcp->instance) return;
static const char* names[] = { "NEW", "UP", "DOWN", "DELETE" };
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[%s] Connection status changed to %s", etcp->log_name,
(status >= 0 && status < (int)(sizeof(names)/sizeof(names[0]))) ? names[status] : "?");
struct etcp_status_cbk_entry* cbe = etcp->instance->conn_status_cbks;
while (cbe) { struct etcp_status_cbk_entry* n = cbe->next; cbe->fn(etcp, status, cbe->arg); cbe = n; }
}
void etcp_cbk_fire(struct ETCP_CONN* conn, int event) {
if (!conn) return;
static const char* names[] = { "INIT", "REINIT", "UP", "DOWN", "NODE_CHANGED" };
int idx = 0, e = event; while (e >>= 1) idx++;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] callback event: %s", conn->log_name, (idx >= 0 && idx < (int)(sizeof(names)/sizeof(names[0]))) ? names[idx] : "?");
conn->callbacks_running = 1;
struct etcp_cbk_entry* cbe = conn->cbks;
while (cbe) { struct etcp_cbk_entry* n = cbe->next; if (cbe->event_mask & event) cbe->fn(conn, event, cbe->arg); cbe = n; }
conn->callbacks_running = 0;
}
static void etcp_on_up(struct ETCP_CONN* etcp) {
int total_links = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) total_links++;
uint64_t elapsed_ms = (get_time_tb() - etcp->setup_start_tb) / 10;
char links_str[256] = {0}; int pp = 0;
for (struct ETCP_LINK* l = etcp->links; l; l = l->next) {
if (l->link_status) {
if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", ");
pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str);
}
}
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection UP (%d/%d links, mtu=%d, reinit=%u)%s%s",
etcp->log_name, etcp->links_up, total_links, etcp->mtu, etcp->reinit_count,
pp > 0 ? " — up: " : "", links_str);
(void)elapsed_ms;
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_ON_UP_BEFORE: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x",
etcp->log_name,
etcp->crypto_ctx.session_key[0], etcp->crypto_ctx.session_key[1],
etcp->crypto_ctx.session_key[2], etcp->crypto_ctx.session_key[3],
etcp->crypto_ctx.peer_public_key[0], etcp->crypto_ctx.peer_public_key[1],
etcp->crypto_ctx.peer_public_key[2], etcp->crypto_ctx.peer_public_key[3]);
etcp_cbk_fire(etcp, ETCP_CBK_EVENT_UP);
etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_UP);
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_ON_UP_AFTER: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x",
etcp->log_name,
etcp->crypto_ctx.session_key[0], etcp->crypto_ctx.session_key[1],
etcp->crypto_ctx.session_key[2], etcp->crypto_ctx.session_key[3],
etcp->crypto_ctx.peer_public_key[0], etcp->crypto_ctx.peer_public_key[1],
etcp->crypto_ctx.peer_public_key[2], etcp->crypto_ctx.peer_public_key[3]);
}
static void etcp_on_down(struct ETCP_CONN* etcp, struct ETCP_LINK* down_link) {
char links_str[256] = {0}; int pp = 0;
for (struct ETCP_LINK* l = etcp->links; l; l = l->next) {
if (l == down_link) continue;
if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", ");
pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str);
}
if (down_link)
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — link %s lost%s%s",
etcp->log_name, sockaddr_storage_to_str(&down_link->remote_addr).str,
pp > 0 ? ", remaining: " : "", links_str);
else
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — closing, links: %s", etcp->log_name, pp ? links_str : "(none)");
etcp_cbk_fire(etcp, ETCP_CBK_EVENT_DOWN);
etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DOWN);
}
// Phase 2 resources cleanup (callable both sync and async via call_soon)
static void etcp_connection_free_resources(struct ETCP_CONN* etcp) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] freeing resources phase 2", etcp->log_name);
// Close links
if (etcp->links) {
struct ETCP_LINK* link = etcp->links;
while (link) { struct ETCP_LINK* next = link->next; etcp_link_close(link); link = next; }
etcp->links = NULL;
}
// Drain and free all queues
drain_and_free_queue(&etcp->input_queue);
drain_and_free_queue(&etcp->output_queue);
drain_and_free_queue(&etcp->input_send_q);
drain_and_free_queue(&etcp->input_wait_ack);
drain_and_free_queue(&etcp->recv_q);
drain_and_free_queue(&etcp->ack_q);
// Free callback chain
{ struct etcp_cbk_entry* cbe = etcp->cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->cbks = NULL; }
if (etcp->inflight_pool) { memory_pool_destroy(etcp->inflight_pool); etcp->inflight_pool = NULL; }
if (etcp->io_pool) { memory_pool_destroy(etcp->io_pool); etcp->io_pool = NULL; }
etcp_connect_cancel_for_conn(etcp->instance, etcp);
u_free(etcp->name);
u_free(etcp);
}
static void etcp_connection_free_deferred(void* arg) {
etcp_connection_free_resources((struct ETCP_CONN*)arg);
}
// Close connection: phase 1 detach + deferred phase 2 cleanup
void etcp_connection_close(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!etcp) return;
if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] already deleted", etcp->log_name); return; }
if (etcp->callbacks_running) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "[%s] FATAL: etcp_connection_close called from inside callback chain — SEGFAULTING to show backtrace",
etcp->log_name);
*(volatile int*)0 = 0;
}
// === PHASE 1: detach from external world ===
if (etcp->links_up != 0) { etcp->links_up = 0; etcp_on_down(etcp, NULL); }
// Cancel active timers
if (etcp->retrans_timer) { uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer); etcp->retrans_timer = NULL; }
if (etcp->ack_resp_timer) { uasync_cancel_timeout(etcp->instance->ua, etcp->ack_resp_timer); etcp->ack_resp_timer = NULL; }
etcp_metrics_stop_timer(etcp);
routing_del_conn(etcp);
etcp_router_transit_queues_destroy(etcp);
if (etcp->normalizer) { pn_deinit((struct PKTNORM*)etcp->normalizer); etcp->normalizer = NULL; }
// Close links (detach from socket scan queues)
if (etcp->links) {
struct ETCP_LINK* link = etcp->links;
while (link) { struct ETCP_LINK* next = link->next; etcp_link_close(link); link = next; }
etcp->links = NULL;
}
// Remove from instance queue
if (etcp->conn_queue && etcp->conn_queue_entry) {
queue_remove_data(etcp->conn_queue, etcp->conn_queue_entry);
queue_entry_free(etcp->conn_queue_entry);
etcp->conn_queue_entry = NULL;
etcp->conn_queue = NULL;
}
etcp->state = 2; // deleted
etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DELETE);
// === PHASE 2: deferred resource cleanup (only if no outstanding refs) ===
if (etcp->ref_count == 0)
uasync_call_soon(etcp->instance->ua, etcp, etcp_connection_free_deferred);
// else: cleanup deferred until last etcp_conn_ref_free()
}
// Take a reference on connection. Returns -1 if already deleted (state==2).
int etcp_conn_ref_take(struct ETCP_CONN* conn) {
if (!conn) return -1;
if (conn->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ref_take on deleted conn", conn->log_name); return -1; }
conn->ref_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ref_take ref_count=%d", conn->log_name, conn->ref_count);
return 0;
}
// Release a reference. If last ref and conn is deleted, schedule deferred free.
void etcp_conn_ref_free(struct ETCP_CONN* conn) {
if (!conn) return;
if (conn->ref_count <= 0) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ref_free underflow ref_count=%d", conn->log_name, conn->ref_count); return; }
conn->ref_count--;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ref_free ref_count=%d state=%d", conn->log_name, conn->ref_count, conn->state);
if (conn->ref_count == 0 && conn->state == 2)
uasync_call_soon(conn->instance->ua, conn, etcp_connection_free_deferred);
}
// Reset connection
void etcp_conn_reset(struct ETCP_CONN* etcp) {
// Reset IDs
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "Resetting ETCP instance [%s]", etcp->log_name);
etcp->next_tx_id = 1;
etcp->last_rx_id = 0;
etcp->last_delivered_id = 0;
etcp->rx_ack_till = 0;
// Устанавливаем флаг ожидания первого пакета
etcp->got_initial_pkt = 0;
// Reset metrics
etcp->unacked_bytes = 0;
etcp->rtt_last = 0;
etcp->rtt_avg_10 = 0;
etcp->jitter = 0;
etcp->last_rtt_cb_time = 0;
etcp->bytes_sent_total = 0;
etcp->retransmissions_count = 0;
etcp->ack_packets_count = 0;
etcp->rx_dup_count = 0;
etcp->tx_dup_count = 0;
// Reset RTT history
memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history));
etcp->rtt_history_idx = 0;
etcp->last_rr_link = NULL;
// Return unconfirmed inflight data to normalizer input (до нормалайзера)
if (etcp->normalizer) {
etcp->normalizer->input->callback_suspended = 1;
etcp_return_inflight_to_normalizer(etcp);
}
// Clear queues (keep queue structures)
clear_queue(etcp->input_queue);
clear_queue(etcp->output_queue);
clear_queue(etcp->input_send_q);
clear_queue(etcp->input_wait_ack);
clear_queue(etcp->recv_q);
clear_queue(etcp->ack_q);
// clear_queue leaves callback_suspended=1, resume to prevent deadlock
queue_resume_callback(etcp->input_queue);
queue_resume_callback(etcp->input_send_q);
queue_resume_callback(etcp->input_wait_ack);
queue_resume_callback(etcp->output_queue);
queue_resume_callback(etcp->ack_q);
if (etcp->normalizer && etcp->normalizer->input) {
clear_queue(etcp->normalizer->input);
queue_resume_callback(etcp->normalizer->input);
}
// В etcp_conn_reset(), после очистки очередей добавьте:
struct ETCP_LINK* l = etcp->links;
while (l) {
l->inflight_bytes = 0;
l->inflight_packets = 0;
l->delivered_bytes = 0;
l->acked_bytes = 0;
l->acked_packets = 0;
l->last_ack_time_tb = get_time_tb();
l = l->next;
}
// Cancel active timers to prevent memory leaks
if (etcp->retrans_timer) {
uasync_cancel_timeout(etcp->instance->ua, etcp->retrans_timer);
etcp->retrans_timer = NULL;
}
if (etcp->ack_resp_timer) {
uasync_cancel_timeout(etcp->instance->ua, etcp->ack_resp_timer);
etcp->ack_resp_timer = NULL;
}
etcp->reset_count++;
// Reset normalizer (packer/unpacker state for reconnection)
if (etcp->normalizer) {
pn_reset(etcp->normalizer);
}
queue_resume_callback(etcp->input_queue);
queue_resume_callback(etcp->output_queue);
queue_resume_callback(etcp->input_send_q);
queue_resume_callback(etcp->input_wait_ack);
queue_resume_callback(etcp->recv_q);
queue_resume_callback(etcp->ack_q);
etcp->tx_state=ETCP_TX_STATE_DATA_WAIT;
int up=0;
struct ETCP_LINK* link = etcp->links;
while (link) {
if (link->link_status) up=1;
link = link->next;
}
if (up) {
if (etcp->links_up==0) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "reset conn: set link up");
etcp->links_up=1;
etcp_on_up(etcp);
} else {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "reset conn: set link down/up");
etcp_on_down(etcp, NULL);
etcp_on_up(etcp);
}
} else {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "links not up - skip on_down/on_up");
if (etcp->links_up) {
etcp->links_up=0;
etcp->setup_start_tb = get_time_tb();
etcp_on_down(etcp, NULL);
}
}
clear_queue(etcp->input_send_q);
clear_queue(etcp->input_wait_ack);
queue_resume_callback(etcp->input_send_q);
queue_resume_callback(etcp->input_wait_ack);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "end");
}
void etcp_links_reset(struct ETCP_CONN* etcp) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение
// Сбрасываем initialized во всех линках
struct ETCP_LINK* link = etcp->links;
while (link) {
link->initialized = 0;
link = link->next;
}
}
void etcp_conn_reinit(struct ETCP_CONN* etcp, const char* reason) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] REINIT: %s (init=%d links=%d reinit=%u tx=%d)",
etcp->log_name, reason, etcp->initialized, etcp->links_up, etcp->reinit_count, etcp->tx_state);
etcp->setup_start_tb = get_time_tb();
etcp->reinit_count++;
etcp->reinit_pending = 1;
etcp->reset_done = 0;
etcp->initialized = 0;// еще раз придёт conn_ready_callback
// Отменяем висящие NAT-ping'и для этого соединения
if (etcp->instance && etcp->instance->nat_det) {
nat_detection_cancel_for_conn(etcp->instance->nat_det, etcp);
}
// Сбрасываем ретрансмиты роутера ДО очистки ETCP-очередей (избегаем гонки с data_pool)
if (etcp->peer_node_id) etcp_router_pause_retrans_for_node(etcp->instance, etcp->peer_node_id);
// Сбрасываем NAT-статус у всех линков этого соединения
struct ETCP_LINK* l = etcp->links;
while (l) {
l->nat_type = NAT_TYPE_UNKNOWN;
l->nat_check_status = NAT_CHECK_NONE;
l = l->next;
}
// Вызываем etcp_conn_reset для сброса состояния
etcp_conn_reset(etcp);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "end");
}
// внутренняя функция. Вызывается один раз когда первый линк готов.
void etcp_conn_ready(struct ETCP_CONN* conn) {
if (!conn) return;
if (conn->initialized) return; // already ready
conn->initialized = 1;
conn->reset_done = 1;
if (conn->tx_state == 0) { conn->tx_state = ETCP_TX_STATE_DATA_WAIT; }
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Connection ready", conn->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_CONN_READY: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x my_pub=%02x%02x%02x%02x links=%p",
conn->log_name,
conn->crypto_ctx.session_key[0], conn->crypto_ctx.session_key[1],
conn->crypto_ctx.session_key[2], conn->crypto_ctx.session_key[3],
conn->crypto_ctx.peer_public_key[0], conn->crypto_ctx.peer_public_key[1],
conn->crypto_ctx.peer_public_key[2], conn->crypto_ctx.peer_public_key[3],
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[0] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[1] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[2] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[3] : 0,
(void*)conn->links);
etcp_conn_queue_set_ready(conn);
}
// Move from pending to indexed connections, fire ready callbacks
void etcp_conn_queue_set_ready(struct ETCP_CONN* conn) {
if (!conn || conn->state != 0) return;
// Reindex: remove old entry (key=0), add new entry with real peer_node_id
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] queue_set_ready: removing old entry=%p key=0 -> new key=0x%llx",
conn->log_name, (void*)conn->conn_queue_entry, (unsigned long long)conn->peer_node_id);
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
}
struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to alloc queue entry for ready", conn->log_name); return; }
struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
ce->peer_node_id = conn->peer_node_id;
ce->conn = conn;
conn->conn_queue_entry = qe;
conn->conn_queue = conn->instance->connections;
conn->state = 1;
queue_data_put_with_index(conn->instance->connections, qe);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] queue_set_ready: new entry=%p put in queue, count=%d",
conn->log_name, (void*)qe, queue_entry_count(conn->instance->connections));
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] moved to ready queue, state=%d peer_node_id=0x%llx",
conn->log_name, conn->state, (unsigned long long)conn->peer_node_id);
etcp_metrics_start_timer(conn);
etcp_cbk_fire(conn, ETCP_CBK_EVENT_INIT);
if (conn->reinit_pending) { conn->reinit_pending = 0; etcp_cbk_fire(conn, ETCP_CBK_EVENT_REINIT); }
if (conn->links_up) etcp_on_up(conn);
}
// Update log_name when peer_node_id becomes known
void etcp_update_log_name(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!etcp || !etcp->instance) return;
uint16_t local_id = etcp->instance->node_id;
uint16_t peer_id = etcp->peer_node_id;
const char* name = etcp->name ? etcp->name : EMPTY_NAME;
if (etcp->instance->topo_groups && peer_id) {
struct TOPO_NODE* ni = topo_node_registry_find(etcp->instance->topo_groups, etcp->peer_node_id);
if (ni && ni->node_name && ni->node_name[0]) name = ni->node_name;
}
snprintf(etcp->log_name, sizeof(etcp->log_name), "%04X->%04X [%s]", local_id, peer_id, name);
}
// ====================================================================== Отправка данных в etcp после нормализации
// Send data through ETCP connection
// Allocates memory from data_pool and places in input queue
// Returns: 0 on success, -1 on failure
int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!etcp || !data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] invalid parameters (etcp=%p, data=%p, len=%zu)", etcp->log_name, etcp, data, len);
return -1;
}
if (!etcp->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] input_queue is NULL for etcp=%p", etcp->log_name, etcp);
return -1;
}
// Check length against maximum packet size
if (len > ETCP_MAX_PAYLOAD_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] packet too large (len=%zu, max=%d)", etcp->log_name, len, ETCP_MAX_PAYLOAD_SIZE);
return -1;
}
// Allocate packet data from data_pool (following ETCP reception pattern)
uint8_t* packet_data = memory_pool_alloc(etcp->instance->data_pool);
if (!packet_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate packet data from data_pool", etcp->log_name);
return -1;
}
// Copy user data to packet buffer
memcpy(packet_data, data, len);
struct ETCP_FRAGMENT* pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate queue entry", etcp->log_name);
memory_pool_free(etcp->instance->data_pool, packet_data);
return -1;
}
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->ll.dgram = packet_data; // Point to data_pool allocation
pkt->ll.len = len; // размер packet_data
pkt->ll.dgram_pool = etcp->instance->data_pool;
pkt->ll.memlen = etcp->instance->data_pool->object_size;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] created PACKET %p with data %p (len=%zu)", etcp->log_name, pkt, packet_data, len);
// Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, (struct ll_entry*)pkt) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name);
queue_dgram_free(&pkt->ll);
queue_entry_free(&pkt->ll);
return -1;
}
return 0;
}
static void input_queue_try_push(struct ETCP_CONN* etcp) {// пробуем протолкнуть при отправке
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
// когда очередь отправки пуста - пробуем взять новый пакет на обработку
size_t wait_ack_bytes = queue_total_bytes(etcp->input_wait_ack);
if (wait_ack_bytes <= etcp->optimal_inflight) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] resume input queue: inflight_bytes=%d, input_len=%d", etcp->log_name, wait_ack_bytes, etcp->input_queue->total_bytes);
if (etcp->input_queue->count > 0) input_queue_cb(etcp->input_queue, etcp);
else queue_resume_callback(etcp->input_queue);// и только когда больше нечего отправлять - забираем новый пакет
}
}
static void input_queue_try_resume(struct ETCP_CONN* etcp) {// при ACK
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
// Сперва отправим всё из очереди отправки
size_t send_q_bytes = queue_total_bytes(etcp->input_send_q);
// queue_resume_callback(etcp->input_send_q);// вызвать лишний раз resume не страшно.
if (send_q_bytes>0) return;
// когда очередь отправки пуста - пробуем взять новый пакет на обработку
size_t wait_ack_bytes = queue_total_bytes(etcp->input_wait_ack);
if (wait_ack_bytes <= etcp->optimal_inflight) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] resume input queue: inflight_bytes=%d, input_len=%d", etcp->log_name, wait_ack_bytes, etcp->input_queue->total_bytes);
queue_resume_callback(etcp->input_queue);// и только когда больше нечего отправлять - забираем новый пакет
}
}
// Called from link-level etcp_link_update_inflight_lim() after inflight_lim_bytes changes.
// Recalculates connection-level optimal_inflight and resumes input_queue if room opened up.
void etcp_conn_on_inflight_lim_changed(struct ETCP_CONN* etcp) {
if (!etcp) return;
uint32_t sum = 0;
for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes;
etcp->optimal_inflight = sum;
// input_queue_try_resume(etcp);
}
void etcp_stats(struct ETCP_CONN* etcp) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] stats for conn=%p:", etcp->log_name, etcp);
// Queue statistics
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Queues:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_queue: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->input_queue), queue_total_bytes(etcp->input_queue));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_send_q: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->input_send_q), queue_total_bytes(etcp->input_send_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_wait_ack: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->input_wait_ack), queue_total_bytes(etcp->input_wait_ack));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_q: %zu pkts", etcp->log_name,
queue_entry_count(etcp->ack_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] recv_q: %zu pkts", etcp->log_name,
queue_entry_count(etcp->recv_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] output_queue: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->output_queue), queue_total_bytes(etcp->output_queue));
// RTT metrics
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RTT metrics:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_last: %u (0.1ms)", etcp->log_name, etcp->rtt_last);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_avg_10: %u (0.1ms)", etcp->log_name, etcp->rtt_avg_10);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] jitter: %u (0.1ms)", etcp->log_name, etcp->jitter);
// Counters
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Counters:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] bytes_sent_total: %u", etcp->log_name, etcp->bytes_sent_total);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] retransmissions_count: %u", etcp->log_name, etcp->retransmissions_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_packets_count: %u", etcp->log_name, etcp->ack_packets_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] unacked_bytes: %u", etcp->log_name, etcp->unacked_bytes);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rx_dup_count: %u", etcp->log_name, etcp->rx_dup_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] tx_dup_count: %u", etcp->log_name, etcp->tx_dup_count);
// IDs
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] IDs:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] next_tx_id: %u", etcp->log_name, etcp->next_tx_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] last_rx_id: %u", etcp->log_name, etcp->last_rx_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] last_delivered_id:%u", etcp->log_name, etcp->last_delivered_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rx_ack_till: %u", etcp->log_name, etcp->rx_ack_till);
}
// Input callback for input_queue (добавление новых кодограмм в стек)
// input_queue -> input_send_q
static void input_queue_cb(struct ll_queue* q, void* arg) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
struct ETCP_FRAGMENT* in_pkt = (struct ETCP_FRAGMENT*)queue_data_get(q);
if (!in_pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot get element (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp);
queue_resume_callback(q);
return;
}
// Create INFLIGHT_PACKET
struct INFLIGHT_PACKET* p = (struct INFLIGHT_PACKET*)queue_entry_new_from_pool(etcp->inflight_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp);
queue_dgram_free(&in_pkt->ll);
queue_entry_free((struct ll_entry*)in_pkt);
queue_resume_callback(q);
return;
}
// Setup inflight packet (based on protocol.txt)
// memset(p, 0, sizeof(*p));
p->seq = etcp->next_tx_id++; // Assign seq
p->state = INFLIGHT_STATE_WAIT_SEND;
p->last_timestamp = 0;
p->ll.dgram = in_pkt->ll.dgram;
p->ll.dgram_pool = in_pkt->ll.dgram_pool;
p->ll.len = in_pkt->ll.len;
// memory_pool_free(etcp->io_pool, in_pkt);// перемещаем из io_pool в inflight_pool
queue_entry_free(&in_pkt->ll);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] TX input -> inflight (seq=%u, len=%u); Qlen: in=%d snd=%d wait=%d", etcp->log_name, p->seq, p->ll.len, etcp->input_queue->count, etcp->input_send_q->count, etcp->input_wait_ack->count);
int len=p->ll.len;// сохраним len
// Add to send queue
if (queue_data_put_with_index(etcp->input_send_q, &p->ll) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet seq=%u to input_send_q", etcp->log_name, p->seq);
queue_dgram_free(&p->ll);
queue_entry_free(&p->ll);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] EXIT (queue put failed)", etcp->log_name);
return;
}
etcp->unacked_bytes += len;
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "successfully moved from input_queue to input_send_q");
// etcp_conn_process_send_queue(etcp);// сразу обработаем этот пакет
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] nextloop, input_queue size=%d ", etcp->log_name, q->count);
}
static void ack_timeout_cb(void* arg);
static void ack_timeout_check(struct ETCP_CONN* etcp) {
uint64_t now = get_time_tb();
int32_t timeout = (etcp->rtt_avg_10 * RETRANS_K1 + etcp->jitter * RETRANS_K2) / 16;
if (timeout<50) timeout=50;
if (timeout>10000) timeout=10000;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "starting check, now=%llu, timeout=%llu, rtt_avg_10=%u, jitter=%u",
(unsigned long long)now, (unsigned long long)timeout, etcp->rtt_avg_10, etcp->jitter);
struct ll_entry* current;
current = etcp->input_wait_ack->head;
while (current) {
struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current;
int64_t elapsed = now - pkt->last_timestamp;
if (elapsed > timeout) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] timeout for seq=%u, elapsed=%lld, now=%llu, timeout=%llu, send_count=%u. Moving: wait_ack -> send_q",
etcp->log_name, pkt->seq, (unsigned long long)elapsed, (unsigned long long)now, (unsigned long long)timeout, pkt->send_count);
// Remove from wait_ack
pkt=(struct INFLIGHT_PACKET*)queue_data_get(etcp->input_wait_ack);
if (!pkt) break;
// Increment counters
pkt->retrans_req_count++; // Optional, if used for retrans request logic
pkt->last_timestamp = now;
// Change state and add to send_q for retransmission
pkt->state = INFLIGHT_STATE_WAIT_SEND;
queue_data_put_with_index(etcp->input_send_q, (struct ll_entry*)pkt);
// Update stats
etcp->retransmissions_count++;
etcp_metrics_add_loss(etcp, 1);
}
else {// не надо до конца сканировать - они уже сортированы по таймстемпу т.к. очередь fifo, а timestamp = время добавления в очередь = время отправки
// shedule timer
int64_t next_timeout=timeout - elapsed;
if (next_timeout<0) next_timeout=0;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] retransmission timer set for %llu units", etcp->log_name, next_timeout);
etcp->retrans_timer = uasync_set_timeout(etcp->instance->ua, next_timeout+10, etcp, ack_timeout_cb, "etcp_retrans");
return;
}
current = etcp->input_wait_ack->head;
}
queue_resume_callback(etcp->input_wait_ack);
}
static void ack_timeout_cb(void* arg) {// сработал таймер переотправки
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
etcp->retrans_timer=NULL;
ack_timeout_check(etcp);// он установит новый таймер или выгребет всё и установит ожидание на очередь
}
static void wait_ack_cb(struct ll_queue* q, void* arg) {// добавили пакет в ожидание ACK
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
if (!etcp->retrans_timer) ack_timeout_check(etcp);// если таймер ретрансмиссий взведен - дожидаемся таймера.
}
static void input_send_q_cb(struct ll_queue* q, void* arg) {// etcp->input_send_q data ready
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_CONN* etcp=(struct ETCP_CONN*)arg;
etcp_conn_process_send_queue(etcp);
if (etcp->tx_state==ETCP_TX_STATE_DATA_WAIT) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "resume input_send_q");
queue_resume_callback(etcp->input_send_q);
} else DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "no resume - link busy");
}
/*
static void send_ack_req_cb(struct ll_queue* q, void* arg) {// etcp->ack_q data ready
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_CONN* etcp=(struct ETCP_CONN*)arg;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] processing", etcp->log_name);
etcp_conn_process_send_queue(etcp);
if (etcp->tx_state==ETCP_TX_STATE_DATA_WAIT) queue_resume_callback(etcp->ack_q);
}
*/
void etcp_on_link_down(struct ETCP_CONN* etcp, struct ETCP_LINK* down_link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
int up=0;
struct ETCP_LINK* link = etcp->links;
while (link) {
if (link->link_status) up=1;
link = link->next;
}
int was_up = etcp->links_up;
etcp->links_up = up;
if (up == 0 && was_up != 0) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All links fall down");
etcp_on_down(etcp, down_link);
}
}
static void ack_response_timer_cb(void* arg) {// проверяем неотправленные ack response и отправляем если надо.
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
etcp->ack_resp_timer=NULL;
if (etcp->ack_q->count==0) return;// нечего отправлять
etcp_conn_process_send_queue(etcp);// проталкиваем (она же должна отправлять только ack если больше ничего нет)
// если ack все еще заняты - обновляем таймаут
if (etcp->ack_q->count) etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb, "etcp_ack_resp");
// else etcp->ack_resp_timer=NULL;
}
static void etcp_link_ready_callback(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!etcp) return;
if (etcp->tx_state==0) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "ETCP not ready, skip link state change");
return;// not initialized
}
if (etcp->links_up==0) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] etcp_link_ready_callback: links_up 0→1, calling etcp_on_up (initialized=%d tx_state=%d)", etcp->log_name, etcp->initialized, etcp->tx_state);
etcp->links_up=1;
etcp_on_up(etcp);
} else DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[%s] etcp_link_ready_callback: links_up=%d already up", etcp->log_name, etcp->links_up);
if (etcp->tx_state!=ETCP_TX_STATE_LINK_WAIT) return;
etcp->tx_state = ETCP_TX_STATE_DATA_WAIT;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "resume input_send_q+ack_q; link_wait->data_wait");
queue_resume_callback(etcp->input_send_q);
// queue_resume_callback(etcp->ack_q);
if (etcp->ack_q->count && etcp->ack_resp_timer == NULL) {
ack_response_timer_cb(etcp);
}
}
// Process packets in send queue and transmit them
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {// вызываем когда есть элемент в send_q или надо отправить ack
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_DGRAM* dgram;
if (etcp->tx_state!=ETCP_TX_STATE_DATA_WAIT) {
char l_status[256]={0};
struct ETCP_LINK* link = etcp->links;
while (link) {
snprintf (l_status+strlen(l_status), 256-strlen(l_status), "L%d%d%s ", link->recv_keepalive, link->remote_keepalive, (link->shaper_timer)?"wait":"rdy");
link = link->next;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX state: %d (link not ready, skip send) %s", etcp->log_name, etcp->tx_state, l_status);
return;
}
dgram = etcp_request_pkt(etcp);
while(dgram) {
etcp_loadbalancer_send(dgram);
dgram = etcp_request_pkt(etcp);
}
if (etcp->tx_state==ETCP_TX_STATE_DATA_WAIT) {
queue_resume_callback(etcp->input_send_q);
}
}
// Подготовить и отправить кодограмму
// вызывается линком когда освобождается или очередью если появляются данные на передачу
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
struct ETCP_LINK* link = NULL;
// если есть активный burst — используем этот линк напрямую
for (struct ETCP_LINK* l = etcp->links; l; l = l->next) {
if (l->burst_active) { link = l; break; }
}
if (!link) link = etcp_loadbalancer_select_link(etcp);
if (!link) {
etcp->tx_state=ETCP_TX_STATE_LINK_WAIT;
etcp->cnt_link_wait++;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no link available", etcp->log_name);
return NULL;// если линков нет - ждём появления свободного
}
etcp->tx_state=ETCP_TX_STATE_DATA_WAIT;
size_t send_q_size = queue_entry_count(etcp->input_send_q);
if (send_q_size == 0) {// сгребаем из input_queue
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_send_q empty, check if avail input_queue -> inflight");
input_queue_try_push(etcp);
}
// First, check if there's a packet in input_send_q (retrans or new)
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "getting packet from input_send_q");
struct INFLIGHT_PACKET* inf_pkt = (struct INFLIGHT_PACKET*)queue_data_get(etcp->input_send_q);
if (inf_pkt) {
uint64_t now=get_time_tb();
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] send_q->wait_ack seq=%d TS=%llu", etcp->log_name, inf_pkt->seq, now);
// === NEW: per-link inflight + send_hist logic ===
// Compute 1-based link number in the linked list (as requested)
uint8_t link_num = 0;
struct ETCP_LINK* tmp = etcp->links;
while (tmp) {
link_num++;
if (tmp == link) break;
tmp = tmp->next;
}
if (link_num == 0) link_num = 255; // safety (should never happen)
// Fill history (send_count is the index before this attempt)
inf_pkt->send_hist[inf_pkt->send_count % 8] = link_num;
// Always subtract from previous last_link (if any) – this handles retransmission
if (inf_pkt->last_link) {
inf_pkt->last_link->total_retransmissions++;
inf_pkt->last_link->bbr_loss_since_ack += inf_pkt->ll.len;
bbr_note_loss(inf_pkt->last_link->bbr);
inf_pkt->last_link->inflight_bytes -= inf_pkt->ll.len;
inf_pkt->last_link->inflight_packets--;
}
// Always add to the CURRENT link (first send or retransmission)
link->inflight_bytes += inf_pkt->ll.len;
link->inflight_packets++;
// BBR: сохраняем снэпшот на момент отправки
inf_pkt->delivered_at_send = link->delivered_bytes;
inf_pkt->inflight_at_send = link->inflight_bytes;
inf_pkt->is_app_limited = (etcp->input_queue->count == 0) ? 1 : 0;
// Update last_link for future ACK/retrans
inf_pkt->last_link = link;
inf_pkt->last_timestamp=now;
inf_pkt->send_count++;
inf_pkt->state=INFLIGHT_STATE_WAIT_ACK;
queue_data_put_with_index(etcp->input_wait_ack, &inf_pkt->ll);// move dgram to wait_ack queue
etcp_metrics_add_sent(etcp, inf_pkt->ll.len);
}
size_t ack_q_size = queue_entry_count(etcp->ack_q);
if (!inf_pkt && ack_q_size == 0 && !link->burst_active) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no data/ack to send", etcp->log_name);
return NULL;
}
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ETCP_DGRAM", etcp->log_name);
return NULL;
}
dgram->link = link;
dgram->noencrypt_len=0;
dgram->timestamp = get_current_timestamp();
// формат ack: [01] [elements count] [4 байта last_delivered_id] [2 байта rx_dup_count] и <[4 байта seq][2 байта recv_ts][2 байта txrx delay ts]> x count
dgram->data[0]=1;// ack
int ptr=2;
dgram->data[ptr++]=etcp->last_delivered_id;
dgram->data[ptr++]=etcp->last_delivered_id>>8;
dgram->data[ptr++]=etcp->last_delivered_id>>16;
dgram->data[ptr++]=etcp->last_delivered_id>>24;
dgram->data[ptr++]=etcp->rx_dup_count;
dgram->data[ptr++]=etcp->rx_dup_count>>8;
int data_len=0;
if (inf_pkt) data_len=inf_pkt->ll.len;
int remain_len = link->mtu
- 28 /*udp headers*/
- SC_NONCE_SIZE - SC_TAG_SIZE - SC_CRC32_SIZE
- (ETCP_ENCRYPTED_HDR_SIZE + ETCP_ACK_BASE_SIZE)
- (1 + sizeof(inf_pkt->seq)) /*payload hdr*/
- data_len;
// DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "remain_len= %d pl=%d", remain_len, data_len);
// добавим опциональные заголовки
struct ACK_PACKET* ack_pkt;
while (remain_len>=8) {
ack_pkt = (struct ACK_PACKET*)queue_data_get(etcp->ack_q);
if (!ack_pkt) break;
remain_len-=8;
// seq 4 байта
dgram->data[ptr++]=ack_pkt->seq;
dgram->data[ptr++]=ack_pkt->seq>>8;
dgram->data[ptr++]=ack_pkt->seq>>16;
dgram->data[ptr++]=ack_pkt->seq>>24;
// ts приема 2 байта
dgram->data[ptr++]=ack_pkt->recv_timestamp;
dgram->data[ptr++]=ack_pkt->recv_timestamp>>8;
// время задержки 2 байта между recv и ack
uint16_t dly=get_current_timestamp()-ack_pkt->recv_timestamp;
dgram->data[ptr++]=dly;
dgram->data[ptr++]=dly>>8;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX Send ACK seq=%d need=%d", etcp->log_name, ack_pkt->seq, etcp->last_delivered_id+1);
queue_entry_free((struct ll_entry*)ack_pkt);
if (inf_pkt && inf_pkt->ll.len+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки)
if (ptr>500) break;
}
dgram->data[1]=(ptr - 8) / 8;
int max_enc = link->mtu - 28 - SC_NONCE_SIZE - SC_TAG_SIZE - SC_CRC32_SIZE - ETCP_ENCRYPTED_HDR_SIZE;
if (max_enc > PACKET_DATA_SIZE) max_enc = PACKET_DATA_SIZE;
if (max_enc < 0) max_enc = 0;
uint8_t burst_flags = 0; int pkt_sz_pos = 0;
// === MEAS_RESP piggyback (если есть готовый ответ и не активен burst) ===
if (link->burst_resp_pending && !link->burst_active && max_enc - ptr >= MEAS_RESP_SECTION_SIZE) {
dgram->data[ptr++] = ETCP_SECTION_MEAS_RESP;
dgram->data[ptr++] = link->burst_recv_id & 0xFF; dgram->data[ptr++] = (link->burst_recv_id >> 8) & 0xFF;
dgram->data[ptr++] = link->burst_resp_valid;
uint32_t v = link->burst_resp_gap_avg;
dgram->data[ptr++] = v & 0xFF; dgram->data[ptr++] = (v >> 8) & 0xFF; dgram->data[ptr++] = (v >> 16) & 0xFF; dgram->data[ptr++] = (v >> 24) & 0xFF;
v = link->burst_resp_gap_min;
dgram->data[ptr++] = v & 0xFF; dgram->data[ptr++] = (v >> 8) & 0xFF; dgram->data[ptr++] = (v >> 16) & 0xFF; dgram->data[ptr++] = (v >> 24) & 0xFF;
dgram->data[ptr++] = link->burst_resp_pkt_count;
link->burst_resp_pending = 0;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX burst resp id=%u valid=%u gap_min=%u", etcp->log_name, link->burst_recv_id, link->burst_resp_valid, link->burst_resp_gap_min);
}
// === MEAS_TS section (burst) ===
if (link->burst_active && max_enc - ptr >= MEAS_TS_SECTION_SIZE) {
burst_flags = 0;
if (!inf_pkt) burst_flags |= MEAS_FLAG_IS_FILLER;
if (link->burst_seq == 0) burst_flags |= MEAS_FLAG_IS_FIRST;
if (link->burst_seq + 1 >= link->burst_count) burst_flags |= MEAS_FLAG_IS_LAST;
dgram->data[ptr++] = ETCP_SECTION_MEAS_TS;
dgram->data[ptr++] = link->burst_id & 0xFF; dgram->data[ptr++] = (link->burst_id >> 8) & 0xFF;
dgram->data[ptr++] = burst_flags;
dgram->data[ptr++] = link->burst_seq;
uint16_t ts_us = (uint16_t)get_time_us();
dgram->data[ptr++] = ts_us & 0xFF; dgram->data[ptr++] = (ts_us >> 8) & 0xFF;
pkt_sz_pos = ptr;
dgram->data[ptr++] = 0; dgram->data[ptr++] = 0;
}
if (link->last_recv_updated && remain_len>=5) {// если есть данные - добавим channel_timestamp
uint64_t now=get_time_tb();
uint64_t dt=now - link->last_recv_local_time;
link->last_recv_updated=0;
if (dt<1000000) {
dgram->data[ptr++]=ETCP_SECTION_TIMESTAMP;
uint16_t t=link->last_recv_timestamp + dt;
dgram->data[ptr++]=t;
dgram->data[ptr++]=t>>8;
t=link->last_recv_local_time - link->last_recv_timestamp;
dgram->data[ptr++]=t;
dgram->data[ptr++]=t>>8;
remain_len-=5;
}
}
// DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "remain_len(2)= %d", remain_len);
if (inf_pkt) {
// === FILLER section before PAYLOAD (burst padding) ===
if (link->burst_active) {
int payload_overhead = 1 + 4 + data_len; // type + seq + data
int buf_avail = PACKET_DATA_SIZE - ptr - payload_overhead;
int filler_data = max_enc - ptr - FILLER_HDR_SIZE - payload_overhead;
if (filler_data > buf_avail) filler_data = buf_avail;
if (filler_data > 0) {
dgram->data[ptr++] = ETCP_SECTION_FILLER;
dgram->data[ptr++] = filler_data & 0xFF;
dgram->data[ptr++] = (filler_data >> 8) & 0xFF;
memset(&dgram->data[ptr], 0, filler_data); ptr += filler_data;
}
}
// фрейм data (0) обязательно в конец
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX DATA: seq=%u len=%u retry=%d; Qlen: in=%d snd=%d wait=%d", etcp->log_name, inf_pkt->seq, inf_pkt->ll.len, inf_pkt->send_count, etcp->input_queue->count, etcp->input_send_q->count, etcp->input_wait_ack->count);
dgram->data[ptr++]=0;// payload
dgram->data[ptr++]=inf_pkt->seq;
dgram->data[ptr++]=inf_pkt->seq>>8;
dgram->data[ptr++]=inf_pkt->seq>>16;
dgram->data[ptr++]=inf_pkt->seq>>24;
if ((int)(ptr + inf_pkt->ll.len) <= PACKET_DATA_SIZE) {
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len;
} else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload overflow: ptr=%d len=%u max=%d", etcp->log_name, ptr, inf_pkt->ll.len, PACKET_DATA_SIZE);
}
else {
if (link->burst_active) {
int buf_avail = PACKET_DATA_SIZE - ptr;
int filler_data = max_enc - ptr - FILLER_HDR_SIZE;
if (filler_data > buf_avail - FILLER_HDR_SIZE) filler_data = buf_avail - FILLER_HDR_SIZE;
if (filler_data > 0) {
dgram->data[ptr++] = ETCP_SECTION_FILLER;
dgram->data[ptr++] = filler_data & 0xFF;
dgram->data[ptr++] = (filler_data >> 8) & 0xFF;
memset(&dgram->data[ptr], 0, filler_data); ptr += filler_data;
}
}
int chk=queue_check_consistency(etcp->ack_q);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] only ACK (size=%d) packet with %d bytes total (chk=%d) rem=%d", etcp->log_name, ack_q_size, ptr, chk, remain_len);
}
if (ptr>=PACKET_DATA_SIZE-50) DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] SIZE ERROR!!! %d", etcp->log_name, ptr);
dgram->data_len=ptr;
// === Fill pkt_sz in MEAS_TS and advance burst state ===
if (link->burst_active && pkt_sz_pos > 0) {
uint16_t pkt_sz = (uint16_t)dgram->data_len;
dgram->data[pkt_sz_pos] = pkt_sz & 0xFF;
dgram->data[pkt_sz_pos + 1] = (pkt_sz >> 8) & 0xFF;
link->burst_pkt_size = pkt_sz;
link->burst_seq++;
if (link->burst_seq >= link->burst_count) etcp_link_burst_finish(link);
}
etcp_dump_pkt_sections(dgram, link, 1);
return dgram;
}
// ====================================================================== Прием данных
void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
// пробуем собрать выходную очередь из фрагментов
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] etcp=%p, last_delivered_id=%u, recv_q_count=%d",
// etcp->log_name, etcp, etcp->last_delivered_id, queue_entry_count(etcp->recv_q));
uint32_t next_expected_id = etcp->last_delivered_id + 1;
int delivered_count = 0;
uint32_t delivered_bytes = 0;
// Look for contiguous packets starting from next_expected_id
while (1) {
struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_find_data_by_index(etcp->recv_q, &next_expected_id);
if (!rx_pkt) {
// No more contiguous packets found
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no packet found for id=%u, stopping", etcp->log_name, next_expected_id);
break;
}
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] assembling packet id=%u (len=%u)", etcp->log_name,
// rx_pkt->seq, rx_pkt->ll.len);
// Simply move ETCP_FRAGMENT from recv_q to output_queue - no data copying needed
// Remove from recv_q first
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "Remove from assembly queue");
queue_remove_data(etcp->recv_q, (struct ll_entry*)rx_pkt);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "move: ETCP -> PN");
// Add to output_queue using the same ETCP_FRAGMENT structure
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] moving packet id=%u to output_queue (qlen=%d)", etcp->log_name,
next_expected_id, etcp->output_queue->count);
uint16_t pkt_len = rx_pkt->ll.len;
if (queue_data_put(etcp->output_queue, (struct ll_entry*)rx_pkt) == 0) {
delivered_bytes += pkt_len;
delivered_count++;
} else {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet id=%u to output_queue", etcp->log_name,
next_expected_id);
// Put it back in recv_q if we can't add to output_queue
queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt);
break;
}
// Update state for next iteration
etcp->last_delivered_id = next_expected_id;
next_expected_id++;
}
if (delivered_count>0) DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] delivered %u contiguous packets (%u bytes), last_delivered_id=%u, asm_queue=%d, output_queue=%d",
etcp->log_name, delivered_count, delivered_bytes, etcp->last_delivered_id, queue_entry_count(etcp->recv_q), queue_entry_count(etcp->output_queue));
}
// Process ACK receipt - remove acknowledged packet from inflight queues
void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t dts) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!etcp) return;
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] processing ACK for seq=%u, ts=%u, dts=%u", etcp->log_name, seq, ts, dts);
// Find the acknowledged packet in the wait_ack queue
struct INFLIGHT_PACKET* acked_pkt;
acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_index(etcp->input_wait_ack, &seq);
if (acked_pkt) {
etcp->cnt_ack_hit_inf++;
queue_remove_data(etcp->input_wait_ack, (struct ll_entry*)acked_pkt);
}
else {
acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_index(etcp->input_send_q, &seq);
if (acked_pkt) {
etcp->cnt_ack_hit_sndq++;
queue_remove_data(etcp->input_send_q, (struct ll_entry*)acked_pkt);
}
else etcp->cnt_ack_miss++;
}
if (!acked_pkt) {
// Packet might be already acknowledged or not found
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "packet seq=%u not found in wait_ack queue", seq);
return;
}
// === subtract inflight from the LAST link the packet was sent on ===
if (acked_pkt->last_link) {
struct ETCP_LINK* link = acked_pkt->last_link;
link->inflight_bytes -= acked_pkt->ll.len;
link->inflight_packets--;
// BBR: собираем rate_sample и вызываем bbr_main
uint64_t now_tb = get_time_tb();
uint32_t interval_us = link->last_ack_time_tb ? (uint32_t)((now_tb - link->last_ack_time_tb) * 100) : 0;
struct bbr_rate_sample rs = {
.delivered = acked_pkt->ll.len,
.interval_us = interval_us,
.rtt_us = (uint32_t)link->rtt_last * 100,
.acked_sacked = acked_pkt->ll.len,
.prior_delivered = (uint32_t)acked_pkt->delivered_at_send,
.tx_in_flight = acked_pkt->inflight_at_send,
.lost = (int)link->bbr_loss_since_ack,
.is_app_limited = acked_pkt->is_app_limited,
};
link->delivered_bytes += rs.delivered;
link->last_ack_time_tb = now_tb;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[L%u] BBR in: seq=%u del=%u iv=%uus rtt=%uus pri=%u tx_infl=%u cur=%u/%u lost=%d app=%d",
link->local_link_id, seq, rs.delivered, rs.interval_us, rs.rtt_us,
rs.prior_delivered, rs.tx_in_flight, link->inflight_bytes, link->inflight_lim_bytes,
rs.lost, rs.is_app_limited);
uint32_t old_cwnd = link->inflight_lim_bytes;
uint32_t old_pacing = link->bbr_pacing_rate;
uint32_t new_cwnd = old_cwnd;
uint32_t new_pacing = old_pacing;
bbr_main(link->bbr, &rs, &new_cwnd, &new_pacing, link->mtu, link->inflight_bytes,
(link->inflight_bytes >= link->inflight_lim_bytes));
etcp_link_update_inflight_lim(link, new_cwnd);
link->bbr_pacing_rate = new_pacing;
link->bandwidth = (uint32_t)((uint64_t)new_pacing * 8 / 1000);
link->bbr_loss_since_ack = 0;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[L%u] BBR out: cwnd %u→%u pace %u→%u bw=%uK mode=%d cyc=%d infl=%u/%u",
link->local_link_id, old_cwnd, new_cwnd, old_pacing, new_pacing,
link->bandwidth, link->bbr->mode, link->bbr->cycle_idx, link->inflight_bytes, new_cwnd);
link->acked_bytes += acked_pkt->ll.len;
link->acked_packets++;
}
// Update connection statistics
if (etcp->unacked_bytes >= acked_pkt->ll.len) etcp->unacked_bytes -= acked_pkt->ll.len;
else {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] unacked_bytes underflow prevented: %u < %u", etcp->log_name, etcp->unacked_bytes, acked_pkt->ll.len);
etcp->unacked_bytes = 0;
}
etcp->bytes_sent_total += acked_pkt->ll.len;
etcp->ack_packets_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX removed packet seq=%u from wait_ack, unacked_bytes now %u total acked=%u", etcp->log_name, seq, etcp->unacked_bytes, etcp->ack_packets_count);
if (acked_pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, acked_pkt->ll.dgram);
}
memory_pool_free(etcp->inflight_pool, acked_pkt);
// Try to resume sending more packets if window space opened up
input_queue_try_resume(etcp);
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] completed for seq=%u", etcp->log_name, seq);
}
// Process incoming decrypted packet
void etcp_conn_input(struct ETCP_DGRAM* pkt) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!pkt) return;
if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt);
volatile int _halt = 1; while (_halt) {}
}
if (!pkt->data_len) {
memory_pool_free(pkt->link->etcp->instance->pkt_pool, pkt);
return;
}
etcp_dump_pkt_sections(pkt, pkt->link, 0);
pkt->link->last_recv_local_time = get_time_tb();
struct ETCP_CONN* etcp = pkt->link->etcp;
uint8_t* data = pkt->data;
uint16_t len = pkt->data_len;
uint16_t ts = pkt->timestamp; // Received timestamp
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX pkt dlen=%d", etcp->log_name, len);
while (len >= 1) {
uint8_t type = data[0];
// Process sections as per protocol.txt
switch (type) {
case ETCP_SECTION_ACK: {
if (len < 2) { len = 0; break; }
int elm_cnt=data[1];
uint32_t till=data[2] | (data[3]<<8) | (data[4]<<16) | (data[5]<<24);
uint16_t rx_dup_count_16 = data[6] | (data[7]<<8);
uint16_t old_tx_dup_16 = etcp->tx_dup_count & 0xFFFF;
int16_t diff = (int16_t)(rx_dup_count_16 - old_tx_dup_16);
etcp->tx_dup_count += diff;
int ack_section_len = 8 + elm_cnt * 8;
if (ack_section_len > len) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] ACK section too large: elm_cnt=%d len=%d", etcp->log_name, elm_cnt, len);
len = 0;
break;
}
data+=ack_section_len;
len-=ack_section_len;
for (int i=0; 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;
etcp_metrics_add_rtt(etcp, new_rtt);
int recv_dt_tx1=data[3] | (data[4]<<8);// localtime удаленной стороны момента принятия пакета - timestamp этого пакета (на стороне отправителя, т.е. у нас)
int recv_dt_rx=cur_ts - pkt->link->rtt_last/2 - 1000 - ts;
int recv_dt_tx=recv_dt_tx1 - pkt->link->rtt_last/2 - 1000;
if (pkt->link->recv_dt_avg_rx==0) pkt->link->recv_dt_avg_rx=recv_dt_rx*65536;
if (pkt->link->recv_dt_avg_tx==0) pkt->link->recv_dt_avg_tx=recv_dt_tx*65536;
pkt->link->recv_dt_avg_rx +=((int32_t)(recv_dt_rx*65536 - pkt->link->recv_dt_avg_rx))/32;
pkt->link->recv_dt_avg_tx +=((int32_t)(recv_dt_tx*65536 - pkt->link->recv_dt_avg_tx))/32;
pkt->link->rt_last = cur_ts - ts - pkt->link->recv_dt_avg_rx/65536;
pkt->link->tt_last = recv_dt_tx1 - pkt->link->recv_dt_avg_tx/65536;
//tts_correction += ((NOW - RTT/2 - TTS) - tts_correction)/16 (инициализируем сразу по 1 пакету)
data+=5; len-=5;
int prev_rtt = pkt->link->rtt_last;
int d_rtt = new_rtt - prev_rtt;
if (d_rtt<0) d_rtt=-d_rtt;
pkt->link->jitter +=((int32_t)(d_rtt*65536 - pkt->link->jitter))/32;
struct ETCP_LINK* c=etcp->links;
int rtt_sum=0;
int rtt_max=0;
int tt_sum=0;
int j_sum=0;
int cnt=0;
while (c) {
rtt_sum += c->rtt_last;
if (rtt_max < c->rtt_last) rtt_max = c->rtt_last;
tt_sum += c->tt_last;
j_sum += c->jitter/32768;
cnt++;
c=c->next;
}
etcp->rtt_last=rtt_sum/cnt;
etcp->rtt_avg_10=rtt_max;
etcp->tt_last=tt_sum/cnt;
etcp->jitter=j_sum/cnt;
if (etcp->rtt_last) {
uint64_t now_tb = get_time_tb();
if (now_tb - etcp->last_rtt_cb_time > RTT_CB_PERIOD_TB) {
topo_node_ping_update_rtt(etcp->instance->topo_groups, etcp->peer_node_id, etcp->rtt_last);
etcp->last_rtt_cb_time = now_tb;
}
}
break;
}
case ETCP_SECTION_PAYLOAD: {
if (len>=5) {
// формируем ACK
uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24);
if (etcp->got_initial_pkt == 0) {
if (seq==1) { etcp->got_initial_pkt = 1; DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Initial packet seq=1 received", etcp->log_name); }
else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Out-of-order pkt before init seq=%d, waiting seq=1",
etcp->log_name, seq);
len=0;
break;
}
}
else {
int32_t d=seq - etcp->last_delivered_id;
if (d>MAX_INFLIGHT_SIZE || d<-MAX_INFLIGHT_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] Received packet out of inflight bounds: seq=%d last delivered=%d", etcp->log_name, seq, etcp->last_delivered_id);
len=0;
break;
}
}
if (queue_find_data_by_index(etcp->ack_q, &seq) == NULL) {
struct ACK_PACKET* p = (struct ACK_PACKET*)queue_entry_new_from_pool(etcp->instance->ack_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ACK_PACKET", etcp->log_name);
len = 0;
break;
}
p->seq=seq;
p->pkt_timestamp=pkt->timestamp;
p->recv_timestamp=get_current_timestamp();
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX add to ack_q seq=%d", etcp->log_name, seq);
queue_data_put_with_index(etcp->ack_q, (struct ll_entry*)p);
if (etcp->ack_resp_timer == NULL) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] set ack_timer for delayed ACK send", etcp->log_name);
etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb, "etcp_ack_resp");
}
}
if (((int32_t)(etcp->last_delivered_id-seq)<0) && (queue_find_data_by_index(etcp->recv_q, &seq)==NULL)) {// проверяем есть ли пакет с этим seq
uint32_t pkt_len=len-5;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] adding packet seq=%u to recv_q (last_delivered_id=%u)", etcp->log_name, seq, etcp->last_delivered_id);
// отправляем пакет в очередь на сборку
uint8_t* payload_data = memory_pool_alloc(etcp->instance->data_pool);
if (!payload_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate payload_data from data_pool", etcp->log_name);
len=0;
break;
}
struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool);
if (!rx_pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate rx_pkt from io_pool", etcp->log_name);
memory_pool_free(etcp->instance->data_pool, payload_data);
len=0;
break;
}
rx_pkt->seq = seq;
rx_pkt->timestamp = pkt->timestamp;
rx_pkt->ll.dgram = payload_data;
rx_pkt->ll.len = pkt_len;
rx_pkt->ll.dgram_pool = etcp->instance->data_pool;
rx_pkt->ll.memlen = etcp->instance->data_pool->object_size;
if (pkt_len > (int)etcp->instance->data_pool->object_size) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload too large for data_pool: %u > %zu", etcp->log_name, pkt_len, etcp->instance->data_pool->object_size);
memory_pool_free(etcp->instance->data_pool, payload_data);
queue_entry_free(&rx_pkt->ll);
len = 0;
break;
}
if (data + len > pkt->data + pkt->data_len) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] payload section bounds overflow: data+len=%p > end=%p", etcp->log_name, (void*)(data + len), (void*)(pkt->data + pkt->data_len));
memory_pool_free(etcp->instance->data_pool, payload_data);
queue_entry_free(&rx_pkt->ll);
len = 0;
break;
}
// Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX seq=%u need=%u asm_len=%d", etcp->log_name, seq, etcp->last_delivered_id+1, etcp->recv_q->count);
queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt);
etcp_metrics_add_rcvd(etcp, pkt_len);
if ((int32_t)(seq - etcp->last_delivered_id) == 1) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов
} else {
etcp->rx_dup_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RX dup: seq=%u need=%u asm_len=%d", etcp->log_name, seq, etcp->last_delivered_id+1, etcp->recv_q->count);
}
} else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "payload len %d < 5", len);
len=0;
break;
}
case ETCP_SECTION_MEAS_TS: {
if (len < MEAS_TS_SECTION_SIZE) { len = 0; break; }
uint16_t b_id = data[1] | (data[2] << 8);
uint8_t flags = data[3];
uint8_t seq = data[4];
uint16_t pkt_sz = data[7] | (data[8] << 8);
struct ETCP_LINK* link = pkt->link;
if (b_id != link->burst_recv_id) {
link->burst_recv_id = b_id;
link->burst_recv_count = BURST_PACKET_COUNT;
link->burst_recv_next_seq = 0;
link->burst_recv_valid = 1;
link->burst_recv_received = 0;
}
if (link->burst_recv_valid) {
if (seq != link->burst_recv_next_seq) {
if (seq > link->burst_recv_next_seq) etcp_metrics_add_loss(etcp, seq - link->burst_recv_next_seq);
link->burst_recv_valid = 0;
}
if (seq < 16) {
link->burst_recv_times[seq] = get_time_us();
link->burst_recv_next_seq = seq + 1;
link->burst_recv_received++;
link->burst_recv_pkt_size = pkt_sz;
}
}
if ((flags & MEAS_FLAG_IS_LAST) || link->burst_recv_received >= link->burst_recv_count) {
uint8_t total = (flags & MEAS_FLAG_IS_LAST) ? link->burst_recv_received : link->burst_recv_count;
if (link->burst_recv_valid && total >= BURST_SKIP_COUNT + 2) {
uint32_t sum_gap = 0, min_gap = UINT32_MAX;
uint8_t measured = 0;
for (uint8_t i = BURST_SKIP_COUNT; i + 1 < total; i++) {
uint64_t dt = link->burst_recv_times[i + 1] - link->burst_recv_times[i];
uint32_t gap = (uint32_t)(dt > 0 ? dt : 1);
sum_gap += gap; measured++;
if (gap < min_gap) min_gap = gap;
}
if (measured > 0) {
link->burst_resp_gap_avg = sum_gap / measured;
link->burst_resp_gap_min = min_gap;
link->burst_resp_pkt_count = measured;
link->burst_resp_valid = 1;
}
}
link->burst_resp_pending = 1;
}
data += MEAS_TS_SECTION_SIZE; len -= MEAS_TS_SECTION_SIZE;
break;
}
case ETCP_SECTION_MEAS_RESP: {
if (len < MEAS_RESP_SECTION_SIZE) { len = 0; break; }
uint16_t b_id = data[1] | (data[2] << 8);
uint8_t valid = data[3];
uint32_t gap_avg = data[4] | (data[5] << 8) | ((uint32_t)data[6] << 16) | ((uint32_t)data[7] << 24);
uint32_t gap_min = data[8] | (data[9] << 8) | ((uint32_t)data[10] << 16) | ((uint32_t)data[11] << 24);
uint8_t pkt_cnt = data[12];
struct ETCP_LINK* link = pkt->link;
if (b_id == link->burst_id - 1 || b_id == link->burst_id) {
if (link->burst_resp_timer) { uasync_cancel_timeout(link->etcp->instance->ua, link->burst_resp_timer); link->burst_resp_timer = NULL; }
if (valid == MEAS_RESP_VALID && gap_min > 0 && link->burst_pkt_size > 0) {
uint64_t bw_kbps = (uint64_t)link->burst_pkt_size * 8000ULL / gap_min;
if (bw_kbps > 10000000) bw_kbps = 10000000;
link->bandwidth = (uint32_t)bw_kbps;
float rtt_sec = (float)link->bbr->min_rtt_us / 1000000.f;
if (rtt_sec < 0.005f) rtt_sec = 0.005f;
link->burst_target_bdp = (uint32_t)((float)bw_kbps * 1000.f / 8.f * rtt_sec);
// BBR: инициализируем верхнюю границу из burst
if (link->bbr->inflight_hi == ~0U) link->bbr->inflight_hi = link->burst_target_bdp;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] burst resp: BW=%u Kbps, BDP=%u gap_min=%u us pkt_cnt=%u",
link->etcp->log_name, (uint32_t)bw_kbps, link->burst_target_bdp, gap_min, pkt_cnt);
}
}
data += MEAS_RESP_SECTION_SIZE; len -= MEAS_RESP_SECTION_SIZE;
break;
}
case ETCP_SECTION_FILLER: {
if (len < FILLER_HDR_SIZE) { len = 0; break; }
uint16_t fill_len = data[1] | (data[2] << 8);
if ((uint16_t)(fill_len + FILLER_HDR_SIZE) > len) { len = 0; break; }
data += FILLER_HDR_SIZE + fill_len; len -= FILLER_HDR_SIZE + fill_len;
break;
}
case ETCP_SECTION_METRICS: {
if (len < 2 + SC_SIGN_SIZE) { len = 0; break; }
uint16_t csv_len = data[1] | (data[2] << 8);
if ((uint32_t)(2 + SC_SIGN_SIZE + csv_len) > len) { len = 0; break; }
uint8_t* sig = data + 2;
uint8_t* csv = data + 2 + SC_SIGN_SIZE;
uint8_t* peer_pubkey = pkt->link->remote_ed25519_pubkey;
struct sc_stream_sign_state verify_state;
int sig_ok = (sc_stream_sign_verify_init(&verify_state, peer_pubkey) == SC_OK
&& sc_stream_sign_update(&verify_state, csv, csv_len) == SC_OK
&& sc_stream_sign_verify(&verify_state, sig, SC_SIGN_SIZE) == SC_OK);
if (!sig_ok && verify_state.initialized) /* update failed, then verify didn't run */
sc_stream_sign_cleanup(&verify_state);
if (sig_ok) {
const char* stats_dir = etcp->instance->stats_dir;
if (stats_dir[0] && etcp->peer_node_id != 0) {
char path[768];
snprintf(path, sizeof(path), "%s/from_%016llx.dat", stats_dir, (unsigned long long)etcp->peer_node_id);
FILE* f = fopen(path, "a");
if (f) {
fprintf(f, "%.*s,", csv_len, csv);
for (int i = 0; i < SC_SIGN_SIZE; i++) fprintf(f, "%02x", sig[i]);
fprintf(f, "\n");
fclose(f);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Received metrics from peer (%u bytes)", etcp->log_name, csv_len);
if (etcp->instance->api_bindings.on_metrics_rcvd)
etcp->instance->api_bindings.on_metrics_rcvd(etcp->instance->api_bindings.metrics_user_ptr, etcp, csv, csv_len, sig);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] Cannot open from_metrics file: %s", etcp->log_name, path);
}
}
} else {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "[%s] Metrics signature verification FAILED", etcp->log_name);
}
data += 2 + SC_SIGN_SIZE + csv_len; len -= 2 + SC_SIGN_SIZE + csv_len;
break;
}
default:
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "unknown section type=0x%02x", type);
len=0;
break;
}
}
if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt);
volatile int _halt = 1; while (_halt) {}
}
memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram
}
void etcp_update_mtu(struct ETCP_CONN* etcp) {
if (!etcp) return;
int new_mtu = PACKET_DATA_MAX_MTU;
int has_links = 0;
struct ETCP_LINK* link = etcp->links;
while (link) {
has_links = 1;
if (link->mtu > 0 && link->mtu < new_mtu) new_mtu = link->mtu;
link = link->next;
}
if (!has_links) new_mtu = ETCP_RFC791_MIN_MTU;
if (new_mtu != etcp->mtu) {
etcp->mtu = new_mtu;
if (etcp->normalizer) etcp->normalizer->frag_size = etcp->mtu - ACK_REZERV - UDP_HDR_SIZE - UDP_SC_HDR_SIZE;
}
}
// ====================================================================== Метрики: гистограммы RTT/потерь + тотальные счётчики
static const uint16_t etcp_metrics_rtt_bounds[] = { 10, 20, 30, 50, 80, 130, 210, 340, 550, 900, 1500 };
#define ETCP_METRICS_RTT_BOUNDS_COUNT (sizeof(etcp_metrics_rtt_bounds) / sizeof(etcp_metrics_rtt_bounds[0]))
void etcp_metrics_init(struct etcp_metrics* m) {
if (!m) return;
memset(m, 0, sizeof(*m));
}
void etcp_metrics_add_rtt(struct ETCP_CONN* etcp, uint16_t rtt_tb) {
if (!etcp) return;
struct etcp_metrics* m = &etcp->metrics;
int i;
for (i = 0; i < (int)ETCP_METRICS_RTT_BOUNDS_COUNT; i++) {
if (rtt_tb < etcp_metrics_rtt_bounds[i]) break;
}
if (i >= ETCP_METRICS_RTT_BUCKETS) i = ETCP_METRICS_RTT_BUCKETS - 1;
m->rtt_hist[i]++;
m->work_samples++;
}
void etcp_metrics_add_sent(struct ETCP_CONN* etcp, uint32_t len) {
if (!etcp) return;
struct etcp_metrics* m = &etcp->metrics;
m->work_sent++;
m->work_bytes_sent += len;
m->total_sent++;
m->total_bytes_sent += len;
}
void etcp_metrics_add_rcvd(struct ETCP_CONN* etcp, uint32_t len) {
if (!etcp) return;
struct etcp_metrics* m = &etcp->metrics;
m->work_rcvd++;
m->work_bytes_rcvd += len;
m->total_rcvd++;
m->total_bytes_rcvd += len;
}
void etcp_metrics_add_loss(struct ETCP_CONN* etcp, uint32_t count) {
if (!etcp) return;
struct etcp_metrics* m = &etcp->metrics;
m->work_lost += count;
m->total_lost += count;
}
static void etcp_metrics_do_snapshot(struct ETCP_CONN* etcp) {
if (!etcp || !etcp->instance) return;
struct etcp_metrics* m = &etcp->metrics;
// Копируем working → final
memcpy(m->rtt_hist_final, m->rtt_hist, sizeof(m->rtt_hist_final));
memcpy(m->loss_hist_final, m->loss_hist, sizeof(m->loss_hist_final));
m->samples_final = m->work_samples;
m->sent_final = m->work_sent;
m->rcvd_final = m->work_rcvd;
m->bytes_sent_final = m->work_bytes_sent;
m->bytes_rcvd_final = m->work_bytes_rcvd;
m->lost_final = m->work_lost;
// Рассчитываем loss-гистограмму
uint32_t total = m->work_sent + m->work_lost;
if (total > 0) {
uint32_t loss_pct = (m->work_lost * 100 + total / 2) / total;
int bucket = (int)loss_pct;
if (bucket >= ETCP_METRICS_LOSS_BUCKETS) bucket = ETCP_METRICS_LOSS_BUCKETS - 1;
m->loss_hist[bucket]++;
}
// Обнуляем рабочую копию
memset(m->rtt_hist, 0, sizeof(m->rtt_hist));
memset(m->loss_hist, 0, sizeof(m->loss_hist));
m->work_samples = 0;
m->work_sent = 0; m->work_rcvd = 0;
m->work_bytes_sent = 0; m->work_bytes_rcvd = 0;
m->work_lost = 0;
m->last_snapshot_tb = get_time_tb();
// Строим CSV строку
char csv_buf[4096];
int pos = 0;
time_t now = time(NULL);
struct tm* tm = localtime(&now);
char date_str[16], time_str[16];
if (tm) { strftime(date_str, sizeof(date_str), "%Y-%m-%d", tm); strftime(time_str, sizeof(time_str), "%H:%M:%S", tm); }
else { strcpy(date_str, "----"); strcpy(time_str, "----"); }
pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, "%s,%s", date_str, time_str);
for (int i = 0; i < ETCP_METRICS_RTT_BUCKETS; i++) pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%u", m->rtt_hist_final[i]);
for (int i = 0; i < ETCP_METRICS_LOSS_BUCKETS; i++) pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%u", m->loss_hist_final[i]);
pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%u,%u,%u,%u,%u,%u", m->samples_final, m->sent_final, m->rcvd_final, m->bytes_sent_final, m->bytes_rcvd_final, m->lost_final);
pos += snprintf(csv_buf + pos, sizeof(csv_buf) - pos, ",%llu,%llu,%llu,%llu,%llu", (unsigned long long)m->total_sent, (unsigned long long)m->total_rcvd, (unsigned long long)m->total_bytes_sent, (unsigned long long)m->total_bytes_rcvd, (unsigned long long)m->total_lost);
uint32_t csv_len = (uint32_t)pos;
// Запись локального файла
const char* stats_dir = etcp->instance->stats_dir;
if (stats_dir[0] && etcp->peer_node_id != 0) {
char path[768];
snprintf(path, sizeof(path), "%s/%016llx.dat", stats_dir, (unsigned long long)etcp->peer_node_id);
FILE* f = fopen(path, "a");
if (f) { fprintf(f, "%s\n", csv_buf); fclose(f); }
else DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] Cannot open metrics file: %s", etcp->log_name, path);
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Metrics snapshot: samples=%u sent=%u rcvd=%u lost=%u", etcp->log_name, m->samples_final, m->sent_final, m->rcvd_final, m->lost_final);
// Подпись Ed25519 и отправка пиру
if (etcp->crypto_ctx.initialized && etcp->initialized) {
struct sc_stream_sign_state sign_state;
if (sc_stream_sign_init(&etcp->crypto_ctx, &sign_state) == SC_OK) {
uint8_t sig[SC_SIGN_SIZE];
size_t sig_len = sizeof(sig);
if (sc_stream_sign_update(&sign_state, (uint8_t*)csv_buf, csv_len) == SC_OK
&& sc_stream_sign_final(&sign_state, sig, &sig_len) == SC_OK) {
// Построить и отправить дграмму с секцией METRICS
struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp);
if (link) {
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (dgram) {
dgram->link = link;
dgram->noencrypt_len = 0;
dgram->timestamp = get_current_timestamp();
int ptr = 0;
// Минимальная ACK-секция
dgram->data[ptr++] = 1; dgram->data[ptr++] = 0;
dgram->data[ptr++] = etcp->last_delivered_id; dgram->data[ptr++] = etcp->last_delivered_id>>8;
dgram->data[ptr++] = etcp->last_delivered_id>>16; dgram->data[ptr++] = etcp->last_delivered_id>>24;
dgram->data[ptr++] = etcp->rx_dup_count; dgram->data[ptr++] = etcp->rx_dup_count>>8;
// METRICS секция
dgram->data[ptr++] = ETCP_SECTION_METRICS;
uint16_t csv16 = (uint16_t)csv_len;
dgram->data[ptr++] = csv16 & 0xFF; dgram->data[ptr++] = (csv16 >> 8) & 0xFF;
memcpy(dgram->data + ptr, sig, SC_SIGN_SIZE); ptr += SC_SIGN_SIZE;
memcpy(dgram->data + ptr, csv_buf, csv_len); ptr += csv_len;
dgram->data_len = ptr;
etcp_encrypt_send(dgram);
memory_pool_free(etcp->instance->pkt_pool, dgram);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Metrics sent to peer (%u bytes)", etcp->log_name, csv_len);
}
}
} else {
sc_stream_sign_cleanup(&sign_state);
}
}
}
}
static void metrics_snapshot_timer_cb(void* arg) {
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
if (!etcp) return;
etcp->metrics.timer = NULL;
if (etcp->metrics.work_sent + etcp->metrics.work_rcvd >= ETCP_METRICS_MIN_PACKETS)
etcp_metrics_do_snapshot(etcp);
// Перезапускаем таймер
etcp->metrics.timer = uasync_set_timeout(etcp->instance->ua, ETCP_METRICS_INTERVAL_TB, etcp, metrics_snapshot_timer_cb, "etcp_metrics");
}
void etcp_metrics_start_timer(struct ETCP_CONN* etcp) {
if (!etcp || etcp->metrics.timer) return;
etcp_metrics_init(&etcp->metrics);
etcp->metrics.timer = uasync_set_timeout(etcp->instance->ua, ETCP_METRICS_INTERVAL_TB, etcp, metrics_snapshot_timer_cb, "etcp_metrics");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Metrics timer started (interval=%ums)", etcp->log_name, ETCP_METRICS_INTERVAL_TB / 10);
}
void etcp_metrics_stop_timer(struct ETCP_CONN* etcp) {
if (!etcp) return;
if (etcp->metrics.timer && etcp->instance) {
uasync_cancel_timeout(etcp->instance->ua, etcp->metrics.timer);
etcp->metrics.timer = NULL;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Metrics timer stopped", etcp->log_name);
}
}