Browse Source

etcp_keepalive: протокол смены keepalive-режима (standby/normal) + событие active mode

- новый модуль etcp_keepalive.c/h: keepalive-пакеты, таймер, адаптивный период,
  согласование при handshake вынесены из etcp_connections.c как есть
- протокол KEEPALIVE_REQ/RESP (0x09/0x0A): мобильный узел при смене active mode
  шлёт желаемый режим; отвечающий считает max (standby при любом standby), не-мобильный
  применяет как есть; инициатор применяет согласованный режим после RESP
- режимы: NORMAL (3с, адаптивный рост), STANDBY (30с, ka_pinned=1 без роста)
- active mode как событие: utun_add/remove_activity_cbk; подписчики topo_group
  (nodeinfo) и etcp_keepalive (рассылка REQ всем линкам)
topo_upd
evgeny 2 months ago
parent
commit
9348706736
  1. 28
      confdefs.h
  2. 2
      conftest.mk
  3. BIN
      conftest.tar
  4. 2
      src/Makefile.am
  5. 21
      src/routing_layer/topo_group.c
  6. 7
      src/routing_layer/topo_group.h
  7. 151
      src/transport_layer/etcp_connections.c
  8. 3
      src/transport_layer/etcp_connections.h
  9. 26
      src/transport_layer/etcp_connections_doc.md
  10. 296
      src/transport_layer/etcp_keepalive.c
  11. 43
      src/transport_layer/etcp_keepalive.h
  12. 61
      src/utun_instance.c
  13. 12
      src/utun_instance.h

28
confdefs.h

@ -1,28 +0,0 @@
/* confdefs.h */
#define PACKAGE_NAME "utun"
#define PACKAGE_TARNAME "utun"
#define PACKAGE_VERSION "2.0.0"
#define PACKAGE_STRING "utun 2.0.0"
#define PACKAGE_BUGREPORT "https://github.com/anomalyco/utun3/issues"
#define PACKAGE_URL ""
#define PACKAGE "utun"
#define VERSION "2.0.0"
#define USE_OPENSSL 1
#define HAVE_ARPA_INET_H 1
#define HAVE_FCNTL_H 1
#define HAVE_LIMITS_H 1
#define HAVE_NETINET_IN_H 1
#define HAVE_STDINT_H 1
#define HAVE_STDLIB_H 1
#define HAVE_STRING_H 1
#define HAVE_SYS_IOCTL_H 1
#define HAVE_SYS_SOCKET_H 1
#define HAVE_UNISTD_H 1
#define HAVE_MALLOC 1
#define HAVE_GETTIMEOFDAY 1
#define HAVE_MEMSET 1
#define HAVE_SOCKET 1
#define HAVE_STRCHR 1
#define HAVE_STRDUP 1
#define HAVE_STRERROR 1
#define HAVE_STRSTR 1

2
conftest.mk

@ -1,2 +0,0 @@
conftest.ts1: conftest.ts2
touch conftest.ts2

BIN
conftest.tar

Binary file not shown.

2
src/Makefile.am

@ -29,6 +29,7 @@ utun_CORE_SOURCES = \
tun_windows.c \
transport_layer/etcp.c \
transport_layer/etcp_connections.c \
transport_layer/etcp_keepalive.c \
transport_layer/etcp_bbr.c \
transport_layer/etcp_loadbalancer.c \
transport_layer/etcp_debug.c \
@ -109,6 +110,7 @@ libutun_a_SOURCES = \
tun_windows.c \
transport_layer/etcp.c \
transport_layer/etcp_connections.c \
transport_layer/etcp_keepalive.c \
transport_layer/etcp_bbr.c \
transport_layer/etcp_loadbalancer.c \
transport_layer/etcp_debug.c \

21
src/routing_layer/topo_group.c

