Browse Source

socket_monitor: extract socket monitoring from etcp_connections, add per-instance context

- New files: transport_layer/socket_monitor.c/h — socket state monitoring with per-instance ctx
- etcp_api: expose get_latency/1s_pkt_count measuring, add etcp_hashrnd for tests
- utun_instance: store socket_monitor per instance, add is_socket_ok check
- topo_node: detect IPv4 addresses with port from listening sockets
- topo_group: use node_name for display in routes
- Makefile.am: add socket_monitor.c to build
topo_upd
evgeny 2 months ago
parent
commit
f5efaa495b
  1. 2
      src/Makefile.am
  2. 3
      src/routing_layer/topo_group.c
  3. 85
      src/routing_layer/topo_node.c
  4. 2
      src/routing_layer/topo_node.h
  5. 27
      src/transport_layer/etcp_api.c
  6. 17
      src/transport_layer/etcp_api.h
  7. 261
      src/transport_layer/socket_monitor.c
  8. 16
      src/transport_layer/socket_monitor.h
  9. 5
      src/utun_instance.c
  10. 2
      src/utun_instance.h

2
src/Makefile.am

@ -44,6 +44,7 @@ utun_CORE_SOURCES = \
transport_layer/etcp_api.c \ transport_layer/etcp_api.c \
transport_layer/etcp_connect.c \ transport_layer/etcp_connect.c \
transport_layer/node_conn_direct.c \ transport_layer/node_conn_direct.c \
transport_layer/socket_monitor.c \
control_server.c \ control_server.c \
firewall.c \ firewall.c \
eim_nat.c \ eim_nat.c \
@ -118,6 +119,7 @@ libutun_a_SOURCES = \
transport_layer/etcp_api.c \ transport_layer/etcp_api.c \
transport_layer/etcp_connect.c \ transport_layer/etcp_connect.c \
transport_layer/node_conn_direct.c \ transport_layer/node_conn_direct.c \
transport_layer/socket_monitor.c \
control_server.c \ control_server.c \
firewall.c \ firewall.c \
eim_nat.c \ eim_nat.c \

3
src/routing_layer/topo_group.c

