From 937b2de41babd1b066d8a67939f1f659e70df935 Mon Sep 17 00:00:00 2001 From: "roberto@maverick64" Date: Wed, 2 Feb 2011 16:38:51 +0100 Subject: [PATCH] the uWSGI fastrouter --- buildconf/default.ini | 4 +- cache.c | 14 +- emperor.c | 96 +++++-- event.c | 31 +++ plugins/cgi/cgi_plugin.c | 139 ++++++++++ plugins/cgi/uwsgiplugin.py | 5 + plugins/fastrouter/fastrouter.c | 415 ++++++++++++++++++++++++++++++ plugins/fastrouter/uwsgiplugin.py | 7 + plugins/python/python_plugin.c | 4 +- plugins/python/uwsgi_pymodule.c | 55 +++- protocol.c | 205 ++++++++++++++- regexp.c | 35 +++ socket.c | 11 + utils.c | 13 +- uwsgi.c | 5 + uwsgi.h | 17 +- uwsgiconfig.py | 42 +-- 17 files changed, 1017 insertions(+), 81 deletions(-) create mode 100644 plugins/cgi/cgi_plugin.c create mode 100644 plugins/cgi/uwsgiplugin.py create mode 100644 plugins/fastrouter/fastrouter.c create mode 100644 plugins/fastrouter/uwsgiplugin.py create mode 100644 regexp.c diff --git a/buildconf/default.ini b/buildconf/default.ini index a404be47..af8197a9 100644 --- a/buildconf/default.ini +++ b/buildconf/default.ini @@ -16,7 +16,7 @@ async = true http = true evdis = false ldap = false -routing = false +pcre = auto stackless = false debug = false unbit = false @@ -24,7 +24,7 @@ xml_implementation = libxml2 plugins = bin_name = uwsgi plugin_dir = . -embedded_plugins = python, ping, proxy, nagios, rpc +embedded_plugins = python, ping, proxy, nagios, rpc, fastrouter locking = auto event = auto diff --git a/cache.c b/cache.c index c00a6248..a98ece18 100644 --- a/cache.c +++ b/cache.c @@ -148,6 +148,8 @@ int uwsgi_cache_request(struct wsgi_request *wsgi_req) { uint16_t vallen = 0; char *value; + char *argv[3]; + uint8_t argc = 0; switch(wsgi_req->uh.modifier2) { case 0: @@ -155,12 +157,22 @@ int uwsgi_cache_request(struct wsgi_request *wsgi_req) { if (wsgi_req->uh.pktsize > 0) { value = uwsgi_cache_get(wsgi_req->buffer, wsgi_req->uh.pktsize, &vallen); if (value && vallen > 0) { - wsgi_req->response_size = write(wsgi_req->poll.fd, value, vallen); + wsgi_req->uh.pktsize = vallen; + wsgi_req->response_size = write(wsgi_req->poll.fd, &wsgi_req->uh, 4); + wsgi_req->response_size += write(wsgi_req->poll.fd, value, vallen); } } break; case 1: // set + if (wsgi_req->uh.pktsize > 0) { + argc = 3; + if (!uwsgi_parse_array(wsgi_req->buffer, wsgi_req->uh.pktsize, argv, &argc)) { + if (argc > 1) { + uwsgi_cache_set(argv[0], strlen(argv[0]), argv[1], strlen(argv[1]), 0); + } + } + } break; case 2: // del diff --git a/emperor.c b/emperor.c index 3b595ec5..9724f4e1 100644 --- a/emperor.c +++ b/emperor.c @@ -1,4 +1,6 @@ #include "uwsgi.h" +#include + extern struct uwsgi_server uwsgi; extern char **environ; @@ -148,7 +150,6 @@ void emperor_add(char *name, time_t born) { uenvs++; } - uwsgi_log("OK\n"); // close the left side of the pipe close(n_ui->pipe[0]); @@ -188,16 +189,28 @@ void emperor_loop() { int waitpid_status; int has_children = 0; int i_am_alone = 0; + int simple_mode = 0; + glob_t g; + int i; + struct dirent *de; memset(&ui_base, 0, sizeof(struct uwsgi_instance)); uwsgi_log("*** starting uWSGI Emperor ***\n"); - if (chdir(uwsgi.emperor_dir)) { - uwsgi_error("chdir()"); + if (!glob(uwsgi.emperor_dir, GLOB_MARK, NULL, &g)) { + if (g.gl_pathc == 1 && g.gl_pathv[0][strlen(g.gl_pathv[0])-1] == '/' ) { + simple_mode = 1; + if (chdir(uwsgi.emperor_dir)) { + uwsgi_error("chdir()"); + exit(1); + } + } + } + else { + uwsgi_error("glob()"); exit(1); } - struct dirent *de; ui = &ui_base; @@ -214,35 +227,72 @@ void emperor_loop() { } } - DIR *dir = opendir("."); - while((de = readdir(dir)) != NULL) { - if (!strcmp(de->d_name+(strlen(de->d_name)-4), ".xml") || - !strcmp(de->d_name+(strlen(de->d_name)-4), ".ini") || - !strcmp(de->d_name+(strlen(de->d_name)-4), ".yml") || - !strcmp(de->d_name+(strlen(de->d_name)-5), ".yaml") - ) { + if (simple_mode) { + DIR *dir = opendir("."); + while((de = readdir(dir)) != NULL) { + if (!strcmp(de->d_name+(strlen(de->d_name)-4), ".xml") || + !strcmp(de->d_name+(strlen(de->d_name)-4), ".ini") || + !strcmp(de->d_name+(strlen(de->d_name)-4), ".yml") || + !strcmp(de->d_name+(strlen(de->d_name)-5), ".yaml") + ) { - if (strlen(de->d_name) >= 0xff) continue; + if (strlen(de->d_name) >= 0xff) continue; - if (stat(de->d_name, &st)) continue; + if (stat(de->d_name, &st)) continue; - if (!S_ISREG(st.st_mode)) continue; + if (!S_ISREG(st.st_mode)) continue; - ui_current = emperor_get(de->d_name); + ui_current = emperor_get(de->d_name); - if (ui_current) { - // check if mtime is changed and the uWSGI instance must be reloaded - if (st.st_mtime > ui_current->last_mod) { - emperor_respawn(ui_current, st.st_mtime); + if (ui_current) { + // check if mtime is changed and the uWSGI instance must be reloaded + if (st.st_mtime > ui_current->last_mod) { + emperor_respawn(ui_current, st.st_mtime); + } + } + else { + emperor_add(de->d_name, st.st_mtime); } } - else { - emperor_add(de->d_name, st.st_mtime); - } + } + closedir(dir); + } + else { + if (glob(uwsgi.emperor_dir, GLOB_MARK, NULL, &g)) { + uwsgi_error("glob()"); + continue; + } + + for(i=0;i<(int)g.gl_pathc;i++) { + if (!strcmp(g.gl_pathv[i]+(strlen(g.gl_pathv[i])-4), ".xml") || + !strcmp(g.gl_pathv[i]+(strlen(g.gl_pathv[i])-4), ".ini") || + !strcmp(g.gl_pathv[i]+(strlen(g.gl_pathv[i])-4), ".yml") || + !strcmp(g.gl_pathv[i]+(strlen(g.gl_pathv[i])-5), ".yaml") + ) { + + + if (strlen(g.gl_pathv[i]) >= 0xff) continue; + + if (stat(g.gl_pathv[i], &st)) continue; + + if (!S_ISREG(st.st_mode)) continue; + + ui_current = emperor_get(g.gl_pathv[i]); + + if (ui_current) { + // check if mtime is changed and the uWSGI instance must be reloaded + if (st.st_mtime > ui_current->last_mod) { + emperor_respawn(ui_current, st.st_mtime); + } + } + else { + emperor_add(g.gl_pathv[i], st.st_mtime); + } + } + } } - closedir(dir); // check for removed instances diff --git a/event.c b/event.c index 57306698..4f99b690 100644 --- a/event.c +++ b/event.c @@ -118,6 +118,37 @@ int event_queue_add_fd_read(int eq, int fd) { return fd; } +int event_queue_del_fd(int eq, int fd) { + + struct epoll_event ee; + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.data.fd = fd; + + if (epoll_ctl(eq, EPOLL_CTL_DEL, fd, &ee)) { + uwsgi_error("epoll_ctl()"); + return -1; + } + + return fd; +} + +int event_queue_add_fd_write(int eq, int fd) { + + struct epoll_event ee; + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.events = EPOLLOUT; + ee.data.fd = fd; + + if (epoll_ctl(eq, EPOLL_CTL_ADD, fd, &ee)) { + uwsgi_error("epoll_ctl()"); + return -1; + } + + return fd; +} + int event_queue_wait(int eq, int timeout, int *interesting_fd) { int ret; diff --git a/plugins/cgi/cgi_plugin.c b/plugins/cgi/cgi_plugin.c new file mode 100644 index 00000000..c6575af1 --- /dev/null +++ b/plugins/cgi/cgi_plugin.c @@ -0,0 +1,139 @@ +#include "../../uwsgi.h" + +#include + +extern struct uwsgi_server uwsgi; + +char *cgi_docroot; + +#define LONG_ARGS_CGI_BASE 17000 + ((9 + 1) * 100) +#define LONG_ARGS_CGI LONG_ARGS_CGI_BASE + 1 + +struct option uwsgi_cgi_options[] = { + + {"cgi", required_argument, 0, LONG_ARGS_CGI}, + {0, 0, 0, 0}, + +}; + +int uwsgi_cgi_init(){ + + + uwsgi_log("initialized CGI engine\n"); + + return 1; + +} + +int uwsgi_cgi_request(struct wsgi_request *wsgi_req) { + + int i; + pid_t cgi_pid; + int waitpid_status; + char *argv[2]; + char full_path[PATH_MAX]; + + /* Standard CGI request */ + if (!wsgi_req->uh.pktsize) { + uwsgi_log("Invalid CGI request. skip.\n"); + return -1; + } + + + if (uwsgi_parse_vars(wsgi_req)) { + uwsgi_log("Invalid CGI request. skip.\n"); + return -1; + } + + // check for file availability (and 'runnability') + + cgi_pid = fork(); + + if (cgi_pid < 0) { + uwsgi_error("fork()"); + return -1; + } + + if (cgi_pid > 0) { + // close input + close(wsgi_req->poll.fd); + wsgi_req->fd_closed = 1; + + // now wait for fd + if (waitpid(cgi_pid, &waitpid_status ,0) > 0) { + uwsgi_log("CGI FINISHED\n"); + } + return 0; + } + + // close all the fd except wsgi_req->poll.fd and 2; + + for(i=0;i< (int)uwsgi.max_fd;i++) { + if (i != wsgi_req->poll.fd && i != 2) { + close(i); + } + } + + // now map wsgi_req->poll.fd to 0 && 1 + if (wsgi_req->poll.fd != 0) { + dup2(wsgi_req->poll.fd, 0); + close(wsgi_req->poll.fd); + } + + dup2(0,1); + + + // fill cgi env + for(i=0;ivar_cnt;i++) { + // no need to free the putenv() memory + if (putenv(uwsgi_concat3n(wsgi_req->hvec[i].iov_base, wsgi_req->hvec[i].iov_len, "=", 1, wsgi_req->hvec[i+1].iov_base, wsgi_req->hvec[i+1].iov_len))) { + uwsgi_error("putenv()"); + } + i++; + } + + char *path_info = uwsgi_concat4n(cgi_docroot, strlen(cgi_docroot), "/", 1,wsgi_req->path_info, wsgi_req->path_info_len, "", 0); + uwsgi_log("requested %s %s\n", path_info, realpath(path_info, full_path)); + + argv[0] = full_path; + argv[1] = NULL; + if (execv(argv[0], argv)) { + uwsgi_error("execv()"); + } + + // never here + exit(1); + + return 0; +} + + +void uwsgi_cgi_after_request(struct wsgi_request *wsgi_req) { + + if (uwsgi.shared->options[UWSGI_OPTION_LOGGING]) + log_request(wsgi_req); +} + +int uwsgi_cgi_manage_options(int i, char *optarg) { + + switch(i) { + case LONG_ARGS_CGI: + cgi_docroot = optarg; + return 1; + } + + return 0; +} + + +struct uwsgi_plugin cgi_plugin = { + + .name = "cgi", + .modifier1 = 9, + .init = uwsgi_cgi_init, + .options = uwsgi_cgi_options, + .manage_opt = uwsgi_cgi_manage_options, + .request = uwsgi_cgi_request, + .after_request = uwsgi_cgi_after_request, + +}; diff --git a/plugins/cgi/uwsgiplugin.py b/plugins/cgi/uwsgiplugin.py new file mode 100644 index 00000000..983040f3 --- /dev/null +++ b/plugins/cgi/uwsgiplugin.py @@ -0,0 +1,5 @@ +NAME='cgi' +CFLAGS = [] +LDFLAGS = [] +LIBS = [] +GCC_LIST = ['cgi_plugin'] diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c new file mode 100644 index 00000000..4b5da9fa --- /dev/null +++ b/plugins/fastrouter/fastrouter.c @@ -0,0 +1,415 @@ +/* + + uWSGI fastrouter + + requires: + + - async + - caching + - pcre (optional) + +*/ + +#include "../../uwsgi.h" + +#define LONG_ARGS_FASTROUTER 150001 + +#define FASTROUTER_STATUS_FREE 0 +#define FASTROUTER_STATUS_CONNECTING 1 +#define FASTROUTER_STATUS_RECV_HDR 2 +#define FASTROUTER_STATUS_RECV_VARS 3 +#define FASTROUTER_STATUS_RESPONSE 4 + +struct uwsgi_fastrouter { + char *socket_name; + int use_cache; +} ufr; + +struct option fastrouter_options[] = { + {"fastrouter", required_argument, 0, LONG_ARGS_FASTROUTER}, + {"fastrouter-use-cache", no_argument, &ufr.use_cache, 1}, + {0, 0, 0, 0}, +}; + +extern struct uwsgi_server uwsgi; + +struct fastrouter_session { + + int fd; + int instance_fd; + int status; + struct uwsgi_header uh; + uint8_t h_pos; + char buffer[0xffff]; + uint16_t pos; + + char *hostname; + uint16_t hostname_len; + + char *instance_address; + uint16_t instance_address_len; + + int pass_fd; +}; + +void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) { + + struct fastrouter_session *fr_session = (struct fastrouter_session *) data; + + //uwsgi_log("%.*s = %.*s\n", keylen, key, vallen, val); + if (!uwsgi_strncmp("SERVER_NAME", 11, key, keylen) && !fr_session->hostname_len) { + fr_session->hostname = val; + fr_session->hostname_len = vallen; + return; + } + + if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen)) { + fr_session->hostname = val; + fr_session->hostname_len = vallen; + return; + } +} + +struct fastrouter_session *alloc_fr_session() { + + return uwsgi_malloc(sizeof(struct fastrouter_session)); +} + +void fastrouter_loop() { + + int fr_queue; + int fr_server; + int nevents; + int interesting_fd; + int new_connection; + ssize_t len; + int i; + + struct msghdr msg; + union { + struct cmsghdr cmsg; + char control [CMSG_SPACE (sizeof (int))]; + } msg_control; + struct cmsghdr *cmsg; + + struct sockaddr_un fr_addr; + socklen_t fr_addr_len = sizeof(struct sockaddr_un); + + struct fastrouter_session *fr_session; + + struct fastrouter_session *fr_table[2048]; + + struct iovec iov[2]; + + int soopt; + socklen_t solen = sizeof(int); + + for(i=0;i<2048;i++) { + fr_table[i] = NULL; + } + + fr_server = bind_to_tcp(ufr.socket_name, uwsgi.listen_queue, ufr.socket_name); + + fr_queue = event_queue_init(); + event_queue_add_fd_read(fr_queue, fr_server); + + for (;;) { + + nevents = event_queue_wait(fr_queue, -1, &interesting_fd); + + if (nevents > 0) { + + //uwsgi_log("interesting_fd: %d\n", interesting_fd); + + if (interesting_fd == fr_server) { + new_connection = accept(fr_server, (struct sockaddr *) &fr_addr, &fr_addr_len); + if (new_connection < 0) { + continue; + } + + fr_table[new_connection] = alloc_fr_session(); + fr_table[new_connection]->fd = new_connection; + fr_table[new_connection]->instance_fd = -1; + fr_table[new_connection]->status = FASTROUTER_STATUS_RECV_HDR; + fr_table[new_connection]->h_pos = 0; + fr_table[new_connection]->pos = 0; + fr_table[new_connection]->instance_address_len = 0; + + event_queue_add_fd_read(fr_queue, new_connection); + + } + else { + fr_session = fr_table[interesting_fd]; + + // something is going wrong... + if (fr_session == NULL) continue; + + switch(fr_session->status) { + + case FASTROUTER_STATUS_RECV_HDR: + len = recv(fr_session->fd, (char *)(&fr_session->uh) + fr_session->h_pos, 4-fr_session->h_pos, 0); + if (len <= 0) { + uwsgi_error("recv()"); + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + free(fr_session); + break; + } + fr_session->h_pos += len; + if (fr_session->h_pos == 4) { +#ifdef UWSGI_DEBUG + uwsgi_log("modifier1: %d pktsize: %d modifier2: %d\n", fr_session->uh.modifier1, fr_session->uh.pktsize, fr_session->uh.modifier2); +#endif + fr_session->status = FASTROUTER_STATUS_RECV_VARS; + } + break; + + + case FASTROUTER_STATUS_RECV_VARS: + len = recv(fr_session->fd, fr_session->buffer + fr_session->pos, fr_session->uh.pktsize - fr_session->pos, 0); + if (len <= 0) { + uwsgi_error("recv()"); + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + free(fr_session); + break; + } + fr_session->pos += len; + if (fr_session->pos == fr_session->uh.pktsize) { + if (uwsgi_hooked_parse(fr_session->buffer, fr_session->uh.pktsize, fr_get_hostname, (void *) fr_session)) { + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + free(fr_session); + break; + } + + if (fr_session->hostname_len == 0) { + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + free(fr_session); + break; + } + + //uwsgi_log("requested domain %.*s\n", fr_session->hostname_len, fr_session->hostname); + if (ufr.use_cache) { + fr_session->instance_address = uwsgi_cache_get(fr_session->hostname, fr_session->hostname_len, &fr_session->instance_address_len); + } + + // no address found + if (!fr_session->instance_address_len) { + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + free(fr_session); + break; + } + + + fr_session->pass_fd = is_unix(fr_session->instance_address, fr_session->instance_address_len); + + fr_session->instance_fd = uwsgi_connectn(fr_session->instance_address, fr_session->instance_address_len, 0, 1); + if (fr_session->instance_fd < 0) { + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + free(fr_session); + break; + } + + + fr_session->status = FASTROUTER_STATUS_CONNECTING; + fr_table[fr_session->instance_fd] = fr_session; + event_queue_add_fd_write(fr_queue, fr_session->instance_fd); + } + break; + + + + case FASTROUTER_STATUS_CONNECTING: + + if (interesting_fd == fr_session->instance_fd) { + + if (getsockopt(fr_session->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&soopt), &solen) < 0) { + uwsgi_error("getsockopt()"); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + + if (soopt) { + uwsgi_log("unable to connect() to uwsgi instance: %s\n", strerror(soopt)); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + + iov[0].iov_base = &fr_session->uh; + iov[0].iov_len = 4; + iov[1].iov_base = fr_session->buffer; + iov[1].iov_len = fr_session->uh.pktsize; + + // fd passing: PERFORMANCE EXTREME BOOST !!! + if (fr_session->pass_fd) { + msg.msg_name = NULL; + msg.msg_namelen = 0; + msg.msg_iov = iov; + msg.msg_iovlen = 2; + msg.msg_flags = 0; + msg.msg_control = &msg_control; + msg.msg_controllen = sizeof (msg_control); + + cmsg = CMSG_FIRSTHDR (&msg); + cmsg->cmsg_len = CMSG_LEN (sizeof (int)); + cmsg->cmsg_level = SOL_SOCKET; + cmsg->cmsg_type = SCM_RIGHTS; + + *((int *) CMSG_DATA (cmsg)) = fr_session->fd; + + if (sendmsg(fr_session->instance_fd, &msg, 0) < 0) { + uwsgi_error("sendmsg()"); + } + + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + + if (writev(fr_session->instance_fd, iov, 2) < 0) { + uwsgi_error("writev()"); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + + event_queue_del_fd(fr_queue, fr_session->instance_fd); + event_queue_add_fd_read(fr_queue, fr_session->instance_fd); + fr_session->status = FASTROUTER_STATUS_RESPONSE; + } + + break; + + case FASTROUTER_STATUS_RESPONSE: + + // data from instance + if (interesting_fd == fr_session->instance_fd) { + len = recv(fr_session->instance_fd, fr_session->buffer, 0xffff, 0); + if (len <= 0) { + if (len < 0) uwsgi_error("recv()"); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + + len = send(fr_session->fd, fr_session->buffer, len, 0); + + if (len <= 0) { + if (len < 0) uwsgi_error("send()"); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + } + // body from client + else if (interesting_fd == fr_session->fd) { + + //uwsgi_log("receiving body...\n"); + len = recv(fr_session->fd, fr_session->buffer, 0xffff, 0); + if (len <= 0) { + if (len < 0) uwsgi_error("recv()"); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + + + len = send(fr_session->instance_fd, fr_session->buffer, len, 0); + + if (len <= 0) { + if (len < 0) uwsgi_error("send()"); + close(fr_session->fd); + close(fr_session->instance_fd); + fr_table[fr_session->fd] = NULL; + fr_table[fr_session->instance_fd] = NULL; + free(fr_session); + break; + } + } + + break; + + + + // fallback to destroy !!! + default: + uwsgi_log("default action\n"); + close(fr_session->fd); + fr_table[fr_session->fd] = NULL; + if (fr_session->instance_fd != -1) { + close(fr_session->instance_fd); + fr_table[fr_session->instance_fd] = NULL; + } + free(fr_session); + break; + + } + } + + } + } +} + +int fastrouter_init() { + + if (ufr.socket_name) { + + if (ufr.use_cache && !uwsgi.cache_max_items) { + uwsgi_log("you need to create a uwsgi cache to use the fastrouter (add --cache )\n"); + exit(1); + } + + if (register_gateway("fastrouter", fastrouter_loop) == NULL) { + uwsgi_log("unable to register the fastrouter gateway\n"); + exit(1); + } + } + + return 0; +} + +int fastrouter_opt(int i, char *optarg) { + + switch(i) { + case LONG_ARGS_FASTROUTER: + ufr.socket_name = optarg; + return 1; + } + return 0; +} + + +struct uwsgi_plugin fastrouter_plugin = { + + .options = fastrouter_options, + .manage_opt = fastrouter_opt, + .init = fastrouter_init, +}; + diff --git a/plugins/fastrouter/uwsgiplugin.py b/plugins/fastrouter/uwsgiplugin.py new file mode 100644 index 00000000..4c5200a3 --- /dev/null +++ b/plugins/fastrouter/uwsgiplugin.py @@ -0,0 +1,7 @@ + +NAME='fastrouter' +CFLAGS = [] +LDFLAGS = [] +LIBS = [] + +GCC_LIST = ['fastrouter'] diff --git a/plugins/python/python_plugin.c b/plugins/python/python_plugin.c index a21c943f..ad8eda49 100644 --- a/plugins/python/python_plugin.c +++ b/plugins/python/python_plugin.c @@ -577,9 +577,7 @@ void init_uwsgi_embedded_module() { init_uwsgi_module_sharedarea(new_uwsgi_module); } - if (uwsgi.cache_max_items > 0) { - init_uwsgi_module_cache(new_uwsgi_module); - } + init_uwsgi_module_cache(new_uwsgi_module); } #endif diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index ae9689b7..1cc1834b 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -1796,6 +1796,12 @@ PyObject *py_uwsgi_send_message(PyObject * self, PyObject * args) { } UWSGI_GET_GIL + + // if it is a fd passing request, return None + if (fd >=0 && cl == -1) { + Py_INCREF(Py_None); + return Py_None; + } // request sent, return the iterator response ui = PyObject_New(uwsgi_Iter, &uwsgi_IterType); if (!ui) { @@ -2348,16 +2354,21 @@ static PyMethodDef uwsgi_sa_methods[] = { PyObject *py_uwsgi_cache_del(PyObject * self, PyObject * args) { char *key; - char *value; + Py_ssize_t keylen = 0; + char *remote = NULL; - if (!PyArg_ParseTuple(args, "s:cache_del", &key, &value)) { + if (!PyArg_ParseTuple(args, "s#|s:cache_del", &key, &keylen, &remote)) { return NULL; } - - if (uwsgi_cache_del(key, strlen(key))) { - Py_INCREF(Py_None); - return Py_None; + if (remote) { + uwsgi_simple_send_string(remote, 111, 2, key, keylen, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); + } + else if (uwsgi.cache_max_items) { + if (uwsgi_cache_del(key, strlen(key))) { + Py_INCREF(Py_None); + return Py_None; + } } Py_INCREF(Py_True); @@ -2371,10 +2382,12 @@ PyObject *py_uwsgi_cache_set(PyObject * self, PyObject * args) { char *key; char *value; Py_ssize_t vallen = 0; + Py_ssize_t keylen = 0; + char *remote = NULL; uint64_t expires = 0; - if (!PyArg_ParseTuple(args, "ss#|i:cache_set", &key, &value, &vallen, &expires)) { + if (!PyArg_ParseTuple(args, "s#s#|is:cache_set", &key, &keylen, &value, &vallen, &expires, &remote)) { return NULL; } @@ -2382,9 +2395,14 @@ PyObject *py_uwsgi_cache_set(PyObject * self, PyObject * args) { return PyErr_Format(PyExc_ValueError, "uWSGI cache items size must be < 64K, requested %d bytes", (int) vallen); } - if (uwsgi_cache_set(key, strlen(key), value, vallen, expires)) { - Py_INCREF(Py_None); - return Py_None; + if (remote) { + uwsgi_simple_send_string2(remote, 111, 1, key, keylen, value, vallen, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); + } + else if (uwsgi.cache_max_items) { + if (uwsgi_cache_set(key, strlen(key), value, vallen, expires)) { + Py_INCREF(Py_None); + return Py_None; + } } Py_INCREF(Py_True); @@ -2415,13 +2433,24 @@ PyObject *py_uwsgi_cache_get(PyObject * self, PyObject * args) { char *key; uint16_t valsize; - char *value; + Py_ssize_t keylen = 0; + char *value = NULL; + char *remote = NULL; + char buffer[0xffff]; - if (!PyArg_ParseTuple(args, "s:cache_get", &key)) { + if (!PyArg_ParseTuple(args, "s#|s:cache_get", &key, &keylen, &remote)) { return NULL; } - value = uwsgi_cache_get(key, strlen(key), &valsize); + if (remote) { + uwsgi_simple_message_string(remote, 111, 0, key, keylen, buffer, &valsize, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); + if (valsize > 0) { + value = buffer; + } + } + else if (uwsgi.cache_max_items) { + value = uwsgi_cache_get(key, strlen(key), &valsize); + } if (value) { return PyString_FromStringAndSize(value, valsize); diff --git a/protocol.c b/protocol.c index 8867fd48..492c5b38 100644 --- a/protocol.c +++ b/protocol.c @@ -127,7 +127,9 @@ int uwsgi_enqueue_message(char *host, int port, uint8_t modifier1, uint8_t modif return uwsgi_poll.fd; } -ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *message, uint16_t size, int pfd, size_t plen, int timeout) { + + +ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *message, uint16_t size, int pfd, ssize_t plen, int timeout) { struct pollfd uwsgi_mpoll; ssize_t cnt; @@ -135,7 +137,13 @@ ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *m char buffer[4096]; ssize_t ret = 0; int pret; - + struct msghdr msg; + struct iovec iov [1]; + union { + struct cmsghdr cmsg; + char control [CMSG_SPACE (sizeof (int))]; + } msg_control; + struct cmsghdr *cmsg; if (!timeout) timeout = uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]; @@ -143,7 +151,33 @@ ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *m uh.pktsize = size; uh.modifier2 = modifier2; - cnt = write(fd, &uh, 4); + if (pfd >= 0 && plen == -1) { + // pass the fd + iov[0].iov_base = &uh; + iov[0].iov_len = 4; + + msg.msg_name = NULL; + msg.msg_namelen = 0; + msg.msg_iov = iov; + msg.msg_iovlen = 1; + msg.msg_flags = 0; + + msg.msg_control = &msg_control; + msg.msg_controllen = sizeof (msg_control); + + cmsg = CMSG_FIRSTHDR (&msg); + cmsg->cmsg_len = CMSG_LEN (sizeof (int)); + cmsg->cmsg_level = SOL_SOCKET; + cmsg->cmsg_type = SCM_RIGHTS; + + *((int *) CMSG_DATA (cmsg)) = pfd; + + uwsgi_log("passing fd\n"); + cnt = sendmsg(fd, &msg, 0); + } + else { + cnt = write(fd, &uh, 4); + } if (cnt != 4) { uwsgi_error("write()"); return -1; @@ -201,6 +235,13 @@ ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *m int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header *uh, char *buffer) { int rlen, i; + struct msghdr msg; + struct iovec iov [1]; + struct cmsghdr *cmsg; + union { + struct cmsghdr cmsg; + char control [CMSG_SPACE (sizeof (int))]; + } msg_control; if (!timeout) timeout = 1; @@ -216,7 +257,20 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header close(upoll->fd); return 0; } - rlen = read(upoll->fd, uh, 4); + + iov [0].iov_base = uh; + iov [0].iov_len = 4; + + msg.msg_name = NULL; + msg.msg_namelen = 0; + msg.msg_iov = iov; + msg.msg_iovlen = 1; + msg.msg_control = &msg_control; + msg.msg_controllen = sizeof (msg_control); + msg.msg_flags = 0; + + //rlen = read(upoll->fd, uh, 4); + rlen = recvmsg(upoll->fd, &msg, 0); if (rlen > 0 && rlen < 4) { i = rlen; while (i < 4) { @@ -263,6 +317,7 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header return 0; } + //uwsgi_log("ready for reading %d bytes\n", wsgi_req.size); i = 0; @@ -291,6 +346,20 @@ int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header return 0; } + cmsg = CMSG_FIRSTHDR (&msg); + while(cmsg != NULL) { + if (cmsg->cmsg_len != CMSG_LEN (sizeof (int)) || + cmsg->cmsg_level != SOL_SOCKET || + cmsg->cmsg_type != SCM_RIGHTS) continue; + + // upgrade connection to the new socket + uwsgi_log("upgrading fd %d to ", upoll->fd); + close(upoll->fd); + upoll->fd = *((int *) CMSG_DATA (cmsg)); + uwsgi_log("%d\n", upoll->fd); + cmsg = CMSG_NXTHDR (&msg, cmsg); + } + return 1; } @@ -298,18 +367,14 @@ int uwsgi_parse_array(char *buffer, uint16_t size, char **argv, uint8_t *argc) { char *ptrbuf, *bufferend; uint16_t strsize = 0; - int i; + uint8_t max = *argc; *argc = 0; ptrbuf = buffer; bufferend = ptrbuf + size; - for(i=0;i> 8) & 0xff); + + strsize2[0] = (uint8_t) (item2_len & 0xff); + strsize2[1] = (uint8_t) ((item2_len >> 8) & 0xff); + + iov[0].iov_base = &uh; + iov[0].iov_len = 4; + + iov[1].iov_base = strsize1; + iov[1].iov_len = 2; + + iov[2].iov_base = item1; + iov[2].iov_len = item1_len; + + iov[3].iov_base = strsize2; + iov[3].iov_len = 2; + + iov[4].iov_base = item2; + iov[4].iov_len = item2_len; + + if (writev(fd, iov, 5) < 0) { + uwsgi_error("writev()"); + } + + close(fd); + + return 0; +} + +int uwsgi_simple_send_string(char *socket_name, uint8_t modifier1, uint8_t modifier2, char *item1, uint16_t item1_len, int timeout) { + + struct uwsgi_header uh; + char strsize1[2]; + + struct iovec iov[3]; + + int fd = uwsgi_connect(socket_name, timeout, 0); + + if (fd < 0) { + return -1; + } + + uh.modifier1 = modifier1; + uh.pktsize = 2+item1_len; + uh.modifier2 = modifier2; + + strsize1[0] = (uint8_t) (item1_len & 0xff); + strsize1[1] = (uint8_t) ((item1_len >> 8) & 0xff); + + iov[0].iov_base = &uh; + iov[0].iov_len = 4; + + iov[1].iov_base = strsize1; + iov[1].iov_len = 2; + + iov[2].iov_base = item1; + iov[2].iov_len = item1_len; + + if (writev(fd, iov, 3) < 0) { + uwsgi_error("writev()"); + } + + close(fd); + + return 0; +} + diff --git a/regexp.c b/regexp.c new file mode 100644 index 00000000..228fbd2f --- /dev/null +++ b/regexp.c @@ -0,0 +1,35 @@ +#ifdef UWSGI_PCRE + +#include "uwsgi.h" + +#include + +/* + +void uwsgi_regexp_match(regexp, what) { + + int ret,i; + + for(i=0;inroutes;i++) { + + ret = pcre_exec(ur->pattern, ur->pattern_extra, wsgi_req->path_info, wsgi_req->path_info_len, 0, 0, wsgi_req->ovector, (ur->args+1)*3 ); + + if (ret >= 0) { + if (ur->action) { + ur->action(uwsgi, wsgi_req, ur); + } + else { + uwsgi_route_action_wsgi(uwsgi, wsgi_req, ur); + } + } + + // TODO check for errors if < 0 && != NO_MATCH + } + + return; +} + +*/ + + +#endif diff --git a/socket.c b/socket.c index 38f248d3..81be899c 100644 --- a/socket.c +++ b/socket.c @@ -218,6 +218,17 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst } #endif +int uwsgi_connectn(char *socket_name, uint16_t len, int timeout, int async) { + + int fd ; + + char *zeroed_socket_name = uwsgi_concat2n(socket_name, len, "", 0); + fd = uwsgi_connect(zeroed_socket_name, timeout, async); + + free(zeroed_socket_name); + return fd; +} + int uwsgi_connect(char *socket_name, int timeout, int async) { int ret; diff --git a/utils.c b/utils.c index b91a5b56..ca71111c 100644 --- a/utils.c +++ b/utils.c @@ -1229,8 +1229,6 @@ char *uwsgi_cheap_string(char *buf, int len) { char *cheap_buf = buf-1; - uwsgi_log("original buf: %.*s\n", len ,buf); - for(i=0;i 1) { if (!getrlimit(RLIMIT_NOFILE, &uwsgi.rl)) { if ((unsigned long) uwsgi.rl.rlim_cur < (unsigned long) uwsgi.async) { @@ -1018,6 +1019,10 @@ int uwsgi_start(void *v_argv) { } } + if (!getrlimit(RLIMIT_NOFILE, &uwsgi.rl)) { + uwsgi.max_fd = uwsgi.rl.rlim_max; + } + uwsgi.wsgi_requests = uwsgi_malloc(sizeof(struct wsgi_request *) * uwsgi.cores); for (i = 0; i < uwsgi.cores; i++) { diff --git a/uwsgi.h b/uwsgi.h index 9079f8cc..dec97d24 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -104,10 +104,6 @@ extern int pivot_root(const char * new_root, const char * put_old); #define UWSGI_PLUGIN_BASE "" #endif -#ifdef UWSGI_ROUTING -#include -#endif - #include #include #include @@ -815,6 +811,8 @@ struct uwsgi_server { pid_t mypid; int mywid; + rlim_t max_fd; + struct timeval start_tv; int abstract_socket; @@ -1124,6 +1122,7 @@ int bind_to_tcp(char *, int, char *); int bind_to_udp(char *, int, int); int timed_connect(struct pollfd *, const struct sockaddr *, int, int, int); int uwsgi_connect(char *, int, int); +int uwsgi_connectn(char *, uint16_t, int, int); int connect_to_tcp(char *, int, int, int); int connect_to_unix(char *, int, int); #ifdef UWSGI_SCTP @@ -1375,7 +1374,7 @@ char *generate_socket_name(char *); #define UMIN(a,b) ((a)>(b)?(b):(a)) -ssize_t uwsgi_send_message(int, uint8_t, uint8_t, char *, uint16_t, int, size_t, int); +ssize_t uwsgi_send_message(int, uint8_t, uint8_t, char *, uint16_t, int, ssize_t, int); char *uwsgi_cluster_best_node(void); @@ -1393,6 +1392,8 @@ inline void *uwsgi_malloc(size_t); int event_queue_init(void); int event_queue_add_fd_read(int, int); +int event_queue_add_fd_write(int, int); +int event_queue_del_fd(int, int); int event_queue_wait(int, int, int *); int event_queue_add_timer(int, int *, int); @@ -1453,3 +1454,9 @@ char *uwsgi_num2str(int); char *magic_sub(char *, int, int *, char *[]); void init_magic_table(char *[]); + +char *uwsgi_simple_message_string(char *, uint8_t, uint8_t, char *, uint16_t, char *, uint16_t *, int); +int uwsgi_simple_send_string2(char *, uint8_t, uint8_t, char *, uint16_t, char *, uint16_t, int); +int uwsgi_simple_send_string(char *, uint8_t, uint8_t, char *, uint16_t, int); + +int is_unix(char *, int); diff --git a/uwsgiconfig.py b/uwsgiconfig.py index 3b1c73d2..de823582 100644 --- a/uwsgiconfig.py +++ b/uwsgiconfig.py @@ -332,6 +332,28 @@ class uConf(object): if self.get('udp'): self.cflags.append("-DUWSGI_UDP") + if self.get('pcre'): + if self.get('pcre') == 'auto': + pcreconf = spcall('pcre-config --libs') + if pcreconf: + self.libs.append(pcreconf) + pcreconf = spcall("pcre-config --cflags") + self.cflags.append(pcreconf) + self.gcc_list.append('regexp') + self.cflags.append("-DUWSGI_PCRE") + + else: + pcreconf = spcall('pcre-config --libs') + if pcreconf is None: + print("*** libpcre headers unavailable. uWSGI build is interrupted. You have to install pcre development package or disable pcre") + sys.exit(1) + else: + self.libs.append(pcreconf) + pcreconf = spcall("pcre-config --cflags") + self.cflags.append(pcreconf) + self.gcc_list.append('regexp') + self.cflags.append("-DUWSGI_PCRE") + if self.get('async'): self.cflags.append("-DUWSGI_ASYNC") self.gcc_list.append('async') @@ -356,26 +378,6 @@ class uConf(object): self.gcc_list.append('ldap') self.libs.append('-lldap') - """ - if ROUTING: - depends_on("ROUTING", ['WEB3', 'XML']) - cflags.append("-DUWSGI_ROUTING") - gcc_list.append('routing') - pcreconf = spcall("pcre-config --cflags") - if pcreconf is None: - print ("*** Unable to locate pcre-config. The uWSGI build has been interrupted. You have to install pcre.") - sys.exit(1) - else: - cflags.append(pcreconf) - - pcreconf = spcall("pcre-config --libs") - if pcreconf is None: - print ("*** Unable to locate pcre-config. The uWSGI build has been interrupted. You have to install pcre.") - sys.exit(1) - else: - libs.append(pcreconf) - """ - if self.get('http'): self.cflags.append("-DUWSGI_HTTP") self.gcc_list.append('http')