@ -338,11 +338,31 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
etcp_add_conn_status_cbk(instance, topo_group_conn_status, g);
etcp_add_socket_cbk(instance, topo_node_on_socket_changed, NULL,
ETCP_SOCKET_EVENT_ADDR_CHANGED | ETCP_SOCKET_EVENT_STATUS_CHANGED);
utun_add_activity_cbk(instance, topo_group_on_activity_change, NULL);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "TOPO groups initialized (with default group)");
return g;
}
void topo_group_on_activity_change(struct UTUN_INSTANCE* instance, int active, void* arg) {
(void)active; (void)arg;
if (!instance || !instance->topo_groups || !instance->topo_groups->group_list) return;
struct ll_entry* ge = instance->topo_groups->group_list->head;
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
topo_group_update_my_nodeinfo(instance, g);
if (g->senders_list) {
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0);
se = se->next;
}
}
ge = ge->next;
}
}
void topo_groups_destroy(struct UTUN_INSTANCE* instance) {
if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "instance is NULL"); return; }
if (!instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "topo_groups is NULL"); return; }
@ -355,6 +375,7 @@ void topo_groups_destroy(struct UTUN_INSTANCE* instance) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 2 conn_cbk_remove done");
etcp_remove_socket_cbk(instance, topo_node_on_socket_changed, NULL);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 3 socket_cbk_remove done");
utun_remove_activity_cbk(instance, topo_group_on_activity_change, NULL);
route_connectivity_cancel_all(instance);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[TOPO_DESTROY] 4 connectivity_cancel done");

7
src/routing_layer/topo_group.h

