diff --git a/async.c b/async.c new file mode 100644 index 00000000..d8a2a6af --- /dev/null +++ b/async.c @@ -0,0 +1,131 @@ +#ifdef UWSGI_ASYNC + +#include "uwsgi.h" + + +#ifdef __linux__ + +#include + +int async_queue_init(int serverfd) { + int epfd ; + struct epoll_event ee; + + epfd = epoll_create(256); + + if (epfd < 0) { + perror("epoll_create()"); + return -1 ; + } + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.events = EPOLLIN; + ee.data.fd = serverfd; + + if (epoll_ctl(epfd, EPOLL_CTL_ADD, serverfd, &ee)) { + perror("epoll_ctl()"); + close(epfd); + return -1; + } + + return epfd; +} + +int async_add(int queuefd, int fd, int etype) { + struct epoll_event ee; + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.events = etype; + ee.data.fd = fd; + + if (epoll_ctl(queuefd, EPOLL_CTL_ADD, fd, &ee)) { + perror("epoll_ctl()"); + return -1; + } + + return 0; +} +#endif + +struct wsgi_request *next_wsgi_req(struct uwsgi_server *uwsgi, struct wsgi_request *wsgi_req) { + + uint8_t *ptr = (uint8_t *) wsgi_req ; + + ptr += sizeof(struct wsgi_request)+(uwsgi->buffer_size-1) ; + + return (struct wsgi_request *) ptr ; +} +struct wsgi_request *find_first_available_wsgi_req(struct uwsgi_server *uwsgi) { + + struct wsgi_request* wsgi_req = uwsgi->wsgi_requests ; + int i ; + + for(i=0;iasync;i++) { + //fprintf(stderr,"request %d fd %d switches %d\n", i, wsgi_req->poll.fd, wsgi_req->async_switches); + if (wsgi_req->async_status == 0) { + return wsgi_req ; + } + wsgi_req = next_wsgi_req(uwsgi, wsgi_req) ; + } + + return NULL ; +} + +struct wsgi_request *find_wsgi_req_by_fd(struct uwsgi_server *uwsgi, int fd, int etype) { + + struct wsgi_request* wsgi_req = uwsgi->wsgi_requests ; + int i ; + + for(i=0;iasync;i++) { + if (wsgi_req->async_waiting_fd == fd && wsgi_req->async_waiting_fd_type & etype) { + return wsgi_req ; + } + wsgi_req = next_wsgi_req(uwsgi, wsgi_req) ; + } + + return NULL ; + +} + + +struct wsgi_request * async_loop(struct uwsgi_server *uwsgi) { + + struct wsgi_request *wsgi_req ; + int i ; + + uwsgi->async_running = -1 ; + wsgi_req = uwsgi->wsgi_requests ; + + + for(i=0;iasync;i++) { + if (wsgi_req->async_status == UWSGI_AGAIN) { + //fprintf(stderr,"REQUEST MONITORED %d %d\n",wsgi_req->async_waiting_fd, wsgi_req->async_waiting_fd_monitored); + if (wsgi_req->async_waiting_fd != -1 && !wsgi_req->async_waiting_fd_monitored) { + // add fd to monitoring + if (async_add(uwsgi->async_queue, wsgi_req->async_waiting_fd, wsgi_req->async_waiting_fd_type)) { + // error adding fd to the async queue, better to close it... + close(wsgi_req->async_waiting_fd); + wsgi_req->async_status = UWSGI_OK ; + return wsgi_req; + } + wsgi_req->async_waiting_fd_monitored = 1; + wsgi_req->async_status = UWSGI_AGAIN; + } + else if (wsgi_req->async_waiting_fd == -1) { + uwsgi->async_running = 0 ; + // st global wsgi_req for python functions + uwsgi->wsgi_req = wsgi_req ; + wsgi_req->async_status = (*uwsgi->shared->hooks[wsgi_req->modifier]) (uwsgi, wsgi_req); + + if (wsgi_req->async_status < UWSGI_AGAIN) { + return wsgi_req; + } + } + } + wsgi_req = next_wsgi_req(uwsgi, wsgi_req) ; + } + + return NULL; + +} +#endif diff --git a/tests/__init__.py b/tests/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/tests/cpubound_async.py b/tests/cpubound_async.py new file mode 100644 index 00000000..18ab1ceb --- /dev/null +++ b/tests/cpubound_async.py @@ -0,0 +1,5 @@ + +def application(env, start_response): + start_response( '200 OK', [ ('Content-Type','text/html') ]) + for i in range(1,100000): + yield "

%s

" % i