Browse Source

Implement etcp_ack_recv function for ETCP packet acknowledgment

- Add etcp_ack_recv function declaration to etcp.h
- Implement complete ACK processing with RTT calculation and statistics
- Handle packet removal from inflight queues with proper memory management
- Update connection state and trigger new packet transmission
- Integrates with existing etcp_conn_input ACK section processing
nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
c1f73584e0
  1. 542
      src/etcp.c
  2. 15
      src/etcp.h

542
src/etcp.c

@ -29,7 +29,10 @@
//#define CONTAINER_OF(ptr, type, member) ((type *)((char *)(ptr) - offsetof(type, member)))
// Forward declarations
static void input_queue_cb(struct ll_queue* q, void* data, void* arg);
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 etcp_conn_process_send_queue(struct ETCP_CONN* etcp);
// Get current time in 0.1ms units
uint64_t get_current_time_units() {
@ -71,14 +74,14 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance) {
etcp->input_send_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE); // Hash for send_q
etcp->input_wait_ack = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE); // Hash for wait_ack
etcp->recv_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE); // Hash for send_q
etcp->ack_queue = queue_new(instance->ua, 0);
etcp->ack_q = queue_new(instance->ua, 0);
etcp->inflight_pool = memory_pool_init(sizeof(struct INFLIGHT_PACKET));
etcp->rx_pool = memory_pool_init(sizeof(struct RX_PACKET));
etcp->rx_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT));
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "etcp_connection_create: queues created - 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);
if (!etcp->input_queue || !etcp->output_queue || !etcp->input_send_q || !etcp->recv_q || !etcp->ack_queue ||
if (!etcp->input_queue || !etcp->output_queue || !etcp->input_send_q || !etcp->recv_q || !etcp->ack_q ||
!etcp->input_wait_ack || !etcp->inflight_pool || !etcp->rx_pool) {
etcp_connection_close(etcp);
return NULL;
@ -94,6 +97,8 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance) {
// 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);
etcp->link_ready_for_send_fn = etcp_link_ready_callback;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connection_create: queues created, mtu=%d, window_size=%u, next_tx_id=%u",
etcp->mtu, etcp->window_size, etcp->next_tx_id);
@ -122,9 +127,9 @@ void etcp_connection_close(struct ETCP_CONN* etcp) {
queue_free(etcp->input_wait_ack);
etcp->input_wait_ack = NULL;
}
if (etcp->ack_queue) {
queue_free(etcp->ack_queue);
etcp->ack_queue = NULL;
if (etcp->ack_q) {
queue_free(etcp->ack_q);
etcp->ack_q = NULL;
}
if (etcp->recv_q) {
queue_free(etcp->recv_q);
@ -169,6 +174,70 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) {
// Clear inflight, rx_list, etc.
}
// ====================================================================== Отправка данных
// Send data through ETCP connection
// Allocates memory from data_pool and places in input queue
// Returns: 0 on success, -1 on failure
int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_send: ENTER etcp=%p, data=%p, len=%zu", etcp, data, len);
if (!etcp || !data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: invalid parameters (etcp=%p, data=%p, len=%zu)", etcp, data, len);
return -1;
}
if (!etcp->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: input_queue is NULL for etcp=%p", etcp);
return -1;
}
// Check length against maximum packet size
if (len > PACKET_DATA_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: packet too large (len=%zu, max=%d)", len, PACKET_DATA_SIZE);
return -1;
}
// Allocate packet data from data_pool (following ETCP reception pattern)
uint8_t* packet_data = memory_pool_alloc(etcp->instance->data_pool);
if (!packet_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate packet data from data_pool");
return -1;
}
// Copy user data to packet buffer
memcpy(packet_data, data, len);
// Create queue entry - this allocates ll_entry + data pointer
struct ETCP_FRAGMENT* pkt = memory_pool_alloc(etcp->rx_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate queue entry");
memory_pool_free(etcp->instance->data_pool, packet_data);
return -1;
}
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->pkt_data = packet_data; // Point to data_pool allocation
pkt->pkt_len = len;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_send: created PACKET %p with data %p (len=%zu)", pkt, packet_data, len);
// Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, pkt, 0) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to add to input queue");
memory_pool_free(etcp->instance->data_pool, packet_data);
memory_pool_free(etcp->rx_pool, pkt);
return -1;
}
return 0;
}
static void input_queue_try_resume(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_try_resume: ENTER etcp=%p", etcp);
@ -185,43 +254,153 @@ static void input_queue_try_resume(struct ETCP_CONN* etcp) {
// Input callback for input_queue (добавление новых кодограмм в стек)
// input_queue -> input_send_q
static void input_queue_cb(struct ll_queue* q, void* data, void* arg) {
static void input_queue_cb(struct ll_queue* q, void* arg) {
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
struct RX_PACKET* in_pkt = (struct RX_PACKET*)data;
struct ETCP_FRAGMENT* in_pkt = queue_data_get(q);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: processing RX_PACKET %p (seq=%u, len=%u, data=%p)",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: processing ETCP_FRAGMENT %p (seq=%u, len=%u, data=%p)",
in_pkt, in_pkt->seq, in_pkt->pkt_len, in_pkt->pkt_data);
// Create INFLIGHT_PACKET
struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot allocate INFLIGHT_PACKET");
queue_data_free(in_pkt); // Free the RX_PACKET
queue_data_free(in_pkt); // Free the ETCP_FRAGMENT
queue_resume_callback(q);
return;
}
// Setup inflight packet (based on protocol.txt)
memset(p, 0, sizeof(*p));
p->pkt = in_pkt; // Note: RX_PACKET is now owned by inflight; don't free here
p->pkt = in_pkt; // Note: ETCP_FRAGMENT is now owned by inflight; don't free here
p->seq = etcp->next_tx_id++; // Assign seq
p->state = INFLIGHT_STATE_WAIT_SEND;
// Add to send queue
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: adding INFLIGHT_PACKET %p to input_send_q (seq=%u)", p, p->seq);
if (queue_data_put(etcp->input_send_q, p, p->seq) != 0) {
memory_pool_free(etcp->inflight_pool, p);
queue_data_free(in_pkt);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_cb: EXIT (queue put failed)");
return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: successfully added to input_send_q");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: successfully moved from input_queue to input_send_q");
input_queue_try_resume(etcp);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: called input_queue_try_resume for etcp=%p", etcp);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_cb: EXIT success");
}
static void input_send_q_cb(struct ll_queue* q, void* arg) {// etcp->input_send_q processing
struct ETCP_CONN* etcp=(struct ETCP_CONN*)arg;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_send_q_cb: ");
etcp_request_pkt(etcp);
}
// Подготовить и отправить кодограмму
// вызывается линком когда освобождается или очередью если появляются данные на передачу
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp);
if (!link) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no link available");
queue_resume_callback(etcp->input_send_q);
return NULL;
}
size_t send_q_bytes = queue_total_bytes(etcp->input_send_q);
if (send_q_bytes == 0) {// сгребаем из других мест
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: input_send_q empty, trying to resume");
input_queue_try_resume(etcp);
return NULL;
}
// First, check if there's a packet in input_send_q (retrans or new)
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: getting packet from input_send_q");
struct INFLIGHT_PACKET* inf_pkt = queue_data_get(etcp->input_send_q);
if (!inf_pkt) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no packet available from input_send_q");
return NULL;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: got INFLIGHT_PACKET %p (seq=%u, len=%u)",
inf_pkt, inf_pkt->seq, inf_pkt->pkt_len);
inf_pkt->last_timestamp=get_current_time_units();
inf_pkt->send_count++;
inf_pkt->state=INFLIGHT_STATE_WAIT_ACK;
queue_data_put(etcp->input_send_q, inf_pkt, inf_pkt->seq);// move dgram to wait_ack queue
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: failed to allocate ETCP_DGRAM");
return NULL;
}
dgram->link = link;
dgram->noencrypt_len=0;
dgram->timestamp=get_current_timestamp();
dgram->data[0]=1;// ack
int ptr=2;
dgram->data[ptr++]=etcp->last_delivered_id;
dgram->data[ptr++]=etcp->last_delivered_id>>8;
dgram->data[ptr++]=etcp->last_delivered_id>>16;
dgram->data[ptr++]=etcp->last_delivered_id>>24;
// тут (потом) добавим опциональные заголовки
struct ACK_PACKET* ack_pkt;
while (ack_pkt = queue_data_get(etcp->ack_q)) {
// seq 4 байта
dgram->data[ptr++]=ack_pkt->seq;
dgram->data[ptr++]=ack_pkt->seq>>8;
dgram->data[ptr++]=ack_pkt->seq>>16;
dgram->data[ptr++]=ack_pkt->seq>>24;
// ts приема 2 байта
dgram->data[ptr++]=ack_pkt->recv_timestamp;
dgram->data[ptr++]=ack_pkt->recv_timestamp>>8;
// время задержки 2 байта между recv и ack
uint16_t dly=get_current_timestamp()-ack_pkt->recv_timestamp;
dgram->data[ptr++]=dly;
dgram->data[ptr++]=dly>>8;
if (inf_pkt->pkt_len+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки)
}
dgram->data[1]=ptr/8;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: building packet data (seq=%u, len=%u)", inf_pkt->seq, inf_pkt->pkt_len);
dgram->data[ptr++]=0;// payload
memcpy(&dgram->data[ptr], inf_pkt->pkt_data, inf_pkt->pkt_len); ptr+=inf_pkt->pkt_len;
dgram->data_len=ptr;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: built packet with %d bytes total", dgram->data_len);
return dgram;
}
// Callback for when a link is ready to send data
static void etcp_link_ready_callback(struct ETCP_CONN* etcp) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_link_ready_callback: processing send queue for etcp=%p", etcp);
etcp_conn_process_send_queue(etcp);
}
// Process packets in send queue and transmit them
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {
struct ETCP_DGRAM* dgram;
while(dgram = etcp_request_pkt(etcp)) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: sending packet");
etcp_loadbalancer_send(dgram);
}
}
// ====================================================================== Прием данных
void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// пробуем собрать выходную очередь из фрагментов
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: etcp=%p, last_delivered_id=%u, recv_q_count=%d",
@ -233,7 +412,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// Look for contiguous packets starting from next_expected_id
while (1) {
struct RX_PACKET* rx_pkt = queue_find_data_by_id(etcp->recv_q, next_expected_id);
struct ETCP_FRAGMENT* rx_pkt = queue_find_data_by_id(etcp->recv_q, next_expected_id);
if (!rx_pkt) {
// No more contiguous packets found
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: no packet found for id=%u, stopping", next_expected_id);
@ -243,11 +422,11 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: assembling packet id=%u (len=%u)",
rx_pkt->seq, rx_pkt->pkt_len);
// Simply move RX_PACKET from recv_q to output_queue - no data copying needed
// Simply move ETCP_FRAGMENT from recv_q to output_queue - no data copying needed
// Remove from recv_q first
queue_remove_data(etcp->recv_q, rx_pkt);
// Add to output_queue using the same RX_PACKET structure
// Add to output_queue using the same ETCP_FRAGMENT structure
if (queue_data_put(etcp->output_queue, rx_pkt, next_expected_id) == 0) {
delivered_bytes += rx_pkt->pkt_len;
delivered_count++;
@ -270,6 +449,75 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
delivered_count, delivered_bytes, etcp->last_delivered_id, queue_entry_count(etcp->output_queue));
}
// Process ACK receipt - remove acknowledged packet from inflight queues
void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t dts) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: processing ACK for seq=%u, ts=%u, dts=%u", seq, ts, dts);
// Find the acknowledged packet in the wait_ack queue
struct INFLIGHT_PACKET* acked_pkt = queue_find_data_by_id(etcp->input_wait_ack, seq);
if (acked_pkt) queue_remove_data(etcp->input_wait_ack, acked_pkt);
else { acked_pkt = queue_find_data_by_id(etcp->input_send_q, seq);
queue_remove_data(etcp->input_send_q, acked_pkt);
}
if (!acked_pkt) {
// Packet might be already acknowledged or not found
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: packet seq=%u not found in wait_ack queue", seq);
return;
}
// Calculate RTT if timestamps are valid
if (ts != (uint16_t)-1 && dts != (uint16_t)-1) {
uint16_t rtt = timestamp_diff(ts, dts);
etcp->rtt_last = rtt;
// Update RTT averages (exponential smoothing)
if (etcp->rtt_avg_10 == 0) {
etcp->rtt_avg_10 = rtt;
etcp->rtt_avg_100 = rtt;
} else {
// RTT average over 10 packets
etcp->rtt_avg_10 = (etcp->rtt_avg_10 * 9 + rtt) / 10;
// RTT average over 100 packets
etcp->rtt_avg_100 = (etcp->rtt_avg_100 * 99 + rtt) / 100;
}
// Update jitter calculation (max - min of last 10 RTT samples)
etcp->rtt_history[etcp->rtt_history_idx] = rtt;
etcp->rtt_history_idx = (etcp->rtt_history_idx + 1) % 10;
uint16_t rtt_min = UINT16_MAX, rtt_max = 0;
for (int i = 0; i < 10; i++) {
if (etcp->rtt_history[i] < rtt_min) rtt_min = etcp->rtt_history[i];
if (etcp->rtt_history[i] > rtt_max) rtt_max = etcp->rtt_history[i];
}
etcp->jitter = rtt_max - rtt_min;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: RTT updated - last=%u, avg_10=%u, avg_100=%u, jitter=%u",
rtt, etcp->rtt_avg_10, etcp->rtt_avg_100, etcp->jitter);
}
// Update connection statistics
etcp->unacked_bytes -= acked_pkt->pkt_len;
etcp->bytes_sent_total += acked_pkt->pkt_len;
etcp->ack_packets_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: removed packet seq=%u from wait_ack, unacked_bytes now %u",
seq, etcp->unacked_bytes);
if (acked_pkt->pkt_data) {
memory_pool_free(etcp->instance->data_pool, acked_pkt->pkt_data);
}
memory_pool_free(etcp->inflight_pool, acked_pkt);
// Try to resume sending more packets if window space opened up
input_queue_try_resume(etcp);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: completed for seq=%u", seq);
}
// Process incoming decrypted packet
void etcp_conn_input(struct ETCP_DGRAM* pkt) {
@ -286,45 +534,21 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
while (len >= 1) {
uint8_t type = data[0];
uint16_t sec_len = 0; // Section length varies by type
// Determine section length based on type
switch (type) {
case ETCP_SECTION_PAYLOAD:
sec_len = len - 1; // Remaining data is payload
break;
case ETCP_SECTION_ACK:
if (len >= 2) {
uint8_t count = data[1];
sec_len = 1 + 1 + (count * 8) + 4 + 4; // type + count + (id+ts)*count + last_delivered + last_rx
}
break;
case ETCP_SECTION_TIMESTAMP:
sec_len = 1 + 2 + 2; // type + ret_ts + recv_ts
break;
default:
// For unknown types, skip just the type byte to avoid infinite loop
sec_len = 1;
break;
}
if (sec_len > len) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: invalid section len=%u (remaining=%u)", sec_len, len);
break;
}
// Process sections as per protocol.txt
switch (type) {
case ETCP_SECTION_ACK: {
// Mark packets as acknowledged (find using queue_find_data_by_id)
// TODO: Parse ACKs, remove from input_wait_ack, update RTT
// If wait_timeout_active, trigger resume
break;
}
case ETCP_SECTION_TIMESTAMP: {
// Update RTT, jitter
// TODO: Calculate RTT = timestamp_diff(get_current_timestamp(), ts)
// Update averages, jitter = max(last 10) - min(last 10)
int elm_cnt=data[1]*8;
uint32_t till=data[0] | (data[1]<<8) | (data[2]<<16) | (data[3]<<24);
data+=6;
for (int i=0; i<elm_cnt; i++) {
uint32_t seq=data[0] | (data[1]<<8) | (data[2]<<16) | (data[3]<<24);
uint16_t ts=data[4] | (data[5]<<8);
uint16_t dts=data[6] | (data[7]<<8);
etcp_ack_recv(etcp, seq, ts, dts);
data+=8;
}
while (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_PAYLOAD: {
@ -335,14 +559,14 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24);
p->seq=seq;
p->last_timestamp=pkt->timestamp;
queue_data_put(etcp->ack_queue, p, p->seq);
queue_data_put(etcp->ack_q, p, p->seq);
if ((int32_t)(etcp->last_delivered_id-seq)<0) if (queue_find_data_by_id(etcp->recv_q, seq)==NULL) {// проверяем есть ли пакет с этим seq
uint32_t pkt_len=len-5;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: adding packet seq=%u to recv_q (last_delivered_id=%u)", seq, etcp->last_delivered_id);
// отправляем пакет в очередь на сборку
uint8_t* payload_data = memory_pool_alloc(etcp->instance->data_pool);
struct RX_PACKET* rx_pkt = memory_pool_alloc(etcp->rx_pool);
struct ETCP_FRAGMENT* rx_pkt = memory_pool_alloc(etcp->rx_pool);
rx_pkt->seq=seq;
rx_pkt->timestamp=pkt->timestamp;
rx_pkt->pkt_data=payload_data;
@ -355,7 +579,6 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
}
}
len=0;
break;
}
@ -364,214 +587,9 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
break;
}
data += sec_len;
len -= sec_len;
}
memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram
}
// Подготовить и отправить кодограмму
// вызывается линком когда освобождается или очередью если появляются данные на передачу
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: ENTER etcp=%p", etcp);
struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp);
if (!link) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no link available");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: EXIT (no link)");
return NULL;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: selected link %p", link);
size_t send_q_bytes = queue_total_bytes(etcp->input_send_q);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: input_send_q bytes=%zu", send_q_bytes);
if (send_q_bytes == 0) {// сгребаем из других мест
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: input_send_q empty, trying to resume");
input_queue_try_resume(etcp);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: EXIT (send_q empty)");
return NULL;
}
// First, check if there's a packet in input_send_q (retrans or new)
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: getting packet from input_send_q");
struct INFLIGHT_PACKET* inf_pkt = queue_data_get(etcp->input_send_q);
if (!inf_pkt) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no packet available from input_send_q");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: EXIT (no packet)");
return NULL;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: got INFLIGHT_PACKET %p (seq=%u, len=%u)",
inf_pkt, inf_pkt->seq, inf_pkt->pkt_len);
inf_pkt->last_timestamp=get_current_time_units();
inf_pkt->send_count++;
inf_pkt->state=INFLIGHT_STATE_WAIT_ACK;
// queue_data_put(etcp->input_send_q, inf_pkt, inf_pkt->seq);// move dgram to wait_ack queue
// Build outgoing dgram (stub: allocate from pkt_pool)
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: allocating ETCP_DGRAM from pkt_pool");
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: failed to allocate ETCP_DGRAM");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: EXIT (allocation failed)");
return NULL;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: allocated ETCP_DGRAM %p", dgram);
dgram->link = link;
dgram->noencrypt_len=0;
dgram->timestamp=get_current_timestamp();
int ptr=0;
// тут (потом) добавим опциональные заголовки
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: building packet data (seq=%u, len=%u)", inf_pkt->seq, inf_pkt->pkt_len);
dgram->data[ptr++]=0;// payload
memcpy(&dgram->data[ptr], inf_pkt->pkt_data, inf_pkt->pkt_len); ptr+=inf_pkt->pkt_len;
dgram->data_len=ptr;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: built packet with %d bytes total", dgram->data_len);
memory_pool_free(etcp->instance->data_pool, inf_pkt->pkt_data);// для теста отсылаем и освобождаем сразу
memory_pool_free(etcp->inflight_pool, inf_pkt);// для теста отсылаем и освобождаем сразу
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: freed packet data and INFLIGHT_PACKET");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: EXIT success (dgram=%p)", dgram);
return dgram;
}
// Forward declarations
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp);
// Callback for when a link is ready to send data
static void etcp_link_ready_callback(struct ETCP_CONN* etcp) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_link_ready_callback: processing send queue for etcp=%p", etcp);
etcp_conn_process_send_queue(etcp);
}
// Process packets in send queue and transmit them
static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {
if (!etcp || !etcp->links) return;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: etcp=%p, unacked_bytes=%u, window_size=%u",
etcp, etcp->unacked_bytes, etcp->window_size);
// Check if we have window space and packets to send
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: starting loop, unacked_bytes=%u, window_size=%u",
etcp->unacked_bytes, etcp->window_size);
while (etcp->unacked_bytes < etcp->window_size) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: loop iteration, unacked_bytes=%u < window_size=%u",
etcp->unacked_bytes, etcp->window_size);
// Try to get a packet from input_send_q
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: getting packet from input_send_q");
void* inflight_data = queue_data_get(etcp->input_send_q);
if (!inflight_data) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: no packets in send queue");
break;
}
// Get the inflight packet structure
struct INFLIGHT_PACKET* inflight = (struct INFLIGHT_PACKET*)((char*)inflight_data - sizeof(struct ll_entry));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: got INFLIGHT_PACKET %p (seq=%u, len=%u, state=%u)",
inflight, inflight->seq, inflight->pkt_len, inflight->state);
// Build and send the packet
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: calling etcp_request_pkt for seq=%u", inflight->seq);
struct ETCP_DGRAM* dgram = etcp_request_pkt(etcp);
if (dgram) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: sending packet seq=%u", inflight->seq);
etcp_loadbalancer_send(dgram);
// Update inflight state
inflight->last_timestamp = get_current_time_units();
inflight->send_count++;
inflight->state = INFLIGHT_STATE_WAIT_ACK;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: updated INFLIGHT_PACKET %p state to WAIT_ACK", inflight);
// Move to wait_ack queue
void* wait_data = (void*)((char*)inflight + sizeof(struct ll_entry));
queue_data_put(etcp->input_wait_ack, wait_data, inflight->seq);
etcp->unacked_bytes += inflight->pkt_len;
etcp->total_packets_sent++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: packet seq=%u sent, unacked_bytes now %u",
inflight->seq, etcp->unacked_bytes);
} else {
// Put the packet back if we couldn't send it
queue_data_put(etcp->input_send_q, inflight_data, inflight->seq);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: could not build packet, putting back");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: EXIT (build failed)");
break;
}
}
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_conn_process_send_queue: EXIT (loop ended)");
}
// Send data through ETCP connection
// Allocates memory from data_pool and places in input queue
// Returns: 0 on success, -1 on failure
int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
printf ("*********************** etcp_send: ENTER %d %d", g_debug_config.level,g_debug_config.categories);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_send: ENTER etcp=%p, data=%p, len=%zu", etcp, data, len);
if (!etcp || !data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: invalid parameters (etcp=%p, data=%p, len=%zu)", etcp, data, len);
return -1;
}
if (!etcp->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: input_queue is NULL for etcp=%p", etcp);
return -1;
}
// Check length against maximum packet size
if (len > PACKET_DATA_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: packet too large (len=%zu, max=%d)", len, PACKET_DATA_SIZE);
return -1;
}
// Allocate packet data from data_pool (following ETCP reception pattern)
uint8_t* packet_data = memory_pool_alloc(etcp->instance->data_pool);
if (!packet_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate packet data from data_pool");
return -1;
}
// Copy user data to packet buffer
memcpy(packet_data, data, len);
// Create queue entry - this allocates ll_entry + data pointer
struct RX_PACKET* pkt = memory_pool_alloc(etcp->rx_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate queue entry");
memory_pool_free(etcp->instance->data_pool, packet_data);
return -1;
}
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->pkt_data = packet_data; // Point to data_pool allocation
pkt->pkt_len = len;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_send: created RX_PACKET %p with data %p (len=%zu)", pkt, packet_data, len);
// Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, pkt, 0) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to add to input queue");
memory_pool_free(etcp->instance->data_pool, packet_data);
memory_pool_free(etcp->rx_pool, pkt);
return -1;
}
return 0;
}