@ -197,6 +197,13 @@ struct TOPO_GROUPS {
*/
struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance);
/**
* @brief Обработчик события смены active mode (подписчик utun_add_activity_cbk).
*
* Обновляет свой nodeinfo и рассылает его по всем группам и активным BGP-пирам.
*/
void topo_group_on_activity_change(struct UTUN_INSTANCE* instance, int active, void* arg);
/**
* @brief Освобождает все группы и контейнер.
*

151
src/transport_layer/etcp_connections.c

@ -25,6 +25,7 @@
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "etcp_loadbalancer.h"
#include "etcp_keepalive.h"
#include <stdlib.h>
#include <time.h>
#include "../lib/mem.h"
@ -107,11 +108,8 @@ static void etcp_link_remove_from_connections(struct ETCP_SOCKET* conn, struct E
static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision);
//static int etcp_link_send_reset(struct ETCP_LINK* link);
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);
static void start_keepalive_timer(struct ETCP_LINK* link);
static int etcp_tcp_send(struct ETCP_DGRAM* dgram);
// === Burst sender functions ===
@ -318,142 +316,6 @@ void etcp_link_enter_reinit(struct ETCP_LINK* link) {
}
// Вычислить keepalive по типу устройств: оба десктоп/сервер → min, иначе (есть mobile) → max
static uint16_t negotiate_keepalive(uint8_t my_type, uint16_t my_ka,
uint8_t peer_type, uint16_t peer_ka) {
int my_lp = (my_type == CLIENT_TYPE_MOBILE);
int peer_lp = (peer_type == CLIENT_TYPE_MOBILE);
uint16_t result = (my_lp || peer_lp) ? (my_ka > peer_ka ? my_ka : peer_ka)
: (my_ka < peer_ka ? my_ka : peer_ka);
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "negotiate: my=%d(type=%d) peer=%d(type=%d) → %d (rule=%s)",
my_ka, my_type, peer_ka, peer_type, result,
(my_lp || peer_lp) ? "max(low_power)" : "min(both_desktop)");
return result;
}
// Send empty keepalive packet (only timestamp, no sections)
static void etcp_link_send_keepalive(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_KEEPALIVE, "");
if (!link || !link->etcp || !link->etcp->instance) return;
struct ETCP_DGRAM* dgram = u_malloc(sizeof(struct ETCP_DGRAM) + 4);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "malloc failed");
return;
}
dgram->link = link;
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->timestamp = get_current_timestamp();
dgram->flag_up = link->is_tcp ? 1 : link->recv_keepalive;
link->keepalive_sent_count++;
etcp_encrypt_send(dgram);
u_free(dgram);
}
// Check if all links for an ETCP_CONN are down
// Returns 1 if all links are down or no links exist, 0 otherwise
static int etcp_all_links_down(struct ETCP_CONN* etcp) {
if (!etcp || !etcp->links) return 1;
struct ETCP_LINK* l = etcp->links;
while (l) {
if (l->link_status == 1) {
return 0; // At least one link is up
}
l = l->next;
}
return 1; // All links are down
}
static void start_keepalive_timer(struct ETCP_LINK* link) {
// Start keepalive timer
if (link->init_timer) {// cancel init timer
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
if (link->keepalive_timer == NULL) {
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Keepalive timer started on link %p (interval=%d ms)", link->etcp->log_name, link, link->keepalive_interval);
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->keepalive_interval * 10, link, keepalive_timer_cb, "link_keepalive");
}
}
// Keepalive timer callback
static void keepalive_timer_cb(void* arg) {
DEBUG_TRACE(DEBUG_CATEGORY_KEEPALIVE, "");
struct ETCP_LINK* link = (struct ETCP_LINK*)arg;
if (!link || !link->etcp || !link->etcp->instance) {
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "KEEPALIVE NULL !!!!!!!!");
return;
}
link->keepalive_timer = NULL;
// Check if all links are down and start recovery if needed (client only)
if (link->is_server == 0 && etcp_all_links_down(link->etcp)) {
if (link->is_tcp) {
if (link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; }
etcp_tcp_link_start_reconnect(link);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] All links are down, starting recovery", link->etcp->log_name);
etcp_link_enter_reinit(link);// keepalive timr после reinit не нужен
return;
}
// Skip if link is not initialized
if (!link->initialized) {
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Keepalive skipped - link not initialized",
link->etcp->log_name);
goto restart_timer;
}
// Check keepalive timeout
uint64_t now = get_time_tb();
uint64_t timeout_units = (uint64_t)link->keepalive_timeout * 10; // ms -> 0.1ms units
uint64_t elapsed = now - link->last_recv_local_time;
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[%s] ka_check: recv=%d elapsed=%llu timeout=%llu",
link->etcp->log_name, link->recv_keepalive, (unsigned long long)elapsed, (unsigned long long)timeout_units);
if (elapsed > timeout_units) {
if (link->recv_keepalive != 0) {
link->recv_keepalive = 0;
int old_link_status = link->link_status;
link->link_status = 0;
etcp_fire_link_status_cbk(link, link->link_state, old_link_status);
etcp_on_link_down(link->etcp, link);
if (link->is_tcp && link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; }
if (old_link_status) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d down: addr=%s ka=%d remote_ka=%d state=%d init=%d tmo: %llu>%llu els=%llums", link->etcp->log_name, link->local_link_id, sockaddr_storage_to_str(&link->remote_addr).str, link->recv_keepalive, link->remote_keepalive, link->link_state, link->initialized, (unsigned long long)timeout_units, (unsigned long long)elapsed, (unsigned long long)(elapsed/10));
}
}
}
// Adaptive keepalive period (only if adaptive enabled)
if (link->pkt_sent_since_keepalive)
link->ka_period_ms = (uint16_t)link->keepalive_interval;
else if (link->etcp->instance && link->etcp->instance->config &&
link->etcp->instance->config->global.keepalive_adaptive) {
uint32_t next = (uint32_t)link->ka_period_ms * 105 / 100 + 1;
link->ka_period_ms = next > KA_PERIOD_MAX_MS ? KA_PERIOD_MAX_MS : (uint16_t)next;
}
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:
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive");
}
static void sockaddr_to_key(struct sockaddr_storage* addr, uint8_t key[LINK_ADDR_KEY_SIZE], int is_tcp) {
memset(key, 0, LINK_ADDR_KEY_SIZE);
if (!addr) return;
@ -1011,6 +873,7 @@ 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->ka_period_ms = (uint16_t)link->keepalive_interval;
link->ka_pinned = 0;
/* inflight_lim_bytes set above */
link->bandwidth = 10000; // начальная оценка 10 Mbps для шейпера
link->burst_id = 0;
@ -2411,13 +2274,7 @@ int etcp_packet_decrypted(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt,
uint8_t pkt_code = pkt->data[0];
if (pkt_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);
if (etcp_keepalive_on_recv(e_sock, pkt, link, pkt_len)) {
return 0;
}
@ -2590,6 +2447,8 @@ int init_connections(struct UTUN_INSTANCE* instance) {
}
}
etcp_keepalive_register(instance);
// Initialize clients via node_conn_direct
struct CFG_CLIENT* client = config->clients;
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections: clients=%p total_conns=%d", config->clients, queue_entry_count(instance->connections));

3
src/transport_layer/etcp_connections.h

