diff --git a/lib/u_async.c b/lib/u_async.c index c2267155..7af7f93c 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -320,12 +320,27 @@ uint64_t get_time_tb(void) { // QueryPerformanceCounter(&count); // Получаем текущее значение счётчика // return (uint64_t)(count.QuadPart * 10000ULL) / (uint64_t)freq.QuadPart; // Преобразуем в требуемые единицы времени //} -#else -uint64_t get_time_tb(void) { - struct timespec ts; - clock_gettime(CLOCK_MONOTONIC, &ts); - return (uint64_t)ts.tv_sec * 10000ULL + (uint64_t)ts.tv_nsec / 100000ULL; // Преобразуем в требуемые единицы времени -} +#else +uint64_t get_time_tb(void) { + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return (uint64_t)ts.tv_sec * 10000ULL + (uint64_t)ts.tv_nsec / 100000ULL; // Преобразуем в требуемые единицы времени +} +#endif + +#ifdef _WIN32 +uint64_t get_time_us(void) { + LARGE_INTEGER freq, count; + QueryPerformanceFrequency(&freq); + QueryPerformanceCounter(&count); + return (uint64_t)(count.QuadPart * 1000000ULL / freq.QuadPart); +} +#else +uint64_t get_time_us(void) { + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return (uint64_t)ts.tv_sec * 1000000ULL + (uint64_t)ts.tv_nsec / 1000ULL; +} #endif diff --git a/lib/u_async.h b/lib/u_async.h index b54268ff..a78b8f8d 100644 --- a/lib/u_async.h +++ b/lib/u_async.h @@ -80,6 +80,8 @@ void uasync_destroy(struct UASYNC* ua, int close_fds); // текущее время (timebase 0.1ms) uint64_t get_time_tb(void); +// текущее время в микросекундах (для burst-измерений) +uint64_t get_time_us(void); // Timeouts, timebase = 0.1 mS void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* user_arg, timeout_callback_t callback, const char* name); diff --git a/src/Makefile.am b/src/Makefile.am index ff86e994..8e68eccc 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -29,7 +29,8 @@ utun_CORE_SOURCES = \ control_server.c \ firewall.c \ eim_nat.c \ - nat_transport.c + nat_transport.c \ + dummynet.c # Platform-specific TUN libs (Windows only) utun_TUN_LIBS = @TUN_LIBS@ diff --git a/src/dummynet.c b/src/dummynet.c index a2fd21bd..18e5f8d1 100644 --- a/src/dummynet.c +++ b/src/dummynet.c @@ -149,7 +149,7 @@ static void dummynet_process_queue(struct dummynet* dn, int dir_idx) { int delay_tb = 0; if (dir->bandwidth_kbps > 0) { - uint64_t delay_calc = (uint64_t)bytes_sent * 80000ULL / (uint64_t)dir->bandwidth_kbps; + uint64_t delay_calc = (uint64_t)bytes_sent * 80ULL / (uint64_t)dir->bandwidth_kbps; delay_tb = (int)delay_calc; if (delay_tb == 0) delay_tb = 1; } @@ -304,7 +304,7 @@ static void dummynet_delay_callback(void* user_arg) { /* Добавляем в очередь */ uint32_t id = (uint32_t)(uintptr_t)entry; - if (queue_data_put(dir->queue, entry, id) != 0) { + if (queue_data_put(dir->queue, entry) != 0) { dir->stats.dropped++; DEBUG_WARN(DEBUG_CATEGORY_DUMMYNET, "Dir %d: queue_data_put failed", dir_idx); queue_entry_free(entry); diff --git a/src/etcp.c b/src/etcp.c index a46c9be3..32639ec9 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -756,7 +756,12 @@ static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {// вызыв // вызывается линком когда освобождается или очередью если появляются данные на передачу struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); - struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp); + struct ETCP_LINK* link = NULL; +// если есть активный burst — используем этот линк напрямую + for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { + if (l->burst_active) { link = l; break; } + } + if (!link) link = etcp_loadbalancer_select_link(etcp); if (!link) { etcp->tx_state=ETCP_TX_STATE_LINK_WAIT; etcp->cnt_link_wait++; @@ -819,7 +824,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { } size_t ack_q_size = queue_entry_count(etcp->ack_q); - if (!inf_pkt && ack_q_size == 0) { + if (!inf_pkt && ack_q_size == 0 && !link->burst_active) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] no data/ack to send", etcp->log_name); return NULL; } @@ -885,6 +890,41 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { dgram->data[1]=(ptr - 8) / 8; + int max_enc = link->mtu - 28 - SC_NONCE_SIZE - SC_TAG_SIZE - SC_CRC32_SIZE - ETCP_ENCRYPTED_HDR_SIZE; + if (max_enc > PACKET_DATA_SIZE) max_enc = PACKET_DATA_SIZE; + if (max_enc < 0) max_enc = 0; + uint8_t burst_flags = 0; int pkt_sz_pos = 0; + +// === MEAS_RESP piggyback (если есть готовый ответ и не активен burst) === + if (link->burst_resp_pending && !link->burst_active && max_enc - ptr >= MEAS_RESP_SECTION_SIZE) { + dgram->data[ptr++] = ETCP_SECTION_MEAS_RESP; + dgram->data[ptr++] = link->burst_recv_id & 0xFF; dgram->data[ptr++] = (link->burst_recv_id >> 8) & 0xFF; + dgram->data[ptr++] = link->burst_resp_valid; + uint32_t v = link->burst_resp_gap_avg; + dgram->data[ptr++] = v & 0xFF; dgram->data[ptr++] = (v >> 8) & 0xFF; dgram->data[ptr++] = (v >> 16) & 0xFF; dgram->data[ptr++] = (v >> 24) & 0xFF; + v = link->burst_resp_gap_min; + dgram->data[ptr++] = v & 0xFF; dgram->data[ptr++] = (v >> 8) & 0xFF; dgram->data[ptr++] = (v >> 16) & 0xFF; dgram->data[ptr++] = (v >> 24) & 0xFF; + dgram->data[ptr++] = link->burst_resp_pkt_count; + link->burst_resp_pending = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX burst resp id=%u valid=%u gap_min=%u", etcp->log_name, link->burst_recv_id, link->burst_resp_valid, link->burst_resp_gap_min); + } + +// === MEAS_TS section (burst) === + if (link->burst_active && max_enc - ptr >= MEAS_TS_SECTION_SIZE) { + burst_flags = 0; + if (!inf_pkt) burst_flags |= MEAS_FLAG_IS_FILLER; + if (link->burst_seq == 0) burst_flags |= MEAS_FLAG_IS_FIRST; + if (link->burst_seq + 1 >= link->burst_count) burst_flags |= MEAS_FLAG_IS_LAST; + dgram->data[ptr++] = ETCP_SECTION_MEAS_TS; + dgram->data[ptr++] = link->burst_id & 0xFF; dgram->data[ptr++] = (link->burst_id >> 8) & 0xFF; + dgram->data[ptr++] = burst_flags; + dgram->data[ptr++] = link->burst_seq; + uint16_t ts_us = (uint16_t)get_time_us(); + dgram->data[ptr++] = ts_us & 0xFF; dgram->data[ptr++] = (ts_us >> 8) & 0xFF; + pkt_sz_pos = ptr; + dgram->data[ptr++] = 0; dgram->data[ptr++] = 0; + } + if (link->last_recv_updated && remain_len>=5) {// если есть данные - добавим channel_timestamp uint64_t now=get_time_tb(); @@ -907,6 +947,19 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { // DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "remain_len(2)= %d", remain_len); if (inf_pkt) { +// === FILLER section before PAYLOAD (burst padding) === + if (link->burst_active) { + int payload_overhead = 1 + 4 + data_len; // type + seq + data + int buf_avail = PACKET_DATA_SIZE - ptr - payload_overhead; + int filler_data = max_enc - ptr - FILLER_HDR_SIZE - payload_overhead; + if (filler_data > buf_avail) filler_data = buf_avail; + if (filler_data > 0) { + dgram->data[ptr++] = ETCP_SECTION_FILLER; + dgram->data[ptr++] = filler_data & 0xFF; + dgram->data[ptr++] = (filler_data >> 8) & 0xFF; + memset(&dgram->data[ptr], 0, filler_data); ptr += filler_data; + } + } // фрейм data (0) обязательно в конец DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX DATA: seq=%u len=%u retry=%d; Qlen: in=%d snd=%d wait=%d", etcp->log_name, inf_pkt->seq, inf_pkt->ll.len, inf_pkt->send_count, etcp->input_queue->count, etcp->input_send_q->count, etcp->input_wait_ack->count); @@ -919,13 +972,34 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len; } else { + if (link->burst_active) { + int buf_avail = PACKET_DATA_SIZE - ptr; + int filler_data = max_enc - ptr - FILLER_HDR_SIZE; + if (filler_data > buf_avail - FILLER_HDR_SIZE) filler_data = buf_avail - FILLER_HDR_SIZE; + if (filler_data > 0) { + dgram->data[ptr++] = ETCP_SECTION_FILLER; + dgram->data[ptr++] = filler_data & 0xFF; + dgram->data[ptr++] = (filler_data >> 8) & 0xFF; + memset(&dgram->data[ptr], 0, filler_data); ptr += filler_data; + } + } int chk=queue_check_consistency(etcp->ack_q); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] only ACK (size=%d) packet with %d bytes total (chk=%d) rem=%d", etcp->log_name, ack_q_size, ptr, chk, remain_len); } - if (ptr>=PACKET_DATA_SIZE-50) DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] SIZE ERROR!!! %d", ptr); + if (ptr>=PACKET_DATA_SIZE-50) DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] SIZE ERROR!!! %d", etcp->log_name, ptr); dgram->data_len=ptr; +// === Fill pkt_sz in MEAS_TS and advance burst state === + if (link->burst_active && pkt_sz_pos > 0) { + uint16_t pkt_sz = (uint16_t)dgram->data_len; + dgram->data[pkt_sz_pos] = pkt_sz & 0xFF; + dgram->data[pkt_sz_pos + 1] = (pkt_sz >> 8) & 0xFF; + link->burst_pkt_size = pkt_sz; + link->burst_seq++; + if (link->burst_seq >= link->burst_count) etcp_link_burst_finish(link); + } + etcp_dump_pkt_sections(dgram, link, 1); return dgram; @@ -1257,6 +1331,93 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) { break; } + case ETCP_SECTION_MEAS_TS: { + if (len < MEAS_TS_SECTION_SIZE) { len = 0; break; } + uint16_t b_id = data[1] | (data[2] << 8); + uint8_t flags = data[3]; + uint8_t seq = data[4]; + uint16_t pkt_sz = data[7] | (data[8] << 8); + struct ETCP_LINK* link = pkt->link; + if (b_id != link->burst_recv_id) { + link->burst_recv_id = b_id; + link->burst_recv_count = BURST_PACKET_COUNT; + link->burst_recv_next_seq = 0; + link->burst_recv_valid = 1; + link->burst_recv_received = 0; + } + if (link->burst_recv_valid) { + if (seq != link->burst_recv_next_seq) link->burst_recv_valid = 0; + if (seq < 16) { + link->burst_recv_times[seq] = get_time_us(); + link->burst_recv_next_seq = seq + 1; + link->burst_recv_received++; + link->burst_recv_pkt_size = pkt_sz; + } + } + if ((flags & MEAS_FLAG_IS_LAST) || link->burst_recv_received >= link->burst_recv_count) { + uint8_t total = (flags & MEAS_FLAG_IS_LAST) ? link->burst_recv_received : link->burst_recv_count; + if (link->burst_recv_valid && total >= BURST_SKIP_COUNT + 2) { + uint32_t sum_gap = 0, min_gap = UINT32_MAX; + uint8_t measured = 0; + for (uint8_t i = BURST_SKIP_COUNT; i + 1 < total; i++) { + uint64_t dt = link->burst_recv_times[i + 1] - link->burst_recv_times[i]; + uint32_t gap = (uint32_t)(dt > 0 ? dt : 1); + sum_gap += gap; measured++; + if (gap < min_gap) min_gap = gap; + } + if (measured > 0) { + link->burst_resp_gap_avg = sum_gap / measured; + link->burst_resp_gap_min = min_gap; + link->burst_resp_pkt_count = measured; + link->burst_resp_valid = 1; + } + } + link->burst_resp_pending = 1; + } + data += MEAS_TS_SECTION_SIZE; len -= MEAS_TS_SECTION_SIZE; + break; + } + + case ETCP_SECTION_MEAS_RESP: { + if (len < MEAS_RESP_SECTION_SIZE) { len = 0; break; } + uint16_t b_id = data[1] | (data[2] << 8); + uint8_t valid = data[3]; + uint32_t gap_avg = data[4] | (data[5] << 8) | ((uint32_t)data[6] << 16) | ((uint32_t)data[7] << 24); + uint32_t gap_min = data[8] | (data[9] << 8) | ((uint32_t)data[10] << 16) | ((uint32_t)data[11] << 24); + uint8_t pkt_cnt = data[12]; + struct ETCP_LINK* link = pkt->link; + if (b_id == link->burst_id - 1 || b_id == link->burst_id) { + if (link->burst_resp_timer) { uasync_cancel_timeout(link->etcp->instance->ua, link->burst_resp_timer); link->burst_resp_timer = NULL; } + if (valid == MEAS_RESP_VALID && gap_min > 0 && link->burst_pkt_size > 0) { + uint64_t bw_kbps = (uint64_t)link->burst_pkt_size * 8000ULL / gap_min; + if (bw_kbps > 10000000) bw_kbps = 10000000; + link->bandwidth = (uint32_t)bw_kbps; + float rtt_sec = (float)link->rtt_min / 10000.f; + if (rtt_sec < 0.005f) rtt_sec = 0.005f; + link->burst_target_bdp = (uint32_t)((float)bw_kbps * 1000.f / 8.f * rtt_sec); + link->slow_start_threshold = link->burst_target_bdp * 9 / 10; + if (link->inflight_phase == INFLIGHT_PHASE_SLOW_START) { + if (link->inflight_lim_bytes > link->slow_start_threshold) { + link->inflight_lim_bytes = link->slow_start_threshold; + etcp_link_update_inflight_lim(link, link->inflight_lim_bytes); + } + } + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] burst resp: BW=%u Kbps, BDP=%u, slow_start_thr=%u gap_min=%u us pkt_cnt=%u", + link->etcp->log_name, (uint32_t)bw_kbps, link->burst_target_bdp, link->slow_start_threshold, gap_min, pkt_cnt); + } + } + data += MEAS_RESP_SECTION_SIZE; len -= MEAS_RESP_SECTION_SIZE; + break; + } + + case ETCP_SECTION_FILLER: { + if (len < FILLER_HDR_SIZE) { len = 0; break; } + uint16_t fill_len = data[1] | (data[2] << 8); + if ((uint16_t)(fill_len + FILLER_HDR_SIZE) > len) { len = 0; break; } + data += FILLER_HDR_SIZE + fill_len; len -= FILLER_HDR_SIZE + fill_len; + break; + } + default: DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_conn_input: unknown section type=0x%02x", type); len=0; diff --git a/src/etcp.h b/src/etcp.h index 8c8379a5..4a438933 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -28,10 +28,32 @@ uint16_t get_current_timestamp(void); #define ETCP_SECTION_PAYLOAD 0x00 // Data payload #define ETCP_SECTION_ACK 0x01 // ACK section #define ETCP_SECTION_TIMESTAMP 0x06 // Channel timestamp (example, adjust if needed) -#define ETCP_SECTION_MEAS_TS 0x07 // Measurement timestamp for bandwidth -#define ETCP_SECTION_MEAS_RESP 0x08 // Measurement response - -// Inflight packet states +#define ETCP_SECTION_MEAS_TS 0x07 // Measurement timestamp for bandwidth (burst packet) +#define ETCP_SECTION_MEAS_RESP 0x08 // Measurement response (burst result) +#define ETCP_SECTION_FILLER 0x09 // Filler/dummy data (discarded by receiver) + +// Burst measurement constants +#define BURST_PACKET_COUNT 12 // число пакетов в burst +#define BURST_SKIP_COUNT 3 // сколько первых пакетов пропустить при замере gap +#define MIN_BURST_INTERVAL_TB 5000 // минимальный интервал между burst (500ms в 0.1ms) +#define BURST_RESP_TIMEOUT_TB 20000 // таймаут ожидания ответа на burst (2s в 0.1ms) + +// Inflight phases +#define INFLIGHT_PHASE_SLOW_START 0 +#define INFLIGHT_PHASE_CONG_AVOIDANCE 1 + +// Section sizes +#define MEAS_TS_SECTION_SIZE 9 // type(1) + burst_id(2) + flags(1) + seq(1) + ts_us(2) + pkt_sz(2) +#define MEAS_RESP_SECTION_SIZE 13 // type(1) + burst_id(2) + valid(1) + gap_avg(4) + gap_min(4) + pkt_count(1) +#define FILLER_HDR_SIZE 3 // type(1) + len(2) + +// MEAS_TS flags +#define MEAS_FLAG_IS_FILLER 0x01 // пакет содержит FILLER вместо PAYLOAD +#define MEAS_FLAG_IS_LAST 0x02 // последний пакет в burst +#define MEAS_FLAG_IS_FIRST 0x04 // первый пакет в burst + +// MEAS_RESP valid +#define MEAS_RESP_VALID 1 #define INFLIGHT_STATE_WAIT_ACK 0 #define INFLIGHT_STATE_WAIT_SEND 1 #define INFLIGHT_INITIAL_HASH_SIZE 1024 diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 2c3c622f..23bf39c9 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -37,7 +37,7 @@ static void etcp_link_init_timer_cbk(void* arg); static void etcp_link_send_keepalive(struct ETCP_LINK* link); static void keepalive_timer_cb(void* arg); static void link_stats_timer_cb(void* arg); - +static void burst_resp_timeout_cb(void* arg); void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) { if (!link) return; @@ -48,6 +48,53 @@ void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) { loadbalancer_link_ready(link); DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_link_update_inflight_lim: unblocked link (lim %u->%u)", old, new_lim); } +// recalc connection-level optimal_inflight + if (link->etcp) { + uint32_t sum = 0; + for (struct ETCP_LINK* l = link->etcp->links; l; l = l->next) sum += l->inflight_lim_bytes; + if (sum < 100000) sum = 100000; + link->etcp->optimal_inflight = sum; + } +} + +// === Burst sender functions === + +void etcp_link_burst_start(struct ETCP_LINK* link) { + if (!link || !link->etcp || !link->etcp->instance) return; + link->burst_active = 1; + link->burst_seq = 0; + link->burst_count = BURST_PACKET_COUNT; + link->burst_id++; + link->burst_last_time_tb = get_time_tb(); + link->burst_pkt_size = 0; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] burst start id=%u count=%u", link->etcp->log_name, link->burst_id, link->burst_count); +// resume send queue чтобы burst-пакеты были немедленно отправлены + if (link->etcp->link_ready_for_send_fn) link->etcp->link_ready_for_send_fn(link->etcp); +} + +void etcp_link_burst_check(struct ETCP_LINK* link) { + if (!link || !link->etcp) return; + if (link->burst_active) return; + if (link->link_status != 1) return; + if (link->inflight_bytes < link->inflight_lim_bytes) return; + uint64_t now = get_time_tb(); + if (now - link->burst_last_time_tb < MIN_BURST_INTERVAL_TB) return; + etcp_link_burst_start(link); +} + +void etcp_link_burst_finish(struct ETCP_LINK* link) { + if (!link || !link->etcp || !link->etcp->instance) return; + link->burst_active = 0; + link->burst_resp_timer = uasync_set_timeout(link->etcp->instance->ua, BURST_RESP_TIMEOUT_TB, link, burst_resp_timeout_cb, "burst_resp"); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] burst end id=%u pkt_sz=%u wait_resp", link->etcp->log_name, link->burst_id, link->burst_pkt_size); + if (link->etcp->link_ready_for_send_fn) link->etcp->link_ready_for_send_fn(link->etcp); +} + +static void burst_resp_timeout_cb(void* arg) { + struct ETCP_LINK* link = (struct ETCP_LINK*)arg; + if (!link) return; + link->burst_resp_timer = NULL; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] burst resp timeout id=%u", link->etcp->log_name, link->burst_id); } #define INIT_TIMEOUT_INITIAL 500 @@ -730,6 +777,14 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn link->keepalive_sent_count = 0; link->keepalive_recv_count = 0; link->inflight_lim_bytes = 30000; + link->inflight_phase = INFLIGHT_PHASE_SLOW_START; + link->slow_start_threshold = 1000000; // будет обновлён после первого burst-замера + link->last_window_update_tb = 0; + link->bandwidth = 10000; // начальная оценка 10 Mbps для шейпера + link->burst_id = 0; + link->burst_active = 0; + link->burst_last_time_tb = 0; + link->burst_target_bdp = 0; // Выделяем свободный local_link_id int free_id = etcp_find_free_local_link_id(etcp); @@ -775,6 +830,12 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn while (l && l->next) l=l->next; if (l) l->next = link; else etcp->links = link; +// пересчитать connection-level optimal_inflight + { uint32_t sum = 0; + for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes; + if (sum < 100000) sum = 100000; + etcp->optimal_inflight = sum; } + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "NEW link initialized on etcp=[%s] link=%p socket=%s id=%d is_server=%d mtu=%d", etcp->log_name, link, conn->name, link->local_link_id, link->is_server, link->mtu); if (is_server == 0) { @@ -867,30 +928,83 @@ static void link_stats_timer_cb(void* arg) { swm_add(link->rtt_swm, link->rtt_avg10); link->rtt_min=swm_get_min(link->rtt_swm); -// пересчитываем BW - float rtt=(float)link->rtt_min/10000.f;// переводим в секунды - if (rtt<0.005f) rtt=0.005f;// 5 ms - if (rtt>1.f) rtt=1.f;// 1 sec - float win_size=link->inflight_lim_bytes*8.f;// переводим в биты - if (win_size<100000.f) win_size=100000.f;// 10kb min - if (win_size>10000000.f) win_size=10000000.f;// 1Mb max - link->bandwidth=(int)(win_size*8.f/rtt/1024.f*1.3f);// результат в кБит/сек - - if (link->inflight_lim_bytes<10000 || link->inflight_bytes > link->inflight_lim_bytes - 2000) {// очередь загружена - int new_lim=link->inflight_lim_bytes; - if (link->rtt_avg10*10 < link->rtt_min*16) {// подъема rtt нет, плавно увеличиваем win - new_lim+=new_lim/64+1; - if (new_lim>1000000) new_lim=1000000; +// === Burst check: запускаем если линк насыщен и прошло время === + etcp_link_burst_check(link); + + uint64_t now_tb = get_time_tb(); + int do_update = 0; +// адаптивно: обновляем окно примерно раз в RTT + if (now_tb - link->last_window_update_tb >= link->rtt_avg10 || link->last_window_update_tb == 0) { + do_update = 1; + link->last_window_update_tb = now_tb; + } + + if (do_update && link->initialized) { + int new_lim = link->inflight_lim_bytes; + uint32_t loss = link->window_retransmissions; + uint32_t sent = link->window_pkt_transmitted; + uint8_t old_phase = link->inflight_phase; + +// === Phase 0: Slow Start === + if (link->inflight_phase == INFLIGHT_PHASE_SLOW_START) { + if (link->slow_start_threshold == 0) link->slow_start_threshold = 30000; + if (loss > 0) { + link->inflight_phase = INFLIGHT_PHASE_CONG_AVOIDANCE; // первая потеря — выход + new_lim -= new_lim / 2; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] slow_start->cong_avoid (loss: %u) lim %u", link->etcp->log_name, loss, new_lim); + } else if (link->rtt_avg10 > link->rtt_min * 15 / 10) { + link->inflight_phase = INFLIGHT_PHASE_CONG_AVOIDANCE; // RTT вырос >50% + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] slow_start->cong_avoid (RTT spike: avg=%u min=%u) lim %u", link->etcp->log_name, link->rtt_avg10, link->rtt_min, new_lim); + } else if (new_lim >= link->slow_start_threshold) { + link->inflight_phase = INFLIGHT_PHASE_CONG_AVOIDANCE; // достигли порога burst + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] slow_start->cong_avoid (threshold %u reached) lim %u", link->etcp->log_name, link->slow_start_threshold, new_lim); + } else { + new_lim *= 2; // удвоение каждый RTT + } + } + +// === Phase 1: Congestion Avoidance === + if (link->inflight_phase == INFLIGHT_PHASE_CONG_AVOIDANCE) { + if (old_phase == INFLIGHT_PHASE_SLOW_START) { +// только что вышли из slow start — не трогаем лимит в этом тике + } +// Loss-based: мультипликативное снижение + else if (loss > 0 && sent > 0) { + int loss_pct = loss * 100 / sent; + int decrease = new_lim * loss_pct / 50; // loss_pct / 50 ≈ ×2 при 1% потерь + if (decrease < link->mtu) decrease = link->mtu; + if (decrease > new_lim / 2) decrease = new_lim / 2; + new_lim -= decrease; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] cong_avoid loss=%u/%u (%d%%) decrease=%d lim %u", + link->etcp->log_name, loss, sent, loss_pct, decrease, new_lim); } - else { - new_lim-=new_lim/64+1; - if (new_lim<10000) new_lim=10000; +// RTT-based: пропорциональное снижение + else if (link->rtt_avg10 > link->rtt_min * 13 / 10) { + int excess = link->rtt_avg10 - link->rtt_min * 13 / 10; + int ratio = excess * 100 / link->rtt_avg10; // 0..100 + int decrease = new_lim * ratio / 200; // ratio/200 ≈ до 50% + if (decrease < link->mtu) decrease = link->mtu; + if (decrease > new_lim / 4) decrease = new_lim / 4; + new_lim -= decrease; + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] cong_avoid RTT spike avg=%u min=%u decrease=%d lim %u", + link->etcp->log_name, link->rtt_avg10, link->rtt_min, decrease, new_lim); + } +// Additive increase (если нет сигналов перегрузки) + else { + new_lim += link->mtu; // +1 пакет за RTT } - etcp_link_update_inflight_lim(link, new_lim); } +// Clamp + if (new_lim < (int)link->mtu * 2) new_lim = link->mtu * 2; + if (new_lim > 1000000) new_lim = 1000000; + +// Burst target cap (если есть burst_target_bdp — не превышаем 1.1×) + if (link->burst_target_bdp > 0 && new_lim > (int)(link->burst_target_bdp * 11 / 10)) + new_lim = link->burst_target_bdp * 11 / 10; -// если окно забито - пробуем его расширить + if (new_lim != (int)link->inflight_lim_bytes) etcp_link_update_inflight_lim(link, (uint32_t)new_lim); + } // DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] stats window updated (win_timebase=%u us, rtt=%u, retrans=%u, transmitted=%u)", @@ -904,11 +1018,11 @@ static void link_stats_timer_cb(void* arg) { link->window_retransmissions = 0; // 4. Плавная подстройка win_timebase под rtt/2 (в микросекундах) -// uint32_t target_us = (uint32_t)link->rtt_avg10 * 50ULL; // rtt_avg10 (0.1 ms) → rtt/2 в us -// if (target_us < 10000) target_us = 10000; // минимум 10 ms -// if (target_us > 500000) target_us = 500000; // максимум 0.5 s + uint32_t target_us = (uint32_t)link->rtt_avg10 * 50ULL; // rtt_avg10 (0.1 ms) → rtt/2 в us + if (target_us < 10000) target_us = 10000; // минимум 10 ms + if (target_us > 500000) target_us = 500000; // максимум 0.5 s - link->win_timebase = 10000;// (link->win_timebase * 7 + target_us) / 8; + link->win_timebase = (link->win_timebase * 7 + target_us) / 8; // 5. Перезапускаем таймер с новым интервалом start_stats_timer(link); diff --git a/src/etcp_connections.h b/src/etcp_connections.h index 53e44b92..a059770c 100644 --- a/src/etcp_connections.h +++ b/src/etcp_connections.h @@ -214,6 +214,37 @@ struct ETCP_LINK { uint32_t keepalive_recv_count; // Счётчик полученных keepalive uint16_t handshake_minsize; // минимальный размер udp при handshake uint16_t handshake_maxsize; // мax размер udp при handshake (выбирает рандом) + + // === Burst-измерение bandwidth === + // Sender-сторона + uint16_t burst_id; // монотонно возрастающий ID burst + uint8_t burst_active; // 1 = burst в процессе отправки + uint8_t burst_seq; // текущий номер пакета в burst (0..burst_count-1) + uint8_t burst_count; // общее число пакетов в burst + uint64_t burst_last_time_tb; // время последнего burst (0.1ms), для MIN_BURST_INTERVAL + void* burst_resp_timer; // таймер ожидания ответа на burst + uint32_t burst_target_bdp; // BDP из последнего успешного burst (байты) + uint16_t burst_pkt_size; // сохранённый размер пакета для вычисления BW + + // Receiver-сторона + uint16_t burst_recv_id; // ID отслеживаемого burst + uint8_t burst_recv_count; // ожидаемое число пакетов в burst + uint8_t burst_recv_next_seq; // следующий ожидаемый seq + uint8_t burst_recv_valid; // 1 = порядок пакетов не нарушен + uint8_t burst_recv_received; // сколько пакетов уже принято + uint64_t burst_recv_times[16]; // время прихода каждого пакета (µs) + uint16_t burst_recv_pkt_size; // размер пакета из MEAS_TS (для ответа) + uint8_t burst_resp_pending; // 1 = есть готовый ответ для piggyback + uint32_t burst_resp_gap_avg; // средний inter-packet gap (µs) + uint32_t burst_resp_gap_min; // минимальный inter-packet gap (µs) + uint8_t burst_resp_pkt_count; // число пакетов в измерении + uint8_t burst_resp_valid; // валидность измерения + + // === Фазы управления inflight === + // 0 = slow_start, 1 = congestion_avoidance + uint8_t inflight_phase; + uint64_t last_window_update_tb; // время последнего обновления окна (0.1ms) + uint32_t slow_start_threshold; // порог выхода из slow start (burst_target_bdp * 0.9) }; // INITIALIZATION @@ -246,6 +277,11 @@ void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim); int etcp_find_free_local_link_id(struct ETCP_CONN* etcp); void start_stats_timer(struct ETCP_LINK* link); +// Burst measurement functions +void etcp_link_burst_start(struct ETCP_LINK* link); +void etcp_link_burst_check(struct ETCP_LINK* link); +void etcp_link_burst_finish(struct ETCP_LINK* link); + int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bin, const struct sockaddr_storage* addr, int timeout_ms, etcp_ping_callback_t cb, void* user_arg, diff --git a/src/etcp_loadbalancer.c b/src/etcp_loadbalancer.c index 28e370f1..0f91da72 100644 --- a/src/etcp_loadbalancer.c +++ b/src/etcp_loadbalancer.c @@ -141,7 +141,9 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) { } // Update shaper after successful send - if (link->bandwidth == 0) { + if (link->burst_active) { + // burst bypasses shaper — не обновляем нагрузку + } else if (link->bandwidth == 0) { DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] unlimited bandwidth, skipping shaper update", etcp->log_name); // Don't return here, let the function continue to free dgram at the end } else { @@ -178,6 +180,7 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) { int loadbalancer_link_can_send(struct ETCP_LINK* link) { if (!link) return 0; + if (link->burst_active) return 1; // burst bypasses all limits int can = 1; if (link->inflight_bytes >= link->inflight_lim_bytes) { link->send_blocked_inflight = 1; diff --git a/tests/Makefile.am b/tests/Makefile.am index 4bd2cc4c..a706d9b2 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -5,6 +5,7 @@ check_PROGRAMS = \ test_etcp_crypto \ test_etcp_two_instances \ test_etcp_simple_traffic \ + test_etcp_congestion \ test_etcp_minimal \ test_etcp_100_packets \ test_pkt_normalizer_etcp \ @@ -160,6 +161,10 @@ test_etcp_simple_traffic_SOURCES = test_etcp_simple_traffic.c test_etcp_simple_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source test_etcp_simple_traffic_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_congestion_SOURCES = test_etcp_congestion.c +test_etcp_congestion_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source +test_etcp_congestion_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_etcp_minimal_SOURCES = test_etcp_minimal.c test_etcp_minimal_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source test_etcp_minimal_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_congestion.c b/tests/test_etcp_congestion.c new file mode 100644 index 00000000..33a6e419 --- /dev/null +++ b/tests/test_etcp_congestion.c @@ -0,0 +1,423 @@ +/** + * @file test_etcp_congestion.c + * @brief Тест адаптивного управления inflight: burst + 3-фазный алгоритм + * + * Архитектура (однопоточная, один uasync): + * Sender(client) ──Link1── Dummynet1 ──Link1── Receiver(server) + * ──Link2── Dummynet2 ──Link2── + * + * Фазы теста: изменяемые BW/delay/loss на dummynet, метрики каждые 100ms в CSV. + */ + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/memory_pool.h" +#include "../lib/debug_config.h" +#include "../lib/platform_compat.h" +#include "../lib/mem.h" +#include "../src/dummynet.h" +#include "../src/config_parser.h" +#include "../src/utun_instance.h" +#include "../src/etcp.h" +#include "../src/etcp_api.h" +#include "../src/etcp_connections.h" +#include "../src/secure_channel.h" +#include "../src/config_updater.h" +#include "../src/routing.h" +#include "../src/crc32.h" + +/* ===== Конфигурация теста ===== */ +#define TEST_DURATION_MS 25000 +#define METRICS_INTERVAL_TB 1000 // 100ms в 0.1ms units +#define SEND_TIMER_TB 1 // 0.1ms — минимальный re-schedule при заполнении +#define PAYLOAD_SIZE 1200 +#define PAYLOAD_PAD 100 // pad для large MTU burst + +#define DN1_PORT 31001 +#define SRV1_PORT 31002 +#define CLI1_PORT 31000 +#define DN2_PORT 32001 +#define SRV2_PORT 32002 +#define CLI2_PORT 32000 + +#define PHASE_MAX 8 + +/* ===== Структуры ===== */ +struct test_phase { + uint32_t dur_ms; + uint32_t bw_kbps; + uint32_t delay_ms; + uint32_t loss_permille; +}; + +struct test_ctx { + struct UASYNC* ua; + struct UTUN_INSTANCE* sender; + struct UTUN_INSTANCE* receiver; + struct dummynet* dn[2]; + + struct test_phase phases[PHASE_MAX]; + int phase_count; + int current_phase; + void* phase_timer; + void* metrics_timer; + + FILE* csv; + uint64_t start_time_us; + uint64_t bytes_received; + uint64_t last_bytes_recv; + uint64_t bytes_sent; + int test_done; + int sending_active; +}; + +/* ===== Helpers ===== */ +static uint64_t now_us(void) { + struct timeval tv; + gettimeofday(&tv, NULL); + return (uint64_t)tv.tv_sec * 1000000ULL + tv.tv_usec; +} + +/* ===== Forward declarations ===== */ +static void metrics_timer_cb(void* arg); +static void phase_timer_cb(void* arg); +static void send_timer_cb(void* arg); +static void on_recv(struct ETCP_CONN* conn, struct ll_entry* entry); +static void restart_send(struct test_ctx* ctx); + +/* ===== Создание инстансов ===== */ +static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id, + const char* priv_hex, const char* pub_hex) { + struct UTUN_INSTANCE* inst = u_calloc(1, sizeof(*inst)); + if (!inst) return NULL; + inst->ua = u; + inst->node_id = node_id; + if (sc_init_local_keys(&inst->my_keys, pub_hex, priv_hex) != SC_OK) { u_free(inst); return NULL; } + inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET)); + inst->data_pool = memory_pool_init(PACKET_DATA_SIZE); + inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE); + if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) { u_free(inst); return NULL; } + struct utun_config* cfg = u_calloc(1, sizeof(*cfg)); + if (!cfg) { u_free(inst); return NULL; } + strncpy(cfg->global.my_public_key_hex, pub_hex, MAX_KEY_LEN - 1); + strncpy(cfg->global.my_private_key_hex, priv_hex, MAX_KEY_LEN - 1); + cfg->global.my_node_id = node_id; + cfg->global.mtu = 1400; + cfg->global.keepalive_timeout = 5000; + cfg->global.keepalive_interval = 500; + cfg->global.allowed_keys_allow_all = 1; + inst->config = cfg; + return inst; +} + +static int add_server(struct UTUN_INSTANCE* inst, const char* name, int port) { + struct CFG_SERVER* srv = u_calloc(1, sizeof(*srv)); + if (!srv) return -1; + strncpy(srv->name, name, MAX_CONN_NAME_LEN - 1); + srv->ip.ss_family = AF_INET; + ((struct sockaddr_in*)&srv->ip)->sin_addr.s_addr = inet_addr("127.0.0.1"); + ((struct sockaddr_in*)&srv->ip)->sin_port = htons(port); + srv->type = CFG_SERVER_TYPE_PUBLIC; + srv->next = inst->config->servers; + inst->config->servers = srv; + return 0; +} + +static int add_client(struct UTUN_INSTANCE* inst, const char* peer_pubkey) { + struct CFG_CLIENT* cli = u_calloc(1, sizeof(*cli)); + if (!cli) return -1; + strncpy(cli->name, "peer", MAX_CONN_NAME_LEN - 1); + strncpy(cli->peer_public_key_hex, peer_pubkey, MAX_KEY_LEN - 1); + cli->keepalive = 1; + cli->next = inst->config->clients; + inst->config->clients = cli; + return 0; +} + +static struct CFG_CLIENT_LINK* add_link_to_client(struct CFG_CLIENT* cli, struct CFG_SERVER* srv, int remote_port) { + struct CFG_CLIENT_LINK* link = u_calloc(1, sizeof(*link)); + if (!link) return NULL; + link->remote_addr.ss_family = AF_INET; + ((struct sockaddr_in*)&link->remote_addr)->sin_addr.s_addr = inet_addr("127.0.0.1"); + ((struct sockaddr_in*)&link->remote_addr)->sin_port = htons(remote_port); + link->local_srv = srv; + struct CFG_CLIENT_LINK** tail = &cli->links; + while (*tail) tail = &(*tail)->next; + *tail = link; + return link; +} + +/* ===== Dummynet helpers ===== */ +static void dummynet_set_both(struct dummynet* dn, uint32_t bw_kbps, uint32_t delay_ms, + uint32_t loss, int forward_port, int backward_port) { + dummynet_set_direction(dn, DUMMYNET_FORWARD, delay_ms, delay_ms / 4, + bw_kbps, 200, loss, "127.0.0.1", forward_port); + dummynet_set_direction(dn, DUMMYNET_BACKWARD, delay_ms, delay_ms / 4, + bw_kbps, 200, loss, "127.0.0.1", backward_port); +} + +/* ===== Приём данных ===== */ +static void on_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { + struct test_ctx* ctx = (struct test_ctx*)conn->instance->etcp_new_conn_arg; + if (entry) { + ctx->bytes_received += entry->len; + queue_entry_free(entry); + } +} + +/* ===== Отправка на максимальной скорости ===== */ +static void send_burst(struct test_ctx* ctx); + +static void send_timer_cb(void* arg) { + struct test_ctx* ctx = (struct test_ctx*)arg; + send_burst(ctx); +} + +static void send_burst(struct test_ctx* ctx) { + if (ctx->test_done) return; + if (!ctx->sender || !ctx->sender->connections) return; + struct ETCP_CONN* conn = ctx->sender->connections; + if (!conn->initialized) return; + + int sent = 0; + while (sent < 64) { + struct ll_entry* e = ll_alloc_lldgram(sizeof(uint8_t) + PAYLOAD_SIZE); + if (!e) break; + e->dgram[0] = ETCP_ID_DATA; + e->len = 1 + PAYLOAD_SIZE; + if (etcp_send(conn, e) == 0) { ctx->bytes_sent += e->len - 1; sent++; } + else { queue_entry_free(e); break; } + } + if (sent > 0 || !ctx->test_done) { +// если очередь не пуста — планируем немедленный re-schedule + ctx->sending_active = 1; + uasync_set_timeout(ctx->ua, SEND_TIMER_TB, ctx, send_timer_cb, "cong_send"); + } +} + +static void restart_send(struct test_ctx* ctx) { + ctx->sending_active = 1; + send_burst(ctx); +} + +/* ===== Метрики ===== */ +static void csv_header(FILE* f) { + fprintf(f, "time_ms,phase,link,inflight_bytes,inflight_lim,rtt_avg10,rtt_min_us," + "bandwidth_kbps,phase_name,burst_active,burst_bdp," + "retrans,tx_bytes,thru_kbps,optimal_inflight,unacked," + "dn_queue,dn_sent,dn_dropped,dn_lost\n"); +} + +static const char* phase_name(uint8_t p) { + return p == INFLIGHT_PHASE_SLOW_START ? "slow_start" : "cong_avoid"; +} + +static void log_metrics(struct test_ctx* ctx) { + if (!ctx->csv) return; + uint64_t now = now_us(); + double time_ms = (double)(now - ctx->start_time_us) / 1000.0; + double dt_s = 0.1; + uint64_t thru = (uint64_t)((double)(ctx->bytes_received - ctx->last_bytes_recv) * 8.0 / dt_s / 1000.0); + ctx->last_bytes_recv = ctx->bytes_received; + + if (!ctx->sender || !ctx->sender->connections) return; + struct ETCP_CONN* sconn = ctx->sender->connections; + struct ETCP_LINK* link = sconn->links; + int idx = 0; + while (link && idx < 2) { + idx++; + fprintf(ctx->csv, "%.1f,%d,%d,%u,%u,%u,%u,%u,%s,%u,%u,%u,%u,%lu,%u,%u,%d,%lu,%lu,%lu\n", + time_ms, ctx->current_phase, idx, + link->inflight_bytes, link->inflight_lim_bytes, + link->rtt_avg10, (uint32_t)(link->rtt_min * 100U), + link->bandwidth, phase_name(link->inflight_phase), + link->burst_active, link->burst_target_bdp, + link->window_retransmissions, link->window_pkt_transmitted, + (unsigned long)thru, + sconn->optimal_inflight, sconn->unacked_bytes, + dummynet_get_queue_size(ctx->dn[idx - 1], DUMMYNET_FORWARD), + (unsigned long)ctx->dn[idx - 1] ? dummynet_get_stats(ctx->dn[idx - 1], DUMMYNET_FORWARD)->sent : 0, + (unsigned long)ctx->dn[idx - 1] ? dummynet_get_stats(ctx->dn[idx - 1], DUMMYNET_FORWARD)->dropped : 0, + (unsigned long)ctx->dn[idx - 1] ? dummynet_get_stats(ctx->dn[idx - 1], DUMMYNET_FORWARD)->lost : 0); + link = link->next; + } +} + +static void metrics_timer_cb(void* arg) { + struct test_ctx* ctx = (struct test_ctx*)arg; + if (ctx->test_done) return; + log_metrics(ctx); + ctx->metrics_timer = uasync_set_timeout(ctx->ua, METRICS_INTERVAL_TB, ctx, metrics_timer_cb, "cong_metrics"); +} + +/* ===== Фазы ===== */ +static void apply_phase(struct test_ctx* ctx, struct test_phase* ph) { + dummynet_set_both(ctx->dn[0], ph->bw_kbps, ph->delay_ms, ph->loss_permille, SRV1_PORT, CLI1_PORT); + if (ctx->csv) fprintf(ctx->csv, "# phase %d: bw=%u delay=%u loss=%u.%u%%\n", + ctx->current_phase, ph->bw_kbps, ph->delay_ms, + ph->loss_permille / 10, ph->loss_permille % 10); + fflush(ctx->csv); + printf(" Phase %d: bw=%u delay=%u loss=%u.%u%%\n", + ctx->current_phase, ph->bw_kbps, ph->delay_ms, + ph->loss_permille / 10, ph->loss_permille % 10); +} + +static void phase_timer_cb(void* arg) { + struct test_ctx* ctx = (struct test_ctx*)arg; + if (ctx->current_phase >= ctx->phase_count) { ctx->test_done = 1; return; } + apply_phase(ctx, &ctx->phases[ctx->current_phase]); + ctx->current_phase++; + uint32_t next_dur = (ctx->current_phase < ctx->phase_count) ? ctx->phases[ctx->current_phase].dur_ms : 0; + if (next_dur > 0) { + ctx->phase_timer = uasync_set_timeout(ctx->ua, next_dur * 10, ctx, phase_timer_cb, "cong_phase"); + } +} + +/* ===== Main ===== */ +int main(void) { + printf("=== ETCP Congestion Control Test ===\n\n"); + srand((unsigned)time(NULL)); + debug_config_init(); + debug_set_level(DEBUG_LEVEL_WARN); + socket_platform_init(); + crc32_init(); + + struct test_ctx ctx; + memset(&ctx, 0, sizeof(ctx)); + + ctx.ua = uasync_create(); + if (!ctx.ua) { printf("uasync failed\n"); return 1; } + + /* Ключи */ + const char* s_priv = "67b705a92b41bcaae105af2d6a17743faa7b26ccebba8b3b9b0af05e9cd1d5fb"; + const char* s_pub = "1c55e4ccae7c4470707759086738b10681bf88b81f198cc2ab54a647d1556e17c65e6b1833e0c771e5a39382c03067c388915a4c732191bc130480f20f8e00b9"; + const char* c_priv = "4813d31d28b7e9829247f488c6be7672f2bdf61b2508333128e386d1759afed2"; + const char* c_pub = "c594f33c91f3a2222795c2c110c527bf214ad1009197ce14556cb13df3c461b3c373bed8f205a8dd1fc0c364f90bf471d7c6f5db49564c33e4235d268569ac71"; + + printf("Creating instances...\n"); + ctx.sender = create_instance(ctx.ua, 0x1111111111111111ULL, c_priv, c_pub); + ctx.receiver = create_instance(ctx.ua, 0x2222222222222222ULL, s_priv, s_pub); + if (!ctx.sender || !ctx.receiver) { printf("Instance failed\n"); return 1; } + ctx.sender->etcp_new_conn_arg = &ctx; + ctx.receiver->etcp_new_conn_arg = &ctx; + + /* Server side: один слушающий сокет */ + if (add_server(ctx.receiver, "srv1", SRV1_PORT) < 0) { + printf("Server config failed\n"); return 1; + } + + /* Client side: один локальный сокет + один link */ + if (add_server(ctx.sender, "cli1", CLI1_PORT) < 0) { + printf("Client server config failed\n"); return 1; + } + if (add_client(ctx.sender, s_pub) < 0) { printf("Client config failed\n"); return 1; } + struct CFG_CLIENT* cli = ctx.sender->config->clients; + add_link_to_client(cli, ctx.sender->config->servers, DN1_PORT); + + printf("Init receiver...\n"); + if (utun_instance_init(ctx.receiver) < 0) { printf("Receiver init failed\n"); return 1; } + printf("Init sender...\n"); + if (utun_instance_init(ctx.sender) < 0) { printf("Sender init failed\n"); return 1; } + + etcp_bind(ctx.receiver, ETCP_ID_DATA, on_recv); + etcp_set_new_conn_cbk(ctx.receiver, NULL, &ctx); + + printf("Creating dummynet...\n"); + ctx.dn[0] = dummynet_create(ctx.ua, "127.0.0.1", DN1_PORT); + if (!ctx.dn[0]) { printf("Dummynet failed\n"); return 1; } + dummynet_set_both(ctx.dn[0], 100000, 5, 0, SRV1_PORT, CLI1_PORT); + + /* Ожидание handshake — один линк */ + printf("Waiting for connection...\n"); + uint64_t t0 = now_us(); + int links_ready = 0; + while ((now_us() - t0) < 10000000ULL) { + uasync_poll(ctx.ua, 1); + if (ctx.sender->connections && ctx.sender->connections->links) { + struct ETCP_LINK* l = ctx.sender->connections->links; + if (l->initialized && l->link_status == 1) { links_ready = 1; break; } + } + } + printf("Links ready, waiting stabilize...\n"); + t0 = now_us(); + while ((now_us() - t0) < 500000ULL) uasync_poll(ctx.ua, 1); + + /* Фазы */ + ctx.phase_count = 5; + ctx.phases[0] = (struct test_phase){5000, 100000, 5, 0}; + ctx.phases[1] = (struct test_phase){5000, 10000, 5, 0}; + ctx.phases[2] = (struct test_phase){5000, 100000, 5, 0}; + ctx.phases[3] = (struct test_phase){5000, 100000, 50, 0}; + ctx.phases[4] = (struct test_phase){5000, 10000, 50, 10}; + + /* CSV */ + ctx.csv = fopen("test_etcp_congestion.csv", "w"); + if (ctx.csv) csv_header(ctx.csv); + + printf("\n=== Starting test (%d phases, %d ms) ===\n", ctx.phase_count, TEST_DURATION_MS); + ctx.start_time_us = now_us(); + ctx.last_bytes_recv = 0; + + /* Запускаем таймеры */ + uasync_set_timeout(ctx.ua, METRICS_INTERVAL_TB, &ctx, metrics_timer_cb, "cong_metrics"); + uasync_set_timeout(ctx.ua, ctx.phases[0].dur_ms * 10, &ctx, phase_timer_cb, "cong_phase"); + ctx.current_phase = 0; + apply_phase(&ctx, &ctx.phases[0]); ctx.current_phase++; + restart_send(&ctx); + + /* Главный цикл */ + uint64_t last_progress = 0; + while (!ctx.test_done) { + uasync_poll(ctx.ua, 10); + uint64_t elapsed = (now_us() - ctx.start_time_us) / 1000; + if (elapsed >= TEST_DURATION_MS) ctx.test_done = 1; + if (elapsed - last_progress >= 1000) { + last_progress = elapsed; + int l1 = ctx.sender->connections && ctx.sender->connections->links ? 1 : 0; + printf(" t=%lus sent=%lu recv=%lu links=%d\n", (unsigned long)elapsed, (unsigned long)ctx.bytes_sent, (unsigned long)ctx.bytes_received, l1); + } + } + + /* Досылаем и ждём доставки */ + printf("\nWaiting for drain...\n"); + ctx.sending_active = 0; + t0 = now_us(); + while ((now_us() - t0) < 1000000ULL) uasync_poll(ctx.ua, 10); + log_metrics(&ctx); + + /* Итоги */ + double dur_s = (double)(now_us() - ctx.start_time_us) / 1000000.0; + printf("\n=== Results ===\n"); + printf("Duration: %.1f s\n", dur_s); + printf("Bytes sent: %lu\n", (unsigned long)ctx.bytes_sent); + printf("Bytes recv: %lu\n", (unsigned long)ctx.bytes_received); + printf("Throughput: %.1f Kbps\n", ctx.bytes_received * 8.0 / dur_s / 1000.0); + if (ctx.bytes_sent > 0) printf("Loss rate: %.2f%%\n", 100.0 * (ctx.bytes_sent - ctx.bytes_received) / ctx.bytes_sent); + for (int i = 0; i < 1; i++) { + const struct dummynet_stats* st = dummynet_get_stats(ctx.dn[i], DUMMYNET_FORWARD); + if (st) printf("DN%d: recv=%lu sent=%lu lost=%lu dropped=%lu\n", i + 1, + (unsigned long)st->recv, (unsigned long)st->sent, + (unsigned long)st->lost, (unsigned long)st->dropped); + } + + int pass = (ctx.bytes_sent > 0 && ctx.bytes_received > 1000); // проверяем что хоть что-то передалось + printf("\nCSV written to test_etcp_congestion.csv\n"); + printf("[%s]\n", pass ? "PASS" : "FAIL"); + + /* Cleanup */ + if (ctx.csv) fclose(ctx.csv); + dummynet_destroy(ctx.dn[0]); + utun_instance_destroy(ctx.sender); + utun_instance_destroy(ctx.receiver); + uasync_destroy(ctx.ua, 1); + return pass ? 0 : 1; +} diff --git a/tests/test_u_async_performance.c b/tests/test_u_async_performance.c index 79b85ad1..6ba4e9c8 100644 --- a/tests/test_u_async_performance.c +++ b/tests/test_u_async_performance.c @@ -19,7 +19,7 @@ #include "../lib/mem.h" /* Performance measurement */ -static uint64_t get_time_us(void) { +static uint64_t perf_get_time_us(void) { struct timeval tv; utun_gettimeofday(&tv, NULL); return (uint64_t)tv.tv_sec * 1000000ULL + tv.tv_usec; @@ -86,7 +86,7 @@ static void test_socket_callback(int fd, void* arg) { printf("Created %d sockets\n", num_sockets); /* Benchmark 1: Add all sockets */ - uint64_t start_time = get_time_us(); + uint64_t start_time = perf_get_time_us(); int sockets_added = 0; for (int i = 0; i < num_sockets; i++) { void* id = uasync_add_socket(ua, sockets[i], test_socket_callback, NULL, NULL, NULL); @@ -98,7 +98,7 @@ static void test_socket_callback(int fd, void* arg) { socket_fds[i] = sockets[i]; // Store the file descriptor instead of pointer sockets_added++; } - uint64_t add_time = get_time_us() - start_time; + uint64_t add_time = perf_get_time_us() - start_time; printf("DEBUG: Total sockets added: %d\n", sockets_added); printf("Add %d sockets: %llu us (%.2f us per socket)\n", @@ -111,7 +111,7 @@ static void test_socket_callback(int fd, void* arg) { printf("SKIPPING POLLING to test corruption\n"); /* Benchmark 3: Remove all sockets using lookup function */ - start_time = get_time_us(); + start_time = perf_get_time_us(); int removed_count = 0; int failed_count = 0; @@ -136,7 +136,7 @@ static void test_socket_callback(int fd, void* arg) { failed_count++; } } - uint64_t remove_time = get_time_us() - start_time; + uint64_t remove_time = perf_get_time_us() - start_time; printf("DEBUG: Actually removed %d sockets, failed %d\n", removed_count, failed_count); @@ -191,7 +191,7 @@ static void benchmark_high_frequency(void) { int cycles = 10000; printf("Testing %d rapid add/remove cycles...\n", cycles); - uint64_t start_time = get_time_us(); + uint64_t start_time = perf_get_time_us(); for (int cycle = 0; cycle < cycles; cycle++) { /* Remove all sockets */ for (int i = 0; i < num_sockets; i++) { @@ -209,7 +209,7 @@ static void benchmark_high_frequency(void) { /* Poll once per cycle */ uasync_poll(ua, 0); } - uint64_t total_time = get_time_us() - start_time; + uint64_t total_time = perf_get_time_us() - start_time; printf("Completed %d cycles in %llu us\n", cycles, (unsigned long long)total_time); printf("Average time per cycle: %.2f us\n", (double)total_time / cycles); @@ -250,11 +250,11 @@ static void benchmark_scalability(void) { } /* Measure poll time */ - uint64_t start_time = get_time_us(); + uint64_t start_time = perf_get_time_us(); for (int iter = 0; iter < 100; iter++) { uasync_poll(ua, 0); } - uint64_t poll_time = get_time_us() - start_time; + uint64_t poll_time = perf_get_time_us() - start_time; printf(" %d sockets: %.2f us per poll (%.2f ns per socket)\n", num_sockets, (double)poll_time / 100,