commit a169fca4394da5ae4cbe6332f716a4835110c8a6
| author | 斟酌 鵬兄 <tgckpg@gmail.com> |
| date | 2026-05-31T01:40:44Z |
| subject | Handle race conditions introduced by threading |
commit a169fca4394da5ae4cbe6332f716a4835110c8a6
Author: 斟酌 鵬兄 <tgckpg@gmail.com>
Date: 2026-05-31T01:40:44Z
Handle race conditions introduced by threading
---
src/datagram_client.c | 62 ++++++++++++++++++++++++-----
src/datagram_listener.c | 55 ++++++++++++++++++++++----
src/datagram_route.c | 64 ++++++++++++++++++++++--------
src/datagram_route.h | 6 ++-
src/klog.c | 34 ++++++++++++++++
src/klog.h | 102 ++++++++++++++++++++++++++++++++++++++++++++++++
src/route_runtime.c | 27 +++++++++++--
src/route_runtime.h | 3 +-
src/stream_route.c | 16 +++++---
src/stream_route.h | 3 +-
src/tinyproxy.c | 8 +++-
src/worker.c | 5 ++-
src/worker.h | 1 -
13 files changed, 339 insertions(+), 47 deletions(-)
diff --git a/src/datagram_client.c b/src/datagram_client.c
index 9c77395..09b65cd 100644
--- a/src/datagram_client.c
+++ b/src/datagram_client.c
@@ -50,6 +50,8 @@ static int sockaddr_equal(
void cleanup_idle_datagram_clients(struct datagram_route_ctx *ctx)
{
+ compat_mutex_lock(&ctx->clients_mu);
+
time_t now = time(NULL);
struct datagram_client **pp = &ctx->clients;
@@ -71,6 +73,8 @@ void cleanup_idle_datagram_clients(struct datagram_route_ctx *ctx)
free_datagram_client(c);
}
+
+ compat_mutex_unlock(&ctx->clients_mu);
}
struct datagram_client *find_datagram_client(
@@ -260,30 +264,31 @@ int send_datagram_payload_to_upstream(
static void upstream_read_cb(evutil_socket_t fd, short events, void *arg)
{
- (void)events;
-
struct datagram_client *c = arg;
- const struct datagram_route_ctx *ctx = c->ctx;
-
+ struct datagram_route_ctx *ctx = c->ctx;
unsigned char buf[UDP_MAX_PACKET];
+ (void)events;
+
+ compat_mutex_lock(&ctx->clients_mu);
+
for (;;) {
ssize_t n = recv(fd, (char *)buf, sizeof(buf), 0);
if (n < 0) {
int err = EVUTIL_SOCKET_ERROR();
if (socket_err_is_retriable(err)) {
- return;
+ goto out;
}
- LOG_ERROR("udp recvfrom failed",
+ LOG_ERROR("udp upstream recv failed",
"err", _LOGV(evutil_socket_error_to_string(err))
);
- return;
+ goto out;
}
if (n == 0) {
- return;
+ goto out;
}
ssize_t sent = sendto(
@@ -301,9 +306,12 @@ static void upstream_read_cb(evutil_socket_t fd, short events, void *arg)
LOG_ERROR("udp send to client failed",
"err", _LOGV(evutil_socket_error_to_string(err))
);
- return;
+ goto out;
}
}
+
+out:
+ compat_mutex_unlock(&ctx->clients_mu);
}
struct datagram_client *create_datagram_client(
@@ -357,10 +365,46 @@ struct datagram_client *create_datagram_client(
c->next = ctx->clients;
ctx->clients = c;
+#ifdef TINYPROXY_DEBUG
+{
+ struct sockaddr_storage local_addr;
+ socklen_t local_addr_len = sizeof(local_addr);
+ char client_buf[128];
+ char local_buf[128];
+
+ memset(&local_addr, 0, sizeof(local_addr));
+
+ if (getsockname(
+ c->fd,
+ (struct sockaddr *)&local_addr,
+ &local_addr_len
+ ) == 0) {
+ LOG_DEBUG("udp upstream socket created",
+ "listen", _LOGV_ENDPOINT(&r->listen),
+ "upstream", _LOGV_ENDPOINT(&r->upstream),
+ "client", _LOGV_SOCKADDR(&c->client_addr, c->client_addr_len,
+ client_buf, sizeof(client_buf)),
+ "upstream_local", _LOGV_SOCKADDR(&local_addr, local_addr_len,
+ local_buf, sizeof(local_buf)),
+ "fd", _LOGV(c->fd)
+ );
+ } else {
+ LOG_DEBUG("udp upstream socket created",
+ "listen", _LOGV_ENDPOINT(&r->listen),
+ "upstream", _LOGV_ENDPOINT(&r->upstream),
+ "client_family", _LOGV(c->client_addr.ss_family),
+ "client_len", _LOGV(c->client_addr_len),
+ "fd", _LOGV(c->fd),
+ "getsockname_err", _LOGV(evutil_socket_error_to_string(EVUTIL_SOCKET_ERROR()))
+ );
+ }
+}
+#else
LOG_INFO("udp client created",
"client_family", _LOGV(c->client_addr.ss_family),
"client_len", _LOGV(c->client_addr_len)
);
+#endif
return c;
}
diff --git a/src/datagram_listener.c b/src/datagram_listener.c
index 574f577..5816acc 100644
--- a/src/datagram_listener.c
+++ b/src/datagram_listener.c
@@ -8,6 +8,40 @@
#include "datagram_client.h"
#include "datagram_builtin.h"
+#ifdef TINYPROXY_DEBUG
+#include <stdlib.h>
+
+#ifndef _WIN32
+#include <unistd.h>
+#endif
+static void tinyproxy_debug_race_sleep(const char *name)
+{
+ const char *want = getenv("TINYPROXY_RACE_SLEEP");
+ const char *delay_s = getenv("TINYPROXY_RACE_SLEEP_US");
+ long delay_us;
+
+ if (!want || strcmp(want, name) != 0) {
+ return;
+ }
+
+ delay_us = delay_s ? strtol(delay_s, NULL, 10) : 1000;
+ if (delay_us <= 0) {
+ delay_us = 1000;
+ }
+
+#ifdef _WIN32
+ Sleep((DWORD)((delay_us + 999) / 1000));
+#else
+ usleep((useconds_t)delay_us);
+#endif
+}
+#else
+static void tinyproxy_debug_race_sleep(const char *name)
+{
+ (void)name;
+}
+#endif
+
static int handle_datagram_builtin_packet(const struct worker_datagram_packet_msg *pkt)
{
struct datagram_client tmp;
@@ -34,6 +68,7 @@ static struct datagram_client *datagram_route_get_or_create_client(
c = find_datagram_client(pkt->ctx, &pkt->peer_addr, pkt->peer_addr_len);
if (c == NULL) {
+ tinyproxy_debug_race_sleep("udp_client_create");
c = create_datagram_client(
pkt->ctx,
w->base,
@@ -53,31 +88,39 @@ int datagram_route_handle_packet(
struct worker *w,
const struct worker_datagram_packet_msg *pkt)
{
+ struct datagram_route_ctx *ctx = pkt->ctx;
struct datagram_client *c;
int rc;
- if (pkt->route->upstream.kind == ENDPOINT_BUILTIN) {
+ if (ctx->route->upstream.kind == ENDPOINT_BUILTIN) {
return handle_datagram_builtin_packet(pkt);
}
+ compat_mutex_lock(&ctx->clients_mu);
+
c = datagram_route_get_or_create_client(w, pkt);
if (!c) {
- return -ENOMEM;
+ rc = -ENOMEM;
+ goto out;
}
+ tinyproxy_debug_race_sleep("udp_before_send");
+
rc = send_datagram_payload_to_upstream(c, pkt->data, pkt->data_len);
+
+out:
+ compat_mutex_unlock(&ctx->clients_mu);
+
if (rc != 0) {
LOG_WARN("failed to send datagram payload upstream",
"err", _LOGV(strerror(-rc)));
- return rc;
}
- return 0;
+ return rc;
}
static int dispatch_datagram_packet(
struct datagram_route_ctx *ctx,
- evutil_socket_t listen_fd,
const struct sockaddr *peer_addr,
socklen_t peer_addr_len,
const unsigned char *data,
@@ -101,7 +144,6 @@ static int dispatch_datagram_packet(
return worker_enqueue_datagram_packet(
w,
ctx,
- listen_fd,
peer_addr,
peer_addr_len,
data,
@@ -147,7 +189,6 @@ static void listen_read_cb(evutil_socket_t fd, short events, void *arg)
rc = dispatch_datagram_packet(
ctx,
- fd,
(struct sockaddr *)&peer_addr,
peer_addr_len,
buf,
diff --git a/src/datagram_route.c b/src/datagram_route.c
index fcf2377..7af85b2 100644
--- a/src/datagram_route.c
+++ b/src/datagram_route.c
@@ -30,10 +30,10 @@ static int prepare_datagram_route(struct datagram_route_ctx *ctx)
}
int start_datagram_route(
- struct event_base *accept_base,
- struct worker_pool *wpool,
- const struct route *r,
- struct datagram_route_ctx *ctx)
+ struct event_base *accept_base,
+ struct worker_pool *wpool,
+ const struct route *r,
+ struct datagram_route_ctx *ctx)
{
int rc;
char opts[128];
@@ -47,17 +47,24 @@ int start_datagram_route(
ctx->base = accept_base;
ctx->worker_pool = wpool;
ctx->route = r;
- ctx->listen_fd = -1;
+ ctx->listen_fd = EVUTIL_INVALID_SOCKET;
- rc = prepare_datagram_route(ctx);
+ rc = compat_mutex_init(&ctx->clients_mu);
if (rc != 0) {
memset(ctx, 0, sizeof(*ctx));
+ ctx->listen_fd = EVUTIL_INVALID_SOCKET;
+ return rc;
+ }
+
+ rc = prepare_datagram_route(ctx);
+ if (rc != 0) {
+ free_datagram_route(ctx);
return rc;
}
rc = bind_datagram_listener(ctx);
if (rc != 0) {
- stop_datagram_route(ctx);
+ free_datagram_route(ctx);
return rc;
}
@@ -121,33 +128,58 @@ static void stop_unix_datagram_route(struct datagram_route_ctx *ctx)
}
#endif
-void stop_datagram_route(struct datagram_route_ctx *ctx)
+void stop_datagram_route_listener(struct datagram_route_ctx *ctx)
{
if (ctx == NULL) {
return;
}
- while (ctx->clients != NULL) {
- struct datagram_client *c = ctx->clients;
+ if (ctx->listen_ev) {
+ event_free(ctx->listen_ev);
+ ctx->listen_ev = NULL;
+ }
+
+ /*
+ * Do NOT close ctx->listen_fd here.
+ * Workers may still use it to send UDP replies.
+ */
+}
+
+void free_datagram_route(struct datagram_route_ctx *ctx)
+{
+ struct datagram_client *c;
+
+ if (!ctx) {
+ return;
+ }
+
+ compat_mutex_lock(&ctx->clients_mu);
+
+ while (ctx->clients) {
+ c = ctx->clients;
ctx->clients = c->next;
c->next = NULL;
free_datagram_client(c);
}
- if (ctx->listen_ev != NULL) {
+ compat_mutex_unlock(&ctx->clients_mu);
+
+ if (ctx->listen_ev) {
event_free(ctx->listen_ev);
ctx->listen_ev = NULL;
}
+ if (ctx->listen_fd != EVUTIL_INVALID_SOCKET) {
+ evutil_closesocket(ctx->listen_fd);
+ ctx->listen_fd = EVUTIL_INVALID_SOCKET;
+ }
+
#ifndef _WIN32
stop_unix_datagram_route(ctx);
#endif
- if (ctx->listen_fd >= 0) {
- evutil_closesocket(ctx->listen_fd);
- ctx->listen_fd = -1;
- }
+ compat_mutex_destroy(&ctx->clients_mu);
memset(ctx, 0, sizeof(*ctx));
- ctx->listen_fd = -1;
+ ctx->listen_fd = EVUTIL_INVALID_SOCKET;
}
diff --git a/src/datagram_route.h b/src/datagram_route.h
index 0dec487..c5f27fc 100644
--- a/src/datagram_route.h
+++ b/src/datagram_route.h
@@ -4,6 +4,7 @@
#include <event2/util.h>
#include "compat.h"
+#include "compat_thread.h"
#include "worker_pool.h"
struct event;
@@ -13,6 +14,8 @@ struct route;
struct datagram_client;
struct datagram_route_ctx {
+ compat_mutex_t clients_mu;
+
struct event_base *base;
struct worker_pool *worker_pool;
const struct route *route;
@@ -32,6 +35,7 @@ int start_datagram_route(
const struct route *r,
struct datagram_route_ctx *ctx);
-void stop_datagram_route(struct datagram_route_ctx *ctx);
+void stop_datagram_route_listener(struct datagram_route_ctx *ctx);
+void free_datagram_route(struct datagram_route_ctx *ctx);
#endif
diff --git a/src/klog.c b/src/klog.c
index 158a2ae..a888944 100644
--- a/src/klog.c
+++ b/src/klog.c
@@ -6,9 +6,36 @@
#include "klog.h"
#include "route.h"
+#include "compat_thread.h"
+static compat_mutex_t log_mu;
+static int log_mu_ready;
static _Thread_local int klog_worker_id = -1;
+int klog_init(void)
+{
+ if (log_mu_ready) {
+ return 0;
+ }
+
+ if (compat_mutex_init(&log_mu) != 0) {
+ return errno ? errno : EINVAL;
+ }
+
+ log_mu_ready = 1;
+ return 0;
+}
+
+void klog_free(void)
+{
+ if (!log_mu_ready) {
+ return;
+ }
+
+ compat_mutex_destroy(&log_mu);
+ log_mu_ready = 0;
+}
+
void klog_set_worker_id(int id)
{
klog_worker_id = id;
@@ -104,6 +131,9 @@ void log_at(char sev, const char *file, int line, const char *msg, ...)
#ifdef FUZZ
return;
#endif
+ if (log_mu_ready) {
+ compat_mutex_lock(&log_mu);
+ }
struct timespec ts;
struct tm tm;
@@ -145,4 +175,8 @@ void log_at(char sev, const char *file, int line, const char *msg, ...)
va_end(ap);
fputc('\n', stderr);
+
+ if (log_mu_ready) {
+ compat_mutex_unlock(&log_mu);
+ }
}
diff --git a/src/klog.h b/src/klog.h
index c0b73f1..8c65f8a 100644
--- a/src/klog.h
+++ b/src/klog.h
@@ -2,7 +2,11 @@
#define KLOG_H
#include <stdbool.h>
+#include <stddef.h>
#include <stdint.h>
+#include <stdio.h>
+
+#include "compat.h"
struct endpoint;
@@ -33,6 +37,8 @@ struct log_value {
void klog_set_worker_id(int id);
+int klog_init(void);
+void klog_free(void);
void log_at(char sev, const char *file, int line, const char *msg, ...);
static inline struct log_value log_value_str(const char *v)
@@ -93,6 +99,102 @@ static inline struct log_value log_value_endpoint(const struct endpoint *v)
#define _LOGV_ENDPOINT(v) log_value_endpoint(v)
+#ifndef TINYPROXY_SOCKADDR_FORMAT_BUFSIZE
+#define TINYPROXY_SOCKADDR_FORMAT_BUFSIZE 128
+#endif
+
+static inline const char *sockaddr_format(
+ const struct sockaddr *addr,
+ socklen_t addr_len,
+ char *buf,
+ size_t buf_len)
+{
+ if (!buf || buf_len == 0) {
+ return "";
+ }
+
+ buf[0] = '\0';
+
+ if (!addr || addr_len == 0) {
+ snprintf(buf, buf_len, "<null>");
+ return buf;
+ }
+
+ switch (addr->sa_family) {
+ case AF_INET: {
+ const struct sockaddr_in *in = (const struct sockaddr_in *)addr;
+ char ip[INET_ADDRSTRLEN];
+
+ if (addr_len < sizeof(*in)) {
+ snprintf(buf, buf_len, "inet:<short:%u>", (unsigned)addr_len);
+ return buf;
+ }
+
+ if (!inet_ntop(AF_INET, &in->sin_addr, ip, sizeof(ip))) {
+ snprintf(buf, buf_len, "inet:<invalid>:%u",
+ (unsigned)ntohs(in->sin_port));
+ return buf;
+ }
+
+ snprintf(buf, buf_len, "%s:%u",
+ ip,
+ (unsigned)ntohs(in->sin_port));
+ return buf;
+ }
+
+#ifdef AF_INET6
+ case AF_INET6: {
+ const struct sockaddr_in6 *in6 = (const struct sockaddr_in6 *)addr;
+ char ip[INET6_ADDRSTRLEN];
+
+ if (addr_len < sizeof(*in6)) {
+ snprintf(buf, buf_len, "inet6:<short:%u>", (unsigned)addr_len);
+ return buf;
+ }
+
+ if (!inet_ntop(AF_INET6, &in6->sin6_addr, ip, sizeof(ip))) {
+ snprintf(buf, buf_len, "inet6:<invalid>:%u",
+ (unsigned)ntohs(in6->sin6_port));
+ return buf;
+ }
+
+ snprintf(buf, buf_len, "[%s]:%u",
+ ip,
+ (unsigned)ntohs(in6->sin6_port));
+ return buf;
+ }
+#endif
+
+#ifndef _WIN32
+ case AF_UNIX: {
+ const struct sockaddr_un *un = (const struct sockaddr_un *)addr;
+
+ if (addr_len < sizeof(un->sun_family)) {
+ snprintf(buf, buf_len, "unix:<short:%u>", (unsigned)addr_len);
+ return buf;
+ }
+
+ if (un->sun_path[0] == '\0') {
+ snprintf(buf, buf_len, "unix:<abstract-or-empty>");
+ return buf;
+ }
+
+ snprintf(buf, buf_len, "unix:%s", un->sun_path);
+ return buf;
+ }
+#endif
+
+ default:
+ snprintf(buf, buf_len, "family:%d len:%u",
+ (int)addr->sa_family,
+ (unsigned)addr_len);
+ return buf;
+ }
+}
+
+#define _LOGV_SOCKADDR(addr, addr_len, buf, buf_len) \
+ _LOGV(sockaddr_format((const struct sockaddr *)(addr), (addr_len), (buf), (buf_len)))
+
/*
* Structured fields must be passed as:
*
diff --git a/src/route_runtime.c b/src/route_runtime.c
index 39dfc67..0504cff 100644
--- a/src/route_runtime.c
+++ b/src/route_runtime.c
@@ -94,7 +94,7 @@ int start_route(
return -EINVAL;
}
-void stop_route(struct route_ctx *ctx)
+void stop_route_listeners(struct route_ctx *ctx)
{
if (ctx == NULL) {
return;
@@ -102,11 +102,32 @@ void stop_route(struct route_ctx *ctx)
switch (ctx->kind) {
case ROUTE_CTX_STREAM:
- stop_stream_route(&ctx->u.stream);
+ stop_stream_route_listener(&ctx->u.stream);
break;
case ROUTE_CTX_DATAGRAM:
- stop_datagram_route(&ctx->u.datagram);
+ stop_datagram_route_listener(&ctx->u.datagram);
+ break;
+
+ case ROUTE_CTX_NONE:
+ default:
+ break;
+ }
+}
+
+void free_route(struct route_ctx *ctx)
+{
+ if (ctx == NULL) {
+ return;
+ }
+
+ switch (ctx->kind) {
+ case ROUTE_CTX_STREAM:
+ free_stream_route(&ctx->u.stream);
+ break;
+
+ case ROUTE_CTX_DATAGRAM:
+ free_datagram_route(&ctx->u.datagram);
break;
case ROUTE_CTX_NONE:
diff --git a/src/route_runtime.h b/src/route_runtime.h
index 22ad642..ff73e27 100644
--- a/src/route_runtime.h
+++ b/src/route_runtime.h
@@ -29,6 +29,7 @@ int start_route(
const struct route *r,
struct route_ctx *ctx);
-void stop_route(struct route_ctx *ctx);
+void stop_route_listeners(struct route_ctx *ctx);
+void free_route(struct route_ctx *ctx);
#endif
diff --git a/src/stream_route.c b/src/stream_route.c
index fb553ad..69670b5 100644
--- a/src/stream_route.c
+++ b/src/stream_route.c
@@ -92,17 +92,23 @@ static void stop_unix_stream_route(struct stream_route_ctx *ctx)
}
}
-void stop_stream_route(struct stream_route_ctx *ctx)
+void stop_stream_route_listener(struct stream_route_ctx *ctx)
{
- if (ctx == NULL) {
- return;
+ if (ctx->listener) {
+ evconnlistener_free(ctx->listener);
+ ctx->listener = NULL;
}
+}
- if (ctx->listener != NULL) {
- evconnlistener_free(ctx->listener);
+void free_stream_route(struct stream_route_ctx *ctx)
+{
+ if (!ctx) {
+ return;
}
+#ifndef _WIN32
stop_unix_stream_route(ctx);
+#endif
memset(ctx, 0, sizeof(*ctx));
}
diff --git a/src/stream_route.h b/src/stream_route.h
index 747ed14..1b67405 100644
--- a/src/stream_route.h
+++ b/src/stream_route.h
@@ -39,6 +39,7 @@ int start_stream_route(
const struct route *r,
struct stream_route_ctx *ctx);
-void stop_stream_route(struct stream_route_ctx *ctx);
+void stop_stream_route_listener(struct stream_route_ctx *ctx);
+void free_stream_route(struct stream_route_ctx *ctx);
#endif
diff --git a/src/tinyproxy.c b/src/tinyproxy.c
index 105571e..3e130ca 100644
--- a/src/tinyproxy.c
+++ b/src/tinyproxy.c
@@ -48,6 +48,7 @@ static void prevent_socket_write_from_killing_process(void)
int main(int argc, char **argv)
{
prevent_socket_write_from_killing_process();
+ klog_init();
klog_set_worker_id(0);
int opt;
@@ -304,7 +305,7 @@ int main(int argc, char **argv)
out:
for (size_t i = 0; i < route_ctx_count; i++) {
- stop_route(&route_ctxs[i]);
+ stop_route_listeners(&route_ctxs[i]);
}
if (worker_pool_started) {
@@ -312,6 +313,10 @@ out:
worker_pool_join(&wpool);
}
+ for (size_t i = 0; i < route_ctx_count; i++) {
+ free_route(&route_ctxs[i]);
+ }
+
if (worker_pool_ready) {
worker_pool_free(&wpool);
}
@@ -328,6 +333,7 @@ out:
free(routes);
free(inline_routes);
+ klog_free();
#ifdef _WIN32
if (wsa_started) {
diff --git a/src/worker.c b/src/worker.c
index 3031111..c13abec 100644
--- a/src/worker.c
+++ b/src/worker.c
@@ -360,6 +360,8 @@ void worker_free(struct worker *w)
event_base_free(w->base);
w->base = NULL;
}
+
+ klog_free();
}
static int worker_enqueue_msg(struct worker *w, struct worker_msg *msg)
@@ -442,7 +444,6 @@ int worker_enqueue_stream_client(
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,
@@ -451,7 +452,7 @@ int worker_enqueue_datagram_packet(
struct worker_msg *msg;
const struct route *route = ctx->route;
- if (!w || !route || !socket_is_valid(listen_fd)) {
+ if (!w || !route || !socket_is_valid(ctx->listen_fd)) {
return EINVAL;
}
diff --git a/src/worker.h b/src/worker.h
index 4bf1943..fe45a4f 100644
--- a/src/worker.h
+++ b/src/worker.h
@@ -81,7 +81,6 @@ int worker_enqueue_stream_client(
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,