Browse Source

congestion: burst-based BW measurement + adaptive inflight (slow start/cong avoidance/loss response)

New features:
- get_time_us() in lib/u_async.c/h for µs-precision timing
- Burst measurement (ETCP_SECTION_MEAS_TS 0x07, MEAS_RESP 0x08, FILLER 0x09)
  - 12 MTU packets sent at line rate, receiver measures inter-packet gaps
  - Burst triggered when link saturated + min 500ms interval
  - Result: measured BW -> bandwidth for shaper, burst_target_bdp for window cap
- 3-phase adaptive inflight in link_stats_timer_cb:
  - Slow start: windows doubles each RTT, exits on loss/RTT spike/threshold
  - Congestion avoidance: additive increase, proportional RTT shrink, multiplicative loss shrink
  - Burst target cap: inflight_lim <= burst_target_bdp * 1.1
- Dynamic optimal_inflight = sum(link->inflight_lim_bytes), recalculated on change
- Adaptive win_timebase = rtt/2 (was hardcoded 10ms)

Fixes:
- dummynet.c: shaper formula fixed (80000 -> 80)
- dummynet.c: queue_data_put call fixed (3 args -> 2)
- etcp.c: SIZE ERROR format string (missing log_name argument)
- etcp_request_pkt: FILLER leaves room for PAYLOAD overhead

Tests:
- test_etcp_congestion.c: single-threaded, 5-phase BW/delay/loss test
  - sender + receiver + dummynet in one uasync
  - CSV metrics output every 100ms
- All 29 existing tests pass
congestion
Evgeny 5 months ago
parent
commit
61fee27478
  1. 27
      lib/u_async.c
  2. 2
      lib/u_async.h
  3. 3
      src/Makefile.am
  4. 4
      src/dummynet.c
  5. 167
      src/etcp.c
  6. 30
      src/etcp.h
  7. 162
      src/etcp_connections.c
  8. 36
      src/etcp_connections.h
  9. 5
      src/etcp_loadbalancer.c
  10. 5
      tests/Makefile.am
  11. 423
      tests/test_etcp_congestion.c
  12. 18
      tests/test_u_async_performance.c

27
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

2
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);

3
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@

4
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);

167
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;

30
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

162
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);

36
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,

5
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;

5
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)

423
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 <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdint.h>
#include <sys/time.h>
#include <unistd.h>
#include <arpa/inet.h>
#include <math.h>
#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;
}

18
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,

Loading…
Cancel
Save