threaded zeromq/mongrel2 support

This commit is contained in:
roberto@maverick64
2011-04-25 07:13:23 +02:00
parent 816b9d0343
commit e2c335ff73
5 changed files with 112 additions and 39 deletions
+53 -20
View File
@@ -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;i<uwsgi.loops_cnt;i++) {
for (i = 0; i < uwsgi.loops_cnt; i++) {
if (!strcmp(name, uwsgi.loops[i].name)) {
return uwsgi.loops[i].loop;
}
@@ -50,15 +54,11 @@ void *simple_loop(void *arg1) {
// block all signals on new threads
sigfillset(&smask);
pthread_sigmask(SIG_BLOCK, &smask, NULL);
for(i=0;i<0xFF;i++) {
for (i = 0; i < 0xFF; i++) {
if (uwsgi.p[i]->init_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);
}
+7 -5
View File
@@ -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;
+1
View File
@@ -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) {
+49 -13
View File
@@ -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);
}
+2 -1
View File
@@ -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