From c65adaaac754f4ea13831b3fde4d55991a1d08ba Mon Sep 17 00:00:00 2001 From: Roberto De Ioris Date: Sat, 20 Oct 2012 13:47:07 +0200 Subject: [PATCH] first (destructive) commit for 1.4 --- core/buffer.c | 17 +- core/event.c | 249 ++++++++++++- core/plugins.c | 6 +- plugins/corerouter/corerouter.c | 304 ++++++++++++++-- plugins/corerouter/cr.h | 62 ++-- plugins/corerouter/cr_common.c | 217 +----------- plugins/fastrouter/fastrouter.c | 293 ++++++++++++++- plugins/fastrouter/fr.h | 14 +- plugins/fastrouter/fr_events.c | 567 +++++++++++++----------------- plugins/fastrouter/uwsgiplugin.py | 2 +- uwsgi.h | 7 + uwsgiconfig.py | 11 + 12 files changed, 1121 insertions(+), 628 deletions(-) diff --git a/core/buffer.c b/core/buffer.c index f3ee79f3..010c2277 100644 --- a/core/buffer.c +++ b/core/buffer.c @@ -13,6 +13,19 @@ struct uwsgi_buffer *uwsgi_buffer_new(size_t len) { } +int uwsgi_buffer_fix(struct uwsgi_buffer *ub, size_t len) { + if (ub->len < len) { + char *new_buf = realloc(ub->buf, len); + if (!new_buf) { + uwsgi_error("uwsgi_buffer_fix()"); + return -1; + } + ub->buf = new_buf; + ub->len = len; + } + return 0; +} + int uwsgi_buffer_append(struct uwsgi_buffer *ub, char *buf, size_t len) { size_t remains = ub->len - ub->pos; @@ -21,7 +34,7 @@ int uwsgi_buffer_append(struct uwsgi_buffer *ub, char *buf, size_t len) { size_t chunk_size = UMAX(len, (size_t) uwsgi.page_size); char *new_buf = realloc(ub->buf, ub->len + chunk_size); if (!new_buf) { - uwsgi_error("realloc()"); + uwsgi_error("uwsgi_buffer_append()"); return -1; } ub->buf = new_buf; @@ -54,7 +67,7 @@ int uwsgi_buffer_send(struct uwsgi_buffer *ub, int fd) { return -1; } else { - uwsgi_error("write()"); + uwsgi_error("uwsgi_buffer_send()"); return -1; } } diff --git a/core/event.c b/core/event.c index 2cd56964..92d81bef 100644 --- a/core/event.c +++ b/core/event.c @@ -28,7 +28,7 @@ int event_queue_del_fd(int eq, int fd, int event) { return -1; } - return fd; + return 0; } int event_queue_fd_write_to_read(int eq, int fd) { @@ -38,7 +38,7 @@ int event_queue_fd_write_to_read(int eq, int fd) { return -1; } - return fd; + return 0; } @@ -49,10 +49,53 @@ int event_queue_fd_read_to_write(int eq, int fd) { return -1; } - return fd; + return 0; } +int event_queue_fd_readwrite_to_read(int eq, int fd) { + + if (port_associate(eq, PORT_SOURCE_FD, fd, POLLIN, NULL)) { + uwsgi_error("port_associate"); + return -1; + } + + return 0; + +} + +int event_queue_fd_readwrite_to_write(int eq, int fd) { + + if (port_associate(eq, PORT_SOURCE_FD, fd, POLLOUT, NULL)) { + uwsgi_error("port_associate"); + return -1; + } + + return 0; + +} + +int event_queue_fd_write_to_readwrite(int eq, int fd) { + + if (port_associate(eq, PORT_SOURCE_FD, fd, POLLIN|POLLOUT, NULL)) { + uwsgi_error("port_associate"); + return -1; + } + + return 0; + +} + +int event_queue_fd_read_to_readwrite(int eq, int fd) { + + if (port_associate(eq, PORT_SOURCE_FD, fd, POLLIN|POLLOUT, NULL)) { + uwsgi_error("port_associate"); + return -1; + } + + return 0; + +} int event_queue_interesting_fd_has_error(void *events, int id) { @@ -63,7 +106,21 @@ int event_queue_interesting_fd_has_error(void *events, int id) { return 0; } +int event_queue_interesting_fd_is_read(void *events, int id) { + port_event_t *pe = (port_event_t *) events; + if (pe[id].portev_events == POLLIN) { + return 1; + } + return 0; +} +int event_queue_interesting_fd_is_write(void *events, int id) { + port_event_t *pe = (port_event_t *) events; + if (pe[id].portev_events == POLLOUT) { + return 1; + } + return 0; +} int event_queue_add_fd_read(int eq, int fd) { @@ -72,7 +129,7 @@ int event_queue_add_fd_read(int eq, int fd) { return -1; } - return fd; + return 0; } int event_queue_add_fd_write(int eq, int fd) { @@ -82,7 +139,7 @@ int event_queue_add_fd_write(int eq, int fd) { return -1; } - return fd; + return 0; } void *event_queue_alloc(int nevents) { @@ -206,7 +263,7 @@ int event_queue_add_fd_read(int eq, int fd) { return -1; } - return fd; + return 0; } int event_queue_fd_write_to_read(int eq, int fd) { @@ -222,7 +279,7 @@ int event_queue_fd_write_to_read(int eq, int fd) { return -1; } - return fd; + return 0; } int event_queue_fd_read_to_write(int eq, int fd) { @@ -238,9 +295,75 @@ int event_queue_fd_read_to_write(int eq, int fd) { return -1; } - return fd; + return 0; } +int event_queue_fd_readwrite_to_read(int eq, int fd) { + + struct epoll_event ee; + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.events = EPOLLIN; + ee.data.fd = fd; + + if (epoll_ctl(eq, EPOLL_CTL_MOD, fd, &ee)) { + uwsgi_error("epoll_ctl()"); + return -1; + } + + return 0; +} + +int event_queue_fd_readwrite_to_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_MOD, fd, &ee)) { + uwsgi_error("epoll_ctl()"); + return -1; + } + + return 0; +} + + +int event_queue_fd_read_to_readwrite(int eq, int fd) { + + struct epoll_event ee; + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.events = EPOLLIN|EPOLLOUT; + ee.data.fd = fd; + + if (epoll_ctl(eq, EPOLL_CTL_MOD, fd, &ee)) { + uwsgi_error("epoll_ctl()"); + return -1; + } + + return 0; +} + +int event_queue_fd_write_to_readwrite(int eq, int fd) { + + struct epoll_event ee; + + memset(&ee, 0, sizeof(struct epoll_event)); + ee.events = EPOLLIN|EPOLLOUT; + ee.data.fd = fd; + + if (epoll_ctl(eq, EPOLL_CTL_MOD, fd, &ee)) { + uwsgi_error("epoll_ctl()"); + return -1; + } + + return 0; +} + + int event_queue_del_fd(int eq, int fd, int event) { @@ -255,7 +378,7 @@ int event_queue_del_fd(int eq, int fd, int event) { return -1; } - return fd; + return 0; } int event_queue_add_fd_write(int eq, int fd) { @@ -271,7 +394,7 @@ int event_queue_add_fd_write(int eq, int fd) { return -1; } - return fd; + return 0; } void *event_queue_alloc(int nevents) { @@ -292,6 +415,23 @@ int event_queue_interesting_fd_has_error(void *events, int id) { return 0; } +int event_queue_interesting_fd_is_read(void *events, int id) { + struct epoll_event *ee = (struct epoll_event *) events; + if (ee[id].events == EPOLLIN) { + return 1; + } + return 0; +} + +int event_queue_interesting_fd_is_write(void *events, int id) { + struct epoll_event *ee = (struct epoll_event *) events; + if (ee[id].events == EPOLLOUT) { + return 1; + } + return 0; +} + + int event_queue_wait_multi(int eq, int timeout, void *events, int nevents) { int ret; @@ -366,7 +506,7 @@ int event_queue_fd_write_to_read(int eq, int fd) { return -1; } - return fd; + return 0; } int event_queue_fd_read_to_write(int eq, int fd) { @@ -385,9 +525,75 @@ int event_queue_fd_read_to_write(int eq, int fd) { return -1; } - return fd; + return 0; } +int event_queue_fd_readwrite_to_read(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_WRITE, EV_DELETE, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} + +int event_queue_fd_readwrite_to_write(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_READ, EV_DELETE, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} + +int event_queue_fd_read_to_readwrite(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_READ, EV_DELETE, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + EV_SET(&kev, fd, EVFILT_READ|EVFILT_WRITE, EV_ADD, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} + + +int event_queue_fd_write_to_readwrite(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_WRITE, EV_DELETE, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + EV_SET(&kev, fd, EVFILT_READ|EVFILT_WRITE, EV_ADD, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} + + int event_queue_del_fd(int eq, int fd, int event) { @@ -399,7 +605,7 @@ int event_queue_del_fd(int eq, int fd, int event) { return -1; } - return fd; + return 0; } int event_queue_add_fd_read(int eq, int fd) { @@ -412,7 +618,7 @@ int event_queue_add_fd_read(int eq, int fd) { return -1; } - return fd; + return 0; } int event_queue_add_fd_write(int eq, int fd) { @@ -472,8 +678,23 @@ int event_queue_interesting_fd_has_error(void *events, int id) { return 0; } +int event_queue_interesting_fd_is_read(void *events, int id) { + struct kevent *ev = (struct kevent *) events; + if ( ev[id].filter == EVFILT_READ ) { + return 1; + } + return 0; +} +int event_queue_interesting_fd_is_write(void *events, int id) { + struct kevent *ev = (struct kevent *) events; + if ( ev[id].filter == EVFILT_WRITE ) { + return 1; + } + return 0; +} + int event_queue_wait(int eq, int timeout, int *interesting_fd) { int ret; diff --git a/core/plugins.c b/core/plugins.c index b40dac79..15562782 100644 --- a/core/plugins.c +++ b/core/plugins.c @@ -7,7 +7,8 @@ static void uwsgi_plugin_parse_section(char *filename) { size_t s_len = 0; char *buf = uwsgi_elf_section(filename, "uwsgi", &s_len); if (buf) { - char *p = strtok(buf, "\n"); + char *ctx = NULL; + char *p = strtok_r(buf, "\n", &ctx); while(p) { char *equal = strchr(p, '='); if (equal) { @@ -16,7 +17,7 @@ static void uwsgi_plugin_parse_section(char *filename) { uwsgi_load_plugin(-1, equal+1, NULL); } } - p = strtok(NULL, "\n"); + p = strtok_r(NULL, "\n", &ctx); } free(buf); } @@ -78,7 +79,6 @@ void *uwsgi_load_plugin(int modifier, char *plugin, char *has_option) { char *plugin_name = plugin; char *plugin_symbol_name_start = plugin; - struct uwsgi_plugin *up; char linkpath_buf[1024], linkpath[1024]; int linkpath_size; diff --git a/plugins/corerouter/corerouter.c b/plugins/corerouter/corerouter.c index c9ccbc84..78a04446 100644 --- a/plugins/corerouter/corerouter.c +++ b/plugins/corerouter/corerouter.c @@ -199,6 +199,7 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se } else if (cr_session->timed_out) { if (cr_session->instance_address_len > 0) { +/* if (cr_session->status == COREROUTER_STATUS_CONNECTING) { if (!ucr->quiet) uwsgi_log("unable to connect() to uwsgi instance \"%.*s\": timeout\n", (int) cr_session->instance_address_len, cr_session->instance_address); @@ -206,6 +207,7 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se else if (cr_session->status == COREROUTER_STATUS_RESPONSE) { uwsgi_log("timeout waiting for instance \"%.*s\"\n", (int) cr_session->instance_address_len, cr_session->instance_address); } +*/ } } @@ -273,7 +275,7 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se ucr->cr_table[cr_session->instance_fd] = cr_session; - cr_session->status = COREROUTER_STATUS_CONNECTING; + //cr_session->status = COREROUTER_STATUS_CONNECTING; ucr->cr_table[cr_session->instance_fd] = cr_session; event_queue_add_fd_write(ucr->queue, cr_session->instance_fd); return; @@ -297,24 +299,17 @@ end: free(cr_session->buf_file_name); } - if (cr_session->write_queue) - free(cr_session->write_queue); - - if (cr_session->instance_write_queue) - free(cr_session->instance_write_queue); - +/* // could be used to free additional resources if (cr_session->close) cr_session->close(ucr, cr_session); - - if (cr_session->keepalive) { - cr_session->keepalive = 0; - return; - } +*/ close(cr_session->fd); ucr->cr_table[cr_session->fd] = NULL; + uwsgi_buffer_destroy(cr_session->buffer); + cr_del_timeout(ucr, cr_session); free(cr_session); } @@ -340,7 +335,10 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) { cr_session->timed_out = 1; if (cr_session->retry) { cr_session->retry = 0; +/* + TODO allows retry ucr->switch_events(ucr, cr_session, -1); +*/ if (cr_session->retry) { cr_del_timeout(ucr, cr_session); cr_session->timeout = cr_add_fake_timeout(ucr, cr_session); @@ -359,23 +357,235 @@ static void corerouter_expire_timeouts(struct uwsgi_corerouter *ucr) { } } + +int uwsgi_cr_hook_read(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) { + + struct uwsgi_corerouter *ucr = cs->corerouter; + + // first check the case of event removal + if (hook == NULL) { + // nothing changed + if (!cs->event_hook_read) goto unchanged; + // if there is a write event defined, le'ts modify it + if (cs->event_hook_write) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_readwrite_to_write() for %d\n", cs->fd); +#endif + if (event_queue_fd_readwrite_to_write(ucr->queue, cs->fd)) return -1; + } + // simply remove the read event + else { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_del_fd() for %d\n", cs->fd); +#endif + if (event_queue_del_fd(ucr->queue, cs->fd, event_queue_read())) return -1; + } + } + else { + // set the hook + // if write is not defined, simply add a single monitor + if (cs->event_hook_write == NULL) { + if (!cs->event_hook_read) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_add_fd_read() for %d\n", cs->fd); +#endif + if (event_queue_add_fd_read(ucr->queue, cs->fd)) return -1; + } + } + else { + if (!cs->event_hook_read) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_write_to_readwrite() for %d\n", cs->fd); +#endif + if (event_queue_fd_write_to_readwrite(ucr->queue, cs->fd)) return -1; + } + } + } + +unchanged: +#ifdef UWSGI_DEBUG + uwsgi_log("event_hook_read set to %p for %d\n", hook, cs->fd); +#endif + cs->event_hook_read = hook; + return 0; +} + +int uwsgi_cr_hook_write(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) { + + struct uwsgi_corerouter *ucr = cs->corerouter; + + // first check the case of event removal + if (hook == NULL) { + // nothing changed + if (!cs->event_hook_write) goto unchanged; + // if there is a read event defined, le'ts modify it + if (cs->event_hook_read) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_readwrite_to_read() for %d\n", cs->fd); +#endif + if (event_queue_fd_readwrite_to_read(ucr->queue, cs->fd)) return -1; + } + // simply remove the write event + else { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_del_fd() for %d\n", cs->fd); +#endif + if (event_queue_del_fd(ucr->queue, cs->fd, event_queue_write())) return -1; + } + } + else { + // set the hook + // if read is not defined, simply add a single monitor + if (cs->event_hook_read == NULL) { + if (!cs->event_hook_write) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_add_fd_write() for %d\n", cs->fd); +#endif + if (event_queue_add_fd_write(ucr->queue, cs->fd)) return -1; + } + } + else { + if (!cs->event_hook_write) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_read_to_readwrite() for %d\n", cs->fd); +#endif + if (event_queue_fd_read_to_readwrite(ucr->queue, cs->fd)) return -1; + } + } + } + +unchanged: +#ifdef UWSGI_DEBUG + uwsgi_log("event_hook_write set to %p for %d\n", hook, cs->fd); +#endif + cs->event_hook_write = hook; + return 0; +} + +int uwsgi_cr_hook_instance_read(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) { + + struct uwsgi_corerouter *ucr = cs->corerouter; + + // first check the case of event removal + if (hook == NULL) { + // nothing changed + if (!cs->event_hook_instance_read) goto unchanged; + // if there is a write event defined, le'ts modify it + if (cs->event_hook_instance_write) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_readwrite_to_write() for %d\n", cs->instance_fd); +#endif + if (event_queue_fd_readwrite_to_write(ucr->queue, cs->instance_fd)) return -1; + } + // simply remove the read event + else { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_del_fd() for %d\n", cs->instance_fd); +#endif + if (event_queue_del_fd(ucr->queue, cs->instance_fd, event_queue_read())) return -1; + } + } + else { + // set the hook + // if write is not defined, simply add a single monitor + if (cs->event_hook_instance_write == NULL) { + if (!cs->event_hook_instance_read) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_add_fd_read() for %d\n", cs->instance_fd); +#endif + if (event_queue_add_fd_read(ucr->queue, cs->instance_fd)) return -1; + } + } + else { + if (!cs->event_hook_instance_read) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_write_to_readwrite() for %d\n", cs->instance_fd); +#endif + if (event_queue_fd_write_to_readwrite(ucr->queue, cs->instance_fd)) return -1; + } + } + } + +unchanged: +#ifdef UWSGI_DEBUG + uwsgi_log("event_hook_instance_read set to %p for %d\n", hook, cs->instance_fd); +#endif + cs->event_hook_instance_read = hook; + return 0; +} + +int uwsgi_cr_hook_instance_write(struct corerouter_session *cs, ssize_t (*hook)(struct corerouter_session *)) { + + struct uwsgi_corerouter *ucr = cs->corerouter; + + // first check the case of event removal + if (hook == NULL) { + // nothing changed + if (!cs->event_hook_instance_write) goto unchanged; + // if there is a read event defined, le'ts modify it + if (cs->event_hook_instance_read) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_readwrite_to_read() for %d\n", cs->instance_fd); +#endif + if (event_queue_fd_readwrite_to_read(ucr->queue, cs->instance_fd)) return -1; + } + // simply remove the write event + else { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_del_fd() for %d\n", cs->instance_fd); +#endif + if (event_queue_del_fd(ucr->queue, cs->instance_fd, event_queue_write())) return -1; + } + } + else { + // set the hook + // if read is not defined, simply add a single monitor + if (cs->event_hook_instance_read == NULL) { + if (!cs->event_hook_instance_write) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_add_fd_write() for %d\n", cs->instance_fd); +#endif + if (event_queue_add_fd_write(ucr->queue, cs->instance_fd)) return -1; + } + } + else { + if (!cs->event_hook_instance_write) { +#ifdef UWSGI_DEBUG + uwsgi_log("event_queue_fd_read_to_readwrite() for %d\n", cs->instance_fd); +#endif + if (event_queue_fd_read_to_readwrite(ucr->queue, cs->instance_fd)) return -1; + } + } + } + +unchanged: +#ifdef UWSGI_DEBUG + uwsgi_log("event_hook_instance_write set to %p for %d\n", hook, cs->instance_fd); +#endif + cs->event_hook_instance_write = hook; + return 0; +} + + + struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, int new_connection, struct sockaddr *cr_addr, socklen_t cr_addr_len) { ucr->cr_table[new_connection] = uwsgi_calloc(ucr->session_size); ucr->cr_table[new_connection]->fd = new_connection; ucr->cr_table[new_connection]->instance_fd = -1; - ucr->cr_table[new_connection]->status = COREROUTER_STATUS_RECV_HDR; - ucr->cr_table[new_connection]->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]); + // map courerouter and socket + ucr->cr_table[new_connection]->corerouter = ucr; ucr->cr_table[new_connection]->ugs = ugs; - ucr->cr_table[new_connection]->recv = uwsgi_cr_simple_recv; - ucr->cr_table[new_connection]->send = uwsgi_cr_simple_send; - ucr->cr_table[new_connection]->instance_recv = uwsgi_cr_simple_instance_recv; - ucr->cr_table[new_connection]->instance_send = uwsgi_cr_simple_instance_send; + // set initial timeout + ucr->cr_table[new_connection]->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]); + // create dynamic buffer + ucr->cr_table[new_connection]->buffer = uwsgi_buffer_new(uwsgi.page_size); + + // here we prepare the real session and set the hooks ucr->alloc_session(ucr, ugs, ucr->cr_table[new_connection], cr_addr, cr_addr_len); - event_queue_add_fd_read(ucr->queue, new_connection); return ucr->cr_table[new_connection]; } @@ -504,6 +714,7 @@ void uwsgi_corerouter_loop(int id, void *data) { for (;;) { + // set timeouts and harakiri min_timeout = uwsgi_min_rb_timer(ucr->timeouts); if (min_timeout == NULL) { delta = -1; @@ -520,6 +731,7 @@ void uwsgi_corerouter_loop(int id, void *data) { ushared->gateways_harakiri[id] = 0; } + // wait for events nevents = event_queue_wait_multi(ucr->queue, delta, events, ucr->nevents); if (uwsgi.master_process && ucr->harakiri > 0) { @@ -532,8 +744,12 @@ void uwsgi_corerouter_loop(int id, void *data) { for (i = 0; i < nevents; i++) { + // get the interesting fd interesting_fd = event_queue_interesting_fd(events, i); + // something bad happened + if (interesting_fd < 0) continue; + // check if the interesting_fd matches a gateway socket struct uwsgi_gateway_socket *ugs = uwsgi.gateway_sockets; int taken = 0; while (ugs) { @@ -571,9 +787,11 @@ void uwsgi_corerouter_loop(int id, void *data) { continue; } + // manage internal subscription if (interesting_fd == ushared->gateways[id].internal_subscription_pipe[1]) { uwsgi_corerouter_manage_internal_subscription(ucr, interesting_fd); } + // manage a stats request else if (interesting_fd == ucr->cr_stats_server) { corerouter_send_stats(ucr); } @@ -584,17 +802,51 @@ void uwsgi_corerouter_loop(int id, void *data) { if (cr_session == NULL) continue; + // on error, destroy the session if (event_queue_interesting_fd_has_error(events, i)) { - corerouter_close_session(ucr, cr_session); - continue; + corerouter_close_session(ucr, cr_session); + continue; } - cr_session->timeout = corerouter_reset_timeout(ucr, cr_session); - - // implementation specific cycle; - ucr->switch_events(ucr, cr_session, interesting_fd); - + // set timeout + cr_session->timeout = corerouter_reset_timeout(ucr, cr_session); + // call event hook + ssize_t (*hook)(struct corerouter_session *) = NULL; + if (interesting_fd == cr_session->fd) { + if (event_queue_interesting_fd_is_read(events, i)) { + hook = cr_session->event_hook_read; + } + else if (event_queue_interesting_fd_is_write(events, i)) { + hook = cr_session->event_hook_write; + } + } + else if (interesting_fd == cr_session->instance_fd) { + if (event_queue_interesting_fd_is_read(events, i)) { + hook = cr_session->event_hook_instance_read; + } + else if (event_queue_interesting_fd_is_write(events, i)) { + hook = cr_session->event_hook_instance_write; + } + } + if (!hook) { + uwsgi_log("[uwsgi-corerouter] BUG, unexpected event received !!!\n"); + corerouter_close_session(ucr, cr_session); + continue; + } + + ssize_t ret = hook(cr_session); + // connection closed + if (ret == 0) { + corerouter_close_session(ucr, cr_session); + continue; + } + else if (ret < 0) { + if (errno == EINPROGRESS) continue; + corerouter_close_session(ucr, cr_session); + continue; + } + } } } diff --git a/plugins/corerouter/cr.h b/plugins/corerouter/cr.h index 64204030..8b5b3379 100644 --- a/plugins/corerouter/cr.h +++ b/plugins/corerouter/cr.h @@ -9,6 +9,12 @@ #define cr_del_check_timeout(x) rb_erase(&x->rbt, timeouts); #define cr_del_timeout(u, x) rb_erase(&x->timeout->rbt, u->timeouts); free(x->timeout); +#define cr_try_again if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) {\ + errno = EINPROGRESS;\ + return -1;\ + } + + struct corerouter_session; struct uwsgi_corerouter { @@ -19,7 +25,6 @@ struct uwsgi_corerouter { void (*alloc_session)(struct uwsgi_corerouter *, struct uwsgi_gateway_socket *, struct corerouter_session *, struct sockaddr *, socklen_t); int (*mapper)(struct uwsgi_corerouter *, struct corerouter_session *); - void (*switch_events)(struct uwsgi_corerouter *, struct corerouter_session *, int); int has_sockets; int has_backends; @@ -86,15 +91,13 @@ struct corerouter_session { int fd; int instance_fd; - int instance_stopped; - int status; - - uint8_t h_pos; - - uint16_t pos; + // corerouter related to this session + struct uwsgi_corerouter *corerouter; + // gateway socket related to this session struct uwsgi_gateway_socket *ugs; + // parsed hostname char *hostname; uint16_t hostname_len; @@ -110,13 +113,10 @@ struct corerouter_session { int soopt; int timed_out; - // used for tracking required event - int fd_state; - int instance_fd_state; - struct uwsgi_rb_timer *timeout; int instance_failed; + // check content_length size_t post_cl; size_t post_remains; @@ -125,30 +125,27 @@ struct corerouter_session { char *buf_file_name; FILE *buf_file; - uint8_t modifier1; - uint8_t modifier2; - char *tmp_socket_name; + // store the client address struct sockaddr_un addr; socklen_t addr_len; - int keepalive; + // async hooks: + // the session is watiting for this fd + ssize_t (*event_hook_read)(struct corerouter_session *); + ssize_t (*event_hook_write)(struct corerouter_session *); + ssize_t (*event_hook_instance_read)(struct corerouter_session *); + ssize_t (*event_hook_instance_write)(struct corerouter_session *); - char *write_queue; - size_t write_queue_len; - int write_queue_close; + struct uwsgi_buffer *buffer; + size_t buffer_len; + off_t buffer_pos; - char *instance_write_queue; - size_t instance_write_queue_len; - - void (*close)(struct uwsgi_corerouter *, struct corerouter_session *); - - ssize_t (*recv)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); - ssize_t (*send)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); - - ssize_t (*instance_recv)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); - ssize_t (*instance_send)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); + struct uwsgi_header uh; + + uint8_t modifier1; + uint8_t modifier2; }; void uwsgi_opt_corerouter(char *, char *, void *); @@ -184,8 +181,7 @@ int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *, struct corerouter_s int uwsgi_courerouter_has_has_backends(struct uwsgi_corerouter *); -ssize_t uwsgi_cr_simple_recv(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); -ssize_t uwsgi_cr_simple_send(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); - -ssize_t uwsgi_cr_simple_instance_recv(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); -ssize_t uwsgi_cr_simple_instance_send(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); +int uwsgi_cr_hook_read(struct corerouter_session *, ssize_t (*)(struct corerouter_session *)); +int uwsgi_cr_hook_write(struct corerouter_session *, ssize_t (*)(struct corerouter_session *)); +int uwsgi_cr_hook_instance_read(struct corerouter_session *, ssize_t (*)(struct corerouter_session *)); +int uwsgi_cr_hook_instance_write(struct corerouter_session *, ssize_t (*)(struct corerouter_session *)); diff --git a/plugins/corerouter/cr_common.c b/plugins/corerouter/cr_common.c index 7693b97f..dd58f7e2 100644 --- a/plugins/corerouter/cr_common.c +++ b/plugins/corerouter/cr_common.c @@ -10,197 +10,6 @@ extern struct uwsgi_server uwsgi; #include "cr.h" -ssize_t uwsgi_cr_simple_recv(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { - ssize_t ret = recv(cs->fd, buf, len, 0); - if (ret < 0) { - if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { - errno = EINPROGRESS; - return -1; - } - uwsgi_error("recv()"); - } - return ret; -} - -ssize_t uwsgi_cr_simple_instance_recv(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { - ssize_t ret = recv(cs->instance_fd, buf, len, 0); - if (ret < 0) { - if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { - errno = EINPROGRESS; - return -1; - } - uwsgi_error("recv()"); - } - return ret; -} - - -ssize_t uwsgi_cr_simple_send(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { - ssize_t ret = -1; - char *tmp_buf; - off_t pos = 0; - - size_t partial_len = len; - off_t partial_pos = 0; - - if (cs->write_queue_len > 0) { - ret = send(cs->fd, cs->write_queue, cs->write_queue_len, 0); - if (ret > 0) { - cs->write_queue_len-=ret; - pos = ret; - if (cs->write_queue_len == 0) { - free(cs->write_queue); - cs->write_queue = NULL; - if (cs->fd_state) { - event_queue_fd_write_to_read(uc->queue, cs->fd); - cs->fd_state = 0; - } - goto next; - } - goto blocking; - } - else if (ret == 0) { - return 0; - } - else { - if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { - goto blocking; - } - uwsgi_error("send()"); - return -1; - } - } - -next: - if (len == 0) goto end; - - ret = send(cs->fd, buf, len, 0); - if (ret > 0) { - if ((size_t)ret == len) return len; - partial_len-=ret; - partial_pos = ret; - goto blocking; - } - - if (ret == 0) { - return 0; - } - - if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { - goto blocking; - } - uwsgi_error("send()"); - return -1; - -end: - if (cs->write_queue_close) { - return 0; - } - return ret; - -blocking: - // wait for write - if (!cs->fd_state) { - event_queue_fd_read_to_write(uc->queue, cs->fd); - cs->fd_state = 1; - } - // add new datas to the buffer - tmp_buf = malloc(cs->write_queue_len+partial_len); - if (!tmp_buf) { - uwsgi_error("malloc()"); - return -1; - } - if (cs->write_queue_len>0) { - memcpy(tmp_buf, cs->write_queue+pos, cs->write_queue_len); - free(cs->write_queue); - } - memcpy(tmp_buf+cs->write_queue_len, buf+partial_pos, partial_len); - cs->write_queue = tmp_buf; - cs->write_queue_len+=partial_len; - errno = EINPROGRESS; - return -1; -} - -ssize_t uwsgi_cr_simple_instance_send(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { - ssize_t ret; - char *tmp_buf; - off_t pos = 0; - - size_t partial_len = len; - off_t partial_pos = 0; - - if (cs->instance_write_queue_len > 0) { - ret = send(cs->instance_fd, cs->instance_write_queue, cs->instance_write_queue_len, 0); - if (ret > 0) { - cs->instance_write_queue_len-=ret; - pos=ret; - if (cs->instance_write_queue_len == 0) { - free(cs->instance_write_queue); - cs->instance_write_queue = NULL; - if (cs->instance_fd_state) { - event_queue_fd_write_to_read(uc->queue, cs->instance_fd); - cs->instance_fd_state = 0; - } - goto next; - } - goto blocking; - } - else if (ret == 0) { - return 0; - } - else { - if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { - goto blocking; - } - uwsgi_error("send()"); - return -1; - } - } - -next: - ret = send(cs->instance_fd, buf, len, 0); - if (ret > 0) { - if ((size_t)ret == len) return len; - partial_len-=ret; - partial_pos = ret; - goto blocking; - } - - if (ret == 0) { - return 0; - } - - if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { - goto blocking; - } - uwsgi_error("send()"); - return -1; - -blocking: - // wait for write - if (!cs->instance_fd_state) { - event_queue_fd_read_to_write(uc->queue, cs->instance_fd); - cs->instance_fd_state = 1; - } - // add new datas to the buffer - tmp_buf = malloc(cs->instance_write_queue_len+partial_len); - if (!tmp_buf) { - uwsgi_error("malloc()"); - return -1; - } - if (cs->instance_write_queue_len>0) { - memcpy(tmp_buf, cs->instance_write_queue+pos, cs->instance_write_queue_len); - free(cs->instance_write_queue); - } - memcpy(tmp_buf+cs->instance_write_queue_len, buf+partial_pos, partial_len); - cs->instance_write_queue = tmp_buf; - cs->instance_write_queue_len+=partial_len; - errno = EINPROGRESS; - return -1; -} - - - void uwsgi_corerouter_setup_sockets(struct uwsgi_corerouter *ucr) { struct uwsgi_gateway_socket *ugs = uwsgi.gateway_sockets; @@ -208,24 +17,24 @@ void uwsgi_corerouter_setup_sockets(struct uwsgi_corerouter *ucr) { if (!strcmp(ucr->name, ugs->owner)) { if (!ugs->subscription) { if (ugs->name[0] == '=') { - int shared_socket = atoi(ugs->name+1); - if (shared_socket >= 0) { - ugs->fd = uwsgi_get_shared_socket_fd_by_num(shared_socket); + int shared_socket = atoi(ugs->name + 1); + if (shared_socket >= 0) { + ugs->fd = uwsgi_get_shared_socket_fd_by_num(shared_socket); ugs->shared = 1; - if (ugs->fd == -1) { - uwsgi_log("unable to use shared socket %d\n", shared_socket); + if (ugs->fd == -1) { + uwsgi_log("unable to use shared socket %d\n", shared_socket); exit(1); - } + } ugs->name = uwsgi_getsockname(ugs->fd); - } + } } - else if (!uwsgi_startswith("fd://", ugs->name, 5 )) { - int fd_socket = atoi(ugs->name+5); + else if (!uwsgi_startswith("fd://", ugs->name, 5)) { + int fd_socket = atoi(ugs->name + 5); if (fd_socket >= 0) { ugs->fd = fd_socket; ugs->name = uwsgi_getsockname(ugs->fd); if (!ugs->name) { - uwsgi_log("unable to use file descriptor %d as socket\n", fd_socket); + uwsgi_log("unable to use file descriptor %d as socket\n", fd_socket); exit(1); } } @@ -327,8 +136,10 @@ void uwsgi_corerouter_manage_subscription(struct uwsgi_corerouter *ucr, int id, if (node && node->len) { #ifdef UWSGI_SSL if (uwsgi.subscriptions_sign_check_dir) { - if (usr.sign_len == 0 || usr.base_len == 0) return; - if (usr.unix_check <= node->unix_check) return ; + if (usr.sign_len == 0 || usr.base_len == 0) + return; + if (usr.unix_check <= node->unix_check) + return; if (!uwsgi_subscription_sign_check(node->slot, &usr)) { return; } diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index cda3c7b6..1efa0f13 100644 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -2,22 +2,24 @@ uWSGI fastrouter - requires: - - - async - - caching - - pcre (optional) - */ #include "../../uwsgi.h" +#include "../corerouter/cr.h" + +struct uwsgi_fastrouter { + struct uwsgi_corerouter cr; +} ufr; extern struct uwsgi_server uwsgi; -#include "fr.h" - -struct uwsgi_fastrouter ufr; - +struct fastrouter_session { + struct corerouter_session crs; + struct uwsgi_buffer *post_buf; + size_t post_buf_max; + size_t post_buf_len; + off_t post_buf_pos; +}; struct uwsgi_option fastrouter_options[] = { {"fastrouter", required_argument, 0, "run the fastrouter on the specified port", uwsgi_opt_corerouter, &ufr, 0}, @@ -54,13 +56,282 @@ struct uwsgi_option fastrouter_options[] = { {0, 0, 0, 0, 0, 0, 0}, }; +ssize_t fr_recv_uwsgi_header(struct corerouter_session *); +ssize_t fr_instance_read_response(struct corerouter_session *); +ssize_t fr_read_body(struct corerouter_session *); + +void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) { + + // here i use directly corerouter_session + struct corerouter_session *cs = (struct corerouter_session *) data; + + //uwsgi_log("%.*s = %.*s\n", keylen, key, vallen, val); + if (!uwsgi_strncmp("SERVER_NAME", 11, key, keylen) && !cs->hostname_len) { + cs->hostname = val; + cs->hostname_len = vallen; + return; + } + + if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen) && !cs->has_key) { + cs->hostname = val; + cs->hostname_len = vallen; + return; + } + + if (!uwsgi_strncmp("UWSGI_FASTROUTER_KEY", 20, key, keylen)) { + cs->has_key = 1; + cs->hostname = val; + cs->hostname_len = vallen; + return; + } + + if (!uwsgi_strncmp("CONTENT_LENGTH", 14, key, keylen)) { + cs->post_cl = uwsgi_str_num(val, vallen); + return; + } +} + +ssize_t fr_write_body(struct corerouter_session *cs) { + struct fastrouter_session *fs = (struct fastrouter_session *) cs; + ssize_t len = write(cs->instance_fd, fs->post_buf->buf + fs->post_buf_pos, fs->post_buf_len - fs->post_buf_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_write_body()"); + return -1; + } + + fs->post_buf_pos += len; + + // the body chunk has been sent, start again reading from client and instance + if (fs->post_buf_pos == fs->post_buf_len) { + uwsgi_cr_hook_instance_write(cs, NULL); + uwsgi_cr_hook_instance_read(cs, fr_instance_read_response); + uwsgi_cr_hook_read(cs, fr_read_body); + } + + return len; +} + + +ssize_t fr_read_body(struct corerouter_session *cs) { + struct fastrouter_session *fs = (struct fastrouter_session *) cs; + ssize_t len = read(cs->fd, fs->post_buf->buf, fs->post_buf_max); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_read_body()"); + return -1; + } + + // connection closed + if (len == 0) return 0; + + fs->post_buf_len = len; + fs->post_buf_pos = 0; + + // ok we have a body, stop reading from the client and the instance and start writing to the instance + uwsgi_cr_hook_read(cs, NULL); + uwsgi_cr_hook_instance_read(cs, NULL); + uwsgi_cr_hook_instance_write(cs, fr_write_body); + + return len; +} + +ssize_t fr_write_response(struct corerouter_session *cs) { + ssize_t len = write(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_write_response()"); + return -1; + } + + cs->buffer_pos += len; + + // ok this response chunk is sent, let's wait for another one + if (cs->buffer_pos == cs->buffer_len) { + uwsgi_cr_hook_write(cs, NULL); + uwsgi_cr_hook_instance_read(cs, fr_instance_read_response); + } + + return len; +} + +ssize_t fr_instance_read_response(struct corerouter_session *cs) { + ssize_t len = read(cs->instance_fd, cs->buffer->buf, cs->buffer->len); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_instance_read_response()"); + return -1; + } + + // end of the response + if (len == 0) { + return 0; + } + + cs->buffer_pos = 0; + cs->buffer_len = len; + // ok stop reading from the instance, and start writing to the client + uwsgi_cr_hook_instance_read(cs, NULL); + uwsgi_cr_hook_write(cs, fr_write_response); + return len; +} + +ssize_t fr_instance_send_request(struct corerouter_session *cs) { + ssize_t len = write(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->uh.pktsize - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_instance_send_request()"); + return -1; + } + + cs->buffer_pos += len; + + // ok the request is sent, we can start sending client body (if any) and we can start waiting + // for response + if (cs->buffer_pos == cs->uh.pktsize) { + cs->buffer_pos = 0; + // stop writing to the instance + uwsgi_cr_hook_instance_write(cs, NULL); + // start reading from the instance + uwsgi_cr_hook_instance_read(cs, fr_instance_read_response); + // re-start reading from the client (for body or connection close) + struct fastrouter_session *fs = (struct fastrouter_session *) cs; + // allocate a buffer for client body (could be delimited or dynamic) + fs->post_buf_max = UMAX16; + if (cs->post_cl > 0) { + fs->post_buf_max = UMIN(UMAX16, cs->post_cl); + } + fs->post_buf = uwsgi_buffer_new(fs->post_buf_max); + if (!fs->post_buf) return -1; + uwsgi_cr_hook_read(cs, fr_read_body); + } + + return len; +} + +ssize_t fr_instance_send_request_header(struct corerouter_session *cs) { + ssize_t len = write(cs->instance_fd, &cs->uh + cs->buffer_pos, 4 - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_instance_send_request_header()"); + return -1; + } + + cs->buffer_pos += len; + + // ok the request is sent, we can start sending client body (if any) and we can start waiting + // for response + if (cs->buffer_pos == 4) { + cs->buffer_pos = 0; + uwsgi_cr_hook_instance_write(cs, fr_instance_send_request); + } + + return len; +} + +ssize_t fr_instance_connected(struct corerouter_session *cs) { + + socklen_t solen = sizeof(int); + + // first check for errors + if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) { + uwsgi_error("fr_instance_connected()/getsockopt()"); + cs->instance_failed = 1; + return -1; + } + + if (cs->soopt) { + cs->instance_failed = 1; + return -1; + } + + cs->buffer_pos = 0; + + // ok instance is connected, wait for write again + uwsgi_cr_hook_instance_write(cs, fr_instance_send_request_header); + // return a value > 0 + return 1; +} + +ssize_t fr_recv_uwsgi_vars(struct corerouter_session *cs) { + // increase buffer if needed + if (uwsgi_buffer_fix(cs->buffer, cs->uh.pktsize)) return -1; + ssize_t len = read(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->uh.pktsize - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_recv_uwsgi_vars()"); + return -1; + } + + cs->buffer_pos += len; + + // headers received, ready to choose the instance + if (cs->buffer_pos == cs->uh.pktsize) { + struct uwsgi_corerouter *ucr = cs->corerouter; + // find the hostname + if (uwsgi_hooked_parse(cs->buffer->buf, cs->uh.pktsize, fr_get_hostname, (void *) cs)) { + return -1; + } + // check the hostname; + if (cs->hostname_len == 0) return -1; + // find an instance using the key + if (cs->corerouter->mapper(cs->corerouter, cs)) return -1; + // check instance + if (cs->instance_address_len == 0) { + // if fallback nodes are configured, trigger them + if (ucr->fallback) { + cs->instance_failed = 1; + } + return -1; + } + + // stop receiving from the client + uwsgi_cr_hook_read(cs, NULL); + + // start async connect + cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1); + if (cs->instance_fd < 0) { + cs->instance_failed = 1; + cs->soopt = errno; + return -1; + } + // map the instance + cs->corerouter->cr_table[cs->instance_fd] = cs; + // wait for connection + uwsgi_cr_hook_instance_write(cs, fr_instance_connected); + } + + return len; +} + +ssize_t fr_recv_uwsgi_header(struct corerouter_session *cs) { + ssize_t len = read(cs->fd, cs->buffer->buf + cs->buffer_pos, 4 - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_recv_uwsgi_header()"); + return -1; + } + + cs->buffer_pos += len; + + // header ready + if (cs->buffer_pos == 4) { + memcpy(&cs->uh, cs->buffer->buf, 4); + cs->buffer_pos = 0; + uwsgi_cr_hook_read(cs, fr_recv_uwsgi_vars); + } + + return len; +} + void fastrouter_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socket *ugs, struct corerouter_session *cs, struct sockaddr *sa, socklen_t s_len) { + // set the first hook + uwsgi_cr_hook_read(cs, fr_recv_uwsgi_header); } int fastrouter_init() { ufr.cr.session_size = sizeof(struct fastrouter_session); - ufr.cr.switch_events = uwsgi_fastrouter_switch_events; ufr.cr.alloc_session = fastrouter_alloc_session; uwsgi_corerouter_init((struct uwsgi_corerouter *) &ufr); diff --git a/plugins/fastrouter/fr.h b/plugins/fastrouter/fr.h index 7f93dcc3..001f00f9 100644 --- a/plugins/fastrouter/fr.h +++ b/plugins/fastrouter/fr.h @@ -1,20 +1,14 @@ #include "../corerouter/cr.h" -#define FASTROUTER_STATUS_RECV_VARS 10 -#define FASTROUTER_STATUS_BUFFERING 11 - struct uwsgi_fastrouter { - struct uwsgi_corerouter cr; - }; struct fastrouter_session { - struct corerouter_session crs; - struct uwsgi_header uh; - char buffer[UMAX16]; + struct uwsgi_buffer *post_buf; + size_t post_buf_max; + size_t post_buf_len; + off_t post_buf_pos; }; -void uwsgi_fastrouter_switch_events(struct uwsgi_corerouter *, struct corerouter_session *, int interesting_fd); - diff --git a/plugins/fastrouter/fr_events.c b/plugins/fastrouter/fr_events.c index 3922a88b..39a62056 100644 --- a/plugins/fastrouter/fr_events.c +++ b/plugins/fastrouter/fr_events.c @@ -8,349 +8,266 @@ extern struct uwsgi_fastrouter ufr; void fr_get_hostname(char *key, uint16_t keylen, char *val, uint16_t vallen, void *data) { // here i use directly corerouter_session - struct corerouter_session *fr_session = (struct corerouter_session *) data; + struct corerouter_session *cs = (struct corerouter_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; + if (!uwsgi_strncmp("SERVER_NAME", 11, key, keylen) && !cs->hostname_len) { + cs->hostname = val; + cs->hostname_len = vallen; return; } - if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen) && !fr_session->has_key) { - fr_session->hostname = val; - fr_session->hostname_len = vallen; + if (!uwsgi_strncmp("HTTP_HOST", 9, key, keylen) && !cs->has_key) { + cs->hostname = val; + cs->hostname_len = vallen; return; } if (!uwsgi_strncmp("UWSGI_FASTROUTER_KEY", 20, key, keylen)) { - fr_session->has_key = 1; - fr_session->hostname = val; - fr_session->hostname_len = vallen; + cs->has_key = 1; + cs->hostname = val; + cs->hostname_len = vallen; return; } - if (ufr.cr.post_buffering > 0) { - if (!uwsgi_strncmp("CONTENT_LENGTH", 14, key, keylen)) { - fr_session->post_cl = uwsgi_str_num(val, vallen); - return; - } + if (!uwsgi_strncmp("CONTENT_LENGTH", 14, key, keylen)) { + cs->post_cl = uwsgi_str_num(val, vallen); + return; } +} + +ssize_t fr_instance_read_response(struct corerouter_session *); +ssize_t fr_read_body(struct corerouter_session *); + +ssize_t fr_write_body(struct corerouter_session *cs) { + struct fastrouter_session *fs = (struct fastrouter_session *) cs; + ssize_t len = write(cs->instance_fd, fs->post_buf->buf + fs->post_buf_pos, fs->post_buf_len - fs->post_buf_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_write_body()"); + return -1; + } + + fs->post_buf_pos += len; + + // the body chunk has been sent, start again reading from client and instance + if (fs->post_buf_pos == fs->post_buf_len) { + uwsgi_cr_hook_instance_write(cs, NULL); + uwsgi_cr_hook_instance_read(cs, fr_instance_read_response); + uwsgi_cr_hook_read(cs, fr_read_body); + } + + return len; +} + + +ssize_t fr_read_body(struct corerouter_session *cs) { + struct fastrouter_session *fs = (struct fastrouter_session *) cs; + ssize_t len = read(cs->fd, fs->post_buf->buf, fs->post_buf_max); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_read_body()"); + return -1; + } + + // connection closed + if (len == 0) return 0; + + fs->post_buf_len = len; + fs->post_buf_pos = 0; + + // ok we have a body, stop reading from the client and the instance and start writing to the instance + uwsgi_cr_hook_read(cs, NULL); + uwsgi_cr_hook_instance_read(cs, NULL); + uwsgi_cr_hook_instance_write(cs, fr_write_body); + + return len; +} + +ssize_t fr_write_response(struct corerouter_session *cs) { + ssize_t len = write(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->buffer_len - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_write_response()"); + return -1; + } + + cs->buffer_pos += len; + + // ok this response chunk is sent, let's wait for another one + if (cs->buffer_pos == cs->buffer_len) { + uwsgi_cr_hook_write(cs, NULL); + uwsgi_cr_hook_instance_read(cs, fr_instance_read_response); + } + + return len; +} + +ssize_t fr_instance_read_response(struct corerouter_session *cs) { + ssize_t len = read(cs->instance_fd, cs->buffer->buf, cs->buffer->len); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_instance_read_response()"); + return -1; + } + + // end of the response + if (len == 0) { + return 0; } -void uwsgi_fastrouter_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_session *cs, int interesting_fd) { + cs->buffer_pos = 0; + cs->buffer_len = len; + // ok stop reading from the instance, and start writing to the client + uwsgi_cr_hook_instance_read(cs, NULL); + uwsgi_cr_hook_write(cs, fr_write_response); + return len; +} - struct fastrouter_session *fr_session = (struct fastrouter_session *) cs; +ssize_t fr_instance_send_request(struct corerouter_session *cs) { + ssize_t len = write(cs->instance_fd, cs->buffer->buf + cs->buffer_pos, cs->uh.pktsize - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_instance_send_request()"); + return -1; + } + + cs->buffer_pos += len; + + // ok the request is sent, we can start sending client body (if any) and we can start waiting + // for response + if (cs->buffer_pos == cs->uh.pktsize) { + cs->buffer_pos = 0; + // stop writing to the instance + uwsgi_cr_hook_instance_write(cs, NULL); + // start reading from the instance + uwsgi_cr_hook_instance_read(cs, fr_instance_read_response); + // re-start reading from the client (for body or connection close) + struct fastrouter_session *fs = (struct fastrouter_session *) cs; + // allocate a buffer for client body (could be delimited or dynamic) + fs->post_buf_max = UMAX16; + if (cs->post_cl > 0) { + fs->post_buf_max = UMIN(UMAX16, cs->post_cl); + } + fs->post_buf = uwsgi_buffer_new(fs->post_buf_max); + if (!fs->post_buf) return -1; + uwsgi_cr_hook_read(cs, fr_read_body); + } + + return len; +} + +ssize_t fr_instance_send_request_header(struct corerouter_session *cs) { + ssize_t len = write(cs->instance_fd, &cs->uh + cs->buffer_pos, 4 - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_instance_send_request_header()"); + return -1; + } + + cs->buffer_pos += len; + + // ok the request is sent, we can start sending client body (if any) and we can start waiting + // for response + if (cs->buffer_pos == 4) { + cs->buffer_pos = 0; + uwsgi_cr_hook_instance_write(cs, fr_instance_send_request); + } + + return len; +} + +ssize_t fr_instance_connected(struct corerouter_session *cs) { socklen_t solen = sizeof(int); - struct iovec iov[2]; - - struct msghdr msg; - union { - struct cmsghdr cmsg; - char control[CMSG_SPACE(sizeof(int))]; - } msg_control; - struct cmsghdr *cmsg; - - ssize_t len; - char *post_tmp_buf[UMAX16]; - - switch (cs->status) { - - case COREROUTER_STATUS_RECV_HDR: - len = recv(cs->fd, (char *) (&fr_session->uh) + cs->h_pos, 4 - cs->h_pos, 0); -#ifdef UWSGI_EVENT_USE_PORT - event_queue_add_fd_read(ucr->queue, cs->fd); -#endif - if (len <= 0) { - if (len < 0) - uwsgi_error("recv()"); - corerouter_close_session(ucr, cs); - break; - } - cs->h_pos += len; - if (cs->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 - cs->status = FASTROUTER_STATUS_RECV_VARS; - } - break; - - - case FASTROUTER_STATUS_RECV_VARS: - - if (interesting_fd == -1) { - goto choose_node; - } - - len = recv(cs->fd, fr_session->buffer + cs->pos, fr_session->uh.pktsize - cs->pos, 0); -#ifdef UWSGI_EVENT_USE_PORT - event_queue_add_fd_read(ucr->queue, cs->fd); -#endif - if (len <= 0) { - uwsgi_error("recv()"); - corerouter_close_session(ucr, cs); - break; - } - cs->pos += len; - if (cs->pos == fr_session->uh.pktsize) { - if (uwsgi_hooked_parse(fr_session->buffer, fr_session->uh.pktsize, fr_get_hostname, (void *) fr_session)) { - corerouter_close_session(ucr, cs); - break; - } - - if (cs->hostname_len == 0) { - corerouter_close_session(ucr, cs); - break; - } - - - // the mapper hook -choose_node: - if (ucr->mapper(ucr, cs)) - break; - - // no address found - if (!cs->instance_address_len) { - // if fallback nodes are configured, trigger them - if (ucr->fallback) { - cs->instance_failed = 1; - } - corerouter_close_session(ucr, cs); - break; - } - - if (ucr->post_buffering > 0 && cs->post_cl > ucr->post_buffering) { - cs->status = FASTROUTER_STATUS_BUFFERING; - cs->buf_file_name = uwsgi_tmpname(ucr->pb_base_dir, "uwsgiXXXXX"); - if (!cs->buf_file_name) { - uwsgi_error("tempnam()"); - corerouter_close_session(ucr, cs); - break; - } - cs->post_remains = cs->post_cl; - - // 2 + UWSGI_POSTFILE + 2 + cs->buf_file_name - if (fr_session->uh.pktsize + (2 + 14 + 2 + strlen(cs->buf_file_name)) > UMAX16) { - uwsgi_log("unable to buffer request body to file %s: not enough space\n", cs->buf_file_name); - corerouter_close_session(ucr, cs); - break; - } - - char *ptr = fr_session->buffer + fr_session->uh.pktsize; - uint16_t bfn_len = strlen(cs->buf_file_name); - *ptr++ = 14; - *ptr++ = 0; - memcpy(ptr, "UWSGI_POSTFILE", 14); - ptr += 14; - *ptr++ = (char) (bfn_len & 0xff); - *ptr++ = (char) ((bfn_len >> 8) & 0xff); - memcpy(ptr, cs->buf_file_name, bfn_len); - fr_session->uh.pktsize += 2 + 14 + 2 + bfn_len; - - - cs->buf_file = fopen(cs->buf_file_name, "w"); - if (!cs->buf_file) { - uwsgi_error_open(cs->buf_file_name); - corerouter_close_session(ucr, cs); - break; - } - - } - - else { - - cs->pass_fd = is_unix(cs->instance_address, cs->instance_address_len); - - cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1); - - if (cs->instance_fd < 0) { - cs->instance_failed = 1; - cs->soopt = errno; - corerouter_close_session(ucr, cs); - break; - } - - - cs->status = COREROUTER_STATUS_CONNECTING; - ucr->cr_table[cs->instance_fd] = cs; - event_queue_add_fd_write(ucr->queue, cs->instance_fd); - } - } - break; - - - - case COREROUTER_STATUS_CONNECTING: - - if (interesting_fd == cs->instance_fd) { - - if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) { - uwsgi_error("getsockopt()"); - cs->instance_failed = 1; - corerouter_close_session(ucr, cs); - break; - } - - if (cs->soopt) { - cs->instance_failed = 1; - corerouter_close_session(ucr, cs); - break; - } - - fr_session->uh.modifier1 = cs->modifier1; - - 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; - - // increment node requests counter - if (cs->un) - cs->un->requests++; - - // fd passing: PERFORMANCE EXTREME BOOST !!! - if (cs->pass_fd && !uwsgi.no_fd_passing) { - 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; - - memcpy(CMSG_DATA(cmsg), &cs->fd, sizeof(int)); - - if (sendmsg(cs->instance_fd, &msg, 0) < 0) { - uwsgi_error("sendmsg()"); - } - - corerouter_close_session(ucr, cs); - break; - } - - if (writev(cs->instance_fd, iov, 2) < 0) { - uwsgi_error("writev()"); - corerouter_close_session(ucr, cs); - break; - } - - event_queue_fd_write_to_read(ucr->queue, cs->instance_fd); - cs->status = COREROUTER_STATUS_RESPONSE; - } - - break; - case COREROUTER_STATUS_RESPONSE: - - // data from instance - if (interesting_fd == cs->instance_fd) { - len = recv(cs->instance_fd, fr_session->buffer, UMAX16, 0); -#ifdef UWSGI_EVENT_USE_PORT - event_queue_add_fd_read(ucr->queue, cs->instance_fd); -#endif - if (len <= 0) { - if (len < 0) - uwsgi_error("recv()"); - corerouter_close_session(ucr, cs); - break; - } - - len = send(cs->fd, fr_session->buffer, len, 0); - - if (len <= 0) { - if (len < 0) - uwsgi_error("send()"); - corerouter_close_session(ucr, cs); - break; - } - - // update transfer statistics - if (cs->un) - cs->un->transferred += len; - } - // body from client - else if (interesting_fd == cs->fd) { - - //uwsgi_log("receiving body...\n"); - len = recv(cs->fd, fr_session->buffer, UMAX16, 0); -#ifdef UWSGI_EVENT_USE_PORT - event_queue_add_fd_read(ucr->queue, cs->fd); -#endif - if (len <= 0) { - if (len < 0) - uwsgi_error("recv()"); - corerouter_close_session(ucr, cs); - break; - } - - - len = send(cs->instance_fd, fr_session->buffer, len, 0); - - if (len <= 0) { - if (len < 0) - uwsgi_error("send()"); - corerouter_close_session(ucr, cs); - break; - } - } - - break; - - case FASTROUTER_STATUS_BUFFERING: - len = recv(cs->fd, post_tmp_buf, UMIN(UMAX16, cs->post_remains), 0); -#ifdef UWSGI_EVENT_USE_PORT - event_queue_add_fd_read(ucr->queue, cs->fd); -#endif - if (len <= 0) { - if (len < 0) - uwsgi_error("recv()"); - corerouter_close_session(ucr, cs); - break; - } - - if (fwrite(post_tmp_buf, len, 1, cs->buf_file) != 1) { - uwsgi_error("fwrite()"); - corerouter_close_session(ucr, cs); - break; - } - - cs->post_remains -= len; - - if (cs->post_remains == 0) { - // close the buf_file ASAP - fclose(cs->buf_file); - cs->buf_file = NULL; - - cs->pass_fd = is_unix(cs->instance_address, cs->instance_address_len); - - cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1); - - if (cs->instance_fd < 0) { - cs->instance_failed = 1; - corerouter_close_session(ucr, cs); - break; - } - - cs->status = COREROUTER_STATUS_CONNECTING; - ucr->cr_table[cs->instance_fd] = cs; - event_queue_add_fd_write(ucr->queue, cs->instance_fd); - } - break; - - - - - // fallback to destroy !!! - default: - uwsgi_log("unknown event: closing session\n"); - corerouter_close_session(ucr, cs); - break; + // first check for errors + if (getsockopt(cs->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&cs->soopt), &solen) < 0) { + uwsgi_error("fr_instance_connected()/getsockopt()"); + cs->instance_failed = 1; + return -1; } + + if (cs->soopt) { + cs->instance_failed = 1; + return -1; + } + + cs->buffer_pos = 0; + + // ok instance is connected, wait for write again + uwsgi_cr_hook_instance_write(cs, fr_instance_send_request_header); + // return a value > 0 + return 1; +} + +ssize_t fr_recv_uwsgi_vars(struct corerouter_session *cs) { + // increase buffer if needed + if (uwsgi_buffer_fix(cs->buffer, cs->uh.pktsize)) return -1; + ssize_t len = read(cs->fd, cs->buffer->buf + cs->buffer_pos, cs->uh.pktsize - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_recv_uwsgi_vars()"); + return -1; + } + + cs->buffer_pos += len; + + // headers received, ready to choose the instance + if (cs->buffer_pos == cs->uh.pktsize) { + struct uwsgi_corerouter *ucr = cs->corerouter; + // find the hostname + if (uwsgi_hooked_parse(cs->buffer->buf, cs->uh.pktsize, fr_get_hostname, (void *) cs)) { + return -1; + } + // check the hostname; + if (cs->hostname_len == 0) return -1; + // find an instance using the key + if (cs->corerouter->mapper(cs->corerouter, cs)) return -1; + // check instance + if (cs->instance_address_len == 0) { + // if fallback nodes are configured, trigger them + if (ucr->fallback) { + cs->instance_failed = 1; + } + return -1; + } + + // stop receiving from the client + uwsgi_cr_hook_read(cs, NULL); + + // start async connect + cs->instance_fd = uwsgi_connectn(cs->instance_address, cs->instance_address_len, 0, 1); + if (cs->instance_fd < 0) { + cs->instance_failed = 1; + cs->soopt = errno; + return -1; + } + // map the instance + cs->corerouter->cr_table[cs->instance_fd] = cs; + // wait for connection + uwsgi_cr_hook_instance_write(cs, fr_instance_connected); + } + + return len; +} + +ssize_t fr_recv_uwsgi_header(struct corerouter_session *cs) { + ssize_t len = read(cs->fd, cs->buffer->buf + cs->buffer_pos, 4 - cs->buffer_pos); + if (len < 0) { + cr_try_again; + uwsgi_error("fr_recv_uwsgi_header()"); + return -1; + } + + cs->buffer_pos += len; + + // header ready + if (cs->buffer_pos == 4) { + memcpy(&cs->uh, cs->buffer->buf, 4); + cs->buffer_pos = 0; + uwsgi_cr_hook_read(cs, fr_recv_uwsgi_vars); + } + + return len; } diff --git a/plugins/fastrouter/uwsgiplugin.py b/plugins/fastrouter/uwsgiplugin.py index c15a27b8..133452b2 100644 --- a/plugins/fastrouter/uwsgiplugin.py +++ b/plugins/fastrouter/uwsgiplugin.py @@ -6,4 +6,4 @@ LIBS = [] REQUIRES = ['corerouter'] -GCC_LIST = ['fastrouter', 'fr_events'] +GCC_LIST = ['fastrouter'] diff --git a/uwsgi.h b/uwsgi.h index 73d8a8e1..5cda78b5 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -2396,6 +2396,12 @@ int event_queue_interesting_fd(void *, int); int event_queue_interesting_fd_has_error(void *, int); int event_queue_fd_write_to_read(int, int); int event_queue_fd_read_to_write(int, int); +int event_queue_fd_readwrite_to_read(int, int); +int event_queue_fd_readwrite_to_write(int, int); +int event_queue_fd_read_to_readwrite(int, int); +int event_queue_fd_write_to_readwrite(int, int); +int event_queue_interesting_fd_is_read(void *, int); +int event_queue_interesting_fd_is_write(void *, int); int event_queue_add_timer(int, int *, int); struct uwsgi_timer *event_queue_ack_timer(int); @@ -3287,6 +3293,7 @@ void uwsgi_set_sockets_protocols(void); struct uwsgi_buffer *uwsgi_buffer_new(size_t); int uwsgi_buffer_append(struct uwsgi_buffer *, char *, size_t); +int uwsgi_buffer_fix(struct uwsgi_buffer *, size_t); void uwsgi_buffer_destroy(struct uwsgi_buffer *); void uwsgi_httpize_var(char *, size_t); diff --git a/uwsgiconfig.py b/uwsgiconfig.py index 56a79689..684d52e9 100644 --- a/uwsgiconfig.py +++ b/uwsgiconfig.py @@ -27,6 +27,16 @@ if not GCC: CPP = os.environ.get('CPP', 'cpp') +CPUCOUNT = 1 +try: + import multiprocessing + CPUCOUNT = multiprocessing.cpu_count() +except: + try: + CPUCOUNT = int(os.sysconf('SC_NPROCESSORS_ONLN')) + except: + pass + binary_list = [] # this is used for reporting (at the end of the build) @@ -182,6 +192,7 @@ def build_uwsgi(uc, print_only=False): print(' '.join(cflags)) sys.exit(0) + print("detected CPU cores: %d" % CPUCOUNT) print("configured CFLAGS: %s" % ' '.join(cflags)) try: