diff --git a/src/Makefile.am b/src/Makefile.am index 2c4ae4fd..69deaa2e 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -44,6 +44,7 @@ utun_CORE_SOURCES = \ transport_layer/etcp_api.c \ transport_layer/etcp_connect.c \ transport_layer/node_conn_direct.c \ + transport_layer/socket_monitor.c \ control_server.c \ firewall.c \ eim_nat.c \ @@ -118,6 +119,7 @@ libutun_a_SOURCES = \ transport_layer/etcp_api.c \ transport_layer/etcp_connect.c \ transport_layer/node_conn_direct.c \ + transport_layer/socket_monitor.c \ control_server.c \ firewall.c \ eim_nat.c \ diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 01f753e1..796a342d 100644 --- a/src/routing_layer/topo_group.c +++ b/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_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)"); return g; @@ -322,6 +324,7 @@ void topo_groups_destroy(struct UTUN_INSTANCE* instance) { etcp_unbind(instance, ETCP_ID_TOPO_ENTRY); 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); struct TOPO_GROUPS* g = instance->topo_groups; diff --git a/src/routing_layer/topo_node.c b/src/routing_layer/topo_node.c index 313aa7a1..4ff60441 100644 --- a/src/routing_layer/topo_node.c +++ b/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; } + +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); +} diff --git a/src/routing_layer/topo_node.h b/src/routing_layer/topo_node.h index 3efb6a3a..f5818d18 100644 --- a/src/routing_layer/topo_node.h +++ b/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); 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); 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); diff --git a/src/transport_layer/etcp_api.c b/src/transport_layer/etcp_api.c index f58f8bcb..e36c3da0 100644 --- a/src/transport_layer/etcp_api.c +++ b/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_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) { if (!conn) return; conn->routing_exchange_active = new_state; diff --git a/src/transport_layer/etcp_api.h b/src/transport_layer/etcp_api.h index 117103a6..56f7b482 100644 --- a/src/transport_layer/etcp_api.h +++ b/src/transport_layer/etcp_api.h @@ -54,6 +54,7 @@ extern "C" { // Forward declarations struct ETCP_CONN; +struct ETCP_SOCKET; struct UTUN_INSTANCE; 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_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); /** diff --git a/src/transport_layer/socket_monitor.c b/src/transport_layer/socket_monitor.c new file mode 100644 index 00000000..17dfa6f8 --- /dev/null +++ b/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 +#include +#include +#include + +#ifndef _WIN32 +#include +#include +#include +#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 +} diff --git a/src/transport_layer/socket_monitor.h b/src/transport_layer/socket_monitor.h new file mode 100644 index 00000000..1ac10bcc --- /dev/null +++ b/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 diff --git a/src/utun_instance.c b/src/utun_instance.c index 05927ebd..9a09dbd5 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -19,6 +19,7 @@ #include "stcp_server.h" #include "control_server.h" #include "transport_layer/node_conn_direct.h" +#include "transport_layer/socket_monitor.h" #include "../lib/u_async.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; } + socket_monitor_init(instance); + // conn_mgr initialized inside topo_group_create (via topo_groups_init) // Initialize firewall @@ -430,6 +433,8 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { } instance->config_conn_handles = NULL; + socket_monitor_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"); diff --git a/src/utun_instance.h b/src/utun_instance.h index 9dc415dc..6e25322b 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -125,6 +125,7 @@ struct UTUN_INSTANCE { // Callback chain for new ETCP connections struct etcp_inst_cbk_entry* new_conn_cbks; 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) struct memory_pool* data_pool;// для входных-выходных данных пакета @@ -135,6 +136,7 @@ struct UTUN_INSTANCE { struct ETCP_SOCKET* etcp_sockets;// linked-list struct TCP_SOCKET* tcp_sockets; // linked-list [TCP-only] 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) struct PING_CONTEXT* pending_pings;