diff --git a/src/etcp.c b/src/etcp.c index 2f3824ba..66aba1da 100644 --- a/src/etcp.c +++ b/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; irx_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; -} \ No newline at end of file diff --git a/src/etcp.h b/src/etcp.h index 0073df26..590a825e 100644 --- a/src/etcp.h +++ b/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