diff --git a/src/etcp.c b/src/etcp.c index ba045162..5c2e8595 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -216,7 +216,8 @@ int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) { // Create queue entry - this allocates ll_entry + data pointer - struct ETCP_FRAGMENT* pkt = memory_pool_alloc(etcp->rx_pool); + + struct ETCP_FRAGMENT* pkt = queue_data_new_from_pool(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); @@ -264,8 +265,16 @@ static void input_queue_cb(struct ll_queue* q, void* arg) { struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg; struct ETCP_FRAGMENT* in_pkt = queue_data_get(q); + if (!in_pkt) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot get element (pool=%p etcp=%p)", etcp->inflight_pool, etcp); + queue_resume_callback(q); + return; + } + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: processing ETCP_FRAGMENT %p (seq=%u, len=%u)", in_pkt, in_pkt->seq, in_pkt->ll.size); + memory_pool_free(etcp->rx_pool, in_pkt);// перемещаем из rx_pool в inflight_pool + // Create INFLIGHT_PACKET struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool); if (!p) { @@ -433,6 +442,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { dgram->data[ptr++]=dly; dgram->data[ptr++]=dly>>8; DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: add ACK N%d dTS=%d", ack_pkt->seq, dly); + queue_data_free(ack_pkt); if (inf_pkt && inf_pkt->ll.size+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки) if (ptr>500) break; @@ -512,7 +522,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) { // 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 ETCP_FRAGMENT structure if (queue_data_put(etcp->output_queue, rx_pkt, next_expected_id) == 0) { delivered_bytes += rx_pkt->ll.size; @@ -591,8 +601,7 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d etcp->bytes_sent_total += acked_pkt->ll.size; 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); + 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); @@ -642,7 +651,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { if (len>=5) { // формируем ACK - struct ACK_PACKET* p = memory_pool_alloc(etcp->instance->ack_pool); + struct ACK_PACKET* p = queue_data_new_from_pool(etcp->instance->ack_pool); uint32_t seq=data[1] | (data[2]<<8) | (data[3]<<16) | (data[4]<<24); p->seq=seq; p->pkt_timestamp=pkt->timestamp; @@ -657,7 +666,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { 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 ETCP_FRAGMENT* rx_pkt = memory_pool_alloc(etcp->rx_pool); + struct ETCP_FRAGMENT* rx_pkt = queue_data_new_from_pool(etcp->rx_pool); rx_pkt->seq=seq; rx_pkt->timestamp=pkt->timestamp; rx_pkt->pkt_data=payload_data; diff --git a/src/etcp.h b/src/etcp.h index fe596215..43f9a0b3 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -77,7 +77,7 @@ struct ETCP_CONN { uint64_t peer_node_id; // Peer node ID // ============ Processing incoming data to be sent by ETCP - struct ll_queue* input_queue; // Incoming packets to send + struct ll_queue* input_queue; // Incoming packets to send (storage: ETCP_FRAGMENT / rx_pool) // Inflight очереди (2 шт) - пока пакет в статусе inflight - к нему прикрепляется struct INFLIGHT_PACKET struct memory_pool* inflight_pool; // память для inflight очередей @@ -90,7 +90,7 @@ struct ETCP_CONN { void (*link_ready_for_send_fn)(struct ETCP_CONN*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные - struct ll_queue* output_queue; // Assembled outgoing packets + struct ll_queue* output_queue; // Assembled outgoing packets (storage: ETCP_FRAGMENT / rx_pool) // IDs and state diff --git a/src/utun_instance.c b/src/utun_instance.c index e03ec1a6..3c794806 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -214,6 +214,13 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { instance->ack_pool = NULL; } + // Cleanup data pool + if (instance->data_pool) { + DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Destroying data pool"); + memory_pool_destroy(instance->data_pool); + instance->data_pool = NULL; + } + // FINALLY destroy uasync (after all resources are cleaned up) if (instance->ua) { DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Destroying uasync instance"); diff --git a/tests/test_etcp_100_packets b/tests/test_etcp_100_packets index 91105b4d..e0f2d801 100755 Binary files a/tests/test_etcp_100_packets and b/tests/test_etcp_100_packets differ