diff --git a/lib/u_async.c b/lib/u_async.c index a001a760..334e19e3 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -297,8 +297,8 @@ static void get_current_time(struct timeval* tv) { ul.HighPart = ft.dwHighDateTime; // Конвертируем 100-наносекундные интервалы в секунды и микросекунды - tv->tv_sec = (long)((ul.QuadPart - 116444736000000000LL) / 10000000); - tv->tv_usec = (long)((ul.QuadPart % 10000000) / 10); + tv->tv_sec = (long)((ul.QuadPart - 116444736000000000ULL) / 10000000ULL); + tv->tv_usec = (long)((ul.QuadPart % 10000000ULL) / 10); #else // Для Linux и других Unix-подобных систем используем gettimeofday gettimeofday(tv, NULL); @@ -308,10 +308,17 @@ static void get_current_time(struct timeval* tv) { #ifdef _WIN32 uint64_t get_time_tb(void) { LARGE_INTEGER freq, count; - QueryPerformanceFrequency(&freq); // Получаем частоту таймера - QueryPerformanceCounter(&count); // Получаем текущее значение счётчика - return (uint64_t)(count.QuadPart * 10000) / (uint64_t)freq.QuadPart; // Преобразуем в требуемые единицы времени + QueryPerformanceFrequency(&freq); + QueryPerformanceCounter(&count); + double t = (double)count.QuadPart * 10000.0 / (double)freq.QuadPart; + return (uint64_t)t; } +//uint64_t get_time_tb(void) { +// LARGE_INTEGER freq, count; +// QueryPerformanceFrequency(&freq); // Получаем частоту таймера +// QueryPerformanceCounter(&count); // Получаем текущее значение счётчика +// return (uint64_t)(count.QuadPart * 10000ULL) / (uint64_t)freq.QuadPart; // Преобразуем в требуемые единицы времени +//} #else uint64_t get_time_tb(void) { struct timespec ts; diff --git a/src/etcp.c b/src/etcp.c index 74b27402..d570f3d3 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -42,54 +42,18 @@ static void send_ack_req_cb(struct ll_queue* q, void* arg); static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp); struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp); -// Helper function to drain and free ETCP_FRAGMENT queue -static void drain_and_free_fragment_queue(struct ETCP_CONN* etcp, struct ll_queue** q) { - if (!*q) return; - struct ETCP_FRAGMENT* pkt; - while ((pkt = (struct ETCP_FRAGMENT*)queue_data_get(*q)) != NULL) { - if (pkt->ll.dgram) { - memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram); - } - memory_pool_free(etcp->io_pool, pkt); - } - queue_free(*q); - *q = NULL; -} - -// Helper function to clear INFLIGHT_PACKET queue (keep queue, free elements) -static void clear_inflight_queue(struct ll_queue* q, struct ETCP_CONN* etcp) { +static void clear_queue(struct ll_queue* q) { if (!q) return; - struct INFLIGHT_PACKET* pkt; - while ((pkt = (struct INFLIGHT_PACKET*)queue_data_get(q)) != NULL) { - if (pkt->ll.dgram) { - memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram); - } - memory_pool_free(etcp->inflight_pool, pkt); + struct ll_entry* pkt; + while ((pkt = queue_data_get(q)) != NULL) { + queue_dgram_free(pkt); + queue_entry_free(pkt); } } -// Helper function to clear ETCP_FRAGMENT queue (keep queue, free elements) -static void clear_fragment_queue(struct ll_queue* q, struct ETCP_CONN* etcp) { - if (!q) return; - struct ETCP_FRAGMENT* pkt; - while ((pkt = (struct ETCP_FRAGMENT*)queue_data_get(q)) != NULL) { - if (pkt->ll.dgram) { - memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram); - } - memory_pool_free(etcp->io_pool, pkt); - } -} - -// Helper function to drain and free INFLIGHT_PACKET queue -static void drain_and_free_inflight_queue(struct ETCP_CONN* etcp, struct ll_queue** q) { +static void drain_and_free_queue(struct ll_queue** q) { if (!*q) return; - struct INFLIGHT_PACKET* pkt; - while ((pkt = (struct INFLIGHT_PACKET*)queue_data_get(*q)) != NULL) { - if (pkt->ll.dgram) { - memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram); - } - memory_pool_free(etcp->inflight_pool, pkt); - } + clear_queue(*q); queue_free(*q); *q = NULL; } @@ -209,11 +173,11 @@ void etcp_connection_close(struct ETCP_CONN* etcp) { } // Drain and free all queues using helper functions - drain_and_free_fragment_queue(etcp, &etcp->input_queue); - drain_and_free_fragment_queue(etcp, &etcp->output_queue); - drain_and_free_inflight_queue(etcp, &etcp->input_send_q); - drain_and_free_inflight_queue(etcp, &etcp->input_wait_ack); - drain_and_free_fragment_queue(etcp, &etcp->recv_q); + 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 ack_q (contains ACK_PACKET from ack_pool - special handling) if (etcp->ack_q) { @@ -283,11 +247,11 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) { etcp->rtt_history_idx = 0; // Clear queues (keep queue structures) - clear_fragment_queue(etcp->input_queue, etcp); - clear_fragment_queue(etcp->output_queue, etcp); - clear_inflight_queue(etcp->input_send_q, etcp); - clear_inflight_queue(etcp->input_wait_ack, etcp); - clear_fragment_queue(etcp->recv_q, etcp); + 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); // Reset timers (just clear the pointers - timers will expire naturally) etcp->retrans_timer = NULL; @@ -389,9 +353,6 @@ int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) { // Copy user data to packet buffer memcpy(packet_data, data, len); - // Create queue entry - this allocates ll_entry + data pointer - - 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); @@ -401,16 +362,18 @@ int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) { 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 = 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) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name); - memory_pool_free(etcp->instance->data_pool, packet_data); - memory_pool_free(etcp->io_pool, pkt); + queue_dgram_free(&pkt->ll); + queue_entry_free(&pkt->ll); return -1; } @@ -419,6 +382,8 @@ int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) { + + static void input_queue_try_push(struct ETCP_CONN* etcp) {// пробуем протолкнуть при отправке DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); @@ -502,10 +467,10 @@ static void input_queue_cb(struct ll_queue* q, void* arg) { } - memory_pool_free(etcp->io_pool, in_pkt);// перемещаем из io_pool в inflight_pool // Create INFLIGHT_PACKET - struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool); + + 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_entry_free((struct ll_entry*)in_pkt); // Free the ETCP_FRAGMENT @@ -522,14 +487,17 @@ static void input_queue_cb(struct ll_queue* q, void* arg) { 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_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input -> inflight (seq=%u, len=%u)", etcp->log_name, p->seq, p->ll.len); int len=p->ll.len;// сохраним len // Add to send queue - if (queue_data_put(etcp->input_send_q, (struct ll_entry*)p, p->seq) != 0) { + if (queue_data_put(etcp->input_send_q, &p->ll, p->seq) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet seq=%u to input_send_q", etcp->log_name, p->seq); - memory_pool_free(etcp->inflight_pool, p); - queue_entry_free((struct ll_entry*)in_pkt); + queue_dgram_free(&p->ll); + queue_entry_free(&p->ll); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] EXIT (queue put failed)", etcp->log_name); return; } @@ -546,28 +514,29 @@ static void ack_timeout_check(struct ETCP_CONN* etcp) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); uint64_t now = get_time_tb(); - uint64_t timeout = 1000;//(uint64_t)(etcp->rtt_avg_10 * RETRANS_K1) + (uint64_t)(etcp->jitter * RETRANS_K2); + int64_t timeout = 1000;//(uint64_t)(etcp->rtt_avg_10 * RETRANS_K1) + (uint64_t)(etcp->jitter * RETRANS_K2); DEBUG_DEBUG(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; while (current = etcp->input_wait_ack->head) { struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current; int64_t elapsed = now - pkt->last_timestamp; if (elapsed > timeout) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] timeout for seq=%u, elapsed=%llu, now=%llu, timeout=%llu, send_count=%u. Moving: wait_ack -> send_q", + DEBUG_INFO(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); + // Increment counters pkt->send_count++; pkt->retrans_req_count++; // Optional, if used for retrans request logic pkt->last_timestamp = now; - // Remove from wait_ack - queue_data_get(etcp->input_wait_ack); - // Change state and add to send_q for retransmission pkt->state = INFLIGHT_STATE_WAIT_SEND; queue_data_put(etcp->input_send_q, (struct ll_entry*)pkt, pkt->seq); @@ -661,14 +630,17 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { 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) { - inf_pkt->last_timestamp=get_time_tb(); + uint64_t now=get_time_tb(); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] send_q->wait_ack seq=%d TS=%llu", etcp->log_name, inf_pkt->seq, now); + inf_pkt->last_timestamp=now; inf_pkt->send_count++; inf_pkt->state=INFLIGHT_STATE_WAIT_ACK; - queue_data_put(etcp->input_wait_ack, (struct ll_entry*)inf_pkt, inf_pkt->seq);// move dgram to wait_ack queue + queue_data_put(etcp->input_wait_ack, &inf_pkt->ll, inf_pkt->seq);// move dgram to wait_ack queue } size_t ack_q_size = queue_entry_count(etcp->ack_q); @@ -685,7 +657,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { dgram->link = link; dgram->noencrypt_len=0; - dgram->timestamp=get_current_timestamp(); + dgram->timestamp = get_current_timestamp(); // формат ack: [01] [elements count] [4 байта last_delivered_id] и <[4 байта seq][2 байта recv_ts][2 байта txrx delay ts]> x count dgram->data[0]=1;// ack @@ -1007,11 +979,12 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { len=0; break; } - rx_pkt->seq=seq; - rx_pkt->timestamp=pkt->timestamp; - rx_pkt->ll.dgram=payload_data; - rx_pkt->ll.dgram_pool=etcp->instance->data_pool; - rx_pkt->ll.len=pkt_len; + 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; // Copy the actual payload data memcpy(payload_data, data + 5, pkt_len); queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, seq); diff --git a/src/etcp_loadbalancer.c b/src/etcp_loadbalancer.c index cf196faa..6a298f5b 100644 --- a/src/etcp_loadbalancer.c +++ b/src/etcp_loadbalancer.c @@ -84,7 +84,6 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) { if (!etcp) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] no ETCP_CONN associated with dgram", etcp->log_name); memory_pool_free(etcp->instance->pkt_pool, dgram); -// u_free(dgram); return; } @@ -94,7 +93,6 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) { if (!dgram->link) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no link available, dropping dgram", etcp->log_name); memory_pool_free(etcp->instance->pkt_pool, dgram); -// u_free(dgram); // Assume free; adjust if pooled return; } } @@ -144,7 +142,6 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) { // (unsigned long long)link->shaper_load_time_tb, (unsigned long long)link->shaper_sub_nanotime, link->shaper_timer==NULL?1:0); memory_pool_free(etcp->instance->pkt_pool, dgram); -// u_free(dgram); // Free the dgram in all cases - we own it } // New: Notify when link is ready (called from timer or external) diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index dac4142a..3a086cd9 100644 --- a/src/pkt_normalizer.c +++ b/src/pkt_normalizer.c @@ -238,9 +238,10 @@ static void pn_send_to_etcp(struct PKTNORM* pn) { frag->seq = 0; frag->timestamp = 0; - frag->ll.dgram = pn->data; - frag->ll.len = pn->data_ptr; - frag->ll.memlen = pn->etcp->instance->data_pool->object_size; + frag->ll.dgram = pn->data; + frag->ll.len = pn->data_ptr; + frag->ll.dgram_pool = pn->etcp->instance->data_pool; + frag->ll.memlen = pn->etcp->instance->data_pool->object_size; DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn->etcp: size=%d memlen=%d frag_size=%d", frag->ll.len, frag->ll.memlen, pn->frag_size); if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->ETCP", pn->data, frag->ll.len);