first (destructive) commit for 1.4

This commit is contained in:
Roberto De Ioris
2012-10-20 13:47:07 +02:00
parent 6ea535a74d
commit c65adaaac7
12 changed files with 1121 additions and 628 deletions
+15 -2
View File
@@ -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;
}
}
+235 -14
View File
@@ -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;
+3 -3
View File
@@ -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;
+278 -26
View File
@@ -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;
}
}
}
}
+29 -33
View File
@@ -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 *));
+14 -203
View File
@@ -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;
}
+282 -11
View File
@@ -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);
+4 -10
View File
@@ -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);
+242 -325
View File
@@ -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;
}
+1 -1
View File
@@ -6,4 +6,4 @@ LIBS = []
REQUIRES = ['corerouter']
GCC_LIST = ['fastrouter', 'fr_events']
GCC_LIST = ['fastrouter']
+7
View File
@@ -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);
+11
View File
@@ -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: