commit ec1f4256b8c796dd6a7289b00e8ff1f6915ba003
| author | 斟酌 鵬兄 <tgckpg@gmail.com> |
| date | 2026-05-30T23:29:08Z |
| subject | Basic UDP thread support |
commit ec1f4256b8c796dd6a7289b00e8ff1f6915ba003
Author: 斟酌 鵬兄 <tgckpg@gmail.com>
Date: 2026-05-30T23:29:08Z
Basic UDP thread support
---
src/datagram_client.c | 13 +++---
src/datagram_client.h | 18 ++------
src/datagram_listener.c | 112 +++++++++++++++++++++++++++++++-----------------
src/datagram_listener.h | 4 ++
src/worker.c | 112 ++++++++++++++++++++++++++++++++++++++++++------
src/worker.h | 25 ++++++++++-
6 files changed, 209 insertions(+), 75 deletions(-)
diff --git a/src/datagram_client.c b/src/datagram_client.c
index 0ad5092..9c77395 100644
--- a/src/datagram_client.c
+++ b/src/datagram_client.c
@@ -74,7 +74,7 @@ void cleanup_idle_datagram_clients(struct datagram_route_ctx *ctx)
}
struct datagram_client *find_datagram_client(
- struct datagram_route_ctx *ctx,
+ const struct datagram_route_ctx *ctx,
const struct sockaddr_storage *addr,
socklen_t addr_len
) {
@@ -203,7 +203,7 @@ int send_datagram_payload_to_upstream(
const unsigned char *payload,
size_t payload_len
) {
- struct datagram_route_ctx *ctx = c->ctx;
+ const struct datagram_route_ctx *ctx = c->ctx;
const struct route *r = ctx->route;
if (!r->opts.proxy_v2) {
@@ -263,7 +263,7 @@ static void upstream_read_cb(evutil_socket_t fd, short events, void *arg)
(void)events;
struct datagram_client *c = arg;
- struct datagram_route_ctx *ctx = c->ctx;
+ const struct datagram_route_ctx *ctx = c->ctx;
unsigned char buf[UDP_MAX_PACKET];
@@ -308,9 +308,10 @@ static void upstream_read_cb(evutil_socket_t fd, short events, void *arg)
struct datagram_client *create_datagram_client(
struct datagram_route_ctx *ctx,
+ struct event_base *base,
const struct sockaddr_storage *client_addr,
- socklen_t client_addr_len
-) {
+ socklen_t client_addr_len)
+{
const struct route *r = ctx->route;
struct datagram_client *c = calloc(1, sizeof(*c));
@@ -340,7 +341,7 @@ struct datagram_client *create_datagram_client(
return NULL;
}
- c->ev = event_new(ctx->base, c->fd, EV_READ | EV_PERSIST, upstream_read_cb, c);
+ c->ev = event_new(base, c->fd, EV_READ | EV_PERSIST, upstream_read_cb, c);
if (c->ev == NULL) {
LOG_ERROR("event_new failed for udp upstream");
free_datagram_client(c);
diff --git a/src/datagram_client.h b/src/datagram_client.h
index a3f238d..3b2c7d7 100644
--- a/src/datagram_client.h
+++ b/src/datagram_client.h
@@ -19,29 +19,17 @@ struct datagram_client {
struct datagram_client *next;
};
-struct datagram_packet {
- const struct route *route;
-
- evutil_socket_t listen_fd;
-
- struct sockaddr_storage peer_addr;
- socklen_t peer_addr_len;
-
- unsigned char data[65535];
- size_t data_len;
-};
-
void cleanup_idle_datagram_clients(struct datagram_route_ctx *ctx);
void free_datagram_client(struct datagram_client *c);
struct datagram_client *create_datagram_client(
struct datagram_route_ctx *ctx,
+ struct event_base *base,
const struct sockaddr_storage *client_addr,
- socklen_t client_addr_len
-);
+ socklen_t client_addr_len);
struct datagram_client *find_datagram_client(
- struct datagram_route_ctx *ctx,
+ const struct datagram_route_ctx *ctx,
const struct sockaddr_storage *addr,
socklen_t addr_len
);
diff --git a/src/datagram_listener.c b/src/datagram_listener.c
index a9c5c31..574f577 100644
--- a/src/datagram_listener.c
+++ b/src/datagram_listener.c
@@ -8,14 +8,12 @@
#include "datagram_client.h"
#include "datagram_builtin.h"
-static int handle_datagram_builtin_packet(
- struct datagram_route_ctx *ctx,
- const struct datagram_packet *pkt)
+static int handle_datagram_builtin_packet(const struct worker_datagram_packet_msg *pkt)
{
struct datagram_client tmp;
memset(&tmp, 0, sizeof(tmp));
- tmp.ctx = ctx;
+ tmp.ctx = pkt->ctx;
tmp.fd = EVUTIL_INVALID_SOCKET;
tmp.client_addr = pkt->peer_addr;
tmp.client_addr_len = pkt->peer_addr_len;
@@ -28,40 +26,49 @@ static int handle_datagram_builtin_packet(
return rc;
}
-static struct datagram_client* datagram_route_get_or_create_client(
- struct datagram_route_ctx *ctx,
- const struct datagram_packet *pkt)
+static struct datagram_client *datagram_route_get_or_create_client(
+ struct worker *w,
+ const struct worker_datagram_packet_msg *pkt)
{
- struct datagram_client *c = find_datagram_client(ctx, &pkt->peer_addr, pkt->peer_addr_len);
+ struct datagram_client *c;
+
+ c = find_datagram_client(pkt->ctx, &pkt->peer_addr, pkt->peer_addr_len);
if (c == NULL) {
- c = create_datagram_client(ctx, &pkt->peer_addr, pkt->peer_addr_len);
+ c = create_datagram_client(
+ pkt->ctx,
+ w->base,
+ &pkt->peer_addr,
+ pkt->peer_addr_len
+ );
if (c == NULL) {
- return c;
+ return NULL;
}
}
+
c->last_seen = time(NULL);
return c;
}
-static int datagram_route_handle_packet(
- struct datagram_route_ctx *ctx,
- const struct datagram_packet *pkt)
+int datagram_route_handle_packet(
+ struct worker *w,
+ const struct worker_datagram_packet_msg *pkt)
{
struct datagram_client *c;
int rc;
if (pkt->route->upstream.kind == ENDPOINT_BUILTIN) {
- return handle_datagram_builtin_packet(ctx, pkt);
+ return handle_datagram_builtin_packet(pkt);
}
- c = datagram_route_get_or_create_client(ctx, pkt);
+ c = datagram_route_get_or_create_client(w, pkt);
if (!c) {
return -ENOMEM;
}
rc = send_datagram_payload_to_upstream(c, pkt->data, pkt->data_len);
if (rc != 0) {
- LOG_WARN("failed to send datagram payload upstream", "err", _LOGV(strerror(-rc)));
+ LOG_WARN("failed to send datagram payload upstream",
+ "err", _LOGV(strerror(-rc)));
return rc;
}
@@ -70,37 +77,59 @@ static int datagram_route_handle_packet(
static int dispatch_datagram_packet(
struct datagram_route_ctx *ctx,
- const struct datagram_packet *pkt
-) {
+ evutil_socket_t listen_fd,
+ const struct sockaddr *peer_addr,
+ socklen_t peer_addr_len,
+ const unsigned char *data,
+ size_t data_len)
+{
+ struct worker *w;
+
+ if (!ctx || !ctx->worker_pool) {
+ return EINVAL;
+ }
+
/*
- * Datagram routes are pinned to ctx->base for now.
- * Later this function can choose a worker/shard.
+ * For now this can be round-robin to test plumbing.
+ * Before real UDP threading, change this to hash by peer_addr.
*/
- return datagram_route_handle_packet(ctx, pkt);
+ w = worker_pool_next(ctx->worker_pool);
+ if (!w) {
+ return EINVAL;
+ }
+
+ return worker_enqueue_datagram_packet(
+ w,
+ ctx,
+ listen_fd,
+ peer_addr,
+ peer_addr_len,
+ data,
+ data_len
+ );
}
static void listen_read_cb(evutil_socket_t fd, short events, void *arg)
{
struct datagram_route_ctx *ctx = arg;
- struct datagram_packet pkt;
+ struct sockaddr_storage peer_addr;
+ socklen_t peer_addr_len;
+ unsigned char buf[65535];
ssize_t n;
+ int rc;
(void)events;
- cleanup_idle_datagram_clients(ctx);
-
- memset(&pkt, 0, sizeof(pkt));
- pkt.route = ctx->route;
- pkt.listen_fd = fd;
- pkt.peer_addr_len = sizeof(pkt.peer_addr);
+ memset(&peer_addr, 0, sizeof(peer_addr));
+ peer_addr_len = sizeof(peer_addr);
n = recvfrom(
fd,
- (char *)pkt.data,
- sizeof(pkt.data),
+ (char *)buf,
+ sizeof(buf),
0,
- (struct sockaddr *)&pkt.peer_addr,
- &pkt.peer_addr_len
+ (struct sockaddr *)&peer_addr,
+ &peer_addr_len
);
if (n < 0) {
@@ -111,18 +140,23 @@ static void listen_read_cb(evutil_socket_t fd, short events, void *arg)
}
LOG_ERROR("udp listen recv failed",
- "err", _LOGV(evutil_socket_error_to_string(err))
- );
+ "err", _LOGV(evutil_socket_error_to_string(err))
+ );
return;
}
- pkt.data_len = (size_t)n;
-
- int rc = dispatch_datagram_packet(ctx, &pkt);
+ rc = dispatch_datagram_packet(
+ ctx,
+ fd,
+ (struct sockaddr *)&peer_addr,
+ peer_addr_len,
+ buf,
+ (size_t)n
+ );
if (rc != 0) {
LOG_WARN("failed to dispatch datagram packet",
- "err", _LOGV(rc)
- );
+ "err", _LOGV(rc)
+ );
}
}
diff --git a/src/datagram_listener.h b/src/datagram_listener.h
index ac876b8..6225a10 100644
--- a/src/datagram_listener.h
+++ b/src/datagram_listener.h
@@ -6,4 +6,8 @@
int bind_datagram_listener(struct datagram_route_ctx *ctx);
void close_datagram_listener(struct datagram_route_ctx *ctx);
+int datagram_route_handle_packet(
+ struct worker *w,
+ const struct worker_datagram_packet_msg *pkt);
+
#endif
diff --git a/src/worker.c b/src/worker.c
index 5d4e5ac..3031111 100644
--- a/src/worker.c
+++ b/src/worker.c
@@ -7,6 +7,8 @@
#include "klog.h"
#include "stream_conn.h"
+#include "datagram_client.h"
+#include "datagram_listener.h"
#include "worker.h"
static void worker_notify_cb(evutil_socket_t fd, short events, void *arg);
@@ -99,6 +101,11 @@ static void worker_process_msg(struct worker *w, struct worker_msg *msg)
worker_adopt_client_fd(w, &msg->payload.stream_client);
break;
+ case WORKER_MSG_DATAGRAM_PACKET:
+ cleanup_idle_datagram_clients(msg->payload.datagram_packet.ctx);
+ datagram_route_handle_packet(w, &msg->payload.datagram_packet);
+ break;
+
default:
LOG_WARN("unknown worker message kind",
"worker", _LOGV(w->id),
@@ -108,25 +115,28 @@ static void worker_process_msg(struct worker *w, struct worker_msg *msg)
}
}
-static void worker_process_pending(struct worker *w)
+static void free_processed_worker_msg(struct worker_msg *msg)
{
- struct worker_msg *msg = worker_take_pending(w);
-
- while (msg) {
- struct worker_msg *next = msg->next;
-
- worker_process_msg(w, msg);
+ if (!msg) {
+ return;
+ }
+ switch (msg->kind) {
+ case WORKER_MSG_STREAM_CLIENT:
/*
- * Important:
- * Do not call free_worker_msg() here.
- *
- * If worker_adopt_client_fd() succeeds, fd ownership moved
- * into stream connection handling.
+ * fd ownership moved to stream connection handling.
+ * Do not close it here.
*/
- free(msg);
- msg = next;
+ break;
+
+ case WORKER_MSG_DATAGRAM_PACKET:
+ free(msg->payload.datagram_packet.data);
+ msg->payload.datagram_packet.data = NULL;
+ msg->payload.datagram_packet.data_len = 0;
+ break;
}
+
+ free(msg);
}
static void *worker_main(void *arg)
@@ -273,6 +283,20 @@ static int socket_is_valid(evutil_socket_t fd)
return fd != (evutil_socket_t)EVUTIL_INVALID_SOCKET;
}
+static void worker_process_pending(struct worker *w)
+{
+ struct worker_msg *msg = worker_take_pending(w);
+
+ while (msg) {
+ struct worker_msg *next = msg->next;
+
+ worker_process_msg(w, msg);
+ free_processed_worker_msg(msg);
+
+ msg = next;
+ }
+}
+
static void free_worker_msg(struct worker_msg *msg)
{
if (!msg) {
@@ -286,6 +310,11 @@ static void free_worker_msg(struct worker_msg *msg)
msg->payload.stream_client.fd = EVUTIL_INVALID_SOCKET;
}
break;
+ case WORKER_MSG_DATAGRAM_PACKET:
+ free(msg->payload.datagram_packet.data);
+ msg->payload.datagram_packet.data = NULL;
+ msg->payload.datagram_packet.data_len = 0;
+ break;
}
free(msg);
@@ -410,6 +439,61 @@ int worker_enqueue_stream_client(
return worker_enqueue_msg(w, msg);
}
+int worker_enqueue_datagram_packet(
+ struct worker *w,
+ struct datagram_route_ctx *ctx,
+ evutil_socket_t listen_fd,
+ const struct sockaddr *peer_addr,
+ socklen_t peer_addr_len,
+ const unsigned char *data,
+ size_t data_len)
+{
+ struct worker_msg *msg;
+ const struct route *route = ctx->route;
+
+ if (!w || !route || !socket_is_valid(listen_fd)) {
+ return EINVAL;
+ }
+
+ if (!peer_addr || peer_addr_len == 0 ||
+ peer_addr_len > sizeof(((struct worker_datagram_packet_msg *)0)->peer_addr)) {
+ return EINVAL;
+ }
+
+ if (!data && data_len > 0) {
+ return EINVAL;
+ }
+
+ msg = calloc(1, sizeof(*msg));
+ if (!msg) {
+ return ENOMEM;
+ }
+
+ msg->kind = WORKER_MSG_DATAGRAM_PACKET;
+ msg->payload.datagram_packet.route = route;
+ msg->payload.datagram_packet.ctx = ctx;
+
+ memcpy(
+ &msg->payload.datagram_packet.peer_addr,
+ peer_addr,
+ peer_addr_len
+ );
+ msg->payload.datagram_packet.peer_addr_len = peer_addr_len;
+
+ if (data_len > 0) {
+ msg->payload.datagram_packet.data = malloc(data_len);
+ if (!msg->payload.datagram_packet.data) {
+ free(msg);
+ return ENOMEM;
+ }
+
+ memcpy(msg->payload.datagram_packet.data, data, data_len);
+ msg->payload.datagram_packet.data_len = data_len;
+ }
+
+ return worker_enqueue_msg(w, msg);
+}
+
static void worker_notify_cb(evutil_socket_t fd, short events, void *arg)
{
struct worker *w = arg;
diff --git a/src/worker.h b/src/worker.h
index 68a5ce0..4bf1943 100644
--- a/src/worker.h
+++ b/src/worker.h
@@ -11,6 +11,8 @@
#include "compat.h"
#include "route.h"
+struct datagram_route_ctx;
+
struct worker {
unsigned int id;
@@ -31,6 +33,7 @@ struct worker {
enum worker_msg_kind {
WORKER_MSG_STREAM_CLIENT,
+ WORKER_MSG_DATAGRAM_PACKET,
};
struct worker_stream_client_msg {
@@ -40,13 +43,24 @@ struct worker_stream_client_msg {
socklen_t peer_addr_len;
};
+struct worker_datagram_packet_msg {
+ const struct route *route;
+ struct datagram_route_ctx *ctx;
+
+ struct sockaddr_storage peer_addr;
+ socklen_t peer_addr_len;
+
+ unsigned char *data;
+ size_t data_len;
+};
+
struct worker_msg {
enum worker_msg_kind kind;
struct worker_msg *next;
union {
struct worker_stream_client_msg stream_client;
-// struct worker_datagram_packet_msg datagram_packet;
+ struct worker_datagram_packet_msg datagram_packet;
} payload;
};
@@ -64,4 +78,13 @@ int worker_enqueue_stream_client(
socklen_t addr_len
);
+int worker_enqueue_datagram_packet(
+ struct worker *w,
+ struct datagram_route_ctx *ctx,
+ evutil_socket_t listen_fd,
+ const struct sockaddr *peer_addr,
+ socklen_t peer_addr_len,
+ const unsigned char *data,
+ size_t data_len);
+
#endif