Browse Source

fix ack timestamp

nodeinfo-routing-update
jeka 7 months ago
parent
commit
65b6b6c150
  1. 17
      lib/u_async.c
  2. 129
      src/etcp.c
  3. 3
      src/etcp_loadbalancer.c
  4. 7
      src/pkt_normalizer.c

17
lib/u_async.c

@ -297,8 +297,8 @@ static void get_current_time(struct timeval* tv) {
ul.HighPart = ft.dwHighDateTime; ul.HighPart = ft.dwHighDateTime;
// Конвертируем 100-наносекундные интервалы в секунды и микросекунды // Конвертируем 100-наносекундные интервалы в секунды и микросекунды
tv->tv_sec = (long)((ul.QuadPart - 116444736000000000LL) / 10000000); tv->tv_sec = (long)((ul.QuadPart - 116444736000000000ULL) / 10000000ULL);
tv->tv_usec = (long)((ul.QuadPart % 10000000) / 10); tv->tv_usec = (long)((ul.QuadPart % 10000000ULL) / 10);
#else #else
// Для Linux и других Unix-подобных систем используем gettimeofday // Для Linux и других Unix-подобных систем используем gettimeofday
gettimeofday(tv, NULL); gettimeofday(tv, NULL);
@ -308,10 +308,17 @@ static void get_current_time(struct timeval* tv) {
#ifdef _WIN32 #ifdef _WIN32
uint64_t get_time_tb(void) { uint64_t get_time_tb(void) {
LARGE_INTEGER freq, count; LARGE_INTEGER freq, count;
QueryPerformanceFrequency(&freq); // Получаем частоту таймера QueryPerformanceFrequency(&freq);
QueryPerformanceCounter(&count); // Получаем текущее значение счётчика QueryPerformanceCounter(&count);
return (uint64_t)(count.QuadPart * 10000) / (uint64_t)freq.QuadPart; // Преобразуем в требуемые единицы времени 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 #else
uint64_t get_time_tb(void) { uint64_t get_time_tb(void) {
struct timespec ts; struct timespec ts;

129
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); static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp);
struct ETCP_DGRAM* etcp_request_pkt(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 clear_queue(struct ll_queue* q) {
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) {
if (!q) return; if (!q) return;
struct INFLIGHT_PACKET* pkt; struct ll_entry* pkt;
while ((pkt = (struct INFLIGHT_PACKET*)queue_data_get(q)) != NULL) { while ((pkt = queue_data_get(q)) != NULL) {
if (pkt->ll.dgram) { queue_dgram_free(pkt);
memory_pool_free(etcp->instance->data_pool, pkt->ll.dgram); queue_entry_free(pkt);
}
memory_pool_free(etcp->inflight_pool, pkt);
} }
} }
// Helper function to clear ETCP_FRAGMENT queue (keep queue, free elements) static void drain_and_free_queue(struct ll_queue** q) {
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) {
if (!*q) return; if (!*q) return;
struct INFLIGHT_PACKET* pkt; clear_queue(*q);
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);
}
queue_free(*q); queue_free(*q);
*q = NULL; *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 all queues using helper functions
drain_and_free_fragment_queue(etcp, &etcp->input_queue); drain_and_free_queue(&etcp->input_queue);
drain_and_free_fragment_queue(etcp, &etcp->output_queue); drain_and_free_queue(&etcp->output_queue);
drain_and_free_inflight_queue(etcp, &etcp->input_send_q); drain_and_free_queue(&etcp->input_send_q);
drain_and_free_inflight_queue(etcp, &etcp->input_wait_ack); drain_and_free_queue(&etcp->input_wait_ack);
drain_and_free_fragment_queue(etcp, &etcp->recv_q); drain_and_free_queue(&etcp->recv_q);
// Drain and free ack_q (contains ACK_PACKET from ack_pool - special handling) // Drain and free ack_q (contains ACK_PACKET from ack_pool - special handling)
if (etcp->ack_q) { if (etcp->ack_q) {
@ -283,11 +247,11 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) {
etcp->rtt_history_idx = 0; etcp->rtt_history_idx = 0;
// Clear queues (keep queue structures) // Clear queues (keep queue structures)
clear_fragment_queue(etcp->input_queue, etcp); clear_queue(etcp->input_queue);
clear_fragment_queue(etcp->output_queue, etcp); clear_queue(etcp->output_queue);
clear_inflight_queue(etcp->input_send_q, etcp); clear_queue(etcp->input_send_q);
clear_inflight_queue(etcp->input_wait_ack, etcp); clear_queue(etcp->input_wait_ack);
clear_fragment_queue(etcp->recv_q, etcp); clear_queue(etcp->recv_q);
// Reset timers (just clear the pointers - timers will expire naturally) // Reset timers (just clear the pointers - timers will expire naturally)
etcp->retrans_timer = NULL; 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 // Copy user data to packet buffer
memcpy(packet_data, data, len); 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); struct ETCP_FRAGMENT* pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool);
if (!pkt) { if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate queue entry", etcp->log_name); 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->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set 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.dgram = packet_data; // Point to data_pool allocation
pkt->ll.len = len; // размер packet_data 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); 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 // Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, (struct ll_entry*)pkt, 0) != 0) { 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); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name);
memory_pool_free(etcp->instance->data_pool, packet_data); queue_dgram_free(&pkt->ll);
memory_pool_free(etcp->io_pool, pkt); queue_entry_free(&pkt->ll);
return -1; 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) {// пробуем протолкнуть при отправке static void input_queue_try_push(struct ETCP_CONN* etcp) {// пробуем протолкнуть при отправке
DEBUG_TRACE(DEBUG_CATEGORY_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 // 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) { if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp); 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 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.dgram_pool = in_pkt->ll.dgram_pool;
p->ll.len = in_pkt->ll.len; 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); 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 int len=p->ll.len;// сохраним len
// Add to send queue // 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); 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_dgram_free(&p->ll);
queue_entry_free((struct ll_entry*)in_pkt); queue_entry_free(&p->ll);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] EXIT (queue put failed)", etcp->log_name); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] EXIT (queue put failed)", etcp->log_name);
return; return;
} }
@ -546,28 +514,29 @@ static void ack_timeout_check(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
uint64_t now = get_time_tb(); 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", 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); (unsigned long long)now, (unsigned long long)timeout, etcp->rtt_avg_10, etcp->jitter);
struct ll_entry* current; struct ll_entry* current;
while (current = etcp->input_wait_ack->head) { while (current = etcp->input_wait_ack->head) {
struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current; struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current;
int64_t elapsed = now - pkt->last_timestamp; int64_t elapsed = now - pkt->last_timestamp;
if (elapsed > timeout) { 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); 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 // Increment counters
pkt->send_count++; pkt->send_count++;
pkt->retrans_req_count++; // Optional, if used for retrans request logic pkt->retrans_req_count++; // Optional, if used for retrans request logic
pkt->last_timestamp = now; pkt->last_timestamp = now;
// Remove from wait_ack
queue_data_get(etcp->input_wait_ack);
// Change state and add to send_q for retransmission // Change state and add to send_q for retransmission
pkt->state = INFLIGHT_STATE_WAIT_SEND; pkt->state = INFLIGHT_STATE_WAIT_SEND;
queue_data_put(etcp->input_send_q, (struct ll_entry*)pkt, pkt->seq); 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); input_queue_try_push(etcp);
} }
// First, check if there's a packet in input_send_q (retrans or new) // 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"); // 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); struct INFLIGHT_PACKET* inf_pkt = (struct INFLIGHT_PACKET*)queue_data_get(etcp->input_send_q);
if (inf_pkt) { 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->send_count++;
inf_pkt->state=INFLIGHT_STATE_WAIT_ACK; 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); 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->link = link;
dgram->noencrypt_len=0; 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 // формат ack: [01] [elements count] [4 байта last_delivered_id] и <[4 байта seq][2 байта recv_ts][2 байта txrx delay ts]> x count
dgram->data[0]=1;// ack dgram->data[0]=1;// ack
@ -1007,11 +979,12 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
len=0; len=0;
break; break;
} }
rx_pkt->seq=seq; rx_pkt->seq = seq;
rx_pkt->timestamp=pkt->timestamp; rx_pkt->timestamp = pkt->timestamp;
rx_pkt->ll.dgram=payload_data; rx_pkt->ll.dgram = payload_data;
rx_pkt->ll.dgram_pool=etcp->instance->data_pool; rx_pkt->ll.len = pkt_len;
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 // Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len); memcpy(payload_data, data + 5, pkt_len);
queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, seq); queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, seq);

