Browse Source

etcp: BBR callback для обновления inflight_lim_bytes + пересчёт optimal_inflight

- BBR: добавлен callback on_cwnd_update(ctx, new_cwnd) в struct bbr, вызывается при изменении cwnd
- etcp_link_update_inflight_lim() — атомарно обновляет inflight_lim_bytes линка и пересчитывает optimal_inflight (сумма по всем линкам)
- BBR теперь сам обновляет лимиты через callback, ручное присваивание в etcp_ack_recv() удалено
- is_server в config_parser.c: печать адресов через sockaddr_storage_to_str (поддержка IPv4/IPv6)
- utun_instance.c: local_sockaddr_equal() — добавлено сравнение AF_INET6
- Убраны отладочные memory_pool_is_freed
tmo
Evgeny 4 months ago
parent
commit
52f32da3f1
  1. 12
      src/config_parser.c
  2. 11
      src/etcp.c
  3. 3
      src/etcp_bbr.c
  4. 4
      src/etcp_bbr.h
  5. 26
      src/etcp_connections.c
  6. 1
      src/etcp_connections.h
  7. 7
      src/utun_instance.c

12
src/config_parser.c

@ -986,9 +986,9 @@ void print_config(const struct utun_config *cfg) {
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Servers:"); DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Servers:");
struct CFG_SERVER *s = cfg->servers; struct CFG_SERVER *s = cfg->servers;
while (s) { while (s) {
struct sockaddr_in *sin = (struct sockaddr_in *)&s->ip; ip_str_t addr_str = sockaddr_storage_to_str(&s->ip);
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, " %s: %s:%d (mark=%d, netif=%u, type=%u, mtu=%d)", DEBUG_INFO(DEBUG_CATEGORY_CONFIG, " %s: %s (mark=%d, netif=%u, type=%u, mtu=%d)",
s->name, ip_to_str(&sin->sin_addr, AF_INET).str, ntohs(sin->sin_port), s->name, addr_str.str,
s->so_mark, s->netif_index, s->type, s->mtu); s->so_mark, s->netif_index, s->type, s->mtu);
s = s->next; s = s->next;
} }
@ -999,9 +999,9 @@ void print_config(const struct utun_config *cfg) {
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, " %s: peer_key=%s, keepalive=%d", c->name, c->peer_public_key_hex, c->keepalive); DEBUG_INFO(DEBUG_CATEGORY_CONFIG, " %s: peer_key=%s, keepalive=%d", c->name, c->peer_public_key_hex, c->keepalive);
struct CFG_CLIENT_LINK *link = c->links; struct CFG_CLIENT_LINK *link = c->links;
while (link) { while (link) {
struct sockaddr_in *sin = (struct sockaddr_in *)&link->remote_addr; ip_str_t addr_str = sockaddr_storage_to_str(&link->remote_addr);
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, " Link: %s:%d (via %s)", DEBUG_INFO(DEBUG_CATEGORY_CONFIG, " Link: %s (via %s)",
ip_to_str(&sin->sin_addr, AF_INET).str, ntohs(sin->sin_port), link->local_srv->name); addr_str.str, link->local_srv->name);
link = link->next; link = link->next;
} }
c = c->next; c = c->next;

11
src/etcp.c

