From 4755bfa9e9260b521e10745132c961ef134596c5 Mon Sep 17 00:00:00 2001 From: aleksisch Date: Thu, 24 Sep 2026 16:20:11 +0300 Subject: [PATCH] fix: data races on loop status, io state flags and hrtime hio_write/hio_close are documented thread-safe, but they read io->ready/closed unlocked while the loop thread rewrites neighbouring bits of the same bitfield word, and hio_write stamps last_write_hrtime from the loop's cur_hrtime. Make those fields atomic, and make loop->status atomic, stored before hloop_stop's wakeup so the woken loop sees STOP. Found with ThreadSanitizer. --- base/hatomic.h | 17 +++++++++++++++++ event/hevent.h | 19 +++++++++++-------- event/hloop.c | 6 +++--- 3 files changed, 31 insertions(+), 11 deletions(-) diff --git a/base/hatomic.h b/base/hatomic.h index 4885c65cf..ff381cce9 100644 --- a/base/hatomic.h +++ b/base/hatomic.h @@ -7,10 +7,14 @@ #include using std::atomic_flag; +using std::atomic_bool; +using std::atomic_int; using std::atomic_long; +using std::atomic_ullong; #define ATOMIC_FLAG_TEST_AND_SET(p) ((p)->test_and_set()) #define ATOMIC_FLAG_CLEAR(p) ((p)->clear()) +#define ATOMIC_EXCHANGE(p, v) ((p)->exchange(v)) #else @@ -22,6 +26,7 @@ using std::atomic_long; #define ATOMIC_FLAG_TEST_AND_SET atomic_flag_test_and_set #define ATOMIC_FLAG_CLEAR atomic_flag_clear +#define ATOMIC_EXCHANGE atomic_exchange #define ATOMIC_ADD atomic_fetch_add #define ATOMIC_SUB atomic_fetch_sub #define ATOMIC_INC(p) ATOMIC_ADD(p, 1) @@ -51,6 +56,7 @@ static inline bool atomic_flag_test_and_set(atomic_flag* p) { return !__sync_bool_compare_and_swap(&p->_Value, 0, 1); } +#define ATOMIC_EXCHANGE(p, v) (__sync_synchronize(), __sync_lock_test_and_set(p, v)) #define ATOMIC_ADD __sync_fetch_and_add #define ATOMIC_SUB __sync_fetch_and_sub #define ATOMIC_INC(p) ATOMIC_ADD(p, 1) @@ -64,6 +70,7 @@ static inline bool atomic_flag_test_and_set(atomic_flag* p) { return InterlockedCompareExchange(&p->_Value, 1, 0); } +#define ATOMIC_EXCHANGE(p, v) InterlockedExchange((volatile LONG*)(p), (v)) #define ATOMIC_ADD InterlockedExchangeAdd #define ATOMIC_SUB(p, n) InterlockedExchangeAdd(p, -(n)) #define ATOMIC_INC(p) InterlockedExchangeAdd(p, 1) @@ -107,6 +114,15 @@ static inline void atomic_flag_clear(atomic_flag* p) { #define ATOMIC_SUB(p, n) (*(p) -= (n)) #endif +#ifndef ATOMIC_EXCHANGE +#define ATOMIC_EXCHANGE atomic_exchange +static inline int atomic_exchange(atomic_int* p, int v) { + int old = *p; + *p = v; + return old; +} +#endif + #ifndef ATOMIC_INC #define ATOMIC_INC(p) ((*(p))++) #endif @@ -126,5 +142,6 @@ typedef atomic_long hatomic_t; #define hatomic_sub ATOMIC_SUB #define hatomic_inc ATOMIC_INC #define hatomic_dec ATOMIC_DEC +#define hatomic_exchange ATOMIC_EXCHANGE #endif // HV_ATOMIC_H_ diff --git a/event/hevent.h b/event/hevent.h index 92c69070e..a4a1ed465 100644 --- a/event/hevent.h +++ b/event/hevent.h @@ -7,6 +7,7 @@ #include "hbuf.h" #include "hmutex.h" +#include "hatomic.h" #include "array.h" #include "list.h" @@ -46,11 +47,11 @@ QUEUE_DECL(hevent_t, event_queue); struct hloop_s { uint32_t flags; - hloop_status_e status; + atomic_int status; // hloop_status_e; storing RUNNING publishes pid/tid to hloop_stop uint64_t start_ms; // ms uint64_t start_hrtime; // us uint64_t end_hrtime; - uint64_t cur_hrtime; + atomic_ullong cur_hrtime; // read by hio_write from other threads uint64_t loop_cnt; long pid; long tid; @@ -131,13 +132,15 @@ struct hperiod_s { }; QUEUE_DECL(offset_buf_t, write_queue); -// sizeof(struct hio_s)=424 on linux-x64 +// sizeof(struct hio_s)=432 on linux-x64 struct hio_s { HEVENT_FIELDS + // read without a lock by hio_write/hio_close/hio_is_opened from other threads, + // so they cannot share a bitfield word with the flags the loop thread rewrites + atomic_bool ready; + atomic_bool connected; + atomic_bool closed; // flags - unsigned ready :1; - unsigned connected :1; - unsigned closed :1; unsigned accept :1; unsigned connect :1; unsigned recv :1; @@ -157,8 +160,8 @@ struct hio_s { int revents; struct sockaddr* localaddr; struct sockaddr* peeraddr; - uint64_t last_read_hrtime; - uint64_t last_write_hrtime; + atomic_ullong last_read_hrtime; + atomic_ullong last_write_hrtime; // written by hio_write from other threads // read fifo_buf_t readbuf; unsigned int read_flags; diff --git a/event/hloop.c b/event/hloop.c index 46ed2eb11..f9ab8c1c7 100644 --- a/event/hloop.c +++ b/event/hloop.c @@ -470,9 +470,9 @@ int hloop_run(hloop_t* loop) { if (loop == NULL) return -1; if (loop->status == HLOOP_STATUS_RUNNING) return -2; - loop->status = HLOOP_STATUS_RUNNING; loop->pid = hv_getpid(); loop->tid = hv_gettid(); + loop->status = HLOOP_STATUS_RUNNING; hlogd("hloop_run tid=%ld", loop->tid); if (loop->intern_nevents == 0) { @@ -523,12 +523,12 @@ int hloop_wakeup(hloop_t* loop) { int hloop_stop(hloop_t* loop) { if (loop == NULL) return -1; - if (loop->status == HLOOP_STATUS_STOP) return -2; + // set the status before waking the loop, or it can wake, see RUNNING and block again + if (hatomic_exchange(&loop->status, HLOOP_STATUS_STOP) == HLOOP_STATUS_STOP) return -2; hlogd("hloop_stop tid=%ld", hv_gettid()); if (hv_gettid() != loop->tid) { hloop_wakeup(loop); } - loop->status = HLOOP_STATUS_STOP; return 0; }