You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

983 lines
41 KiB

// utun_instance.c - Root instance implementation
#include "utun_instance.h"
#include "etcp_router.h"
#include "config_parser.h"
#include "config_updater.h"
#include "tun_if.h"
#include "tun_route.h"
#include "topo_node.h"
#include "route_lib.h"
#include "routing.h"
#include "topo_group.h"
#include "nat_detection.h"
#include "etcp_connections.h"
#include "etcp.h"
#include "conn_mgr.h"
#include "chat/db_sync.h"
#include "chat/chat_core.h"
#include "chat/chat_sync.h"
#include "stcp_server.h"
#include "control_server.h"
#include "transport_layer/node_conn_direct.h"
#include "transport_layer/socket_monitor.h"
#include "transport_layer/auto_socket.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include <stdlib.h>
#include <stdio.h>
#include <string.h>
#include <errno.h>
#ifndef _WIN32
#include <unistd.h>
#endif
#include <unistd.h>
#include <sys/stat.h>
#include "../lib/platform_compat.h"
#include "../lib/mem.h"
#include "../lib/ll_queue.h"
#include <openssl/sha.h>
// Forward declarations
static uint32_t get_dest_ip(const uint8_t *packet, size_t len);
// Global instance for signal handlers
static struct UTUN_INSTANCE *g_instance = NULL;
// Global flag to control TUN initialization (disabled by default)
static int g_tun_init_enabled = 0;
// Function to control TUN initialization
void utun_instance_set_tun_init_enabled(int enabled) {
g_tun_init_enabled = enabled ? 1 : 0;
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN initialization %s", enabled ? "enabled" : "disabled");
}
// Global flag to control topo_group initialization (enabled by default)
static int g_topo_group_enabled = 1;
void utun_instance_set_topo_group_enabled(int enabled) {
g_topo_group_enabled = enabled ? 1 : 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_group initialization %s", enabled ? "enabled" : "disabled");
}
static int local_sockaddr_equal(const struct sockaddr_storage *a, const struct sockaddr_storage *b) {
if (!a || !b || a->ss_family != b->ss_family) return 0;
if (a->ss_family == AF_INET) {
const struct sockaddr_in *ia = (const struct sockaddr_in *)a;
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;
}
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)
// Returns 0 on success, -1 on error (instance is NOT freed on error - caller must handle)
static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* ua, struct utun_config* config) {
// Initialize basic fields
instance->running = 0;
instance->ua = ua;
instance->config = config;
// Set name from config
if (config->global.name[0] != '\0') {
strncpy(instance->name, config->global.name, sizeof(instance->name) - 1);
instance->name[sizeof(instance->name) - 1] = '\0';
} else {
instance->name[0] = '\0';
}
// Set node_id from config
instance->node_id = config->global.my_node_id;
// Set my keys
if (sc_init_local_keys(&instance->my_keys, config->global.my_public_key_hex, config->global.my_private_key_hex) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "Failed to initialize local keys");
return -1;
}
/* Always verify node_id matches pubkey, fix if mismatched */
{
uint8_t pub_bin[SC_PUBKEY_SIZE];
sc_compute_public_key_from_private(instance->my_keys.private_key, pub_bin);
uint64_t derived = sc_derive_node_id_from_pubkey(pub_bin);
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "node_id check: config=%016llx derived=%016llx match=%d",
(unsigned long long)instance->node_id, (unsigned long long)derived,
(instance->node_id == derived));
if (!instance->node_id || instance->node_id != derived) {
if (instance->node_id)
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "node_id mismatch: config=%016llx derived=%016llx — fixing",
(unsigned long long)instance->node_id, (unsigned long long)derived);
instance->node_id = derived;
config->global.my_node_id = derived;
}
}
if (sc_derive_ed25519_pubkey(instance->my_keys.private_key, instance->my_ed25519_pubkey) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "Failed to derive Ed25519 pubkey");
return -1;
}
{
uint8_t ed_priv_hash[SHA512_DIGEST_LENGTH];
SHA512(instance->my_keys.private_key, SC_PRIVKEY_SIZE, ed_priv_hash);
memcpy(instance->my_ed25519_privkey, ed_priv_hash, SC_PRIVKEY_SIZE);
}
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "Ed25519 pubkey derived successfully");
// Initialize networks queue (indexed by 56-bit network id)
instance->networks = queue_new(ua, 16, offsetof(struct NETWORK_ENTRY, id), 7, "networks");
if (!instance->networks) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Failed to create networks queue");
return -1;
}
// Initialize connection queue
instance->connections = queue_new(ua, 256, 0, 8, "connections");
if (!instance->connections) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create connections queue"); return -1; }
instance->tcp_connections = queue_new(ua, sizeof(struct tcp_conn_entry), sizeof(uint64_t), 8, "tcpconn");
if (!instance->tcp_connections) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create tcp_connections queue"); return -1; }
struct CFG_NETWORK *net = config->networks;
while (net) {
struct ll_entry *entry = queue_entry_new(sizeof(struct NETWORK_ENTRY));
if (!entry) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "Failed to allocate network entry for %s", net->name);
net = net->next;
continue;
}
struct NETWORK_ENTRY *ne = (struct NETWORK_ENTRY*)entry->data;
ne->id = net->id & 0x00FFFFFFFFFFFFFFULL;
for (int i = 0; i < SC_PUBKEY_SIZE; i++) {
unsigned int b; sscanf(net->pubkey_hex + i * 2, "%2x", &b); ne->pubkey[i] = (uint8_t)b;
}
for (int i = 0; i < SC_PRIVKEY_SIZE; i++) {
unsigned int b; sscanf(net->signing_key_hex + i * 2, "%2x", &b); ne->signing_key[i] = (uint8_t)b;
}
strncpy(ne->name, net->name, sizeof(ne->name) - 1);
ne->name[sizeof(ne->name) - 1] = '\0';
queue_data_put_with_index(instance->networks, entry);
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Network added: [%s] id=%012llx", ne->name, (unsigned long long)ne->id);
net = net->next;
}
// Create memory pools
instance->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
instance->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
instance->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
// Create routing module
if (routing_create(instance) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "Failed to create routing module");
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_ROUTING, "Routing module created");
if (g_tun_init_enabled && config->global.tun_enabled) {
instance->tun = tun_init(ua, config);
if (!instance->tun) {
DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to initialize TUN device");
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN interface initialized: %s", instance->tun->ifname);
if (config->route_subnets) {
int added = tun_route_add_all(instance->tun->ifindex, instance->tun->ifname, config->route_subnets);
DEBUG_INFO(DEBUG_CATEGORY_TUN, "Added %d system routes for TUN interface", added);
instance->route_subnets = config->route_subnets;
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN initialization disabled - skipping TUN device setup");
instance->tun = NULL;
}
int socket_result = init_sockets(instance);
instance->socket_init_status = socket_result;
if (socket_result != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Socket initialization failed: result=%d", socket_result);
return -1;
}
instance->nat_det = nat_detection_create(instance);
if (!instance->nat_det) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "NAT detection creation failed"); }
if (g_topo_group_enabled) {
instance->topo_groups = topo_groups_init(instance);
if (!instance->topo_groups) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "Failed to initialize BGP module");
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP module initialized");
etcp_router_bind(instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler);
etcp_bind(instance, ETCP_RT_ID_CONN_MGR, conn_mgr_direct_recv_handler);
struct TOPO_GROUP* g = topo_groups_get_default(instance->topo_groups);
if (instance->rt && g && g->local_node) {
topo_group_update_my_nodeinfo(instance, g);
}
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_group initialization disabled, skipping BGP module");
instance->topo_groups = NULL;
}
socket_monitor_init(instance);
auto_socket_init(instance);
// conn_mgr initialized inside topo_group_create (via topo_groups_init)
// Initialize firewall
fw_init(&instance->fw);
if (fw_load_rules(&instance->fw, &config->global) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "Failed to load firewall rules");
} else {
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Firewall initialized: %d rules, bypass_all=%d",
instance->fw.count, instance->fw.bypass_all);
}
// etcp_router — сервисная маршрутизация (после BGP, до TCP proxy/NAT/DATA)
if (etcp_router_init(instance) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to initialize etcp_router");
return -1;
}
// Bind DATA handler via etcp_router (after etcp_router_init)
if (routing_bind(instance) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "Failed to bind DATA via etcp_router");
return -1;
}
// TCP proxy server (exit node, optional) — must be before tcp_proxy_client so
// tcp_proxy_client_create can overwrite the handler if it has remote mappings
#ifndef _WIN32
if (tcp_proxy_server_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "Failed to initialize tcp_proxy_server (non-fatal)");
}
#endif
// TCP proxy client (from [tcp_proxy] config section)
#ifndef _WIN32
if (config->global.tcp_proxy_client_enabled || config->global.tcp_proxy_client_socks_enabled || config->global.tcp_proxy_client_http_proxy_enabled) {
const char* tun_name = config->global.tcp_proxy_client_tun_name;
const char* tun_ip = config->global.tcp_proxy_client_tun_ip;
int mtu = config->global.tcp_proxy_client_mtu;
instance->tcp_proxy_client = tcp_proxy_client_create(instance, ua, tun_name, tun_ip, mtu, g_tun_init_enabled ? 0 : 1,
config->global.tcp_proxy_client_mappings, config->global.tcp_proxy_client_mapping_count,
config->global.tcp_proxy_client_via_node_id,
config->global.tcp_proxy_client_socks_enabled, config->global.tcp_proxy_client_socks_addr,
config->global.tcp_proxy_client_http_proxy_enabled, config->global.tcp_proxy_client_http_proxy_addr);
if (instance->tcp_proxy_client) {
DEBUG_INFO(DEBUG_CATEGORY_TUN, "TCP proxy client enabled: TUN=%s IP=%s MTU=%d socks=%s http=%s mappings=%d",
tun_name, tun_ip, mtu,
config->global.tcp_proxy_client_socks_enabled ? config->global.tcp_proxy_client_socks_addr : "off",
config->global.tcp_proxy_client_http_proxy_enabled ? config->global.tcp_proxy_client_http_proxy_addr : "off",
config->global.tcp_proxy_client_mapping_count);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to create TCP proxy client");
return -1;
}
} else {
instance->tcp_proxy_client = NULL;
}
#endif /* _WIN32 */
return 0;
}
// Create and initialize root instance from config file
struct UTUN_INSTANCE* utun_instance_create(struct UASYNC* ua, const char *config_file) {
// Ensure keys and node_id exist in config (generates them if missing)
if (config_ensure_keys_and_node_id(config_file) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Failed to ensure keys and node_id in config: %s", config_file);
return NULL;
}
// Load configuration
struct utun_config* config = parse_config(config_file);
if (!config) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Failed to load config from %s", config_file);
return NULL;
}
// Open log file only if not using global debug system output
if (config->global.log_file[0]) {
debug_set_output_file(config->global.log_file);
}
// Apply debug level from config
if (config->global.debug_level[0]) {
debug_apply_global_level(config->global.debug_level);
}
// Apply per-category debug levels from [debug] section
for (int i = 0; i < config->global.debug_levels.count; i++) {
debug_apply_category_config(
config->global.debug_levels.category[i],
config->global.debug_levels.level[i]
);
}
// Allocate instance
struct UTUN_INSTANCE *instance = u_calloc(1, sizeof(struct UTUN_INSTANCE));
if (!instance) {
free_config(config);
return NULL;
}
instance->etcp_connect_timeout_tb = 20000;
// Init stats_dir from config file path
{
const char* last_sep = NULL;
for (const char* p = config_file; *p; p++) if (*p == '/' || *p == '\\') last_sep = p;
size_t dir_len = last_sep ? (size_t)(last_sep - config_file) : 1;
if (last_sep) {
memcpy(instance->stats_dir, config_file, dir_len);
instance->stats_dir[dir_len] = '\0';
} else {
instance->stats_dir[0] = '.'; instance->stats_dir[1] = '\0';
}
int remaining = (int)(sizeof(instance->stats_dir) - dir_len - 8);
if (remaining > 0) snprintf(instance->stats_dir + dir_len, remaining, "/stats");
utun_mkdir(instance->stats_dir, 0755);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Metrics stats dir: %s", instance->stats_dir);
}
// Initialize using common function
if (instance_init_common(instance, ua, config) != 0) {
// Cleanup on error
if (instance->config) free_config(instance->config);
u_free(instance);
return NULL;
}
// Log instance info
if (instance->name[0] != '\0') {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "uTun instance '%s' created, node_id=0x%llx",
instance->name, (unsigned long long)instance->node_id);
} else {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "uTun instance created, node_id=0x%llx",
(unsigned long long)instance->node_id);
}
return instance;
}
// Create instance from existing config structure (config ownership transfers to instance)
struct UTUN_INSTANCE* utun_instance_create_from_config(struct UASYNC* ua, struct utun_config* config) {
if (!config) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "NULL config");
return NULL;
}
// Allocate instance
struct UTUN_INSTANCE *instance = u_calloc(1, sizeof(struct UTUN_INSTANCE));
if (!instance) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "Failed to allocate UTUN_INSTANCE");
return NULL;
}
instance->etcp_connect_timeout_tb = 20000;
// Initialize using common function (config ownership transferred to instance)
if (instance_init_common(instance, ua, config) != 0) {
// Cleanup on error - caller still owns config since we failed
u_free(instance);
return NULL;
}
return instance;
}
// Create instance from config text string (writes to temp file, parses, cleans up)
struct UTUN_INSTANCE* utun_instance_create_from_str(struct UASYNC* ua, const char* config_text) {
if (!config_text) return NULL;
char tmp_path[] = "/tmp/utun_cfg_XXXXXX";
int fd = mkstemp(tmp_path);
if (fd < 0) return NULL;
size_t len = strlen(config_text);
if (write(fd, config_text, len) != (ssize_t)len) { close(fd); unlink(tmp_path); return NULL; }
close(fd);
struct UTUN_INSTANCE* inst = utun_instance_create(ua, tmp_path);
unlink(tmp_path);
return inst;
}
// Destroy instance and cleanup resources
void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
if (!instance) return;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Starting cleanup for instance %p", instance);
// Диагностика ресурсов ДО cleanup
utun_instance_diagnose_leaks(instance, "BEFORE_CLEANUP");
if (instance->ua) uasync_print_resources(instance->ua, "INSTANCE_DESTROY_BEFORE");
// Stop running if not already
instance->running = 0;
// Cancel NTP timer
ntp_time_destroy(instance);
ntp_node_time_destroy(instance);
// Shutdown control server first
if (instance->control_srv) {
DEBUG_INFO(DEBUG_CATEGORY_CONTROL, "Shutting down control server");
control_server_shutdown(instance->control_srv);
u_free(instance->control_srv);
instance->control_srv = NULL;
}
// Close config-based connection handles (sends CLOSE before sockets die)
struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles;
while (ch) {
struct CONFIG_CONN_HANDLE* next = ch->next;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[INSTANCE_DESTROY] closing config handle for node=0x%016llx", (unsigned long long)ch->node_id);
node_conn_direct_close(ch->handle);
u_free(ch);
ch = next;
}
instance->config_conn_handles = NULL;
socket_monitor_destroy(instance);
auto_socket_destroy(instance);
// Cleanup BGP module BEFORE sockets (needs live conn_mgr for recovery cleanup)
if (instance->topo_groups) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Destroying BGP module");
etcp_router_unbind(instance, ETCP_RT_ID_CONN_MGR);
etcp_unbind(instance, ETCP_RT_ID_CONN_MGR);
topo_groups_destroy(instance);
}
// Cleanup conn_mgr BEFORE sockets (unbinds from etcp_router, closes NCD handles while conns alive)
// Cleanup ETCP sockets and connections FIRST (before destroying uasync)
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Cleaning up ETCP sockets and connections");
struct ETCP_SOCKET* sock = instance->etcp_sockets;
while (sock) {
struct ETCP_SOCKET* next = sock->next;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Removing socket %p, fd=%d", sock, sock->fd);
etcp_socket_remove(sock); // Полный cleanup сокета
sock = next;
}
instance->etcp_sockets = NULL;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] ETCP sockets cleanup complete");
// Cleanup chat before db_sync (chat releases db_sync instances)
chat_sync_destroy(instance);
chat_core_destroy(instance);
// Cleanup db_sync before connections cleanup (db_sync_destroy iterates connections)
db_sync_destroy(instance);
// Cleanup ETCP connections (phase 1 detach + deferred phase 2 via call_soon)
{
struct ll_entry* entry = instance->connections->head;
while (entry) {
struct ll_entry* next = entry->next;
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (ce && ce->conn) etcp_connection_close(ce->conn);
entry = next;
}
}
queue_free(instance->connections); instance->connections = NULL;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] ETCP connections cleanup complete");
// Cleanup TCP connections (stcp_links)
if (instance->tcp_connections) {
struct ll_entry* entry = instance->tcp_connections->head;
while (entry) {
struct ll_entry* next = entry->next;
struct tcp_conn_entry* te = (struct tcp_conn_entry*)entry->data;
if (te && te->link) stcp_link_close(te->link);
entry = next;
}
}
queue_free(instance->tcp_connections); instance->tcp_connections = NULL;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] TCP connections cleanup complete");
// Cleanup TCP sockets
while (instance->tcp_sockets) tcp_socket_remove(instance->tcp_sockets);
// Wait for all deferred callbacks (etcp_connection_free_deferred)
// to finish before destroying pools
while (instance->ua && instance->ua->immediate_queue_head)
uasync_poll(instance->ua, 0);
struct PING_CONTEXT* p = instance->pending_pings;
while (p) {
struct PING_CONTEXT* next = p->next;
if (p->timeout_timer) uasync_cancel_timeout(instance->ua, p->timeout_timer);
u_free(p);
p = next;
}
instance->pending_pings = NULL;
// Cleanup TUN
if (instance->tun) {
DEBUG_INFO(DEBUG_CATEGORY_TUN, "Closing TUN interface: %s", instance->tun->ifname);
// Delete system routes added at startup
if (instance->tun->ifindex && instance->route_subnets) {
int deleted = tun_route_del_all(instance->tun->ifindex, instance->tun->ifname, instance->route_subnets);
DEBUG_INFO(DEBUG_CATEGORY_TUN, "Deleted %d system routes for TUN interface", deleted);
}
tun_close(instance->tun);
instance->tun = NULL;
}
// Cleanup TCP proxy client module
#ifndef _WIN32
if (instance->tcp_proxy_client) {
DEBUG_INFO(DEBUG_CATEGORY_TUN, "Destroying TCP proxy client module");
tcp_proxy_client_destroy(instance->tcp_proxy_client);
instance->tcp_proxy_client = NULL;
}
// Cleanup TCP proxy server
tcp_proxy_server_destroy(instance);
#endif
// Cleanup routing module (unbinds from etcp_router before etcp_router_destroy)
routing_destroy(instance);
// Cleanup media delivery (unbinds from etcp_router before etcp_router_destroy)
if (instance->md.initialized) {
media_delivery_destroy(instance);
}
// Cleanup media async engine
media_async_destroy(instance->media_async);
instance->media_async = NULL;
// Cleanup NAT (unbinds from etcp_router before etcp_router_destroy)
if (instance->nat_tr.initialized) {
nat_transport_destroy(instance);
}
// Cleanup etcp_router
etcp_router_destroy(instance);
// Cleanup NAT detection
if (instance->nat_det) {
nat_detection_destroy(instance->nat_det);
instance->nat_det = NULL;
}
// Cleanup firewall
fw_free(&instance->fw);
// Cleanup TCP servers
stcp_server_list_destroy_all(instance);
// Cleanup networks queue
if (instance->networks) {
struct ll_entry *entry;
while ((entry = queue_data_get(instance->networks)) != NULL) {
queue_entry_free(entry);
}
queue_free(instance->networks);
instance->networks = NULL;
}
// Cleanup config
if (instance->config) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Freeing configuration");
free_config(instance->config);
instance->config = NULL;
}
// Cleanup packet pool (ensure no leak if stop wasn't called)
if (instance->pkt_pool) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Destroying packet pool");
memory_pool_destroy(instance->pkt_pool);
instance->pkt_pool = NULL;
}
// Cleanup packet pool (ensure no leak if stop wasn't called)
if (instance->ack_pool) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Destroying ack pool");
memory_pool_destroy(instance->ack_pool);
instance->ack_pool = NULL;
}
// Cleanup data pool
if (instance->data_pool) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Destroying data pool");
memory_pool_destroy(instance->data_pool);
instance->data_pool = NULL;
}
if (instance->ua) uasync_print_resources(instance->ua, "INSTANCE_DESTROY_AFTER");
// Note: uasync is NOT destroyed here - caller must destroy it separately
// This allows sharing uasync between multiple instances
instance->ua = NULL;
// Clear global instance
if (g_instance == instance) {
g_instance = NULL;
}
// Free the instance memory
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Freeing instance memory");
// Cleanup nodeinfo callback chain
while (instance->nodeinfo_cbks) {
struct nodeinfo_cbk_entry* r = instance->nodeinfo_cbks;
instance->nodeinfo_cbks = r->next; u_free(r);
}
u_free(instance);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[INSTANCE_DESTROY] Instance destroyed completely");
}
// Stop instance
void utun_instance_stop(struct UTUN_INSTANCE *instance) {
if (!instance) return;
instance->running = 0;
// Wakeup main loop using built-in uasync wakeup
if (instance->ua) {
memory_pool_destroy(instance->pkt_pool);
uasync_wakeup(instance->ua);
}
}
static int mkdir_recursive(const char *path) {
char tmp[512]; int len = snprintf(tmp, sizeof(tmp), "%s", path);
if (len <= 0 || (size_t)len >= sizeof(tmp)) return -1;
for (char *p = tmp + 1; *p; p++) {
if (*p == '/') { *p = '\0'; utun_mkdir(tmp, 0755); *p = '/'; }
}
return utun_mkdir(tmp, 0755);
}
int utun_instance_init(struct UTUN_INSTANCE *instance) {
if (!instance) return -1;
// db_sync — распределённая таблица с репликацией
db_sync_init(instance);
// chat — сообщения, каналы, P2P-синхронизация
if (instance->config->global.chatserver_enabled) {
if (!instance->config->global.db_path[0])
snprintf(instance->config->global.db_path, sizeof(instance->config->global.db_path), "/var/lib/utun");
mkdir_recursive(instance->config->global.db_path);
char db_file[512];
snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path);
if (chat_core_init(instance, db_file) == 0)
chat_sync_init(instance);
}
// Set TUN interface in routing module
if (instance->tun) {
routing_set_tun(instance);
DEBUG_INFO(DEBUG_CATEGORY_ROUTING, "TUN interface registered in routing module");
}
// Initialize NAT (after routing_set_tun, before connections)
if (instance->config->global.nat_enabled) {
int nat_ret = nat_transport_init(instance);
if (nat_ret != 0) {
DEBUG_WARN(DEBUG_CATEGORY_NAT, "NAT transport init failed (non-fatal)");
}
}
// Initialize media delivery
if (media_delivery_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "media_delivery init failed (non-fatal)");
}
// Initialize media async engine
instance->media_async = media_async_create();
if (!instance->media_async) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "media_async create failed");
}
// Note: TUN socket is already registered in tun_init()
// Initialize connections
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "utun_instance_init() calling init_connections() for instance %p", instance);
int conn_result = init_connections(instance);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "init_connections() returned: %d", conn_result);
if (conn_result < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to initialize connections, error=%d", conn_result);
return -1;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully, count=%d", queue_entry_count(instance->connections));
// Initialize control server if configured
if (instance->config->global.control_sock.ss_family != 0) {
instance->control_srv = (struct control_server*)u_calloc(1, sizeof(struct control_server));
if (instance->control_srv) {
if (control_server_init(instance->control_srv, instance->ua, instance,
&instance->config->global.control_sock, 8) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CONTROL, "Failed to initialize control server, continuing without monitoring");
u_free(instance->control_srv);
instance->control_srv = NULL;
} else {
DEBUG_INFO(DEBUG_CATEGORY_CONTROL, "Control server initialized successfully");
}
}
}
// Start the main loop
instance->running = 1;
// Initialize NTP time sync (non-fatal if fails)
if (ntp_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP time sync init failed (non-fatal)");
}
// Initialize NTP node time exchange (non-fatal)
if (ntp_node_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP node time exchange init failed (non-fatal)");
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully");
return 0;
}
// Диагностическая функция для анализа утечек
void utun_instance_diagnose_leaks(struct UTUN_INSTANCE *instance, const char *phase) {
if (!instance) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "[DIAGNOSE] NULL instance for phase: %s", phase);
return;
}
struct {
int etcp_sockets_count;
int etcp_connections_count;
int etcp_links_count;
} report = {0};
// Подсчёт ETCP сокетов
struct ETCP_SOCKET *sock = instance->etcp_sockets;
while (sock) {
report.etcp_sockets_count++;
report.etcp_links_count += queue_entry_count(sock->links_queue);
sock = sock->next;
}
// Подсчёт ETCP соединений
{
struct ll_entry* entry = instance->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
report.etcp_connections_count++;
struct ETCP_LINK *link = ce->conn->links;
while (link) { report.etcp_links_count++; link = link->next; }
entry = entry->next;
}
}
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[LEAK DIAGNOSIS] Phase: %s", phase);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "Instance: %p, Node ID: %llu, UA: %p, Running: %d",
instance, (unsigned long long)instance->node_id, instance->ua, instance->running);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "STRUCTURE COUNTS - ETCP Sockets: %d, ETCP Connections: %d, ETCP Links: %d",
report.etcp_sockets_count, report.etcp_connections_count, report.etcp_links_count);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "RESOURCE STATUS - Memory Pool: %s, TUN: %s, TUN FD: %d",
instance->pkt_pool ? "ALLOCATED" : "NULL",
instance->tun ? instance->tun->ifname : "NULL",
instance->tun ? instance->tun->fd : -1);
if (instance->pkt_pool) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "POTENTIAL LEAK: Memory Pool not freed");
}
if (instance->tun) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "POTENTIAL LEAK: TUN interface not closed");
}
if (report.etcp_sockets_count > 0) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "POTENTIAL LEAK: %d ETCP sockets still allocated", report.etcp_sockets_count);
}
if (report.etcp_connections_count > 0) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "POTENTIAL LEAK: %d ETCP connections still allocated", report.etcp_connections_count);
}
if (report.etcp_links_count > 0) {
DEBUG_INFO(DEBUG_CATEGORY_SYS, "POTENTIAL LEAK: %d ETCP links still allocated", report.etcp_links_count);
}
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[LEAK DIAGNOSIS] Recommendations: pkt_pool=%d, tun=%d, sockets=%d, connections=%d",
instance->pkt_pool ? 1 : 0,
instance->tun ? 1 : 0,
report.etcp_sockets_count,
report.etcp_connections_count);
}
/* Variant B: selective reload. Compares settings for [server:], [client:], link=. Unchanged untouched (no close, no timers, no queues). Changed/deleted: target only that object. New: add. Routing/fw always. BGP self-restarts. */
struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struct UASYNC *ua, const char *config_file) {
if (!instance || !ua || !config_file) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "invalid arguments");
return NULL;
}
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "SIGHUP selective reload from %s", config_file);
debug_reopen_log();
struct utun_config *new_config = parse_config(config_file);
if (!new_config) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Failed to parse new config");
return instance;
}
// Reset + apply debug categories from [debug] section (always, for partial and full paths).
// Mirrors utun_instance_create exactly. Reset ensures removed categories fall back to global.
for (int i = 0; i < DEBUG_CATEGORY_COUNT; i++) {
g_debug_config.category_levels[i] = DEBUG_LEVEL_NONE;
}
if (new_config->global.debug_level[0]) {
debug_apply_global_level(new_config->global.debug_level);
}
for (int i = 0; i < new_config->global.debug_levels.count; i++) {
debug_apply_category_config(
new_config->global.debug_levels.category[i],
new_config->global.debug_levels.level[i]
);
}
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Applied [debug] section (%d categories)", new_config->global.debug_levels.count);
int needs_full = 0;
if (strcmp(instance->config->global.my_public_key_hex, new_config->global.my_public_key_hex) != 0 ||
strcmp(instance->config->global.my_private_key_hex, new_config->global.my_private_key_hex) != 0 ||
instance->config->global.my_node_id != new_config->global.my_node_id ||
strcmp(instance->config->global.tun_ifname, new_config->global.tun_ifname) != 0) {
needs_full = 1;
}
if (needs_full) {
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Incompatible global - full reload");
utun_instance_destroy(instance);
struct UTUN_INSTANCE *new_instance = utun_instance_create(ua, config_file);
if (!new_instance || utun_instance_init(new_instance) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "Full reload failed");
if (new_instance) utun_instance_destroy(new_instance);
free_config(new_config);
return NULL;
}
new_instance->running = 1;
free_config(new_config);
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Full reload completed");
return new_instance;
}
// Selective Variant B (settings compare, only affected updated)
// Servers [server:]
for (struct CFG_SERVER *ns = new_config->servers; ns; ns = ns->next) {
int matched = 0;
for (struct ETCP_SOCKET *os = instance->etcp_sockets; os; os = os->next) {
if (strcmp(ns->name, os->name) == 0 && local_sockaddr_equal(&ns->ip, &os->local_addr)) {
matched = 1;
if (ns->mtu != os->mtu) {
os->mtu = ns->mtu;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Updated socket settings for %s", ns->name);
}
break;
}
}
if (!matched) {
etcp_socket_add(instance, ns);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Added new socket %s", ns->name);
}
}
struct ETCP_SOCKET *os = instance->etcp_sockets;
while (os) {
struct ETCP_SOCKET *next = os->next;
int found = 0;
for (struct CFG_SERVER *ns = new_config->servers; ns; ns = ns->next) {
if (strcmp(ns->name, os->name) == 0) {
found = 1; break;
}
}
if (!found) {
etcp_socket_remove(os);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Removed deleted socket %s", os->name);
}
os = next;
}
// Clients [client:] — close old config handles, rebuild from new config
struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles;
while (ch) { struct CONFIG_CONN_HANDLE* next = ch->next; node_conn_direct_close(ch->handle); u_free(ch); ch = next; }
instance->config_conn_handles = NULL;
for (struct CFG_CLIENT *nc = new_config->clients; nc; nc = nc->next) {
if (strlen(nc->peer_public_key_hex) == 0) { DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "client %s no peer pubkey", nc->name); continue; }
uint8_t pubkey_bin[SC_PUBKEY_SIZE];
if (sc_hex_to_binary(nc->peer_public_key_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid pubkey hex client %s", nc->name); continue; }
uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey_bin);
struct NODE_CONN_DIRECT* handle = NULL;
struct ETCP_CONN* conn = NULL;
for (struct CFG_CLIENT_LINK *nl = nc->links; nl; nl = nl->next) {
struct ETCP_SOCKET* sock = NULL;
for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next)
if (nl->local_srv && strcmp(nl->local_srv->name, s->name) == 0) { sock = s; break; }
if (!sock) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no socket for server '%s' client %s", nl->local_srv ? nl->local_srv->name : "?", nc->name); continue; }
if (!handle) {
struct TOPO_ADDR4 v4_addr; struct TOPO_ADDR6 v6_addr;
struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp));
ni_tmp.node_id = node_id; ni_tmp.node_name = nc->name;
memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE);
if (nl->remote_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&nl->remote_addr;
v4_addr = (struct TOPO_ADDR4){0}; memcpy(v4_addr.addr, &sin->sin_addr.s_addr, 4);
v4_addr.port = ntohs(sin->sin_port); v4_addr.type = TOPO_ADDR_INTERFACE; v4_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v4_addrs = &v4_addr;
} else if (nl->remote_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&nl->remote_addr;
v6_addr = (struct TOPO_ADDR6){0}; memcpy(v6_addr.addr, &sin6->sin6_addr, 16);
v6_addr.port = ntohs(sin6->sin6_port); v6_addr.type = TOPO_ADDR_INTERFACE; v6_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v6_addrs = &v6_addr;
}
int rc = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, sock);
if (rc == NCD_ERR || !handle) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ncd reload open failed client %s", nc->name); break; }
conn = node_conn_direct_get_conn(handle);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "reload: client %s ncd %s node=0x%016llx", nc->name, rc == NCD_NEW ? "NEW" : "REUSED", (unsigned long long)node_id);
} else {
etcp_link_new(conn, sock, &nl->remote_addr, 0);
}
}
if (handle) {
ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, nc->name, MAX_CONN_NAME_LEN - 1);
ch->handle = handle; ch->next = instance->config_conn_handles; instance->config_conn_handles = ch; }
}
}
// Reload networks: clear and repopulate
if (instance->networks) {
struct ll_entry *entry;
while ((entry = queue_data_get(instance->networks)) != NULL) {
queue_entry_free(entry);
}
struct CFG_NETWORK *net = new_config->networks;
while (net) {
struct ll_entry *new_entry = queue_entry_new(sizeof(struct NETWORK_ENTRY));
if (!new_entry) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "Failed to allocate network entry for reload: %s", net->name);
net = net->next;
continue;
}
struct NETWORK_ENTRY *ne = (struct NETWORK_ENTRY*)new_entry->data;
ne->id = net->id & 0x00FFFFFFFFFFFFFFULL;
for (int i = 0; i < SC_PUBKEY_SIZE; i++) {
unsigned int b; sscanf(net->pubkey_hex + i * 2, "%2x", &b); ne->pubkey[i] = (uint8_t)b;
}
for (int i = 0; i < SC_PRIVKEY_SIZE; i++) {
unsigned int b; sscanf(net->signing_key_hex + i * 2, "%2x", &b); ne->signing_key[i] = (uint8_t)b;
}
strncpy(ne->name, net->name, sizeof(ne->name) - 1);
ne->name[sizeof(ne->name) - 1] = '\0';
queue_data_put_with_index(instance->networks, new_entry);
net = net->next;
}
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Networks reloaded");
}
fw_free(&instance->fw);
fw_init(&instance->fw);
fw_load_rules(&instance->fw, &new_config->global);
// routing_destroy(instance); // commented per request - BGP self-restarts on conn events
// routing_create(instance); // commented per request
if (instance->config) free_config(instance->config);
instance->config = new_config;
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Selective partial reload completed (only changed updated, unchanged untouched, [debug] applied)");
return instance;
}