mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-25 12:43:09 +00:00
Check quit flag before pthread_cond_wait
This commit is contained in:
@@ -22,6 +22,7 @@ struct monitor {
|
|||||||
pthread_cond_t cond;
|
pthread_cond_t cond;
|
||||||
pthread_mutex_t mutex;
|
pthread_mutex_t mutex;
|
||||||
int sleep;
|
int sleep;
|
||||||
|
int quit;
|
||||||
};
|
};
|
||||||
|
|
||||||
struct worker_parm {
|
struct worker_parm {
|
||||||
@@ -49,7 +50,7 @@ wakeup(struct monitor *m, int busy) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
static void *
|
static void *
|
||||||
_socket(void *p) {
|
thread_socket(void *p) {
|
||||||
struct monitor * m = p;
|
struct monitor * m = p;
|
||||||
skynet_initthread(THREAD_SOCKET);
|
skynet_initthread(THREAD_SOCKET);
|
||||||
for (;;) {
|
for (;;) {
|
||||||
@@ -79,7 +80,7 @@ free_monitor(struct monitor *m) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
static void *
|
static void *
|
||||||
_monitor(void *p) {
|
thread_monitor(void *p) {
|
||||||
struct monitor * m = p;
|
struct monitor * m = p;
|
||||||
int i;
|
int i;
|
||||||
int n = m->count;
|
int n = m->count;
|
||||||
@@ -99,7 +100,7 @@ _monitor(void *p) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
static void *
|
static void *
|
||||||
_timer(void *p) {
|
thread_timer(void *p) {
|
||||||
struct monitor * m = p;
|
struct monitor * m = p;
|
||||||
skynet_initthread(THREAD_TIMER);
|
skynet_initthread(THREAD_TIMER);
|
||||||
for (;;) {
|
for (;;) {
|
||||||
@@ -111,12 +112,14 @@ _timer(void *p) {
|
|||||||
// wakeup socket thread
|
// wakeup socket thread
|
||||||
skynet_socket_exit();
|
skynet_socket_exit();
|
||||||
// wakeup all worker thread
|
// wakeup all worker thread
|
||||||
|
// we don't need lock m before set m->quit, because it can only set once.
|
||||||
|
m->quit = 1;
|
||||||
pthread_cond_broadcast(&m->cond);
|
pthread_cond_broadcast(&m->cond);
|
||||||
return NULL;
|
return NULL;
|
||||||
}
|
}
|
||||||
|
|
||||||
static void *
|
static void *
|
||||||
_worker(void *p) {
|
thread_worker(void *p) {
|
||||||
struct worker_parm *wp = p;
|
struct worker_parm *wp = p;
|
||||||
int id = wp->id;
|
int id = wp->id;
|
||||||
int weight = wp->weight;
|
int weight = wp->weight;
|
||||||
@@ -124,13 +127,14 @@ _worker(void *p) {
|
|||||||
struct skynet_monitor *sm = m->m[id];
|
struct skynet_monitor *sm = m->m[id];
|
||||||
skynet_initthread(THREAD_WORKER);
|
skynet_initthread(THREAD_WORKER);
|
||||||
struct message_queue * q = NULL;
|
struct message_queue * q = NULL;
|
||||||
for (;;) {
|
while (!m->quit) {
|
||||||
q = skynet_context_message_dispatch(sm, q, weight);
|
q = skynet_context_message_dispatch(sm, q, weight);
|
||||||
if (q == NULL) {
|
if (q == NULL) {
|
||||||
if (pthread_mutex_lock(&m->mutex) == 0) {
|
if (pthread_mutex_lock(&m->mutex) == 0) {
|
||||||
++ m->sleep;
|
++ m->sleep;
|
||||||
// "spurious wakeup" is harmless,
|
// "spurious wakeup" is harmless,
|
||||||
// because skynet_context_message_dispatch() can be call at any time.
|
// because skynet_context_message_dispatch() can be call at any time.
|
||||||
|
if (!m->quit)
|
||||||
pthread_cond_wait(&m->cond, &m->mutex);
|
pthread_cond_wait(&m->cond, &m->mutex);
|
||||||
-- m->sleep;
|
-- m->sleep;
|
||||||
if (pthread_mutex_unlock(&m->mutex)) {
|
if (pthread_mutex_unlock(&m->mutex)) {
|
||||||
@@ -139,13 +143,12 @@ _worker(void *p) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
CHECK_ABORT
|
|
||||||
}
|
}
|
||||||
return NULL;
|
return NULL;
|
||||||
}
|
}
|
||||||
|
|
||||||
static void
|
static void
|
||||||
_start(int thread) {
|
start(int thread) {
|
||||||
pthread_t pid[thread+3];
|
pthread_t pid[thread+3];
|
||||||
|
|
||||||
struct monitor *m = skynet_malloc(sizeof(*m));
|
struct monitor *m = skynet_malloc(sizeof(*m));
|
||||||
@@ -167,9 +170,9 @@ _start(int thread) {
|
|||||||
exit(1);
|
exit(1);
|
||||||
}
|
}
|
||||||
|
|
||||||
create_thread(&pid[0], _monitor, m);
|
create_thread(&pid[0], thread_monitor, m);
|
||||||
create_thread(&pid[1], _timer, m);
|
create_thread(&pid[1], thread_timer, m);
|
||||||
create_thread(&pid[2], _socket, m);
|
create_thread(&pid[2], thread_socket, m);
|
||||||
|
|
||||||
static int weight[] = {
|
static int weight[] = {
|
||||||
-1, -1, -1, -1, 0, 0, 0, 0,
|
-1, -1, -1, -1, 0, 0, 0, 0,
|
||||||
@@ -185,7 +188,7 @@ _start(int thread) {
|
|||||||
} else {
|
} else {
|
||||||
wp[i].weight = 0;
|
wp[i].weight = 0;
|
||||||
}
|
}
|
||||||
create_thread(&pid[i+3], _worker, &wp[i]);
|
create_thread(&pid[i+3], thread_worker, &wp[i]);
|
||||||
}
|
}
|
||||||
|
|
||||||
for (i=0;i<thread+3;i++) {
|
for (i=0;i<thread+3;i++) {
|
||||||
@@ -231,7 +234,7 @@ skynet_start(struct skynet_config * config) {
|
|||||||
|
|
||||||
bootstrap(ctx, config->bootstrap);
|
bootstrap(ctx, config->bootstrap);
|
||||||
|
|
||||||
_start(config->thread);
|
start(config->thread);
|
||||||
|
|
||||||
// harbor_exit may call socket send, so it should exit before socket_free
|
// harbor_exit may call socket send, so it should exit before socket_free
|
||||||
skynet_harbor_exit();
|
skynet_harbor_exit();
|
||||||
|
|||||||
Reference in New Issue
Block a user