Browse Source

etcp: адаптивный keepalive — отдельная команда ETCP_KEEPALIVE с period_ms

- ETCP_KEEPALIVE=0x08, формат: [0x08][period_lo][period_hi]
- link->ka_period_ms: 200ms при трафике, ×1.05+1 в idle, потолок 10s
- keepalive_timeout = peer_period * KA_TIMEOUT_MULT(10)
- приём KA обновляет keepalive_timeout адаптивно (2000ms→100000ms)
- флаг recv_keepalive по-прежнему в flag_up для быстрого сигнала
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
6d661360b1
  1. 45
      src/etcp_connections.c
  2. 7
      src/etcp_connections.h

45
src/etcp_connections.c

@ -239,12 +239,13 @@ static void etcp_link_send_keepalive(struct ETCP_LINK* link) {
} }
dgram->link = link; dgram->link = link;
dgram->data_len = 0; // Empty packet - only timestamp in header dgram->data[0] = ETCP_KEEPALIVE;
dgram->data[1] = link->ka_period_ms & 0xFF;
dgram->data[2] = link->ka_period_ms >> 8;
dgram->data_len = 3;
dgram->noencrypt_len = 0; dgram->noencrypt_len = 0;
dgram->timestamp = get_current_timestamp(); dgram->timestamp = get_current_timestamp();
dgram->flag_up = link->recv_keepalive;
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Sending keepalive on link %p (local_id=%d)",
link->etcp->log_name, link, link->local_link_id);
link->keepalive_sent_count++; link->keepalive_sent_count++;
@ -320,17 +321,21 @@ static void keepalive_timer_cb(void* arg) {
} }
} }
// Send keepalive only if no packets were sent since last tick // Adaptive keepalive period
if (!link->pkt_sent_since_keepalive) { if (link->pkt_sent_since_keepalive)
if (link->is_server) { link->ka_period_ms = KA_PERIOD_MIN_MS; // data flowing → reset to 200ms
if (link->recv_keepalive) etcp_link_send_keepalive(link);// сервер прекращает слать keepalive если линк потерян (ждём keepalive клиента) else {
} uint32_t next = (uint32_t)link->ka_period_ms * 105 / 100 + 1;
else etcp_link_send_keepalive(link); link->ka_period_ms = next > KA_PERIOD_MAX_MS ? KA_PERIOD_MAX_MS : (uint16_t)next;
} }
link->pkt_sent_since_keepalive = 0; link->pkt_sent_since_keepalive = 0;
// Send keepalive (server stops if link lost, client always sends)
if (!link->is_server || link->recv_keepalive)
etcp_link_send_keepalive(link);
restart_timer: restart_timer:
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->keepalive_interval * 10, link, keepalive_timer_cb, "link_keepalive"); link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive");
} }
static uint32_t sockaddr_hash(struct sockaddr_storage* addr) { static uint32_t sockaddr_hash(struct sockaddr_storage* addr) {
@ -781,6 +786,7 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
if (link->keepalive_interval < 10) link->keepalive_interval = 10; if (link->keepalive_interval < 10) link->keepalive_interval = 10;
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->inflight_lim_bytes = link->mtu * 4; // BBR init_cwnd (~4 packets) 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;
@ -1590,16 +1596,19 @@ process_decrypted:
// Count decrypted bytes // Count decrypted bytes
link->total_decrypted += pkt->data_len; link->total_decrypted += pkt->data_len;
// Count received keepalive packets (empty packets with no payload)
if (pkt->data_len == 0) {
link->keepalive_recv_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Received keepalive on link %p (local_id=%d)",
link->etcp->log_name, link, link->local_link_id);
}
size_t offset = 0; size_t offset = 0;
uint8_t code = pkt->data[offset++]; uint8_t code = pkt->data[offset++];
if (code == ETCP_KEEPALIVE) {
if (pkt->data_len >= 3) {
uint16_t peer_period = pkt->data[1] | ((uint16_t)pkt->data[2] << 8);
link->keepalive_timeout = (uint32_t)peer_period * KA_TIMEOUT_MULT;
}
link->keepalive_recv_count++;
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return; // KA handled, nothing more to process
}
if (code == ETCP_INIT_RESPONSE || code == ETCP_INIT_RESPONSE_NOINIT) { if (code == ETCP_INIT_RESPONSE || code == ETCP_INIT_RESPONSE_NOINIT) {
if (pkt_len < 22) { errorcode = 46; DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "INIT_RESPONSE too short: pkt_len=%zu", pkt_len); goto ec_fr; } if (pkt_len < 22) { errorcode = 46; DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "INIT_RESPONSE too short: pkt_len=%zu", pkt_len); goto ec_fr; }
// ETCP_INIT_RESPONSE (0x03) - reset entire ETCP_CONN // ETCP_INIT_RESPONSE (0x03) - reset entire ETCP_CONN

7
src/etcp_connections.h

@ -26,6 +26,12 @@
#define ETCP_INIT_RESPONSE_NOINIT 0x05 #define ETCP_INIT_RESPONSE_NOINIT 0x05
#define ETCP_PING 0x06 #define ETCP_PING 0x06
#define ETCP_PONG 0x07 #define ETCP_PONG 0x07
#define ETCP_KEEPALIVE 0x08
/* Адаптивный keepalive */
#define KA_PERIOD_MIN_MS 200
#define KA_PERIOD_MAX_MS 10000
#define KA_TIMEOUT_MULT 10 /* timeout = period * mult */
#pragma pack(push, 1) #pragma pack(push, 1)
@ -194,6 +200,7 @@ struct ETCP_LINK {
// Keepalive state // Keepalive state
void* keepalive_timer; // Таймер для отправки keepalive пакетов void* keepalive_timer; // Таймер для отправки keepalive пакетов
uint32_t keepalive_timeout; // таймаут (ms) uint32_t keepalive_timeout; // таймаут (ms)
uint16_t ka_period_ms; // адаптивный период отправки keepalive (200→10000->200)
uint8_t pkt_sent_since_keepalive; // Флаг: был ли отправлен пакет с последнего keepalive тика uint8_t pkt_sent_since_keepalive; // Флаг: был ли отправлен пакет с последнего keepalive тика
uint32_t keepalive_sent_count; // Счётчик отправленных keepalive uint32_t keepalive_sent_count; // Счётчик отправленных keepalive
uint32_t keepalive_recv_count; // Счётчик полученных keepalive uint32_t keepalive_recv_count; // Счётчик полученных keepalive

Loading…
Cancel
Save