@ -310,6 +310,8 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
etcp_bind(instance, ETCP_ID_TOPO_ENTRY, topo_group_receive_cbk); etcp_bind(instance, ETCP_ID_TOPO_ENTRY, topo_group_receive_cbk);
etcp_add_conn_status_cbk(instance, topo_group_conn_status, g); 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);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "TOPO groups initialized (with default group)"); DEBUG_INFO(DEBUG_CATEGORY_BGP, "TOPO groups initialized (with default group)");
return g; return g;
@ -322,6 +324,7 @@ void topo_groups_destroy(struct UTUN_INSTANCE* instance) {
etcp_unbind(instance, ETCP_ID_TOPO_ENTRY); etcp_unbind(instance, ETCP_ID_TOPO_ENTRY);
etcp_remove_conn_status_cbk(instance, topo_group_conn_status, instance->topo_groups); etcp_remove_conn_status_cbk(instance, topo_group_conn_status, instance->topo_groups);
etcp_remove_socket_cbk(instance, topo_node_on_socket_changed, NULL);
route_connectivity_cancel_all(instance); route_connectivity_cancel_all(instance);
struct TOPO_GROUPS* g = instance->topo_groups; struct TOPO_GROUPS* g = instance->topo_groups;

85
src/routing_layer/topo_node.c

@ -775,3 +775,88 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
} }
return vc; return vc;
} }
void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return; }
struct TOPO_GROUP* default_group = topo_groups_get_default(instance->topo_groups);
if (!default_group || !default_group->local_node) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "no local_node"); return; }
struct TOPO_NODE* ni = topo_node_registry_find(instance->topo_groups, instance->node_id);
if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "my node not in registry"); return; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "updating my addresses, old ver=%d", ni->ver);
free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, ni->v4_sock_meta); ni->v4_sock_meta = NULL;
free_v4_addr_list(instance->topo_groups->v4_addr_pool, ni->v4_addrs); ni->v4_addrs = NULL;
free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, ni->v6_sock_meta); ni->v6_sock_meta = NULL;
free_v6_addr_list(instance->topo_groups->v6_addr_pool, ni->v6_addrs); ni->v6_addrs = NULL;
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock) {
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; }
if (e_sock->local_addr.ss_family == AF_INET) {
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(instance->topo_groups->v4_sock_meta_pool);
sm->id = e_sock->sock_id; sm->config_type = e_sock->type; sm->nat_type = e_sock->nat_type;
sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; }
struct sockaddr_in* local_sin = (struct sockaddr_in*)&e_sock->local_addr;
struct sockaddr_in* if_sin = (struct sockaddr_in*)&e_sock->interface_addr;
struct sockaddr_in* nat_sin = (struct sockaddr_in*)&e_sock->nat_addr;
int use_local = (local_sin->sin_addr.s_addr != 0);
uint32_t ref_ip = use_local ? local_sin->sin_addr.s_addr : if_sin->sin_addr.s_addr;
uint16_t ref_port = ntohs(use_local ? local_sin->sin_port : if_sin->sin_port);
int nat_verified = (e_sock->nat_addr.ss_family == AF_INET && nat_sin->sin_addr.s_addr != 0 && e_sock->nat_type != NAT_VERIFIED_STRICT);
int nat_differs = (nat_verified && (nat_sin->sin_addr.s_addr != ref_ip || ntohs(nat_sin->sin_port) != ref_port));
{ struct TOPO_ADDR4* a = memory_pool_alloc(instance->topo_groups->v4_addr_pool);
memcpy(a->addr, use_local ? &local_sin->sin_addr.s_addr : &if_sin->sin_addr.s_addr, 4);
a->port = ntohs(use_local ? local_sin->sin_port : if_sin->sin_port);
a->type = (!nat_differs && nat_verified) ? TOPO_ADDR_NAT : TOPO_ADDR_INTERFACE;
a->socket_id = (a->type == TOPO_ADDR_NAT) ? (e_sock->sock_id | 1) : e_sock->sock_id;
a->protocol = TOPO_PROTO_UDP;
a->next = ni->v4_addrs; ni->v4_addrs = a; }
if (nat_differs) {
struct TOPO_ADDR4* a = memory_pool_alloc(instance->topo_groups->v4_addr_pool);
memcpy(a->addr, &nat_sin->sin_addr.s_addr, 4); a->port = ntohs(nat_sin->sin_port);
a->type = TOPO_ADDR_NAT; a->socket_id = e_sock->sock_id | 1; a->protocol = TOPO_PROTO_UDP;
a->next = ni->v4_addrs; ni->v4_addrs = a;
}
} else if (e_sock->local_addr.ss_family == AF_INET6) {
{ struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(instance->topo_groups->v6_sock_meta_pool);
sm6->id = e_sock->sock_id; sm6->config_type = e_sock->type; sm6->nat_type = e_sock->nat_type;
sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; }
struct sockaddr_in6* if_sin6 = (struct sockaddr_in6*)&e_sock->interface_addr;
{ struct TOPO_ADDR6* a6 = memory_pool_alloc(instance->topo_groups->v6_addr_pool);
memcpy(a6->addr, &if_sin6->sin6_addr, 16); a6->port = ntohs(if_sin6->sin6_port);
a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = e_sock->sock_id; a6->protocol = TOPO_PROTO_UDP;
a6->next = ni->v6_addrs; ni->v6_addrs = a6; }
}
e_sock = e_sock->next;
}
ni->ver = (ni->ver % 255) + 1;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "my addresses updated, new ver=%d", ni->ver);
struct ll_entry* ge = instance->topo_groups->group_list->head;
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
if (g->local_node) {
g->local_node->dirty = 1;
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_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg) {
(void)arg;
if (!sock || !sock->instance) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "socket %s changed (event=0x%x), updating nodeinfo", sock->name, event);
topo_node_update_my_addresses(sock->instance);
if (sock->instance->topo_sqlite_db)
topo_node_sqlite_nodeinfo_updated(sock->instance->topo_sqlite_db, sock->instance->node_id);
}