@ -1239,7 +1239,6 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
uint32_t new_pacing = link->bbr_pacing_rate; uint32_t new_pacing = link->bbr_pacing_rate;
bbr_main(link->bbr, &rs, &new_cwnd, &new_pacing, link->mtu, link->inflight_bytes, bbr_main(link->bbr, &rs, &new_cwnd, &new_pacing, link->mtu, link->inflight_bytes,
(link->inflight_bytes >= link->inflight_lim_bytes)); (link->inflight_bytes >= link->inflight_lim_bytes));
link->inflight_lim_bytes = new_cwnd;
link->bbr_pacing_rate = new_pacing; link->bbr_pacing_rate = new_pacing;
link->bandwidth = (uint32_t)((uint64_t)new_pacing * 8 / 1000); link->bandwidth = (uint32_t)((uint64_t)new_pacing * 8 / 1000);
@ -1278,10 +1277,6 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
void etcp_conn_input(struct ETCP_DGRAM* pkt) { void etcp_conn_input(struct ETCP_DGRAM* pkt) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, ""); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!pkt) return; if (!pkt) return;
if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt);
volatile int _halt = 1; while (_halt) {}
}
if (!pkt->data_len) { if (!pkt->data_len) {
memory_pool_free(pkt->link->etcp->instance->pkt_pool, pkt); memory_pool_free(pkt->link->etcp->instance->pkt_pool, pkt);
return; return;
@ -1549,12 +1544,6 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
} }
if (memory_pool_is_freed(pkt->link->etcp->instance->pkt_pool, pkt)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt);
volatile int _halt = 1; while (_halt) {}
}
memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram memory_pool_free(etcp->instance->pkt_pool, pkt); // Free the incoming dgram
} }

3
src/etcp_bbr.c

@ -862,6 +862,8 @@ void bbr_init(struct bbr* bbr)
bbr->rounds_since_probe = 0; bbr->rounds_since_probe = 0;
bbr->bw_probe_samples = 0; bbr->bw_probe_samples = 0;
bbr->prev_probe_too_high = 0; bbr->prev_probe_too_high = 0;
bbr->on_cwnd_update = NULL;
bbr->cwnd_update_ctx = NULL;
} }
void bbr_main(struct bbr* bbr, const struct bbr_rate_sample* rs, void bbr_main(struct bbr* bbr, const struct bbr_rate_sample* rs,
@ -906,6 +908,7 @@ void bbr_main(struct bbr* bbr, const struct bbr_rate_sample* rs,
out: out:
bbr_advance_latest_delivery_signals(bbr, rs, sample_bw); bbr_advance_latest_delivery_signals(bbr, rs, sample_bw);
bbr->loss_in_cycle |= (rs->lost > 0); bbr->loss_in_cycle |= (rs->lost > 0);
if (bbr->on_cwnd_update) bbr->on_cwnd_update(bbr->cwnd_update_ctx, cwnd);
*cwnd_out = cwnd; *cwnd_out = cwnd;
} }

4
src/etcp_bbr.h

@ -45,6 +45,8 @@ struct bbr_rate_sample {
int is_app_limited; int is_app_limited;
}; };
typedef void (*bbr_cwnd_update_fn)(void* ctx, uint32_t new_cwnd);
struct bbr { struct bbr {
uint32_t min_rtt_us; uint32_t min_rtt_us;
uint32_t bw_lo; uint32_t bw_lo;
@ -97,6 +99,8 @@ struct bbr {
uint8_t full_bw_now : 1; uint8_t full_bw_now : 1;
uint8_t pad_unused : 6; uint8_t pad_unused : 6;
uint64_t now_tb; // 0 = real get_time_tb(); >0 = test-controlled time uint64_t now_tb; // 0 = real get_time_tb(); >0 = test-controlled time
bbr_cwnd_update_fn on_cwnd_update;
void* cwnd_update_ctx;
}; };
void bbr_init(struct bbr* bbr); void bbr_init(struct bbr* bbr);

26
src/etcp_connections.c

@ -746,6 +746,8 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn) {
} }
static void bbr_cwnd_updated(void* ctx, uint32_t new_cwnd);
struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server) { struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!remote_addr) return NULL; if (!remote_addr) return NULL;
@ -787,7 +789,6 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
link->keepalive_sent_count = 0; link->keepalive_sent_count = 0;
link->keepalive_recv_count = 0; link->keepalive_recv_count = 0;
link->ka_period_ms = KA_PERIOD_MIN_MS; link->ka_period_ms = KA_PERIOD_MIN_MS;
link->inflight_lim_bytes = link->mtu * 4; // BBR init_cwnd (~4 packets)
link->bandwidth = 10000; // начальная оценка 10 Mbps для шейпера link->bandwidth = 10000; // начальная оценка 10 Mbps для шейпера
link->burst_id = 0; link->burst_id = 0;
link->burst_active = 0; link->burst_active = 0;
@ -832,10 +833,9 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
while (l && l->next) l=l->next; while (l && l->next) l=l->next;
if (l) l->next = link; else etcp->links = link; if (l) l->next = link; else etcp->links = link;
// пересчитать connection-level optimal_inflight etcp_link_update_inflight_lim(link, link->mtu * 4);
{ uint32_t sum = 0; link->bbr->on_cwnd_update = bbr_cwnd_updated;
for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes; link->bbr->cwnd_update_ctx = link;
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); 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);
@ -847,6 +847,18 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
return link; return link;
} }
static void bbr_cwnd_updated(void* ctx, uint32_t new_cwnd) {
etcp_link_update_inflight_lim((struct ETCP_LINK*)ctx, new_cwnd);
}
void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) {
link->inflight_lim_bytes = new_lim;
struct ETCP_CONN* etcp = link->etcp;
uint32_t sum = 0;
for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes;
etcp->optimal_inflight = sum;
}
void etcp_link_close(struct ETCP_LINK* link) { void etcp_link_close(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!link) return; if (!link) return;
@ -1731,10 +1743,6 @@ process_decrypted:
// log_dump("RECV decrypted:", pkt->data, pkt->data_len, link); // log_dump("RECV decrypted:", pkt->data, pkt->data_len, link);
if (link->link_state == 3) { if (link->link_state == 3) {
if (memory_pool_is_freed(e_sock->instance->pkt_pool, pkt)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt);
volatile int _halt = 1; while (_halt) {}
}
etcp_conn_input(pkt); etcp_conn_input(pkt);
} else memory_pool_free(e_sock->instance->pkt_pool, pkt); } else memory_pool_free(e_sock->instance->pkt_pool, pkt);
return; return;