@ -276,6 +276,7 @@ struct ETCP_LINK {
void* keepalive_timer; // Таймер для отправки keepalive пакетов
uint32_t keepalive_timeout; // таймаут (ms)
uint16_t ka_period_ms; // адаптивный период отправки keepalive (200→10000->200)
uint8_t ka_pinned; // 1 = период зафиксирован протоколом (standby-режим мобильного), адаптивный рост отключён
uint8_t pkt_sent_since_keepalive; // Флаг: был ли отправлен пакет с последнего keepalive тика
uint32_t keepalive_sent_count; // Счётчик отправленных keepalive
uint32_t keepalive_recv_count; // Счётчик полученных keepalive
@ -368,6 +369,8 @@ 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);
void etcp_link_enter_init(struct ETCP_LINK *link);
void etcp_link_enter_reinit(struct ETCP_LINK *link);
void etcp_link_enter_ready_tcp(struct ETCP_LINK *link);
void etcp_tcp_link_start_connect(struct ETCP_LINK *link, struct sockaddr_storage *addr, uint16_t port);
void etcp_tcp_link_start_reconnect(struct ETCP_LINK *link);

26
src/transport_layer/etcp_connections_doc.md

@ -85,6 +85,30 @@ void my_ping_callback(int success, uint16_t rtt, void* arg, uint64_t nonce,
- Таймаут = `period × KA_TIMEOUT_MULT` (×10). При превышении → линк падает
- Сервер не шлёт keepalive при `recv_keepalive=0` (линк уже мёртв), клиент шлёт всегда
### Протокол смены keepalive-режима (standby/normal, мобильные узлы)
Модуль `etcp_keepalive.c/h`. Мобильный узел (Android) при смене active mode
(передний/фоновый план) сообщает пиру желаемый keepalive-режим через кодограмму
`ETCP_KEEPALIVE_REQ` (0x09); пир применяет согласованный режим и отвечает
`ETCP_KEEPALIVE_RESP` (0x0A). Формат: `data[0]=code`, `data[1]=режим`.
Режимы:
- `KA_MODE_NORMAL` (0) — активен: интервал `KA_MOBILE_ACTIVE_MS` (3с), адаптивный
(рост 3с → 10с, сброс на трафике).
- `KA_MODE_STANDBY` (1) — спит: интервал `KA_MOBILE_STANDBY_MS` (30с), период
зафиксирован (`ka_pinned=1`), адаптивный рост отключён.
Согласование:
- Мобильный отвечающий считает `agreed = standby`, если его или запрошенный режим
standby (эквивалент `max` интервалов), иначе `normal`.
- Не-мобильный отвечающий применяет запрошенный режим как есть.
- Инициатор применяет согласованный режим только после получения RESP.
- Применение: `keepalive_interval = ka_period_ms = ms`, `keepalive_timeout = ms * KA_TIMEOUT_MULT`,
`ka_pinned` по режиму, рестарт keepalive-таймера.
Смена active mode рассылается через событие `utun_add_activity_cbk`
(подписчики: `topo_group_on_activity_change`, `etcp_keepalive_on_activity`).
## 3. API
### Ключевые структуры
@ -142,6 +166,8 @@ void my_ping_callback(int success, uint16_t rtt, void* arg, uint64_t nonce,
| `ETCP_PING` | 0x06 | One-shot пробник |
| `ETCP_PONG` | 0x07 | Ответ на пробник |
| `ETCP_KEEPALIVE` | 0x08 | keepalive-пакет |
| `ETCP_KEEPALIVE_REQ` | 0x09 | Запрос смены keepalive-режима (standby/normal) |
| `ETCP_KEEPALIVE_RESP` | 0x0A | Ответ: согласованный режим |
### NAT-типы

296
src/transport_layer/etcp_keepalive.c

@ -0,0 +1,296 @@
#include "etcp_keepalive.h"
#include "etcp.h"
#include "etcp_api.h"
#include "stcp_link.h"
#include "topo_node.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/memory_pool.h"
// forward declarations (static)
static void keepalive_timer_cb(void* arg);
static int etcp_all_links_down(struct ETCP_CONN* etcp);
// === перенесено из etcp_connections.c без изменений ===
// Вычислить keepalive по типу устройств: оба десктоп/сервер → min, иначе (есть mobile) → max
uint16_t negotiate_keepalive(uint8_t my_type, uint16_t my_ka,
uint8_t peer_type, uint16_t peer_ka) {
int my_lp = (my_type == CLIENT_TYPE_MOBILE);
int peer_lp = (peer_type == CLIENT_TYPE_MOBILE);
uint16_t result = (my_lp || peer_lp) ? (my_ka > peer_ka ? my_ka : peer_ka)
: (my_ka < peer_ka ? my_ka : peer_ka);
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "negotiate: my=%d(type=%d) peer=%d(type=%d) → %d (rule=%s)",
my_ka, my_type, peer_ka, peer_type, result,
(my_lp || peer_lp) ? "max(low_power)" : "min(both_desktop)");
return result;
}
// Send empty keepalive packet (only timestamp, no sections)
void etcp_link_send_keepalive(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_KEEPALIVE, "");
if (!link || !link->etcp || !link->etcp->instance) return;
struct ETCP_DGRAM* dgram = u_malloc(sizeof(struct ETCP_DGRAM) + 4);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "malloc failed");
return;
}
dgram->link = link;
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->timestamp = get_current_timestamp();
dgram->flag_up = link->is_tcp ? 1 : link->recv_keepalive;
link->keepalive_sent_count++;
etcp_encrypt_send(dgram);
u_free(dgram);
}
// Check if all links for an ETCP_CONN are down
// Returns 1 if all links are down or no links exist, 0 otherwise
static int etcp_all_links_down(struct ETCP_CONN* etcp) {
if (!etcp || !etcp->links) return 1;
struct ETCP_LINK* l = etcp->links;
while (l) {
if (l->link_status == 1) {
return 0; // At least one link is up
}
l = l->next;
}
return 1; // All links are down
}
void start_keepalive_timer(struct ETCP_LINK* link) {
// Start keepalive timer
if (link->init_timer) {// cancel init timer
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
if (link->keepalive_timer == NULL) {
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Keepalive timer started on link %p (interval=%d ms)", link->etcp->log_name, link, link->keepalive_interval);
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->keepalive_interval * 10, link, keepalive_timer_cb, "link_keepalive");
}
}
// Keepalive timer callback
static void keepalive_timer_cb(void* arg) {
DEBUG_TRACE(DEBUG_CATEGORY_KEEPALIVE, "");
struct ETCP_LINK* link = (struct ETCP_LINK*)arg;
if (!link || !link->etcp || !link->etcp->instance) {
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "KEEPALIVE NULL !!!!!!!!");
return;
}
link->keepalive_timer = NULL;
// Check if all links are down and start recovery if needed (client only)
if (link->is_server == 0 && etcp_all_links_down(link->etcp)) {
if (link->is_tcp) {
if (link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; }
etcp_tcp_link_start_reconnect(link);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] All links are down, starting recovery", link->etcp->log_name);
etcp_link_enter_reinit(link);// keepalive timr после reinit не нужен
return;
}
// Skip if link is not initialized
if (!link->initialized) {
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Keepalive skipped - link not initialized",
link->etcp->log_name);
goto restart_timer;
}
// Check keepalive timeout
uint64_t now = get_time_tb();
uint64_t timeout_units = (uint64_t)link->keepalive_timeout * 10; // ms -> 0.1ms units
uint64_t elapsed = now - link->last_recv_local_time;
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[%s] ka_check: recv=%d elapsed=%llu timeout=%llu",
link->etcp->log_name, link->recv_keepalive, (unsigned long long)elapsed, (unsigned long long)timeout_units);
if (elapsed > timeout_units) {
if (link->recv_keepalive != 0) {
link->recv_keepalive = 0;
int old_link_status = link->link_status;
link->link_status = 0;
etcp_fire_link_status_cbk(link, link->link_state, old_link_status);
etcp_on_link_down(link->etcp, link);
if (link->is_tcp && link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; }
if (old_link_status) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d down: addr=%s ka=%d remote_ka=%d state=%d init=%d tmo: %llu>%llu els=%llums", link->etcp->log_name, link->local_link_id, sockaddr_storage_to_str(&link->remote_addr).str, link->recv_keepalive, link->remote_keepalive, link->link_state, link->initialized, (unsigned long long)timeout_units, (unsigned long long)elapsed, (unsigned long long)(elapsed/10));
}
}
}
// Adaptive keepalive period (only if adaptive enabled)
if (link->ka_pinned)
link->ka_period_ms = (uint16_t)link->keepalive_interval;
else if (link->pkt_sent_since_keepalive)
link->ka_period_ms = (uint16_t)link->keepalive_interval;
else if (link->etcp->instance && link->etcp->instance->config &&
link->etcp->instance->config->global.keepalive_adaptive) {
uint32_t next = (uint32_t)link->ka_period_ms * 105 / 100 + 1;
link->ka_period_ms = next > KA_PERIOD_MAX_MS ? KA_PERIOD_MAX_MS : (uint16_t)next;
}
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:
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive");
}
// === протокол смены keepalive-режима (standby/normal, p2p) ===
static uint8_t mobile_desired_ka_mode(struct UTUN_INSTANCE* inst) {
return (inst->client_activity == CLIENT_ACTIVITY_ACTIVE) ? KA_MODE_NORMAL : KA_MODE_STANDBY;
}
static void restart_keepalive_timer(struct ETCP_LINK* link) {
if (!link || !link->etcp || !link->etcp->instance) return;
if (link->keepalive_timer) {
uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer);
link->keepalive_timer = NULL;
}
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive");
}
static void apply_link_ka_mode(struct ETCP_LINK* link, uint8_t mode) {
uint16_t ms = (mode == KA_MODE_STANDBY) ? KA_MOBILE_STANDBY_MS : KA_MOBILE_ACTIVE_MS;
link->keepalive_interval = ms;
link->ka_period_ms = ms;
link->keepalive_timeout = (uint32_t)ms * KA_TIMEOUT_MULT;
link->ka_pinned = (mode == KA_MODE_STANDBY) ? 1 : 0;
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] keepalive mode applied: mode=%s interval=%dms timeout=%dms pinned=%d",
link->etcp->log_name, mode == KA_MODE_STANDBY ? "standby" : "normal",
ms, link->keepalive_timeout, link->ka_pinned);
restart_keepalive_timer(link);
}
static void etcp_keepalive_send_mode_pkt(struct ETCP_LINK* link, uint8_t code, uint8_t mode) {
if (!link || !link->etcp || !link->etcp->instance) return;
struct ETCP_DGRAM* dgram = u_malloc(sizeof(struct ETCP_DGRAM) + 4);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "malloc failed");
return;
}
dgram->link = link;
dgram->data[0] = code;
dgram->data[1] = mode;
dgram->data_len = 2;
dgram->noencrypt_len = 0;
dgram->timestamp = get_current_timestamp();
dgram->flag_up = link->is_tcp ? 1 : link->recv_keepalive;
etcp_encrypt_send(dgram);
u_free(dgram);
}
static void handle_keepalive_req(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, struct ETCP_LINK* link) {
uint8_t req_mode = (pkt->data_len >= 2) ? pkt->data[1] : KA_MODE_NORMAL;
uint8_t agreed;
if (e_sock->instance->client_type == CLIENT_TYPE_MOBILE) {
uint8_t my_mode = mobile_desired_ka_mode(e_sock->instance);
agreed = (my_mode == KA_MODE_STANDBY || req_mode == KA_MODE_STANDBY) ? KA_MODE_STANDBY : KA_MODE_NORMAL;
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] keepalive REQ recv: req=%s my=%s -> agreed=%s",
link->etcp->log_name,
req_mode == KA_MODE_STANDBY ? "standby" : "normal",
my_mode == KA_MODE_STANDBY ? "standby" : "normal",
agreed == KA_MODE_STANDBY ? "standby" : "normal");
} else {
agreed = req_mode;
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] keepalive REQ recv: req=%s (non-mobile) -> agreed=%s",
link->etcp->log_name,
req_mode == KA_MODE_STANDBY ? "standby" : "normal",
agreed == KA_MODE_STANDBY ? "standby" : "normal");
}
apply_link_ka_mode(link, agreed);
etcp_keepalive_send_mode_pkt(link, ETCP_KEEPALIVE_RESP, agreed);
}
static void handle_keepalive_resp(struct ETCP_DGRAM* pkt, struct ETCP_LINK* link) {
uint8_t agreed = (pkt->data_len >= 2) ? pkt->data[1] : KA_MODE_NORMAL;
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] keepalive RESP recv: agreed=%s",
link->etcp->log_name, agreed == KA_MODE_STANDBY ? "standby" : "normal");
apply_link_ka_mode(link, agreed);
}
int etcp_keepalive_on_recv(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt,
struct ETCP_LINK* link, size_t pkt_len) {
(void)pkt_len;
if (!e_sock || !pkt || !link) return 0;
uint8_t code = pkt->data[0];
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 1;
}
if (code == ETCP_KEEPALIVE_REQ) {
handle_keepalive_req(e_sock, pkt, link);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 1;
}
if (code == ETCP_KEEPALIVE_RESP) {
handle_keepalive_resp(pkt, link);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 1;
}
return 0;
}
static void etcp_keepalive_on_activity(struct UTUN_INSTANCE* inst, int active, void* arg) {
(void)active; (void)arg;
if (!inst || inst->client_type != CLIENT_TYPE_MOBILE) return;
uint8_t mode = mobile_desired_ka_mode(inst);
int sent = 0;
if (inst->connections) {
for (struct ll_entry* e = inst->connections->head; e; e = e->next) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
if (!ce || !ce->conn) continue;
for (struct ETCP_LINK* l = ce->conn->links; l; l = l->next) {
if (l->initialized) { etcp_keepalive_send_mode_pkt(l, ETCP_KEEPALIVE_REQ, mode); sent++; }
}
}
}
if (inst->tcp_connections) {
for (struct ll_entry* e = inst->tcp_connections->head; e; e = e->next) {
struct tcp_conn_entry* te = (struct tcp_conn_entry*)e->data;
if (!te || !te->etcp_conn) continue;
for (struct ETCP_LINK* l = te->etcp_conn->links; l; l = l->next) {
if (l->initialized) { etcp_keepalive_send_mode_pkt(l, ETCP_KEEPALIVE_REQ, mode); sent++; }
}
}
}
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "keepalive activity change: mode=%s -> REQ sent to %d links (node=%016llx)",
mode == KA_MODE_STANDBY ? "standby" : "normal", sent, (unsigned long long)inst->node_id);
}
void etcp_keepalive_register(struct UTUN_INSTANCE* inst) {
if (!inst) return;
utun_add_activity_cbk(inst, etcp_keepalive_on_activity, NULL);
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "keepalive activity callback registered");
}

