Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions base/hatomic.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,14 @@
#include <atomic>

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

Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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_
19 changes: 11 additions & 8 deletions event/hevent.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@

#include "hbuf.h"
#include "hmutex.h"
#include "hatomic.h"

#include "array.h"
#include "list.h"
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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;
Expand Down
6 changes: 3 additions & 3 deletions event/hloop.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
}

Expand Down
Loading