1
src/etcp_connections.h

@ -262,6 +262,7 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn);
// connection functions // connection functions
// создает новый канал связи для etcp подключения (ETCP_CONN) // создает новый канал связи для etcp подключения (ETCP_CONN)
struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server); struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server);
void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim);
void etcp_link_close(struct ETCP_LINK* link); void etcp_link_close(struct ETCP_LINK* link);
//int etcp_input_cbk(struct packet_buffer* pkt, struct ETCP_SOCKET* conn);// получает расшифрованный пакет //int etcp_input_cbk(struct packet_buffer* pkt, struct ETCP_SOCKET* conn);// получает расшифрованный пакет
int etcp_encrypt_send(struct ETCP_DGRAM* dgram);// зашифровывает и отправляет пакет int etcp_encrypt_send(struct ETCP_DGRAM* dgram);// зашифровывает и отправляет пакет

7
src/utun_instance.c

@ -49,7 +49,12 @@ static int local_sockaddr_equal(const struct sockaddr_storage *a, const struct s
const struct sockaddr_in *ib = (const struct sockaddr_in *)b; const struct sockaddr_in *ib = (const struct sockaddr_in *)b;
return ia->sin_addr.s_addr == ib->sin_addr.s_addr && ia->sin_port == ib->sin_port; return ia->sin_addr.s_addr == ib->sin_addr.s_addr && ia->sin_port == ib->sin_port;
} }
return 0; // IPv6 stub if (a->ss_family == AF_INET6) {
const struct sockaddr_in6 *ia = (const struct sockaddr_in6 *)a;
const struct sockaddr_in6 *ib = (const struct sockaddr_in6 *)b;
return memcmp(&ia->sin6_addr, &ib->sin6_addr, 16) == 0 && ia->sin6_port == ib->sin6_port;
}
return 0;
} }
// Common initialization function (called by both create functions) // Common initialization function (called by both create functions)

Loading…
Cancel
Save