43
src/transport_layer/etcp_keepalive.h

@ -0,0 +1,43 @@
#ifndef ETCP_KEEPALIVE_H
#define ETCP_KEEPALIVE_H
#ifdef __cplusplus
extern "C" {
#endif
// подмодуль ETCP: keepalive-пакеты, таймер, адаптивный период и протокол смены
// keepalive-режима (standby/normal) при смене active mode у мобильных узлов.
#include "etcp_connections.h"
// Типы кодограмм keepalive
#define ETCP_KEEPALIVE_REQ 0x09 // запрос смены keepalive-режима (data[1] = режим)
#define ETCP_KEEPALIVE_RESP 0x0A // ответ: согласованный режим (data[1] = режим)
// Режимы keepalive (мобильные узлы)
#define KA_MODE_NORMAL 0 // активен: интервал 3с, адаптивный
#define KA_MODE_STANDBY 1 // спит: интервал 30с, пиннится
#define KA_MOBILE_ACTIVE_MS 3000
#define KA_MOBILE_STANDBY_MS 30000
// --- перенесены из etcp_connections.c (имена сохранены) ---
uint16_t negotiate_keepalive(uint8_t my_type, uint16_t my_ka,
uint8_t peer_type, uint16_t peer_ka);
void etcp_link_send_keepalive(struct ETCP_LINK* link);
void start_keepalive_timer(struct ETCP_LINK* link);
// --- новое ---
// Обработка keepalive-пакетов (ETCP_KEEPALIVE/REQ/RESP) из etcp_packet_decrypted.
// Возвращает 1 если пакет обработан (и освобождён), 0 — если это не keepalive.
int etcp_keepalive_on_recv(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt,
struct ETCP_LINK* link, size_t pkt_len);
// Подписка на смену active mode (рассылка KEEPALIVE_REQ всем линкам).
// Вызывается из init_connections.
void etcp_keepalive_register(struct UTUN_INSTANCE* inst);
#ifdef __cplusplus
}
#endif
#endif // ETCP_KEEPALIVE_H