15
src/etcp.h

@ -11,6 +11,8 @@
extern "C" {
#endif
//!!!!!!!!!!! надо переделать ll_queue чтобы возвращала не дата а свою структуру
// Forward declarations
struct UTUN_INSTANCE;
struct UASYNC;
@ -46,7 +48,8 @@ struct INFLIGHT_PACKET {
};
// Список пакетов для сборки. собирается в ll_queue (используем быстрый поиск с хешем)
struct RX_PACKET {
struct ETCP_FRAGMENT {
struct ll_entry ll;
uint32_t seq;
uint16_t timestamp;
uint8_t* pkt_data;
@ -54,6 +57,7 @@ struct RX_PACKET {
};
struct ACK_PACKET {
struct ll_entry ll;
uint32_t seq;// sequence number
uint16_t pkt_timestamp;// timestamp пакета (часы уладенной стороны)
uint32_t recv_timestamp;// время приема (локальное)
@ -83,9 +87,9 @@ struct ETCP_CONN {
struct memory_pool* rx_pool; // память для rx очередей
struct ll_queue* input_send_q; // очередь на отправку (с элементами struct INFLIGHT_PACKET)
struct ll_queue* input_wait_ack; // очередь ожидающих подтверждение (с элементами struct INFLIGHT_PACKET)
struct ll_queue* ack_queue; // неотправленные подтверждения приема пакетов
struct ll_queue* ack_q; // неотправленные подтверждения приема пакетов
struct ll_queue* recv_q; // очередь на сборку (с элементами struct RX_PACKET)
struct ll_queue* recv_q; // очередь на сборку (с элементами struct ETCP_FRAGMENT)
void (*link_ready_for_send_fn)(struct ETCP_CONN*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные
@ -97,6 +101,8 @@ struct ETCP_CONN {
uint32_t last_rx_id; // Last received ID
uint32_t last_delivered_id; // Last delivered to output_queue
uint32_t rx_ack_till;// из ack пакета - по какой пакет получено и собрано на уданенной стороне
// Metrics (RTT, jitter, etc.)
uint16_t rtt_last;
uint16_t rtt_avg_10;
@ -152,6 +158,9 @@ int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len);
// Request next packet for load balancer
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp);
// 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);
#ifdef __cplusplus
}
#endif

Loading…
Cancel
Save