|
|
|
@ -39,38 +39,66 @@ |
|
|
|
|
|
|
|
|
|
#include "src/core/iomgr/alarm_internal.h" |
|
|
|
|
#include "src/core/iomgr/iomgr_internal.h" |
|
|
|
|
#include "src/core/iomgr/iocp_windows.h" |
|
|
|
|
#include "src/core/iomgr/pollset.h" |
|
|
|
|
#include "src/core/iomgr/pollset_windows.h" |
|
|
|
|
|
|
|
|
|
static void remove_worker(grpc_pollset *p, grpc_pollset_worker *worker) { |
|
|
|
|
worker->prev->next = worker->next; |
|
|
|
|
worker->next->prev = worker->prev; |
|
|
|
|
gpr_mu grpc_polling_mu; |
|
|
|
|
static grpc_pollset_worker *g_active_poller; |
|
|
|
|
static grpc_pollset_worker g_global_root_worker; |
|
|
|
|
|
|
|
|
|
void grpc_pollset_global_init() { |
|
|
|
|
gpr_mu_init(&grpc_polling_mu); |
|
|
|
|
g_active_poller = NULL; |
|
|
|
|
g_global_root_worker.links[GRPC_POLLSET_WORKER_LINK_GLOBAL].next = |
|
|
|
|
g_global_root_worker.links[GRPC_POLLSET_WORKER_LINK_GLOBAL].prev = |
|
|
|
|
&g_global_root_worker; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void grpc_pollset_global_shutdown() { |
|
|
|
|
gpr_mu_destroy(&grpc_polling_mu); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void remove_worker(grpc_pollset_worker *worker,
|
|
|
|
|
grpc_pollset_worker_link_type type) { |
|
|
|
|
worker->links[type].prev->links[type].next = worker->links[type].next; |
|
|
|
|
worker->links[type].next->links[type].prev = worker->links[type].prev; |
|
|
|
|
worker->links[type].next = worker->links[type].prev = worker; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static int has_workers(grpc_pollset *p) { |
|
|
|
|
return p->root_worker.next != &p->root_worker; |
|
|
|
|
static int has_workers(grpc_pollset_worker *root, grpc_pollset_worker_link_type type) { |
|
|
|
|
return root->links[type].next != root; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static grpc_pollset_worker *pop_front_worker(grpc_pollset *p) { |
|
|
|
|
if (has_workers(p)) { |
|
|
|
|
grpc_pollset_worker *w = p->root_worker.next; |
|
|
|
|
remove_worker(p, w); |
|
|
|
|
static grpc_pollset_worker *pop_front_worker( |
|
|
|
|
grpc_pollset_worker *root, grpc_pollset_worker_link_type type) { |
|
|
|
|
if (has_workers(root, type)) { |
|
|
|
|
grpc_pollset_worker *w = root->links[type].next; |
|
|
|
|
remove_worker(w, type); |
|
|
|
|
return w; |
|
|
|
|
} else { |
|
|
|
|
return NULL; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void push_back_worker(grpc_pollset *p, grpc_pollset_worker *worker) { |
|
|
|
|
worker->next = &p->root_worker; |
|
|
|
|
worker->prev = worker->next->prev; |
|
|
|
|
worker->prev->next = worker->next->prev = worker; |
|
|
|
|
static void push_back_worker(grpc_pollset_worker *root,
|
|
|
|
|
grpc_pollset_worker_link_type type,
|
|
|
|
|
grpc_pollset_worker *worker) { |
|
|
|
|
worker->links[type].next = root; |
|
|
|
|
worker->links[type].prev = worker->links[type].next->links[type].prev; |
|
|
|
|
worker->links[type].prev->links[type].next =
|
|
|
|
|
worker->links[type].next->links[type].prev =
|
|
|
|
|
worker; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void push_front_worker(grpc_pollset *p, grpc_pollset_worker *worker) { |
|
|
|
|
worker->prev = &p->root_worker; |
|
|
|
|
worker->next = worker->prev->next; |
|
|
|
|
worker->prev->next = worker->next->prev = worker; |
|
|
|
|
static void push_front_worker(grpc_pollset_worker *root,
|
|
|
|
|
grpc_pollset_worker_link_type type,
|
|
|
|
|
grpc_pollset_worker *worker) { |
|
|
|
|
worker->links[type].prev = root; |
|
|
|
|
worker->links[type].next = worker->links[type].prev->links[type].next; |
|
|
|
|
worker->links[type].prev->links[type].next =
|
|
|
|
|
worker->links[type].next->links[type].prev =
|
|
|
|
|
worker; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* There isn't really any such thing as a pollset under Windows, due to the
|
|
|
|
@ -80,69 +108,122 @@ static void push_front_worker(grpc_pollset *p, grpc_pollset_worker *worker) { |
|
|
|
|
|
|
|
|
|
void grpc_pollset_init(grpc_pollset *pollset) { |
|
|
|
|
memset(pollset, 0, sizeof(*pollset)); |
|
|
|
|
gpr_mu_init(&pollset->mu); |
|
|
|
|
pollset->root_worker.next = pollset->root_worker.prev = &pollset->root_worker; |
|
|
|
|
pollset->kicked_without_pollers = 0; |
|
|
|
|
pollset->root_worker.links[GRPC_POLLSET_WORKER_LINK_POLLSET].next =
|
|
|
|
|
pollset->root_worker.links[GRPC_POLLSET_WORKER_LINK_POLLSET].prev =
|
|
|
|
|
&pollset->root_worker; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void grpc_pollset_shutdown(grpc_exec_ctx *exec_ctx, grpc_pollset *pollset, |
|
|
|
|
grpc_closure *closure) { |
|
|
|
|
gpr_mu_lock(&pollset->mu); |
|
|
|
|
gpr_mu_lock(&grpc_polling_mu); |
|
|
|
|
pollset->shutting_down = 1; |
|
|
|
|
grpc_pollset_kick(pollset, GRPC_POLLSET_KICK_BROADCAST); |
|
|
|
|
gpr_mu_unlock(&pollset->mu); |
|
|
|
|
grpc_exec_ctx_enqueue(exec_ctx, closure, 1); |
|
|
|
|
if (!pollset->is_iocp_worker) { |
|
|
|
|
grpc_exec_ctx_enqueue(exec_ctx, closure, 1); |
|
|
|
|
} else { |
|
|
|
|
pollset->on_shutdown = closure; |
|
|
|
|
} |
|
|
|
|
gpr_mu_unlock(&grpc_polling_mu); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void grpc_pollset_destroy(grpc_pollset *pollset) { |
|
|
|
|
gpr_mu_destroy(&pollset->mu); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void grpc_pollset_work(grpc_exec_ctx *exec_ctx, grpc_pollset *pollset, |
|
|
|
|
grpc_pollset_worker *worker, gpr_timespec now, |
|
|
|
|
gpr_timespec deadline) { |
|
|
|
|
int added_worker = 0; |
|
|
|
|
worker->next = worker->prev = NULL; |
|
|
|
|
worker->links[GRPC_POLLSET_WORKER_LINK_POLLSET].next =
|
|
|
|
|
worker->links[GRPC_POLLSET_WORKER_LINK_POLLSET].prev =
|
|
|
|
|
worker->links[GRPC_POLLSET_WORKER_LINK_GLOBAL].next = |
|
|
|
|
worker->links[GRPC_POLLSET_WORKER_LINK_GLOBAL].prev = |
|
|
|
|
NULL; |
|
|
|
|
worker->kicked = 0; |
|
|
|
|
worker->pollset = pollset; |
|
|
|
|
gpr_cv_init(&worker->cv); |
|
|
|
|
if (grpc_alarm_check(exec_ctx, now, &deadline)) { |
|
|
|
|
goto done; |
|
|
|
|
} |
|
|
|
|
if (!pollset->kicked_without_pollers && !pollset->shutting_down) { |
|
|
|
|
push_front_worker(pollset, worker); |
|
|
|
|
if (g_active_poller == NULL) { |
|
|
|
|
grpc_pollset_worker *next_worker; |
|
|
|
|
/* become poller */ |
|
|
|
|
pollset->is_iocp_worker = 1; |
|
|
|
|
g_active_poller = worker; |
|
|
|
|
gpr_mu_unlock(&grpc_polling_mu); |
|
|
|
|
grpc_iocp_work(exec_ctx, deadline); |
|
|
|
|
grpc_exec_ctx_flush(exec_ctx); |
|
|
|
|
gpr_mu_lock(&grpc_polling_mu); |
|
|
|
|
pollset->is_iocp_worker = 0; |
|
|
|
|
g_active_poller = NULL; |
|
|
|
|
/* try to get a worker from this pollsets worker list */ |
|
|
|
|
next_worker = pop_front_worker(&pollset->root_worker, GRPC_POLLSET_WORKER_LINK_POLLSET); |
|
|
|
|
/* try to get a worker from the global list */ |
|
|
|
|
next_worker = pop_front_worker(&g_global_root_worker, GRPC_POLLSET_WORKER_LINK_GLOBAL); |
|
|
|
|
if (next_worker != NULL) { |
|
|
|
|
next_worker->kicked = 1; |
|
|
|
|
gpr_cv_signal(&next_worker->cv); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
if (pollset->shutting_down && pollset->on_shutdown != NULL) { |
|
|
|
|
grpc_exec_ctx_enqueue(exec_ctx, pollset->on_shutdown, 1); |
|
|
|
|
pollset->on_shutdown = NULL; |
|
|
|
|
} |
|
|
|
|
goto done; |
|
|
|
|
} |
|
|
|
|
push_front_worker(&g_global_root_worker, GRPC_POLLSET_WORKER_LINK_GLOBAL, worker); |
|
|
|
|
push_front_worker(&pollset->root_worker, GRPC_POLLSET_WORKER_LINK_POLLSET, worker); |
|
|
|
|
added_worker = 1; |
|
|
|
|
gpr_cv_wait(&worker->cv, &pollset->mu, deadline); |
|
|
|
|
while (!worker->kicked) { |
|
|
|
|
if (gpr_cv_wait(&worker->cv, &grpc_polling_mu, deadline)) { |
|
|
|
|
break; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
pollset->kicked_without_pollers = 0; |
|
|
|
|
} |
|
|
|
|
done: |
|
|
|
|
if (!grpc_closure_list_empty(exec_ctx->closure_list)) { |
|
|
|
|
gpr_mu_unlock(&pollset->mu); |
|
|
|
|
gpr_mu_unlock(&grpc_polling_mu); |
|
|
|
|
grpc_exec_ctx_flush(exec_ctx); |
|
|
|
|
gpr_mu_lock(&pollset->mu); |
|
|
|
|
gpr_mu_lock(&grpc_polling_mu); |
|
|
|
|
} |
|
|
|
|
gpr_cv_destroy(&worker->cv); |
|
|
|
|
if (added_worker) { |
|
|
|
|
remove_worker(pollset, worker); |
|
|
|
|
remove_worker(worker, GRPC_POLLSET_WORKER_LINK_GLOBAL); |
|
|
|
|
remove_worker(worker, GRPC_POLLSET_WORKER_LINK_POLLSET); |
|
|
|
|
} |
|
|
|
|
gpr_cv_destroy(&worker->cv); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void grpc_pollset_kick(grpc_pollset *p, grpc_pollset_worker *specific_worker) { |
|
|
|
|
if (specific_worker != NULL) { |
|
|
|
|
if (specific_worker == GRPC_POLLSET_KICK_BROADCAST) { |
|
|
|
|
for (specific_worker = p->root_worker.next; |
|
|
|
|
for (specific_worker = p->root_worker.links[GRPC_POLLSET_WORKER_LINK_POLLSET].next; |
|
|
|
|
specific_worker != &p->root_worker; |
|
|
|
|
specific_worker = specific_worker->next) { |
|
|
|
|
specific_worker = specific_worker->links[GRPC_POLLSET_WORKER_LINK_POLLSET].next) { |
|
|
|
|
specific_worker->kicked = 1; |
|
|
|
|
gpr_cv_signal(&specific_worker->cv); |
|
|
|
|
} |
|
|
|
|
p->kicked_without_pollers = 1; |
|
|
|
|
if (p->is_iocp_worker) { |
|
|
|
|
grpc_iocp_kick(); |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
gpr_cv_signal(&specific_worker->cv); |
|
|
|
|
if (p->is_iocp_worker) { |
|
|
|
|
if (g_active_poller == specific_worker) { |
|
|
|
|
grpc_iocp_kick(); |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
specific_worker->kicked = 1; |
|
|
|
|
gpr_cv_signal(&specific_worker->cv); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
specific_worker = pop_front_worker(p); |
|
|
|
|
specific_worker = pop_front_worker(&p->root_worker, GRPC_POLLSET_WORKER_LINK_POLLSET); |
|
|
|
|
if (specific_worker != NULL) { |
|
|
|
|
push_back_worker(p, specific_worker); |
|
|
|
|
gpr_cv_signal(&specific_worker->cv); |
|
|
|
|
grpc_pollset_kick(p, specific_worker); |
|
|
|
|
} else if (p->is_iocp_worker) { |
|
|
|
|
grpc_iocp_kick(); |
|
|
|
|
} else { |
|
|
|
|
p->kicked_without_pollers = 1; |
|
|
|
|
} |
|
|
|
|