61
src/utun_instance.c

@ -1027,6 +1027,35 @@ struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struc
return instance;
}
void utun_add_activity_cbk(struct UTUN_INSTANCE* instance, utun_activity_cbk_fn fn, void* arg) {
if (!instance || !fn) return;
struct utun_activity_cbk_entry* e = u_malloc(sizeof(struct utun_activity_cbk_entry));
if (!e) return;
e->fn = fn; e->arg = arg; e->next = instance->activity_cbks;
instance->activity_cbks = e;
}
void utun_remove_activity_cbk(struct UTUN_INSTANCE* instance, utun_activity_cbk_fn fn, void* arg) {
if (!instance || !fn) return;
struct utun_activity_cbk_entry** p = &instance->activity_cbks;
while (*p) {
if ((*p)->fn == fn && (*p)->arg == arg) {
struct utun_activity_cbk_entry* rm = *p;
*p = rm->next; u_free(rm); return;
}
p = &(*p)->next;
}
}
static void utun_fire_activity_cbk(struct UTUN_INSTANCE* instance, int active) {
struct utun_activity_cbk_entry* e = instance->activity_cbks;
while (e) {
struct utun_activity_cbk_entry* n = e->next;
e->fn(instance, active, e->arg);
e = n;
}
}
static void client_activity_timeout_cb(void* arg) {
struct UTUN_INSTANCE* instance = (struct UTUN_INSTANCE*)arg;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "client_activity_timer fired: node=%016llx type=%u -> standby",
@ -1034,21 +1063,7 @@ static void client_activity_timeout_cb(void* arg) {
instance->client_activity_timer = NULL;
if (instance->client_type == CLIENT_TYPE_SERVER) return;
instance->client_activity = CLIENT_ACTIVITY_STANDBY;
if (!instance->topo_groups || !instance->topo_groups->group_list) return;
struct ll_entry* ge = instance->topo_groups->group_list->head;
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
topo_group_update_my_nodeinfo(instance, g);
if (g->senders_list) {
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0);
se = se->next;
}
}
ge = ge->next;
}
utun_fire_activity_cbk(instance, CLIENT_ACTIVITY_STANDBY);
}
void utun_set_client_activity(struct UTUN_INSTANCE* instance, int active) {
@ -1068,19 +1083,5 @@ void utun_set_client_activity(struct UTUN_INSTANCE* instance, int active) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "client_activity set STANDBY node=%016llx type=%u",
(unsigned long long)instance->node_id, instance->client_type);
}
if (!instance->topo_groups || !instance->topo_groups->group_list) return;
struct ll_entry* ge = instance->topo_groups->group_list->head;
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
topo_group_update_my_nodeinfo(instance, g);
if (g->senders_list) {
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0);
se = se->next;
}
}
ge = ge->next;
}
utun_fire_activity_cbk(instance, active);
}