3
src/etcp_loadbalancer.c

@ -84,7 +84,6 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) {
if (!etcp) { if (!etcp) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] no ETCP_CONN associated with dgram", etcp->log_name); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] no ETCP_CONN associated with dgram", etcp->log_name);
memory_pool_free(etcp->instance->pkt_pool, dgram); memory_pool_free(etcp->instance->pkt_pool, dgram);
// u_free(dgram);
return; return;
} }
@ -94,7 +93,6 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) {
if (!dgram->link) { if (!dgram->link) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no link available, dropping dgram", etcp->log_name); DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no link available, dropping dgram", etcp->log_name);
memory_pool_free(etcp->instance->pkt_pool, dgram); memory_pool_free(etcp->instance->pkt_pool, dgram);
// u_free(dgram); // Assume free; adjust if pooled
return; 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); // (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); 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) // New: Notify when link is ready (called from timer or external)

7
src/pkt_normalizer.c

@ -238,9 +238,10 @@ static void pn_send_to_etcp(struct PKTNORM* pn) {
frag->seq = 0; frag->seq = 0;
frag->timestamp = 0; frag->timestamp = 0;
frag->ll.dgram = pn->data; frag->ll.dgram = pn->data;
frag->ll.len = pn->data_ptr; frag->ll.len = pn->data_ptr;
frag->ll.memlen = pn->etcp->instance->data_pool->object_size; 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); 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); if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->ETCP", pn->data, frag->ll.len);

Loading…
Cancel
Save