2
src/routing_layer/topo_node.h

@ -181,6 +181,8 @@ int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t
uint16_t* out_cumulative_rtt); uint16_t* out_cumulative_rtt);
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group); int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group);
void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance);
void topo_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg);
void topo_node_dump_all(struct TOPO_GROUP* group); void topo_node_dump_all(struct TOPO_GROUP* group);
int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size); int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size);
int topo_node_ping_request_cbk(struct TOPO_GROUPS* groups, uint64_t node_id); int topo_node_ping_request_cbk(struct TOPO_GROUPS* groups, uint64_t node_id);

27
src/transport_layer/etcp_api.c

@ -82,6 +82,33 @@ static void etcp_status_cbk_remove_chain(struct etcp_status_cbk_entry** head, et
} }
void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { if (inst) etcp_status_cbk_add_chain(&inst->conn_status_cbks, fn, arg); } void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { if (inst) etcp_status_cbk_add_chain(&inst->conn_status_cbks, fn, arg); }
void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { if (inst) etcp_status_cbk_remove_chain(&inst->conn_status_cbks, fn, arg); } void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg) { if (inst) etcp_status_cbk_remove_chain(&inst->conn_status_cbks, fn, arg); }
static void etcp_socket_cbk_add_chain(struct etcp_socket_cbk_entry** head, etcp_socket_cbk_fn fn, void* arg, int event_mask) {
if (!head || !fn) return;
struct etcp_socket_cbk_entry* e = u_malloc(sizeof(struct etcp_socket_cbk_entry));
if (!e) return;
e->fn = fn; e->arg = arg; e->event_mask = event_mask; e->next = *head;
*head = e;
}
static void etcp_socket_cbk_remove_chain(struct etcp_socket_cbk_entry** head, etcp_socket_cbk_fn fn, void* arg) {
if (!head || !fn) return;
struct etcp_socket_cbk_entry** p = head;
while (*p) {
if ((*p)->fn == fn && (*p)->arg == arg) { struct etcp_socket_cbk_entry* rm = *p; *p = rm->next; u_free(rm); return; }
p = &(*p)->next;
}
}
void etcp_add_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, void* arg, int event_mask) { if (inst) etcp_socket_cbk_add_chain(&inst->socket_cbks, fn, arg, event_mask); }
void etcp_remove_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, void* arg) { if (inst) etcp_socket_cbk_remove_chain(&inst->socket_cbks, fn, arg); }
void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event) {
if (!sock || !sock->instance) return;
static const char* names[] = { "ADDR_CHANGED", "STATUS_CHANGED" };
int idx = 0, e = event; while (e >>= 1) idx++;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket %s event: %s", sock->name, (idx >= 0 && idx < (int)(sizeof(names)/sizeof(names[0]))) ? names[idx] : "?");
struct etcp_socket_cbk_entry* cbe = sock->instance->socket_cbks;
while (cbe) { struct etcp_socket_cbk_entry* n = cbe->next; if (cbe->event_mask & event) cbe->fn(sock, event, cbe->arg); cbe = n; }
}
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) { void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) {
if (!conn) return; if (!conn) return;
conn->routing_exchange_active = new_state; conn->routing_exchange_active = new_state;

17
src/transport_layer/etcp_api.h

@ -54,6 +54,7 @@ extern "C" {
// Forward declarations // Forward declarations
struct ETCP_CONN; struct ETCP_CONN;
struct ETCP_SOCKET;
struct UTUN_INSTANCE; struct UTUN_INSTANCE;
typedef void (*etcp_conn_status_fn)(struct ETCP_CONN* conn, int status, void* arg); typedef void (*etcp_conn_status_fn)(struct ETCP_CONN* conn, int status, void* arg);
@ -188,6 +189,22 @@ int etcp_unbind(struct UTUN_INSTANCE* inst, uint8_t id);
void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); void etcp_add_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg);
void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg); void etcp_remove_conn_status_cbk(struct UTUN_INSTANCE* inst, etcp_conn_status_fn fn, void* arg);
/* Socket property change events (instance-level) */
#define ETCP_SOCKET_EVENT_ADDR_CHANGED (1 << 0)
#define ETCP_SOCKET_EVENT_STATUS_CHANGED (1 << 1)
typedef void (*etcp_socket_cbk_fn)(struct ETCP_SOCKET* sock, int event, void* arg);
struct etcp_socket_cbk_entry {
etcp_socket_cbk_fn fn;
void* arg;
int event_mask;
struct etcp_socket_cbk_entry* next;
};
void etcp_add_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, void* arg, int event_mask);
void etcp_remove_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, void* arg);
void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event);
void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state); void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state);
/** /**

261
src/transport_layer/socket_monitor.c

@ -0,0 +1,261 @@
#include "socket_monitor.h"
#include "etcp_connections.h"
#include "etcp_api.h"
#include "../lib/socket_compat.h"
#include "../lib/platform_compat.h"
#include "../lib/getmyip.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "utun_instance.h"
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <errno.h>
#ifndef _WIN32
#include <linux/netlink.h>
#include <linux/rtnetlink.h>
#include <net/if.h>
#endif
struct SOCKET_MONITOR {
struct UTUN_INSTANCE* instance;
socket_t nl_sock;
void* uasync_handle;
int seq;
uint32_t addr_changes;
uint32_t link_changes;
uint32_t route_changes;
uint32_t socket_updates;
};
static void socket_monitor_read_cb(int fd, void* arg);
#ifndef _WIN32
static void socket_monitor_update_if_addr(struct ETCP_SOCKET* es) {
if (es->local_addr.ss_family == AF_INET) {
if (es->netif_index > 0) {
uint32_t if_ip = get_interface_ip_by_index(es->netif_index);
if (if_ip != 0) {
struct sockaddr_in* sin = (struct sockaddr_in*)&es->interface_addr;
sin->sin_family = AF_INET;
sin->sin_addr.s_addr = if_ip;
}
} else {
struct sockaddr_in* sin = (struct sockaddr_in*)&es->local_addr;
if (sin->sin_addr.s_addr == 0) {
struct sockaddr_storage remote;
memset(&remote, 0, sizeof(remote));
struct sockaddr_in* rem_sin = (struct sockaddr_in*)&remote;
rem_sin->sin_family = AF_INET;
rem_sin->sin_port = htons(53);
inet_pton(AF_INET, "8.8.8.8", &rem_sin->sin_addr);
get_outgoing_local_ip(&es->local_addr, &remote, &es->interface_addr);
}
}
} else if (es->local_addr.ss_family == AF_INET6) {
if (es->netif_index > 0) {
uint8_t v6addr[16];
int got = get_interface_ipv6_addr_nl(es->netif_index, 1, v6addr);
if (got != 0) got = get_interface_ipv6_addr_nl(es->netif_index, 0, v6addr);
if (got == 0) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&es->interface_addr;
sin6->sin6_family = AF_INET6;
memcpy(&sin6->sin6_addr, v6addr, 16);
}
} else {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&es->local_addr;
if (memcmp(&sin6->sin6_addr, &in6addr_any, 16) == 0) {
uint16_t v6if = get_default_route_netif_index(AF_INET6);
if (v6if > 0) {
uint8_t v6addr[16];
int got = get_interface_ipv6_addr_nl(v6if, 1, v6addr);
if (got != 0) got = get_interface_ipv6_addr_nl(v6if, 0, v6addr);
if (got == 0) {
struct sockaddr_in6* if6 = (struct sockaddr_in6*)&es->interface_addr;
if6->sin6_family = AF_INET6;
memcpy(&if6->sin6_addr, v6addr, 16);
}
}
if (es->interface_addr.ss_family == 0) {
struct sockaddr_storage remote;
memset(&remote, 0, sizeof(remote));
struct sockaddr_in6* rem_sin6 = (struct sockaddr_in6*)&remote;
rem_sin6->sin6_family = AF_INET6;
rem_sin6->sin6_port = htons(53);
inet_pton(AF_INET6, "2001:4860:4860::8888", &rem_sin6->sin6_addr);
get_outgoing_local_ip(&es->local_addr, &remote, &es->interface_addr);
}
}
}
}
if (es->interface_addr.ss_family == es->local_addr.ss_family && es->local_addr.ss_family != 0) {
if (es->local_addr.ss_family == AF_INET) {
struct sockaddr_in* sin_cfg = (struct sockaddr_in*)&es->local_addr;
struct sockaddr_in* sin_if = (struct sockaddr_in*)&es->interface_addr;
sin_if->sin_port = sin_cfg->sin_port;
} else if (es->local_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6_cfg = (struct sockaddr_in6*)&es->local_addr;
struct sockaddr_in6* sin6_if = (struct sockaddr_in6*)&es->interface_addr;
sin6_if->sin6_port = sin6_cfg->sin6_port;
}
}
}
static void socket_monitor_check_and_fire(struct SOCKET_MONITOR* sm, struct ETCP_SOCKET* es, uint32_t ifindex) {
if (es->netif_index != ifindex) return;
struct sockaddr_storage old_if_addr = es->interface_addr;
socket_monitor_update_if_addr(es);
if (memcmp(&old_if_addr, &es->interface_addr, sizeof(old_if_addr)) != 0) {
sm->socket_updates++;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket %s interface_addr changed: %s -> %s",
es->name, sockaddr_storage_to_str(&old_if_addr).str, sockaddr_storage_to_str(&es->interface_addr).str);
etcp_socket_cbk_fire(es, ETCP_SOCKET_EVENT_ADDR_CHANGED);
}
}
static void socket_monitor_handle_addr(struct SOCKET_MONITOR* sm, struct ifaddrmsg* ifa, int msg_type) {
sm->addr_changes++;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "addr event: type=%s ifindex=%u family=%d",
msg_type == RTM_NEWADDR ? "NEW" : "DEL", ifa->ifa_index, ifa->ifa_family);
struct ETCP_SOCKET* es = sm->instance->etcp_sockets;
while (es) { socket_monitor_check_and_fire(sm, es, ifa->ifa_index); es = es->next; }
}
static void socket_monitor_handle_link(struct SOCKET_MONITOR* sm, struct ifinfomsg* ifi, int msg_type) {
sm->link_changes++;
int is_up = (ifi->ifi_flags & IFF_UP) != 0;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "link event: type=%s ifindex=%u flags=0x%x up=%d",
msg_type == RTM_NEWLINK ? "NEW" : "DEL", ifi->ifi_index, ifi->ifi_flags, is_up);
struct ETCP_SOCKET* es = sm->instance->etcp_sockets;
while (es) {
if (es->netif_index == ifi->ifi_index) {
etcp_socket_cbk_fire(es, ETCP_SOCKET_EVENT_STATUS_CHANGED);
if (is_up) socket_monitor_check_and_fire(sm, es, ifi->ifi_index);
}
es = es->next;
}
}
static void socket_monitor_handle_route(struct SOCKET_MONITOR* sm, struct rtmsg* rtm, int msg_type) {
if (rtm->rtm_dst_len != 0) return;
sm->route_changes++;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "default route event: type=%s family=%d",
msg_type == RTM_NEWROUTE ? "NEW" : "DEL", rtm->rtm_family);
struct ETCP_SOCKET* es = sm->instance->etcp_sockets;
while (es) {
if (es->netif_index != 0 || es->local_addr.ss_family == 0) { es = es->next; continue; }
int is_any = 0;
if (es->local_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&es->local_addr;
is_any = (sin->sin_addr.s_addr == 0);
} else if (es->local_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&es->local_addr;
is_any = (memcmp(&sin6->sin6_addr, &in6addr_any, 16) == 0);
}
if (is_any) {
struct sockaddr_storage old_if_addr = es->interface_addr;
socket_monitor_update_if_addr(es);
if (memcmp(&old_if_addr, &es->interface_addr, sizeof(old_if_addr)) != 0) {
sm->socket_updates++;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket %s interface_addr changed via default route: %s -> %s",
es->name, sockaddr_storage_to_str(&old_if_addr).str, sockaddr_storage_to_str(&es->interface_addr).str);
etcp_socket_cbk_fire(es, ETCP_SOCKET_EVENT_ADDR_CHANGED);
}
}
es = es->next;
}
}
static void socket_monitor_read_cb(int fd, void* arg) {
struct SOCKET_MONITOR* sm = (struct SOCKET_MONITOR*)arg;
char buf[8192];
struct sockaddr_nl nladdr;
socklen_t nladdr_len = sizeof(nladdr);
for (;;) {
ssize_t n = recvfrom(fd, buf, sizeof(buf), MSG_DONTWAIT, (struct sockaddr*)&nladdr, &nladdr_len);
if (n < 0) {
if (errno == EAGAIN || errno == EWOULDBLOCK) break;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "netlink recvfrom fd=%d failed: %s", fd, strerror(errno));
break;
}
if (n == 0) break;
size_t remaining = (size_t)n;
struct nlmsghdr* nlh = (struct nlmsghdr*)buf;
while (NLMSG_OK(nlh, remaining)) {
if (nlh->nlmsg_type == RTM_NEWADDR || nlh->nlmsg_type == RTM_DELADDR) {
socket_monitor_handle_addr(sm, (struct ifaddrmsg*)NLMSG_DATA(nlh), nlh->nlmsg_type);
} else if (nlh->nlmsg_type == RTM_NEWLINK || nlh->nlmsg_type == RTM_DELLINK) {
socket_monitor_handle_link(sm, (struct ifinfomsg*)NLMSG_DATA(nlh), nlh->nlmsg_type);
} else if (nlh->nlmsg_type == RTM_NEWROUTE || nlh->nlmsg_type == RTM_DELROUTE) {
socket_monitor_handle_route(sm, (struct rtmsg*)NLMSG_DATA(nlh), nlh->nlmsg_type);
}
if (nlh->nlmsg_type == NLMSG_DONE) break;
nlh = NLMSG_NEXT(nlh, remaining);
}
}
}
#endif
int socket_monitor_init(struct UTUN_INSTANCE* instance) {
if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "instance is NULL"); return -1; }
if (instance->socket_monitor) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socket monitor already initialized"); return -1; }
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "initializing socket monitor");
#ifdef _WIN32
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket monitor not available on Windows");
return 0;
#else
struct SOCKET_MONITOR* sm = u_calloc(1, sizeof(struct SOCKET_MONITOR));
if (!sm) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "alloc failed"); return -1; }
sm->instance = instance;
sm->nl_sock = socket(AF_NETLINK, SOCK_RAW, NETLINK_ROUTE);
if (sm->nl_sock < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "failed to create netlink socket: %s", strerror(errno));
u_free(sm);
return -1;
}
struct sockaddr_nl nladdr;
memset(&nladdr, 0, sizeof(nladdr));
nladdr.nl_family = AF_NETLINK;
nladdr.nl_groups = RTMGRP_LINK | RTMGRP_IPV4_IFADDR | RTMGRP_IPV6_IFADDR
| RTMGRP_IPV4_ROUTE | RTMGRP_IPV6_ROUTE;
if (bind(sm->nl_sock, (struct sockaddr*)&nladdr, sizeof(nladdr)) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "failed to bind netlink socket: %s", strerror(errno));
close(sm->nl_sock);
u_free(sm);
return -1;
}
if (socket_set_nonblocking(sm->nl_sock) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "failed to set netlink socket non-blocking: %s", strerror(errno));
}
sm->uasync_handle = uasync_add_socket(instance->ua, sm->nl_sock, socket_monitor_read_cb, NULL, NULL, sm);
if (!sm->uasync_handle) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "failed to register netlink socket fd=%d with uasync", sm->nl_sock);
close(sm->nl_sock);
u_free(sm);
return -1;
}
instance->socket_monitor = sm;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket monitor initialized: nl_sock=%d", sm->nl_sock);
return 0;
#endif
}
void socket_monitor_destroy(struct UTUN_INSTANCE* instance) {
if (!instance) return;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "destroying socket monitor");
#ifdef _WIN32
return;
#else
struct SOCKET_MONITOR* sm = (struct SOCKET_MONITOR*)instance->socket_monitor;
if (!sm) return;
instance->socket_monitor = NULL;
uasync_remove_socket(instance->ua, sm->uasync_handle);
close(sm->nl_sock);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket monitor destroyed: nl_sock=%d addr_changes=%u link_changes=%u route_changes=%u socket_updates=%u",
sm->nl_sock, sm->addr_changes, sm->link_changes, sm->route_changes, sm->socket_updates);
u_free(sm);
#endif
}

16
src/transport_layer/socket_monitor.h

@ -0,0 +1,16 @@
#ifndef SOCKET_MONITOR_H
#define SOCKET_MONITOR_H
#ifdef __cplusplus
extern "C" {
#endif
struct UTUN_INSTANCE;
int socket_monitor_init(struct UTUN_INSTANCE* instance);
void socket_monitor_destroy(struct UTUN_INSTANCE* instance);
#ifdef __cplusplus
}
#endif
#endif

5
src/utun_instance.c

@ -19,6 +19,7 @@
#include "stcp_server.h" #include "stcp_server.h"
#include "control_server.h" #include "control_server.h"
#include "transport_layer/node_conn_direct.h" #include "transport_layer/node_conn_direct.h"
#include "transport_layer/socket_monitor.h"
#include "../lib/u_async.h" #include "../lib/u_async.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
@ -218,6 +219,8 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
instance->topo_groups = NULL; instance->topo_groups = NULL;
} }
socket_monitor_init(instance);
// conn_mgr initialized inside topo_group_create (via topo_groups_init) // conn_mgr initialized inside topo_group_create (via topo_groups_init)
// Initialize firewall // Initialize firewall
@ -430,6 +433,8 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
} }
instance->config_conn_handles = NULL; instance->config_conn_handles = NULL;
socket_monitor_destroy(instance);
// Cleanup BGP module BEFORE sockets (needs live conn_mgr for recovery cleanup) // Cleanup BGP module BEFORE sockets (needs live conn_mgr for recovery cleanup)
if (instance->topo_groups) { if (instance->topo_groups) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Destroying BGP module"); DEBUG_INFO(DEBUG_CATEGORY_BGP, "Destroying BGP module");

2
src/utun_instance.h

@ -125,6 +125,7 @@ struct UTUN_INSTANCE {
// Callback chain for new ETCP connections // Callback chain for new ETCP connections
struct etcp_inst_cbk_entry* new_conn_cbks; struct etcp_inst_cbk_entry* new_conn_cbks;
struct etcp_status_cbk_entry* conn_status_cbks; // instance-level: NEW/UP/DOWN/DELETE struct etcp_status_cbk_entry* conn_status_cbks; // instance-level: NEW/UP/DOWN/DELETE
struct etcp_socket_cbk_entry* socket_cbks; // instance-level: socket ADDR/STATUS changes
void* test_user_ptr; // Generic user pointer (used by tests) void* test_user_ptr; // Generic user pointer (used by tests)
struct memory_pool* data_pool;// для входных-выходных данных пакета struct memory_pool* data_pool;// для входных-выходных данных пакета
@ -135,6 +136,7 @@ struct UTUN_INSTANCE {
struct ETCP_SOCKET* etcp_sockets;// linked-list struct ETCP_SOCKET* etcp_sockets;// linked-list
struct TCP_SOCKET* tcp_sockets; // linked-list [TCP-only] struct TCP_SOCKET* tcp_sockets; // linked-list [TCP-only]
struct stcp_server *stcp_server; // TCP server (single for now) struct stcp_server *stcp_server; // TCP server (single for now)
void* socket_monitor; // SOCKET_MONITOR* (opaque, transport_layer/socket_monitor.c)
// Pending one-shot pings (for callback on PONG or timeout) // Pending one-shot pings (for callback on PONG or timeout)
struct PING_CONTEXT* pending_pings; struct PING_CONTEXT* pending_pings;

Loading…
Cancel
Save