12
src/utun_instance.h

@ -80,6 +80,15 @@ struct tcp_conn_entry {
struct ETCP_CONN* etcp_conn; // = stcp_link_get_etcp_conn(link) — for etcp_send() compat
};
// Подписка на смену active mode (client_activity). Событие рассылается
// подписчикам (topo_group, keepalive и т.д.) при каждом изменении активности.
typedef void (*utun_activity_cbk_fn)(struct UTUN_INSTANCE* instance, int active, void* arg);
struct utun_activity_cbk_entry {
utun_activity_cbk_fn fn;
void* arg;
struct utun_activity_cbk_entry* next;
};
// uTun instance configuration
struct UTUN_INSTANCE {
// Identification
@ -189,6 +198,7 @@ struct UTUN_INSTANCE {
uint16_t keepalive_interval; // желаемый keepalive (ms), из конфига. для handshake
uint8_t client_activity; // CLIENT_ACTIVITY_STANDBY/ACTIVE
void* client_activity_timer; // uasync timer handle for inactivity timeout
struct utun_activity_cbk_entry* activity_cbks; // подписки на смену client_activity
// TCP proxy server (exit node)
struct tcp_proxy_server tcp_proxy_server;
@ -224,6 +234,8 @@ void utun_instance_stop(struct UTUN_INSTANCE *instance);
void utun_instance_set_tun_init_enabled(int enabled);
void utun_instance_set_topo_group_enabled(int enabled);
void utun_set_client_activity(struct UTUN_INSTANCE* instance, int active);
void utun_add_activity_cbk(struct UTUN_INSTANCE* instance, utun_activity_cbk_fn fn, void* arg);
void utun_remove_activity_cbk(struct UTUN_INSTANCE* instance, utun_activity_cbk_fn fn, void* arg);
// Diagnostic function for memory leak analysis
void utun_instance_diagnose_leaks(struct UTUN_INSTANCE* instance, const char* phase);

Loading…
Cancel
Save