diff --git a/confdefs.h b/confdefs.h deleted file mode 100644 index ea8e1156..00000000 --- a/confdefs.h +++ /dev/null @@ -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 diff --git a/conftest.mk b/conftest.mk deleted file mode 100644 index 771e2547..00000000 --- a/conftest.mk +++ /dev/null @@ -1,2 +0,0 @@ -conftest.ts1: conftest.ts2 - touch conftest.ts2 diff --git a/conftest.tar b/conftest.tar deleted file mode 100644 index 2b0a5bf8..00000000 Binary files a/conftest.tar and /dev/null differ diff --git a/src/Makefile.am b/src/Makefile.am index 9170e490..5fe99e0b 100644 --- a/src/Makefile.am +++ b/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 \ diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 419566e9..aee079dd 100644 --- a/src/routing_layer/topo_group.c +++ b/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"); diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index d6cccb81..fa934ec7 100644 --- a/src/routing_layer/topo_group.h +++ b/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 Освобождает все группы и контейнер. * diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index f840ac3d..b68c4507 100644 --- a/src/transport_layer/etcp_connections.c +++ b/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 #include #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)); diff --git a/src/transport_layer/etcp_connections.h b/src/transport_layer/etcp_connections.h index 38da781d..e7378a20 100644 --- a/src/transport_layer/etcp_connections.h +++ b/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); diff --git a/src/transport_layer/etcp_connections_doc.md b/src/transport_layer/etcp_connections_doc.md index 4f5b575a..35e73617 100644 --- a/src/transport_layer/etcp_connections_doc.md +++ b/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-типы diff --git a/src/transport_layer/etcp_keepalive.c b/src/transport_layer/etcp_keepalive.c new file mode 100644 index 00000000..c9741d94 --- /dev/null +++ b/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"); +} diff --git a/src/transport_layer/etcp_keepalive.h b/src/transport_layer/etcp_keepalive.h new file mode 100644 index 00000000..345be775 --- /dev/null +++ b/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 diff --git a/src/utun_instance.c b/src/utun_instance.c index 0d8159bb..363f3e85 100644 --- a/src/utun_instance.c +++ b/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); } diff --git a/src/utun_instance.h b/src/utun_instance.h index 8f412e44..c0df27c1 100644 --- a/src/utun_instance.h +++ b/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);