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.
 
 
 
 
 
 

263 lines
12 KiB

// udp_proxy.c — UDP датаграммный прокси: клиент ↔ exit через etcp_router
#include "udp_proxy.h"
#include "etcp.h"
#include "etcp_api.h"
#include "etcp_router.h"
#include "utun_instance.h"
#include "tun_if.h"
#include "../routing_layer/topo_node.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/ll_queue.h"
#include "../lib/mem.h"
#include <stdlib.h>
#include <string.h>
#include <errno.h>
#ifndef _WIN32
#include <unistd.h>
#include <netinet/in.h>
#include <netinet/ip.h>
#include <arpa/inet.h>
#include <sys/socket.h>
#endif
#define UDP_FLOW_TIMEOUT_TB 600000 // 60с
struct udp_proxy_ctx* g_udp_ctx = NULL;
static void flow_expire_timer_cb(void* arg);
static void flow_expire(struct udp_proxy_ctx* ctx);
static struct udp_flow* flow_find(struct udp_flow* head, uint64_t client_node_id,
uint32_t src_ip, uint16_t src_port,
uint32_t dst_ip, uint16_t dst_port) {
struct udp_flow* f;
for (f = head; f; f = f->next)
if (f->client_node_id == client_node_id && f->src_ip == src_ip &&
f->src_port == src_port && f->dst_ip == dst_ip && f->dst_port == dst_port) return f;
return NULL;
}
static void flow_read_cb(socket_t sock, void* arg) {
(void)sock; struct udp_flow* f = (struct udp_flow*)arg;
if (!f || f->sock == SOCKET_INVALID || !g_udp_ctx) return;
uint8_t buf[1600];
struct sockaddr_in from; socklen_t flen = sizeof(from);
ssize_t n = recvfrom(f->sock, buf, sizeof(buf), 0, (struct sockaddr*)&from, &flen);
if (n <= 0) return;
f->last_activity_tb = get_time_tb();
// Обернуть ответ: меняем src/dst для восстановления на стороне клиента
struct ll_entry* e = queue_entry_new(0);
if (!e) return;
e->dgram = u_malloc(UDP_PROXY_HDR_SIZE + n);
if (!e->dgram) { queue_entry_free(e); return; }
e->dgram[0] = ETCP_RT_ID_UDP_PROXY;
e->dgram[1] = UDP_PROXY_SUBCMD_DATA;
// Ответ src = оригинальный dst_ip:dest_port
memcpy(e->dgram + 2, &f->dst_ip, 4);
memcpy(e->dgram + 6, &f->dst_port, 2);
// Ответ dst = оригинальный src_ip:src_port
memcpy(e->dgram + 8, &f->src_ip, 4);
memcpy(e->dgram + 12, &f->src_port, 2);
memcpy(e->dgram + UDP_PROXY_HDR_SIZE, buf, n);
e->len = UDP_PROXY_HDR_SIZE + n;
etcp_route_send(g_udp_ctx->inst, TOPO_GROUP_UTUN, f->client_node_id, e, 0, 0);
}
// ====================================================================
// Exit узел: принять REQUEST, создать сокет, переслать
// ====================================================================
static void exit_handle_data(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_udp_ctx ? g_udp_ctx->inst : NULL);
if (!inst || !g_udp_ctx || entry->len < UDP_PROXY_RECV_HDR_SIZE + 1) goto drop;
uint64_t client_node_id; memcpy(&client_node_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8);
uint32_t src_ip; memcpy(&src_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4);
uint16_t src_port; memcpy(&src_port, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 5, 2);
uint32_t dst_ip; memcpy(&dst_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 7, 4);
uint16_t dst_port; memcpy(&dst_port, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 11, 2);
uint8_t* payload = entry->dgram + UDP_PROXY_RECV_HDR_SIZE;
size_t payload_len = entry->len - UDP_PROXY_RECV_HDR_SIZE;
struct udp_flow* f = flow_find(g_udp_ctx->flows, client_node_id, src_ip, src_port, dst_ip, dst_port);
if (!f) {
f = u_calloc(1, sizeof(struct udp_flow));
if (!f) goto drop;
f->client_node_id = client_node_id; f->src_ip = src_ip; f->src_port = src_port;
f->dst_ip = dst_ip; f->dst_port = dst_port;
f->ua = g_udp_ctx->ua; f->created_tb = get_time_tb(); f->last_activity_tb = f->created_tb;
f->sock = socket(AF_INET, SOCK_DGRAM, 0);
if (f->sock == SOCKET_INVALID) { u_free(f); DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "udp_proxy: socket failed"); goto drop; }
socket_set_nonblocking(f->sock);
struct sockaddr_in bind_addr = {.sin_family = AF_INET, .sin_addr = {.s_addr = INADDR_ANY}, .sin_port = 0};
bind(f->sock, (struct sockaddr*)&bind_addr, sizeof(bind_addr));
f->read_id = uasync_add_socket_t(g_udp_ctx->ua, f->sock, flow_read_cb, NULL, NULL, "udp_flow", f);
if (!f->read_id) { socket_close_wrapper(f->sock); u_free(f); goto drop; }
f->next = g_udp_ctx->flows; g_udp_ctx->flows = f; g_udp_ctx->flow_count++;
if (!g_udp_ctx->expire_timer)
g_udp_ctx->expire_timer = uasync_set_timeout(g_udp_ctx->ua, g_udp_ctx->flow_timeout_tb, g_udp_ctx, flow_expire_timer_cb, "udp_expire");
}
f->last_activity_tb = get_time_tb();
struct sockaddr_in addr = {.sin_family = AF_INET, .sin_addr = {.s_addr = dst_ip}, .sin_port = dst_port};
sendto(f->sock, payload, payload_len, 0, (struct sockaddr*)&addr, sizeof(addr));
drop:
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
// Сторона клиента: принять REPLY, доставить в TUN
// ====================================================================
static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (entry->len < UDP_PROXY_RECV_HDR_SIZE + 1) { queue_dgram_free(entry); queue_entry_free(entry); return; }
const uint8_t* d = entry->dgram;
uint32_t src_ip, dst_ip;
uint16_t src_port, dst_port;
// В сообщении: src_ip:src_port = оригинальный dst, dst_ip:dst_port = оригинальный src
memcpy(&src_ip, d + ROUTER_SVC_PAYLOAD_OFF + 1, 4);
memcpy(&src_port, d + ROUTER_SVC_PAYLOAD_OFF + 5, 2);
memcpy(&dst_ip, d + ROUTER_SVC_PAYLOAD_OFF + 7, 4);
memcpy(&dst_port, d + ROUTER_SVC_PAYLOAD_OFF + 11, 2);
uint8_t* payload = d + UDP_PROXY_RECV_HDR_SIZE;
size_t payload_len = entry->len - UDP_PROXY_RECV_HDR_SIZE;
struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_udp_ctx ? g_udp_ctx->inst : NULL);
udp_proxy_deliver_reply(inst, src_ip, src_port, dst_ip, dst_port, payload, payload_len);
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
// Единый etcp_router коллбэк
// ====================================================================
void udp_proxy_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || !entry->dgram || entry->len < ROUTER_SVC_HDR_SIZE) {
if (entry) { queue_dgram_free(entry); queue_entry_free(entry); }
return;
}
uint8_t subcmd = entry->dgram[ROUTER_SVC_PAYLOAD_OFF];
if (subcmd == UDP_PROXY_SUBCMD_DATA) {
if (g_udp_ctx && g_udp_ctx->is_exit) exit_handle_data(conn, entry);
else client_handle_reply(conn, entry);
return;
}
queue_dgram_free(entry); queue_entry_free(entry);
}
// ====================================================================
// Сторона клиента: отправить UDP датаграмму в exit узел
// ====================================================================
int udp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id,
uint32_t src_ip, uint16_t src_port,
uint32_t dst_ip, uint16_t dst_port,
const uint8_t* payload, size_t payload_len) {
if (!inst || !g_udp_ctx) return -1;
struct ll_entry* e = queue_entry_new(0);
if (!e) return -1;
e->dgram = u_malloc(UDP_PROXY_HDR_SIZE + payload_len);
if (!e->dgram) { queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_RT_ID_UDP_PROXY;
e->dgram[1] = UDP_PROXY_SUBCMD_DATA;
memcpy(e->dgram + 2, &src_ip, 4);
memcpy(e->dgram + 6, &src_port, 2);
memcpy(e->dgram + 8, &dst_ip, 4);
memcpy(e->dgram + 12, &dst_port, 2);
if (payload_len > 0) memcpy(e->dgram + UDP_PROXY_HDR_SIZE, payload, payload_len);
e->len = UDP_PROXY_HDR_SIZE + payload_len;
return etcp_route_send(inst, TOPO_GROUP_UTUN, exit_node_id, e, 0, 0);
}
// ====================================================================
// Сторона клиента: доставить UDP ответ в TUN
// ====================================================================
int udp_proxy_deliver_reply(struct UTUN_INSTANCE* inst,
uint32_t src_ip, uint16_t src_port,
uint32_t dst_ip, uint16_t dst_port,
const uint8_t* payload, size_t payload_len) {
if (!inst || !inst->tcp_proxy_client || !inst->tcp_proxy_client->tun) return -1;
// Собрать IP/UDP ответный пакет
size_t ip_len = 20 + 8 + payload_len;
uint8_t* pkt = u_malloc(ip_len);
if (!pkt) return -1;
memset(pkt, 0, 20);
pkt[0] = 0x45; pkt[8] = 64; pkt[9] = IPPROTO_UDP;
memcpy(pkt + 12, &src_ip, 4); memcpy(pkt + 16, &dst_ip, 4);
uint16_t total_len = htons(20 + 8 + payload_len);
memcpy(pkt + 2, &total_len, 2);
{ uint32_t ip_sum = 0; for (int i = 0; i < 10; i++) { uint16_t w; memcpy(&w, pkt + i*2, 2); ip_sum += w; }
ip_sum = (ip_sum & 0xFFFF) + (ip_sum >> 16); ip_sum += (ip_sum >> 16);
uint16_t cs = (uint16_t)(~ip_sum); memcpy(pkt + 10, &cs, 2); }
uint16_t udp_len = htons(8 + payload_len);
memcpy(pkt + 24, &udp_len, 2);
memcpy(pkt + 20, &src_port, 2);
memcpy(pkt + 22, &dst_port, 2);
if (payload_len > 0) memcpy(pkt + 28, payload, payload_len);
struct ll_entry* e = queue_entry_new(0);
if (!e) { u_free(pkt); return -1; }
e->dgram = u_malloc(1 + ip_len);
e->dgram[0] = 4; memcpy(e->dgram + 1, pkt, ip_len); e->len = 1 + ip_len;
u_free(pkt);
queue_data_put(inst->tcp_proxy_client->tun->input_queue, e);
return 0;
}
static void flow_expire_timer_cb(void* arg) {
struct udp_proxy_ctx* ctx = (struct udp_proxy_ctx*)arg;
flow_expire(ctx);
ctx->expire_timer = NULL;
if (ctx->flows)
ctx->expire_timer = uasync_set_timeout(ctx->ua, ctx->flow_timeout_tb, ctx, flow_expire_timer_cb, "udp_expire");
}
static void flow_expire(struct udp_proxy_ctx* ctx) {
uint64_t now = get_time_tb(); struct udp_flow** prev = &ctx->flows;
while (*prev) {
struct udp_flow* f = *prev;
if (now - f->last_activity_tb > ctx->flow_timeout_tb) {
*prev = f->next; ctx->flow_count--;
if (f->read_id) { uasync_remove_socket_t(f->ua, f->sock); f->read_id = NULL; }
socket_close_wrapper(f->sock); u_free(f);
} else prev = &f->next;
}
}
// ====================================================================
// Публичные init/destroy
// ====================================================================
int udp_proxy_init(struct UTUN_INSTANCE* inst, struct UASYNC* ua) {
if (!inst) return -1;
struct udp_proxy_ctx* ctx = u_calloc(1, sizeof(struct udp_proxy_ctx));
if (!ctx) return -1;
ctx->inst = inst; ctx->ua = ua;
ctx->flow_timeout_tb = UDP_FLOW_TIMEOUT_TB;
ctx->expire_timer = NULL;
ctx->is_exit = inst->tcp_proxy_server.enabled;
g_udp_ctx = ctx;
etcp_router_bind(inst, ETCP_RT_ID_UDP_PROXY, udp_proxy_recv_cb);
ctx->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_PROXY, "udp_proxy initialized (exit=%d)", ctx->is_exit);
return 0;
}
void udp_proxy_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !g_udp_ctx) return;
etcp_router_unbind(inst, ETCP_RT_ID_UDP_PROXY);
if (g_udp_ctx->expire_timer) { uasync_cancel_timeout(g_udp_ctx->ua, g_udp_ctx->expire_timer); g_udp_ctx->expire_timer = NULL; }
struct udp_flow* f = g_udp_ctx->flows;
while (f) { struct udp_flow* n = f->next;
if (f->read_id) { uasync_remove_socket_t(g_udp_ctx->ua, f->sock); f->read_id = NULL; }
socket_close_wrapper(f->sock); u_free(f); f = n; }
u_free(g_udp_ctx); g_udp_ctx = NULL;
DEBUG_INFO(DEBUG_CATEGORY_PROXY, "udp_proxy destroyed");
}