From e2c335ff731eed28000152473128d445c5e4a2e9 Mon Sep 17 00:00:00 2001 From: "roberto@maverick64" Date: Mon, 25 Apr 2011 07:13:23 +0200 Subject: [PATCH] threaded zeromq/mongrel2 support --- loop.c | 73 ++++++++++++++++++++++++++++++++++++-------------- proto/zeromq.c | 12 +++++---- utils.c | 1 + uwsgi.c | 62 +++++++++++++++++++++++++++++++++--------- uwsgi.h | 3 ++- 5 files changed, 112 insertions(+), 39 deletions(-) diff --git a/loop.c b/loop.c index 8ed12743..e02265d4 100644 --- a/loop.c +++ b/loop.c @@ -2,12 +2,16 @@ extern struct uwsgi_server uwsgi; -struct wsgi_request* threaded_current_wsgi_req() { return pthread_getspecific(uwsgi.tur_key); } -struct wsgi_request* simple_current_wsgi_req() { return uwsgi.wsgi_req; } +struct wsgi_request *threaded_current_wsgi_req() { + return pthread_getspecific(uwsgi.tur_key); +} +struct wsgi_request *simple_current_wsgi_req() { + return uwsgi.wsgi_req; +} void uwsgi_register_loop(char *name, void *loop) { - + if (uwsgi.loops_cnt >= MAX_LOOPS) { uwsgi_log("you can define %d loops at max\n", MAX_LOOPS); exit(1); @@ -21,8 +25,8 @@ void uwsgi_register_loop(char *name, void *loop) { void *uwsgi_get_loop(char *name) { int i; - - for(i=0;iinit_thread) { uwsgi.p[i]->init_thread(core_id); } } - /* - pts = PyThreadState_New(uwsgi.main_thread->interp); - pthread_setspecific(uwsgi.ut_save_key, (void *) pts); - */ } } #endif @@ -85,37 +85,70 @@ void *simple_loop(void *arg1) { #ifdef UWSGI_THREADING pthread_exit(NULL); #endif - + //never here return NULL; } #ifdef UWSGI_ZEROMQ void *zeromq_loop(void *arg1) { + sigset_t smask; + int i; long core_id = (long) arg1; - struct wsgi_request *wsgi_req = uwsgi.wsgi_requests[core_id]; + struct wsgi_request *wsgi_req = uwsgi.wsgi_requests[core_id]; uwsgi.zeromq_recv_flag = 0; + if (uwsgi.threads > 1) { + + pthread_setspecific(uwsgi.tur_key, (void *) wsgi_req); + + if (core_id > 0) { + // block all signals on new threads + sigfillset(&smask); + pthread_sigmask(SIG_BLOCK, &smask, NULL); + for (i = 0; i < 0xFF; i++) { + if (uwsgi.p[i]->init_thread) { + uwsgi.p[i]->init_thread(core_id); + } + } + + + void *tmp_zmq_pull = zmq_socket(uwsgi.zmq_context, ZMQ_PULL); + if (tmp_zmq_pull == NULL) { + uwsgi_error("zmq_socket()"); + exit(1); + } + if (zmq_connect(tmp_zmq_pull, uwsgi.zmq_receiver) < 0) { + uwsgi_error("zmq_connect()"); + exit(1); + } + + pthread_setspecific(uwsgi.zmq_pull, tmp_zmq_pull); + } + + } + + while (uwsgi.workers[uwsgi.mywid].manage_next_request) { - UWSGI_CLEAR_STATUS; + UWSGI_CLEAR_STATUS; - wsgi_req_setup(wsgi_req, core_id, -1); + wsgi_req_setup(wsgi_req, core_id, -1); uwsgi.edge_triggered = 1; int socket_id = uwsgi.zmq_socket; wsgi_req->socket = &uwsgi.sockets[socket_id]; - wsgi_req->poll.fd = wsgi_req->socket->proto_accept(wsgi_req, uwsgi.sockets[socket_id].fd); + wsgi_req->poll.fd = wsgi_req->socket->proto_accept(wsgi_req, uwsgi.sockets[socket_id].fd); - if (wsgi_req->poll.fd >= 0) { - wsgi_req_recv(wsgi_req); - } + if (wsgi_req->poll.fd >= 0) { + wsgi_req_recv(wsgi_req); + } - uwsgi_close_request(wsgi_req); - } + uwsgi_close_request(wsgi_req); + } exit(0); } diff --git a/proto/zeromq.c b/proto/zeromq.c index c41c7b2d..ec5eb2c7 100644 --- a/proto/zeromq.c +++ b/proto/zeromq.c @@ -235,7 +235,7 @@ int uwsgi_proto_zeromq_accept(struct wsgi_request *wsgi_req, int fd) { if (uwsgi.edge_triggered == 0) { - if (zmq_getsockopt(uwsgi.zmq_pull, ZMQ_EVENTS, &events, &events_len) < 0) { + if (zmq_getsockopt(pthread_getspecific(uwsgi.zmq_pull), ZMQ_EVENTS, &events, &events_len) < 0) { uwsgi_error("zmq_getsockopt()"); uwsgi.edge_triggered = 0; return -1; @@ -246,9 +246,7 @@ int uwsgi_proto_zeromq_accept(struct wsgi_request *wsgi_req, int fd) { wsgi_req->do_not_add_to_async_queue = 1; wsgi_req->proto_parser_status = 0; zmq_msg_init(&message); - if (uwsgi.threads > 1) pthread_mutex_lock(&uwsgi.zmq_lock); - if (zmq_recv(uwsgi.zmq_pull, &message, uwsgi.zeromq_recv_flag) < 0) { - if (uwsgi.threads > 1) pthread_mutex_unlock(&uwsgi.zmq_lock); + if (zmq_recv(pthread_getspecific(uwsgi.zmq_pull), &message, uwsgi.zeromq_recv_flag) < 0) { if (errno == EAGAIN) { uwsgi.edge_triggered = 0; } @@ -258,7 +256,6 @@ int uwsgi_proto_zeromq_accept(struct wsgi_request *wsgi_req, int fd) { zmq_msg_close(&message); return -1; } - if (uwsgi.threads > 1) pthread_mutex_unlock(&uwsgi.zmq_lock); uwsgi.edge_triggered = 1; message_size = zmq_msg_size(&message); //uwsgi_log("%.*s\n", (int) wsgi_req->proto_parser_pos, zmq_msg_data(&message)); @@ -370,9 +367,11 @@ void uwsgi_proto_zeromq_close(struct wsgi_request *wsgi_req) { if (!wsgi_req->proto_parser_pos) return; zmq_msg_init_data(&reply, wsgi_req->proto_parser_buf, wsgi_req->proto_parser_pos, uwsgi_proto_zeromq_free, NULL); + pthread_mutex_lock(&uwsgi.zmq_lock); if (zmq_send(uwsgi.zmq_pub, &reply, 0)) { uwsgi_error("zmq_send()"); } + pthread_mutex_unlock(&uwsgi.zmq_lock); zmq_msg_close(&reply); if (wsgi_req->async_post && wsgi_req->body_as_file) { @@ -413,11 +412,14 @@ ssize_t uwsgi_proto_zeromq_write(struct wsgi_request *wsgi_req, char *buf, size_ //uwsgi_log("|%.*s|\n", (int)wsgi_req->proto_parser_pos+len, zmq_body); zmq_msg_init_data(&reply, zmq_body, wsgi_req->proto_parser_pos+len, uwsgi_proto_zeromq_free, NULL); + pthread_mutex_lock(&uwsgi.zmq_lock); if (zmq_send(uwsgi.zmq_pub, &reply, 0)) { uwsgi_error("zmq_send()"); + pthread_mutex_unlock(&uwsgi.zmq_lock); zmq_msg_close(&reply); return -1; } + pthread_mutex_unlock(&uwsgi.zmq_lock); zmq_msg_close(&reply); return len; diff --git a/utils.c b/utils.c index ebe2ddc3..74e42771 100644 --- a/utils.c +++ b/utils.c @@ -574,6 +574,7 @@ int wsgi_req_accept(struct wsgi_request *wsgi_req) { */ polling: + uwsgi.edge_triggered = 1; ret = poll(uwsgi.sockets_poll, uwsgi.sockets_cnt + uwsgi.master_process, uwsgi.edge_triggered - 1); if (ret < 0) { diff --git a/uwsgi.c b/uwsgi.c index d45f20f0..263e2c0e 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -1903,7 +1903,7 @@ int uwsgi_start(void *v_argv) { #ifdef UWSGI_ZEROMQ if (uwsgi.zmq_receiver && uwsgi.zmq_responder) { - uwsgi.zmq_context = zmq_init(4); + uwsgi.zmq_context = zmq_init(1); if (uwsgi.zmq_context == NULL) { uwsgi_error("zmq_init()"); exit(1); @@ -1913,16 +1913,6 @@ int uwsgi_start(void *v_argv) { pthread_mutex_init(&uwsgi.zmq_lock, NULL); } - uwsgi.zmq_pull = zmq_socket(uwsgi.zmq_context, ZMQ_PULL); - if (uwsgi.zmq_pull == NULL) { - uwsgi_error("zmq_socket()"); - exit(1); - } - if (zmq_connect(uwsgi.zmq_pull, uwsgi.zmq_receiver) < 0) { - uwsgi_error("zmq_connect()"); - exit(1); - } - uwsgi.zmq_pub = zmq_socket(uwsgi.zmq_context, ZMQ_PUB); if (uwsgi.zmq_pub == NULL) { uwsgi_error("zmq_socket()"); @@ -1961,12 +1951,31 @@ int uwsgi_start(void *v_argv) { uwsgi.sockets[uwsgi.zmq_socket].proto_sendfile = uwsgi_proto_zeromq_sendfile; uwsgi.sockets[uwsgi.zmq_socket].edge_trigger = 1; + + if (pthread_key_create(&uwsgi.zmq_pull, NULL)) { + uwsgi_error("pthread_key_create()"); + exit(1); + } + + void *tmp_zmq_pull = zmq_socket(uwsgi.zmq_context, ZMQ_PULL); + if (tmp_zmq_pull == NULL) { + uwsgi_error("zmq_socket()"); + exit(1); + } + if (zmq_connect(tmp_zmq_pull, uwsgi.zmq_receiver) < 0) { + uwsgi_error("zmq_connect()"); + exit(1); + } + + pthread_setspecific(uwsgi.zmq_pull, tmp_zmq_pull); + size_t zmq_socket_len = sizeof(int); - if (zmq_getsockopt(uwsgi.zmq_pull, ZMQ_FD, &uwsgi.sockets[uwsgi.zmq_socket].fd, &zmq_socket_len) < 0) { + if (zmq_getsockopt(pthread_getspecific(uwsgi.zmq_pull), ZMQ_FD, &uwsgi.sockets[uwsgi.zmq_socket].fd, &zmq_socket_len) < 0) { uwsgi_error("zmq_getsockopt()"); exit(1); } + uwsgi.sockets_poll[uwsgi.zmq_socket].fd = uwsgi.sockets[uwsgi.zmq_socket].fd; uwsgi.sockets_poll[uwsgi.zmq_socket].events = POLLIN; uwsgi.sockets[uwsgi.zmq_socket].bound = 1; @@ -2085,7 +2094,34 @@ int uwsgi_start(void *v_argv) { } else { #ifdef UWSGI_ZEROMQ - if (uwsgi.zeromq && uwsgi.cores < 2 && uwsgi.sockets_cnt == 1) { + if (uwsgi.zeromq && uwsgi.async < 2 && uwsgi.sockets_cnt == 1) { + + pthread_attr_t pa; + pthread_t *a_thread; + int ret; + + if (uwsgi.threads > 1) { + ret = pthread_attr_init(&pa); + if (ret) { + uwsgi_log("pthread_attr_init() = %d\n", ret); + exit(1); + } + ret = pthread_attr_setdetachstate(&pa, PTHREAD_CREATE_DETACHED); + if (ret) { + uwsgi_log("pthread_attr_setdetachstate() = %d\n", ret); + exit(1); + } + if (pthread_key_create(&uwsgi.tur_key, NULL)) { + uwsgi_error("pthread_key_create()"); + exit(1); + } + for (i = 1; i < uwsgi.threads; i++) { + long j = i; + a_thread = uwsgi_malloc(sizeof(pthread_t)); + pthread_create(a_thread, &pa, zeromq_loop, (void *) j); + } + } + long y = 0; zeromq_loop((void *) y); } diff --git a/uwsgi.h b/uwsgi.h index 690c4f9c..f13194c2 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -1076,10 +1076,11 @@ struct uwsgi_server { char *zmq_responder; int zmq_socket; void *zmq_context; - void *zmq_pull; + //void *zmq_pull; void *zmq_pub; int zeromq_recv_flag; pthread_mutex_t zmq_lock; + pthread_key_t zmq_pull; #endif struct uwsgi_socket sockets[MAX_SOCKETS]; // leave a slot for no-orphan mode