diff --git a/gateway.c b/gateway.c
index 2fc94a72..df577957 100644
--- a/gateway.c
+++ b/gateway.c
@@ -140,3 +140,16 @@ struct uwsgi_gateway_socket *uwsgi_new_gateway_socket_from_fd(int fd, char *owne
}
+void uwsgi_gateway_go_cheap(char *gw_id, int queue, int *i_am_cheap) {
+
+ uwsgi_log("[%s pid %d] no more nodes available. Going cheap...\n", gw_id, (int) uwsgi.mypid);
+ struct uwsgi_gateway_socket *ugs = uwsgi.gateway_sockets;
+ while (ugs) {
+ if (!strcmp(ugs->owner, gw_id) && !ugs->subscription) {
+ event_queue_del_fd(queue, ugs->fd, event_queue_read());
+ }
+ ugs = ugs->next;
+ }
+ *i_am_cheap = 1;
+}
+
diff --git a/hello_world.py b/hello_world.py
index 8b2301a2..7f4a4145 100644
--- a/hello_world.py
+++ b/hello_world.py
@@ -2,5 +2,6 @@ import uwsgi
#uwsgi.cache_set('/', "HTTP/1.1 200 OK\r\nContent-Type: text/html\r\n\r\nHello World from cache")
def application(env, start_response):
start_response('200 OK', [('Content-Type', 'text/html')])
+ yield "foobar"
yield str(env['wsgi.input'].fileno())
yield "
Hello World
"
diff --git a/lib/corerouter.h b/lib/corerouter.h
index 0ca30daa..0ae99e91 100644
--- a/lib/corerouter.h
+++ b/lib/corerouter.h
@@ -91,19 +91,6 @@ static void *uwsgi_corerouter_setup_event_queue(char *gw_id, int id, int nevents
return event_queue_alloc(nevents);
}
-static void __attribute__ ((unused)) uwsgi_corerouter_go_cheap(char *gw_id, int queue, int *i_am_cheap) {
-
- uwsgi_log("[%s pid %d] no more nodes available. Going cheap...\n", gw_id, (int) uwsgi.mypid);
- struct uwsgi_gateway_socket *ugs = uwsgi.gateway_sockets;
- while (ugs) {
- if (!strcmp(ugs->owner, gw_id) && !ugs->subscription) {
- event_queue_del_fd(queue, ugs->fd, event_queue_read());
- }
- ugs = ugs->next;
- }
- *i_am_cheap = 1;
-}
-
static void __attribute__ ((unused)) uwsgi_corerouter_manage_subscription(char *gw_id, int id, struct uwsgi_gateway_socket *ugs, int queue, struct uwsgi_subscribe_slot **subscriptions, int regexp, void (*parse_hook) (char *, uint16_t, char *, uint16_t, void *), int cheap, int *i_am_cheap) {
int i;
@@ -145,7 +132,7 @@ static void __attribute__ ((unused)) uwsgi_corerouter_manage_subscription(char *
uwsgi_remove_subscribe_node(subscriptions, node);
}
if (*subscriptions == NULL && cheap && !*i_am_cheap) {
- uwsgi_corerouter_go_cheap(gw_id, queue, i_am_cheap);
+ uwsgi_gateway_go_cheap(gw_id, queue, i_am_cheap);
}
}
}
@@ -205,7 +192,7 @@ static void __attribute__ ((unused)) uwsgi_corerouter_manage_internal_subscripti
uwsgi_remove_subscribe_node(subscriptions, node);
}
if (*subscriptions == NULL && cheap && !*i_am_cheap) {
- uwsgi_corerouter_go_cheap(gw_id, queue, i_am_cheap);
+ uwsgi_gateway_go_cheap(gw_id, queue, i_am_cheap);
}
}
}
diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c
index 945ca5a7..19a1e8d7 100644
--- a/plugins/fastrouter/fastrouter.c
+++ b/plugins/fastrouter/fastrouter.c
@@ -17,23 +17,10 @@ extern struct uwsgi_server uwsgi;
#include "../../lib/corerouter.h"
-#define FASTROUTER_STATUS_FREE 0
-#define FASTROUTER_STATUS_CONNECTING 1
-#define FASTROUTER_STATUS_RECV_HDR 2
-#define FASTROUTER_STATUS_RECV_VARS 3
-#define FASTROUTER_STATUS_RESPONSE 4
-#define FASTROUTER_STATUS_BUFFERING 5
-
#ifdef UWSGI_SCTP
-#define FASTROUTER_STATUS_SCTP_NODE_FREE 6
-#define FASTROUTER_STATUS_SCTP_RESPONSE 7
+extern struct uwsgi_fr_sctp_node *uwsgi_fastrouter_sctp_nodes;
#endif
-#define add_timeout(x) uwsgi_add_rb_timer(ufr.timeouts, time(NULL)+ufr.socket_timeout, x)
-#define add_check_timeout(x) uwsgi_add_rb_timer(timeouts, time(NULL)+x, NULL)
-#define del_check_timeout(x) rb_erase(&x->rbt, timeouts);
-#define del_timeout(x) rb_erase(&x->timeout->rbt, ufr.timeouts); free(x->timeout);
-
void fastrouter_send_stats(int);
#include "fr.h"
@@ -203,54 +190,9 @@ void fastrouter_manage_subscription(char *key, uint16_t keylen, char *val, uint1
}
}
-struct fastrouter_session {
-
- int fd;
- int instance_fd;
- int status;
- struct uwsgi_header uh;
- uint8_t h_pos;
- uint16_t pos;
-
- char *hostname;
- uint16_t hostname_len;
-
- int has_key;
-#ifdef UWSGI_SCTP
- int persistent;
-#endif
-
- char *instance_address;
- uint64_t instance_address_len;
-
- struct uwsgi_subscribe_node *un;
- struct uwsgi_string_list *static_node;
- int pass_fd;
- int soopt;
- int timed_out;
-
- struct uwsgi_rb_timer *timeout;
- int instance_failed;
-
- size_t post_cl;
- size_t post_remains;
-
- struct uwsgi_string_list *fallback;
-
- char *buf_file_name;
- FILE *buf_file;
-
- uint8_t modifier1;
- uint8_t modifier2;
-
- char *tmp_socket_name;
-
- char buffer[0xffff];
-};
-
static struct uwsgi_rb_timer *reset_timeout(struct fastrouter_session *);
-static void close_session(struct fastrouter_session *fr_session) {
+void close_session(struct fastrouter_session *fr_session) {
if (fr_session->instance_fd != -1) {
@@ -302,7 +244,7 @@ static void close_session(struct fastrouter_session *fr_session) {
uwsgi_remove_subscribe_node(&ufr.subscriptions, fr_session->un);
}
if (ufr.subscriptions == NULL && ufr.cheap && !ufr.i_am_cheap && !ufr.fallback) {
- uwsgi_corerouter_go_cheap("uWSGI fastrouter", ufr.queue, &ufr.i_am_cheap);
+ uwsgi_gateway_go_cheap("uWSGI fastrouter", ufr.queue, &ufr.i_am_cheap);
}
}
@@ -403,7 +345,20 @@ static void expire_timeouts() {
if (urbt->key <= current) {
fr_session = (struct fastrouter_session *) urbt->data;
fr_session->timed_out = 1;
- close_session(fr_session);
+ if (fr_session->retry) {
+ fr_session->retry = 0;
+ uwsgi_fastrouter_switch_events(fr_session, -1, ufr.magic_table);
+ if (fr_session->retry) {
+ del_timeout(fr_session);
+ fr_session->timeout = add_fake_timeout(fr_session);
+ }
+ else {
+ fr_session->timeout = reset_timeout(fr_session);
+ }
+ }
+ else {
+ close_session(fr_session);
+ }
continue;
}
@@ -514,27 +469,14 @@ void fastrouter_loop(int id) {
time_t delta;
- char *post_tmp_buf[0xffff];
- int tmp_socket_name_len;
-
-
struct uwsgi_rb_timer *min_timeout;
- struct msghdr msg;
- union {
- struct cmsghdr cmsg;
- char control[CMSG_SPACE(sizeof(int))];
- } msg_control;
- struct cmsghdr *cmsg;
-
int interesting_fd;
int new_connection;
- ssize_t len;
- char *magic_table[0xff];
if (ufr.pattern) {
- init_magic_table(magic_table);
+ init_magic_table(ufr.magic_table);
}
struct sockaddr_un fr_addr;
@@ -542,9 +484,6 @@ void fastrouter_loop(int id) {
struct fastrouter_session *fr_session;
- struct iovec iov[2];
-
- socklen_t solen = sizeof(int);
ufr.timeouts = uwsgi_init_rb_timer();
@@ -687,8 +626,14 @@ void fastrouter_loop(int id) {
continue;
if (event_queue_interesting_fd_has_error(events, i)) {
- close_session(fr_session);
- continue;
+#ifdef UWSGI_SCTP
+ if (!fr_session->persistent) {
+#endif
+ close_session(fr_session);
+ continue;
+#ifdef UWSGI_SCTP
+ }
+#endif
}
#ifdef UWSGI_SCTP
@@ -698,503 +643,10 @@ void fastrouter_loop(int id) {
#ifdef UWSGI_SCTP
}
#endif
-
- switch (fr_session->status) {
-
- case FASTROUTER_STATUS_RECV_HDR:
- len = recv(fr_session->fd, (char *) (&fr_session->uh) + fr_session->h_pos, 4 - fr_session->h_pos, 0);
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("recv()");
- close_session(fr_session);
- break;
- }
- fr_session->h_pos += len;
- if (fr_session->h_pos == 4) {
-#ifdef UWSGI_DEBUG
- uwsgi_log("modifier1: %d pktsize: %d modifier2: %d\n", fr_session->uh.modifier1, fr_session->uh.pktsize, fr_session->uh.modifier2);
-#endif
- fr_session->status = FASTROUTER_STATUS_RECV_VARS;
- }
- break;
-
-
- case FASTROUTER_STATUS_RECV_VARS:
- len = recv(fr_session->fd, fr_session->buffer + fr_session->pos, fr_session->uh.pktsize - fr_session->pos, 0);
- if (len <= 0) {
- uwsgi_error("recv()");
- close_session(fr_session);
- break;
- }
- fr_session->pos += len;
- if (fr_session->pos == fr_session->uh.pktsize) {
- if (uwsgi_hooked_parse(fr_session->buffer, fr_session->uh.pktsize, fr_get_hostname, (void *) fr_session)) {
- close_session(fr_session);
- break;
- }
-
- if (fr_session->hostname_len == 0) {
- close_session(fr_session);
- break;
- }
-
-#ifdef UWSGI_DEBUG
- //uwsgi_log("requested domain %.*s\n", fr_session->hostname_len, fr_session->hostname);
-#endif
- if (ufr.use_cache) {
- fr_session->instance_address = uwsgi_cache_get(fr_session->hostname, fr_session->hostname_len, &fr_session->instance_address_len);
- char *cs_mod = uwsgi_str_contains(fr_session->instance_address, fr_session->instance_address_len, ',');
- if (cs_mod) {
- fr_session->modifier1 = uwsgi_str_num(cs_mod + 1, (fr_session->instance_address_len - (cs_mod - fr_session->instance_address)) - 1);
- fr_session->instance_address_len = (cs_mod - fr_session->instance_address);
- }
- }
- else if (ufr.pattern) {
- magic_table['s'] = uwsgi_concat2n(fr_session->hostname, fr_session->hostname_len, "", 0);
- fr_session->tmp_socket_name = magic_sub(ufr.pattern, ufr.pattern_len, &tmp_socket_name_len, magic_table);
- free(magic_table['s']);
- fr_session->instance_address_len = tmp_socket_name_len;
- fr_session->instance_address = fr_session->tmp_socket_name;
- }
- else if (ufr.has_subscription_sockets) {
- fr_session->un = uwsgi_get_subscribe_node(&ufr.subscriptions, fr_session->hostname, fr_session->hostname_len, ufr.subscription_regexp);
- if (fr_session->un && fr_session->un->len) {
- fr_session->instance_address = fr_session->un->name;
- fr_session->instance_address_len = fr_session->un->len;
- fr_session->modifier1 = fr_session->un->modifier1;
- }
- else if (ufr.subscriptions == NULL && ufr.cheap && !ufr.i_am_cheap) {
- uwsgi_corerouter_go_cheap("uWSGI fastrouter", ufr.queue, &ufr.i_am_cheap);
- }
- }
- else if (ufr.base) {
- fr_session->tmp_socket_name = uwsgi_concat2nn(ufr.base, ufr.base_len, fr_session->hostname, fr_session->hostname_len, &tmp_socket_name_len);
- fr_session->instance_address_len = tmp_socket_name_len;
- fr_session->instance_address = fr_session->tmp_socket_name;
- }
- else if (ufr.code_string_code && ufr.code_string_function) {
- if (uwsgi.p[ufr.code_string_modifier1]->code_string) {
- fr_session->instance_address = uwsgi.p[ufr.code_string_modifier1]->code_string("uwsgi_fastrouter", ufr.code_string_code, ufr.code_string_function, fr_session->hostname, fr_session->hostname_len);
- if (fr_session->instance_address) {
- fr_session->instance_address_len = strlen(fr_session->instance_address);
- char *cs_mod = uwsgi_str_contains(fr_session->instance_address, fr_session->instance_address_len, ',');
- if (cs_mod) {
- fr_session->modifier1 = uwsgi_str_num(cs_mod + 1, (fr_session->instance_address_len - (cs_mod - fr_session->instance_address)) - 1);
- fr_session->instance_address_len = (cs_mod - fr_session->instance_address);
- }
- }
- }
- }
- else if (ufr.to_socket) {
- fr_session->instance_address = ufr.to_socket->name;
- fr_session->instance_address_len = ufr.to_socket->name_len;
- }
- else if (ufr.static_nodes) {
- if (!ufr.current_static_node) {
- ufr.current_static_node = ufr.static_nodes;
- }
-
- fr_session->static_node = ufr.current_static_node;
-
- // is it a dead node ?
- if (fr_session->static_node->custom > 0) {
-
- // gracetime passed ?
- if (fr_session->static_node->custom + ufr.static_node_gracetime <= (uint64_t) uwsgi_now()) {
- fr_session->static_node->custom = 0;
- }
- else {
- struct uwsgi_string_list *tmp_node = fr_session->static_node;
- struct uwsgi_string_list *next_node = fr_session->static_node->next;
- fr_session->static_node = NULL;
- // needed for 1-node only setups
- if (!next_node) next_node = ufr.static_nodes;
-
- while(tmp_node != next_node) {
- if (!next_node) {
- next_node = ufr.static_nodes;
- }
-
- if (tmp_node == next_node) break;
-
- if (next_node->custom == 0) {
- fr_session->static_node = next_node;
- break;
- }
- next_node = next_node->next;
- }
- }
- }
-
- if (fr_session->static_node) {
-
- fr_session->instance_address = fr_session->static_node->value;
- fr_session->instance_address_len = fr_session->static_node->len;
- // set the next one
- ufr.current_static_node = fr_session->static_node->next;
- }
- else {
- // set the next one
- ufr.current_static_node = ufr.current_static_node->next;
- }
-
- }
-#ifdef UWSGI_SCTP
- else if (ufr.has_sctp_sockets > 0) {
-
- struct uwsgi_fr_sctp_node *ufsn = uwsgi_fastrouter_sctp_nodes;
- int choosen_fd = -1;
- while(ufsn) {
- if (ufr.fr_table[ufsn->fd]->status == FASTROUTER_STATUS_SCTP_NODE_FREE) {
- choosen_fd = ufsn->fd;
- break;
- }
- if (ufsn->next == uwsgi_fastrouter_sctp_nodes) {
- break;
- }
-
- ufsn = ufsn->next;
- }
-
- // no nodes available
- if (choosen_fd == -1) break;
-
- struct sctp_sndrcvinfo sinfo;
- memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
- sinfo.sinfo_stream = 0;
- memcpy(&sinfo.sinfo_ppid, &fr_session->uh, sizeof(uint32_t));
- len = sctp_send(choosen_fd, fr_session->buffer, fr_session->uh.pktsize, &sinfo, 0);
- fr_session->instance_fd = choosen_fd;
- fr_session->status = FASTROUTER_STATUS_SCTP_RESPONSE;
- ufr.fr_table[fr_session->instance_fd]->status = FASTROUTER_STATUS_SCTP_RESPONSE;
- ufr.fr_table[fr_session->instance_fd]->fd = fr_session->fd;
- break;
- }
-#endif
-
- // no address found
- if (!fr_session->instance_address_len) {
- // if fallback nodes are configured, trigger them
- if (ufr.fallback) {
- fr_session->instance_failed = 1;
- }
- close_session(fr_session);
- break;
- }
-
- if (ufr.post_buffering > 0 && fr_session->post_cl > ufr.post_buffering) {
- fr_session->status = FASTROUTER_STATUS_BUFFERING;
- fr_session->buf_file_name = uwsgi_tmpname(ufr.pb_base_dir, "uwsgiXXXXX");
- if (!fr_session->buf_file_name) {
- uwsgi_error("tempnam()");
- close_session(fr_session);
- break;
- }
- fr_session->post_remains = fr_session->post_cl;
-
- // 2 + UWSGI_POSTFILE + 2 + fr_session->buf_file_name
- if (fr_session->uh.pktsize + (2 + 14 + 2 + strlen(fr_session->buf_file_name)) > 0xffff) {
- uwsgi_log("unable to buffer request body to file %s: not enough space\n", fr_session->buf_file_name);
- close_session(fr_session);
- break;
- }
-
- char *ptr = fr_session->buffer + fr_session->uh.pktsize;
- uint16_t bfn_len = strlen(fr_session->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, fr_session->buf_file_name, bfn_len);
- fr_session->uh.pktsize += 2 + 14 + 2 + bfn_len;
-
-
- fr_session->buf_file = fopen(fr_session->buf_file_name, "w");
- if (!fr_session->buf_file) {
- uwsgi_error_open(fr_session->buf_file_name);
- close_session(fr_session);
- break;
- }
-
- }
-
- else {
-
- fr_session->pass_fd = is_unix(fr_session->instance_address, fr_session->instance_address_len);
-
- fr_session->instance_fd = uwsgi_connectn(fr_session->instance_address, fr_session->instance_address_len, 0, 1);
-
- if (fr_session->instance_fd < 0) {
- fr_session->instance_failed = 1;
- fr_session->soopt = errno;
- close_session(fr_session);
- break;
- }
-
-
- fr_session->status = FASTROUTER_STATUS_CONNECTING;
- ufr.fr_table[fr_session->instance_fd] = fr_session;
- event_queue_add_fd_write(ufr.queue, fr_session->instance_fd);
- }
- }
- break;
-
-
-
- case FASTROUTER_STATUS_CONNECTING:
-
- if (interesting_fd == fr_session->instance_fd) {
-
- if (getsockopt(fr_session->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&fr_session->soopt), &solen) < 0) {
- uwsgi_error("getsockopt()");
- fr_session->instance_failed = 1;
- close_session(fr_session);
- break;
- }
-
- if (fr_session->soopt) {
- fr_session->instance_failed = 1;
- close_session(fr_session);
- break;
- }
-
- fr_session->uh.modifier1 = fr_session->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 (fr_session->un)
- fr_session->un->requests++;
-
- // fd passing: PERFORMANCE EXTREME BOOST !!!
- if (fr_session->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), &fr_session->fd, sizeof(int));
-
- if (sendmsg(fr_session->instance_fd, &msg, 0) < 0) {
- uwsgi_error("sendmsg()");
- }
-
- close_session(fr_session);
- break;
- }
-
- if (writev(fr_session->instance_fd, iov, 2) < 0) {
- uwsgi_error("writev()");
- close_session(fr_session);
- break;
- }
-
- event_queue_fd_write_to_read(ufr.queue, fr_session->instance_fd);
- fr_session->status = FASTROUTER_STATUS_RESPONSE;
- }
-
- break;
-#ifdef UWSGI_SCTP
- case FASTROUTER_STATUS_SCTP_NODE_FREE:
-
- {
- struct sctp_sndrcvinfo sinfo;
- int msg_flags;
-
- memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
- len = sctp_recvmsg(fr_session->instance_fd, fr_session->buffer, 0xffff, NULL, NULL, &sinfo, &msg_flags);
- }
- // remove the SCTP node
- uwsgi_fr_sctp_add_node(fr_session->instance_fd);
- ufr.fr_table[fr_session->instance_fd] = NULL;
- free(fr_session);
- close(interesting_fd);
- break;
- case FASTROUTER_STATUS_SCTP_RESPONSE:
+ uwsgi_fastrouter_switch_events(fr_session, interesting_fd, ufr.magic_table);
+
- // data from instance
- if (interesting_fd == fr_session->instance_fd) {
- struct sctp_sndrcvinfo sinfo;
- int msg_flags;
- memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
- len = sctp_recvmsg(fr_session->instance_fd, fr_session->buffer, 0xffff, NULL, NULL, &sinfo, &msg_flags);
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("recv()");
- if (!msg_flags) {
- // REMOVE THE NODE
- uwsgi_fr_sctp_add_node(fr_session->instance_fd);
- ufr.fr_table[fr_session->instance_fd] = NULL;
- free(fr_session);
- close(interesting_fd);
- }
- close_session(ufr.fr_table[fr_session->fd]);
- break;
- }
-
- // check for close packet
- if (sinfo.sinfo_stream == 2) {
- uwsgi_log("C L O S I N G\n");
- fr_session->status = FASTROUTER_STATUS_SCTP_NODE_FREE;
- close_session(ufr.fr_table[fr_session->fd]);
- break;
- }
-
- len = send(fr_session->fd, fr_session->buffer, len, 0);
-
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("send()");
- close_session(ufr.fr_table[fr_session->fd]);
- break;
- }
-
- // update transfer statistics
- if (fr_session->un)
- fr_session->un->transferred += len;
-
- }
- // body from client
- else if (interesting_fd == fr_session->fd) {
-
- //uwsgi_log("receiving body...\n");
- len = recv(fr_session->fd, fr_session->buffer, 0xffff, 0);
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("recv()");
- close_session(fr_session);
- break;
- }
-
- struct sctp_sndrcvinfo sinfo;
- memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
- // stream 1 is for BODY
- sinfo.sinfo_stream = 1;
-
- len = sctp_send(fr_session->instance_fd, fr_session->buffer, len, &sinfo, 0);
-
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("send()");
- close_session(fr_session);
- break;
- }
- }
-
- break;
-#endif
- case FASTROUTER_STATUS_RESPONSE:
-
- // data from instance
- if (interesting_fd == fr_session->instance_fd) {
- len = recv(fr_session->instance_fd, fr_session->buffer, 0xffff, 0);
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("recv()");
- close_session(fr_session);
- break;
- }
-
- len = send(fr_session->fd, fr_session->buffer, len, 0);
-
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("send()");
- close_session(fr_session);
- break;
- }
-
- // update transfer statistics
- if (fr_session->un)
- fr_session->un->transferred += len;
- }
- // body from client
- else if (interesting_fd == fr_session->fd) {
-
- //uwsgi_log("receiving body...\n");
- len = recv(fr_session->fd, fr_session->buffer, 0xffff, 0);
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("recv()");
- close_session(fr_session);
- break;
- }
-
-
- len = send(fr_session->instance_fd, fr_session->buffer, len, 0);
-
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("send()");
- close_session(fr_session);
- break;
- }
- }
-
- break;
-
- case FASTROUTER_STATUS_BUFFERING:
- len = recv(fr_session->fd, post_tmp_buf, UMIN(0xffff, fr_session->post_remains), 0);
- if (len <= 0) {
- if (len < 0)
- uwsgi_error("recv()");
- close_session(fr_session);
- break;
- }
-
- if (fwrite(post_tmp_buf, len, 1, fr_session->buf_file) != 1) {
- uwsgi_error("fwrite()");
- close_session(fr_session);
- break;
- }
-
- fr_session->post_remains -= len;
-
- if (fr_session->post_remains == 0) {
- // close the buf_file ASAP
- fclose(fr_session->buf_file);
- fr_session->buf_file = NULL;
-
- fr_session->pass_fd = is_unix(fr_session->instance_address, fr_session->instance_address_len);
-
- fr_session->instance_fd = uwsgi_connectn(fr_session->instance_address, fr_session->instance_address_len, 0, 1);
-
- if (fr_session->instance_fd < 0) {
- fr_session->instance_failed = 1;
- close_session(fr_session);
- break;
- }
-
- fr_session->status = FASTROUTER_STATUS_CONNECTING;
- ufr.fr_table[fr_session->instance_fd] = fr_session;
- event_queue_add_fd_write(ufr.queue, fr_session->instance_fd);
- }
- break;
-
-
-
-
- // fallback to destroy !!!
- default:
- uwsgi_log("unknown event: closing session\n");
- close_session(fr_session);
- break;
-
- }
}
}
}
diff --git a/plugins/fastrouter/fr.h b/plugins/fastrouter/fr.h
index 773f2791..ee352c36 100644
--- a/plugins/fastrouter/fr.h
+++ b/plugins/fastrouter/fr.h
@@ -1,3 +1,22 @@
+#define FASTROUTER_STATUS_FREE 0
+#define FASTROUTER_STATUS_CONNECTING 1
+#define FASTROUTER_STATUS_RECV_HDR 2
+#define FASTROUTER_STATUS_RECV_VARS 3
+#define FASTROUTER_STATUS_RESPONSE 4
+#define FASTROUTER_STATUS_BUFFERING 5
+
+#ifdef UWSGI_SCTP
+#define FASTROUTER_STATUS_SCTP_NODE_FREE 6
+#define FASTROUTER_STATUS_SCTP_RESPONSE 7
+#endif
+
+#define add_timeout(x) uwsgi_add_rb_timer(ufr.timeouts, time(NULL)+ufr.socket_timeout, x)
+#define add_fake_timeout(x) uwsgi_add_rb_timer(ufr.timeouts, time(NULL)+1, x)
+#define add_check_timeout(x) uwsgi_add_rb_timer(timeouts, time(NULL)+x, NULL)
+#define del_check_timeout(x) rb_erase(&x->rbt, timeouts);
+#define del_timeout(x) rb_erase(&x->timeout->rbt, ufr.timeouts); free(x->timeout);
+
+
struct uwsgi_fastrouter {
int has_sockets;
@@ -14,6 +33,8 @@ struct uwsgi_fastrouter {
int use_cache;
int nevents;
+ char *magic_table[0xff];
+
int queue;
char *pattern;
@@ -70,5 +91,58 @@ struct uwsgi_fr_sctp_node {
};
struct uwsgi_fr_sctp_node *uwsgi_fr_sctp_add_node(int);
+void uwsgi_fr_sctp_del_node(int);
void uwsgi_opt_fastrouter_sctp(char *, char *, void *);
+
#endif
+
+struct fastrouter_session {
+
+ int fd;
+ int instance_fd;
+ int status;
+ struct uwsgi_header uh;
+ uint8_t h_pos;
+ uint16_t pos;
+
+ char *hostname;
+ uint16_t hostname_len;
+
+ int has_key;
+ int retry;
+#ifdef UWSGI_SCTP
+ int persistent;
+#endif
+
+ char *instance_address;
+ uint64_t instance_address_len;
+
+ struct uwsgi_subscribe_node *un;
+ struct uwsgi_string_list *static_node;
+ int pass_fd;
+ int soopt;
+ int timed_out;
+
+ struct uwsgi_rb_timer *timeout;
+ int instance_failed;
+
+ size_t post_cl;
+ size_t post_remains;
+
+ struct uwsgi_string_list *fallback;
+
+ char *buf_file_name;
+ FILE *buf_file;
+
+ uint8_t modifier1;
+ uint8_t modifier2;
+
+ char *tmp_socket_name;
+
+ char buffer[0xffff];
+};
+
+
+void uwsgi_fastrouter_switch_events(struct fastrouter_session *, int intersting_fd, char **);
+void close_session(struct fastrouter_session *);
+void fr_get_hostname(char *, uint16_t, char *, uint16_t, void *);
diff --git a/plugins/fastrouter/fr_events.c b/plugins/fastrouter/fr_events.c
new file mode 100644
index 00000000..4b74fa0a
--- /dev/null
+++ b/plugins/fastrouter/fr_events.c
@@ -0,0 +1,547 @@
+#include "../../uwsgi.h"
+
+#include "fr.h"
+
+extern struct uwsgi_server uwsgi;
+extern struct uwsgi_fastrouter ufr;
+
+#ifdef UWSGI_SCTP
+extern struct uwsgi_fr_sctp_node *uwsgi_fastrouter_sctp_nodes;
+#endif
+
+void uwsgi_fastrouter_switch_events(struct fastrouter_session *fr_session, int interesting_fd, char **magic_table) {
+
+ 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[0xffff];
+
+ int tmp_socket_name_len;
+
+ switch (fr_session->status) {
+
+ case FASTROUTER_STATUS_RECV_HDR:
+ len = recv(fr_session->fd, (char *) (&fr_session->uh) + fr_session->h_pos, 4 - fr_session->h_pos, 0);
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("recv()");
+ close_session(fr_session);
+ break;
+ }
+ fr_session->h_pos += len;
+ if (fr_session->h_pos == 4) {
+#ifdef UWSGI_DEBUG
+ uwsgi_log("modifier1: %d pktsize: %d modifier2: %d\n", fr_session->uh.modifier1, fr_session->uh.pktsize, fr_session->uh.modifier2);
+#endif
+ fr_session->status = FASTROUTER_STATUS_RECV_VARS;
+ }
+ break;
+
+
+ case FASTROUTER_STATUS_RECV_VARS:
+
+ if (interesting_fd == -1) goto choose_node;
+
+ len = recv(fr_session->fd, fr_session->buffer + fr_session->pos, fr_session->uh.pktsize - fr_session->pos, 0);
+ if (len <= 0) {
+ uwsgi_error("recv()");
+ close_session(fr_session);
+ break;
+ }
+ fr_session->pos += len;
+ if (fr_session->pos == fr_session->uh.pktsize) {
+ if (uwsgi_hooked_parse(fr_session->buffer, fr_session->uh.pktsize, fr_get_hostname, (void *) fr_session)) {
+ close_session(fr_session);
+ break;
+ }
+
+ if (fr_session->hostname_len == 0) {
+ close_session(fr_session);
+ break;
+ }
+
+#ifdef UWSGI_DEBUG
+ //uwsgi_log("requested domain %.*s\n", fr_session->hostname_len, fr_session->hostname);
+#endif
+
+choose_node:
+ if (ufr.use_cache) {
+ fr_session->instance_address = uwsgi_cache_get(fr_session->hostname, fr_session->hostname_len, &fr_session->instance_address_len);
+ char *cs_mod = uwsgi_str_contains(fr_session->instance_address, fr_session->instance_address_len, ',');
+ if (cs_mod) {
+ fr_session->modifier1 = uwsgi_str_num(cs_mod + 1, (fr_session->instance_address_len - (cs_mod - fr_session->instance_address)) - 1);
+ fr_session->instance_address_len = (cs_mod - fr_session->instance_address);
+ }
+ }
+ else if (ufr.pattern) {
+ magic_table['s'] = uwsgi_concat2n(fr_session->hostname, fr_session->hostname_len, "", 0);
+ fr_session->tmp_socket_name = magic_sub(ufr.pattern, ufr.pattern_len, &tmp_socket_name_len, magic_table);
+ free(magic_table['s']);
+ fr_session->instance_address_len = tmp_socket_name_len;
+ fr_session->instance_address = fr_session->tmp_socket_name;
+ }
+ else if (ufr.has_subscription_sockets) {
+ fr_session->un = uwsgi_get_subscribe_node(&ufr.subscriptions, fr_session->hostname, fr_session->hostname_len, ufr.subscription_regexp);
+ if (fr_session->un && fr_session->un->len) {
+ fr_session->instance_address = fr_session->un->name;
+ fr_session->instance_address_len = fr_session->un->len;
+ fr_session->modifier1 = fr_session->un->modifier1;
+ }
+ else if (ufr.subscriptions == NULL && ufr.cheap && !ufr.i_am_cheap) {
+ uwsgi_gateway_go_cheap("uWSGI fastrouter", ufr.queue, &ufr.i_am_cheap);
+ }
+ }
+ else if (ufr.base) {
+ fr_session->tmp_socket_name = uwsgi_concat2nn(ufr.base, ufr.base_len, fr_session->hostname, fr_session->hostname_len, &tmp_socket_name_len);
+ fr_session->instance_address_len = tmp_socket_name_len;
+ fr_session->instance_address = fr_session->tmp_socket_name;
+ }
+ else if (ufr.code_string_code && ufr.code_string_function) {
+ if (uwsgi.p[ufr.code_string_modifier1]->code_string) {
+ fr_session->instance_address = uwsgi.p[ufr.code_string_modifier1]->code_string("uwsgi_fastrouter", ufr.code_string_code, ufr.code_string_function, fr_session->hostname, fr_session->hostname_len);
+ if (fr_session->instance_address) {
+ fr_session->instance_address_len = strlen(fr_session->instance_address);
+ char *cs_mod = uwsgi_str_contains(fr_session->instance_address, fr_session->instance_address_len, ',');
+ if (cs_mod) {
+ fr_session->modifier1 = uwsgi_str_num(cs_mod + 1, (fr_session->instance_address_len - (cs_mod - fr_session->instance_address)) - 1);
+ fr_session->instance_address_len = (cs_mod - fr_session->instance_address);
+ }
+ }
+ }
+ }
+ else if (ufr.to_socket) {
+ fr_session->instance_address = ufr.to_socket->name;
+ fr_session->instance_address_len = ufr.to_socket->name_len;
+ }
+ else if (ufr.static_nodes) {
+ if (!ufr.current_static_node) {
+ ufr.current_static_node = ufr.static_nodes;
+ }
+
+ fr_session->static_node = ufr.current_static_node;
+
+ // is it a dead node ?
+ if (fr_session->static_node->custom > 0) {
+
+ // gracetime passed ?
+ if (fr_session->static_node->custom + ufr.static_node_gracetime <= (uint64_t) uwsgi_now()) {
+ fr_session->static_node->custom = 0;
+ }
+ else {
+ struct uwsgi_string_list *tmp_node = fr_session->static_node;
+ struct uwsgi_string_list *next_node = fr_session->static_node->next;
+ fr_session->static_node = NULL;
+ // needed for 1-node only setups
+ if (!next_node) next_node = ufr.static_nodes;
+
+ while(tmp_node != next_node) {
+ if (!next_node) {
+ next_node = ufr.static_nodes;
+ }
+
+ if (tmp_node == next_node) break;
+
+ if (next_node->custom == 0) {
+ fr_session->static_node = next_node;
+ break;
+ }
+ next_node = next_node->next;
+ }
+ }
+ }
+
+ if (fr_session->static_node) {
+
+ fr_session->instance_address = fr_session->static_node->value;
+ fr_session->instance_address_len = fr_session->static_node->len;
+ // set the next one
+ ufr.current_static_node = fr_session->static_node->next;
+ }
+ else {
+ // set the next one
+ ufr.current_static_node = ufr.current_static_node->next;
+ }
+
+ }
+#ifdef UWSGI_SCTP
+ else if (ufr.has_sctp_sockets > 0) {
+
+ struct uwsgi_fr_sctp_node *ufsn = uwsgi_fastrouter_sctp_nodes;
+ int choosen_fd = -1;
+ // find the first available server
+ while(ufsn) {
+ if (ufr.fr_table[ufsn->fd]->status == FASTROUTER_STATUS_SCTP_NODE_FREE) {
+ choosen_fd = ufsn->fd;
+ break;
+ }
+ if (ufsn->next == uwsgi_fastrouter_sctp_nodes) {
+ break;
+ }
+
+ ufsn = ufsn->next;
+ }
+
+ // no nodes available
+ if (choosen_fd == -1) { fr_session->retry = 1; break; }
+
+ struct sctp_sndrcvinfo sinfo;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ memcpy(&sinfo.sinfo_ppid, &fr_session->uh, sizeof(uint32_t));
+ sinfo.sinfo_stream = fr_session->fd;
+ len = sctp_send(choosen_fd, fr_session->buffer, fr_session->uh.pktsize, &sinfo, 0);
+ fr_session->instance_fd = choosen_fd;
+ fr_session->status = FASTROUTER_STATUS_SCTP_RESPONSE;
+ ufr.fr_table[fr_session->instance_fd]->status = FASTROUTER_STATUS_SCTP_RESPONSE;
+ ufr.fr_table[fr_session->instance_fd]->fd = fr_session->fd;
+ break;
+ }
+#endif
+
+ // no address found
+ if (!fr_session->instance_address_len) {
+ // if fallback nodes are configured, trigger them
+ if (ufr.fallback) {
+ fr_session->instance_failed = 1;
+ }
+ close_session(fr_session);
+ break;
+ }
+
+ if (ufr.post_buffering > 0 && fr_session->post_cl > ufr.post_buffering) {
+ fr_session->status = FASTROUTER_STATUS_BUFFERING;
+ fr_session->buf_file_name = uwsgi_tmpname(ufr.pb_base_dir, "uwsgiXXXXX");
+ if (!fr_session->buf_file_name) {
+ uwsgi_error("tempnam()");
+ close_session(fr_session);
+ break;
+ }
+ fr_session->post_remains = fr_session->post_cl;
+
+ // 2 + UWSGI_POSTFILE + 2 + fr_session->buf_file_name
+ if (fr_session->uh.pktsize + (2 + 14 + 2 + strlen(fr_session->buf_file_name)) > 0xffff) {
+ uwsgi_log("unable to buffer request body to file %s: not enough space\n", fr_session->buf_file_name);
+ close_session(fr_session);
+ break;
+ }
+
+ char *ptr = fr_session->buffer + fr_session->uh.pktsize;
+ uint16_t bfn_len = strlen(fr_session->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, fr_session->buf_file_name, bfn_len);
+ fr_session->uh.pktsize += 2 + 14 + 2 + bfn_len;
+
+
+ fr_session->buf_file = fopen(fr_session->buf_file_name, "w");
+ if (!fr_session->buf_file) {
+ uwsgi_error_open(fr_session->buf_file_name);
+ close_session(fr_session);
+ break;
+ }
+
+ }
+
+ else {
+
+ fr_session->pass_fd = is_unix(fr_session->instance_address, fr_session->instance_address_len);
+
+ fr_session->instance_fd = uwsgi_connectn(fr_session->instance_address, fr_session->instance_address_len, 0, 1);
+
+ if (fr_session->instance_fd < 0) {
+ fr_session->instance_failed = 1;
+ fr_session->soopt = errno;
+ close_session(fr_session);
+ break;
+ }
+
+
+ fr_session->status = FASTROUTER_STATUS_CONNECTING;
+ ufr.fr_table[fr_session->instance_fd] = fr_session;
+ event_queue_add_fd_write(ufr.queue, fr_session->instance_fd);
+ }
+ }
+ break;
+
+
+
+ case FASTROUTER_STATUS_CONNECTING:
+
+ if (interesting_fd == fr_session->instance_fd) {
+
+ if (getsockopt(fr_session->instance_fd, SOL_SOCKET, SO_ERROR, (void *) (&fr_session->soopt), &solen) < 0) {
+ uwsgi_error("getsockopt()");
+ fr_session->instance_failed = 1;
+ close_session(fr_session);
+ break;
+ }
+
+ if (fr_session->soopt) {
+ fr_session->instance_failed = 1;
+ close_session(fr_session);
+ break;
+ }
+
+ fr_session->uh.modifier1 = fr_session->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 (fr_session->un)
+ fr_session->un->requests++;
+
+ // fd passing: PERFORMANCE EXTREME BOOST !!!
+ if (fr_session->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), &fr_session->fd, sizeof(int));
+
+ if (sendmsg(fr_session->instance_fd, &msg, 0) < 0) {
+ uwsgi_error("sendmsg()");
+ }
+
+ close_session(fr_session);
+ break;
+ }
+
+ if (writev(fr_session->instance_fd, iov, 2) < 0) {
+ uwsgi_error("writev()");
+ close_session(fr_session);
+ break;
+ }
+
+ event_queue_fd_write_to_read(ufr.queue, fr_session->instance_fd);
+ fr_session->status = FASTROUTER_STATUS_RESPONSE;
+ }
+
+ break;
+#ifdef UWSGI_SCTP
+ case FASTROUTER_STATUS_SCTP_NODE_FREE:
+
+ {
+ struct sctp_sndrcvinfo sinfo;
+ int msg_flags = 0;
+
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ len = sctp_recvmsg(interesting_fd, fr_session->buffer, 0xffff, NULL, NULL, &sinfo, &msg_flags);
+ // remove the SCTP node
+ uwsgi_log("removing SCTP node %d flags = %d len = %d\n", interesting_fd, msg_flags, len);
+ uwsgi_fr_sctp_del_node(interesting_fd);
+ if (ufr.fr_table[interesting_fd]->timeout) {
+ del_timeout(ufr.fr_table[interesting_fd]);
+ }
+ free(ufr.fr_table[interesting_fd]);
+ ufr.fr_table[interesting_fd] = NULL;
+ close(interesting_fd);
+ }
+
+ break;
+ case FASTROUTER_STATUS_SCTP_RESPONSE:
+
+ // data from instance
+ if (interesting_fd == fr_session->instance_fd) {
+ struct sctp_sndrcvinfo sinfo;
+ struct uwsgi_header *uh;
+ int msg_flags =0 ;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ len = sctp_recvmsg(fr_session->instance_fd, fr_session->buffer, 0xffff, NULL, NULL, &sinfo, &msg_flags);
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("recv()");
+ close_session(ufr.fr_table[fr_session->fd]);
+ // REMOVE THE NODE
+ uwsgi_log("removing SCTP node %d flags = %d len = %d\n", interesting_fd, msg_flags, len);
+ uwsgi_fr_sctp_del_node(interesting_fd);
+ if (ufr.fr_table[interesting_fd]->timeout) {
+ del_timeout(ufr.fr_table[interesting_fd]);
+ }
+ free(ufr.fr_table[interesting_fd]);
+ ufr.fr_table[interesting_fd] = NULL;
+ close(interesting_fd);
+ break;
+ }
+
+ if (sinfo.sinfo_stream != fr_session->fd) {
+ uwsgi_log("INVALID STREAM !!!\n");
+ close_session(ufr.fr_table[fr_session->fd]);
+ break;
+ }
+
+ uh = (struct uwsgi_header *) &sinfo.sinfo_ppid ;
+
+ // check for close packet
+ if (uh->modifier1 == 200) {
+ fr_session->status = FASTROUTER_STATUS_SCTP_NODE_FREE;
+ close_session(ufr.fr_table[fr_session->fd]);
+ break;
+ }
+
+ len = send(fr_session->fd, fr_session->buffer, len, 0);
+
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("send()");
+ close_session(ufr.fr_table[fr_session->fd]);
+ break;
+ }
+
+ // update transfer statistics
+ if (fr_session->un)
+ fr_session->un->transferred += len;
+
+ }
+ // body from client
+ else if (interesting_fd == fr_session->fd) {
+
+ uwsgi_log("BODy FROM CLIENT\n");
+
+ //uwsgi_log("receiving body...\n");
+ len = recv(fr_session->fd, fr_session->buffer, 0xffff, 0);
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("recv()");
+ close_session(fr_session);
+ break;
+ }
+
+ struct sctp_sndrcvinfo sinfo;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ // map the stream id to the file descriptor
+ sinfo.sinfo_stream = fr_session->fd;
+
+ len = sctp_send(fr_session->instance_fd, fr_session->buffer, len, &sinfo, 0);
+
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("send()");
+ close_session(fr_session);
+ break;
+ }
+ }
+
+ break;
+#endif
+ case FASTROUTER_STATUS_RESPONSE:
+
+ // data from instance
+ if (interesting_fd == fr_session->instance_fd) {
+ len = recv(fr_session->instance_fd, fr_session->buffer, 0xffff, 0);
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("recv()");
+ close_session(fr_session);
+ break;
+ }
+
+ len = send(fr_session->fd, fr_session->buffer, len, 0);
+
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("send()");
+ close_session(fr_session);
+ break;
+ }
+
+ // update transfer statistics
+ if (fr_session->un)
+ fr_session->un->transferred += len;
+ }
+ // body from client
+ else if (interesting_fd == fr_session->fd) {
+
+ //uwsgi_log("receiving body...\n");
+ len = recv(fr_session->fd, fr_session->buffer, 0xffff, 0);
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("recv()");
+ close_session(fr_session);
+ break;
+ }
+
+
+ len = send(fr_session->instance_fd, fr_session->buffer, len, 0);
+
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("send()");
+ close_session(fr_session);
+ break;
+ }
+ }
+
+ break;
+
+ case FASTROUTER_STATUS_BUFFERING:
+ len = recv(fr_session->fd, post_tmp_buf, UMIN(0xffff, fr_session->post_remains), 0);
+ if (len <= 0) {
+ if (len < 0)
+ uwsgi_error("recv()");
+ close_session(fr_session);
+ break;
+ }
+
+ if (fwrite(post_tmp_buf, len, 1, fr_session->buf_file) != 1) {
+ uwsgi_error("fwrite()");
+ close_session(fr_session);
+ break;
+ }
+
+ fr_session->post_remains -= len;
+
+ if (fr_session->post_remains == 0) {
+ // close the buf_file ASAP
+ fclose(fr_session->buf_file);
+ fr_session->buf_file = NULL;
+
+ fr_session->pass_fd = is_unix(fr_session->instance_address, fr_session->instance_address_len);
+
+ fr_session->instance_fd = uwsgi_connectn(fr_session->instance_address, fr_session->instance_address_len, 0, 1);
+
+ if (fr_session->instance_fd < 0) {
+ fr_session->instance_failed = 1;
+ close_session(fr_session);
+ break;
+ }
+
+ fr_session->status = FASTROUTER_STATUS_CONNECTING;
+ ufr.fr_table[fr_session->instance_fd] = fr_session;
+ event_queue_add_fd_write(ufr.queue, fr_session->instance_fd);
+ }
+ break;
+
+
+
+
+ // fallback to destroy !!!
+ default:
+ uwsgi_log("unknown event: closing session\n");
+ close_session(fr_session);
+ break;
+
+ }
+}
diff --git a/plugins/fastrouter/fr_sctp.c b/plugins/fastrouter/fr_sctp.c
index ae8dd1b9..0d80453d 100644
--- a/plugins/fastrouter/fr_sctp.c
+++ b/plugins/fastrouter/fr_sctp.c
@@ -40,6 +40,37 @@ struct uwsgi_fr_sctp_node *uwsgi_fr_sctp_add_node(int fd) {
}
+void uwsgi_fr_sctp_del_node(int fd) {
+
+ struct uwsgi_fr_sctp_node *ufsn = uwsgi_fastrouter_sctp_nodes;
+ while(ufsn) {
+
+ if (ufsn->fd == fd) {
+ struct uwsgi_fr_sctp_node *prev = ufsn->prev;
+ struct uwsgi_fr_sctp_node *next = ufsn->next;
+
+ prev->next = next;
+ next->prev = prev;
+
+ // check: am i the only node ?
+ if ( ufsn == prev || ufsn == next ) {
+ free(uwsgi_fastrouter_sctp_nodes);
+ uwsgi_fastrouter_sctp_nodes = NULL;
+ break;
+ }
+
+ free(ufsn);
+ break;
+ }
+
+ if (ufsn->next == uwsgi_fastrouter_sctp_nodes) {
+ break;
+ }
+
+ ufsn = ufsn->next;
+ }
+}
+
void uwsgi_opt_fastrouter_sctp(char *opt, char *value, void *foobar) {
struct uwsgi_gateway_socket *ugs = uwsgi_new_gateway_socket(value, "uWSGI fastrouter");
diff --git a/plugins/fastrouter/uwsgiplugin.py b/plugins/fastrouter/uwsgiplugin.py
index a2e91713..903c6f27 100644
--- a/plugins/fastrouter/uwsgiplugin.py
+++ b/plugins/fastrouter/uwsgiplugin.py
@@ -4,4 +4,4 @@ CFLAGS = []
LDFLAGS = []
LIBS = []
-GCC_LIST = ['fr_sctp', 'fastrouter']
+GCC_LIST = ['fr_sctp', 'fastrouter', 'fr_events']
diff --git a/proto/sctp.c b/proto/sctp.c
index eacee320..6b9793ae 100644
--- a/proto/sctp.c
+++ b/proto/sctp.c
@@ -4,10 +4,45 @@
extern struct uwsgi_server uwsgi;
+#ifdef __linux__
+ssize_t sctp_sendv(int s, struct iovec *iov, size_t iov_len,
+ const struct sctp_sndrcvinfo *sinfo, int flags)
+{
+ struct msghdr outmsg;
+
+ outmsg.msg_name = NULL;
+ outmsg.msg_namelen = 0;
+ outmsg.msg_iov = iov;
+ outmsg.msg_iovlen = iov_len;
+ outmsg.msg_controllen = 0;
+
+ if (sinfo) {
+ char outcmsg[CMSG_SPACE(sizeof(struct sctp_sndrcvinfo))];
+ struct cmsghdr *cmsg;
+
+ outmsg.msg_control = outcmsg;
+ outmsg.msg_controllen = sizeof(outcmsg);
+ outmsg.msg_flags = 0;
+
+ cmsg = CMSG_FIRSTHDR(&outmsg);
+ cmsg->cmsg_level = IPPROTO_SCTP;
+ cmsg->cmsg_type = SCTP_SNDRCV;
+ cmsg->cmsg_len = CMSG_LEN(sizeof(struct sctp_sndrcvinfo));
+
+ outmsg.msg_controllen = cmsg->cmsg_len;
+ memcpy(CMSG_DATA(cmsg), sinfo, sizeof(struct sctp_sndrcvinfo));
+ }
+
+ return sendmsg(s, &outmsg, flags);
+}
+
+#endif
+
int uwsgi_proto_sctp_parser(struct wsgi_request *wsgi_req) {
struct sctp_sndrcvinfo sinfo;
- int msg_flags;
+ memset(&sinfo, 0, sizeof(sinfo));
+ int msg_flags = 0;
ssize_t len = sctp_recvmsg(wsgi_req->socket->fd, wsgi_req->buffer, uwsgi.buffer_size, NULL, NULL, &sinfo, &msg_flags);
@@ -17,24 +52,31 @@ int uwsgi_proto_sctp_parser(struct wsgi_request *wsgi_req) {
// connection lost, retrigger it
close(wsgi_req->socket->fd);
wsgi_req->socket->fd = connect_to_sctp(wsgi_req->socket->name, wsgi_req->socket->queue);
+ // avoid closing connection
+ wsgi_req->fd_closed = 1;
}
return -1;
}
else if (len == 0) {
- uwsgi_log("lost connection with the SCTP server\n");
+ uwsgi_log("lost connection with the SCTP server %d\n", msg_flags);
// connection lost, retrigger it
close(wsgi_req->socket->fd);
wsgi_req->socket->fd = connect_to_sctp(wsgi_req->socket->name, wsgi_req->socket->queue);
+ // avoid closing connection
+ wsgi_req->fd_closed = 1;
return -2;
}
- // check for a request stream
- if (sinfo.sinfo_stream != 0) {
- uwsgi_log("invalid SCTP stream id (must be 0)\n");
+ // get the uwsgi 4 bytes header from ppid
+ memcpy(&wsgi_req->uh, &sinfo.sinfo_ppid, sizeof(uint32_t));
+
+ // check for invalid modifiers
+ if (wsgi_req->uh.modifier1 == 199 || wsgi_req->uh.modifier1 == 200) {
+ uwsgi_log("invalid SCTP uwsgi modifier1: %d\n", wsgi_req->uh.modifier1);
return -1;
}
- memcpy(&wsgi_req->uh, &sinfo.sinfo_ppid, sizeof(uint32_t));
+ wsgi_req->stream_id = sinfo.sinfo_stream;
/* big endian ? */
#ifdef __BIG_ENDIAN__
@@ -56,25 +98,37 @@ int uwsgi_proto_sctp_parser(struct wsgi_request *wsgi_req) {
}
ssize_t uwsgi_proto_sctp_writev_header(struct wsgi_request * wsgi_req, struct iovec * iovec, size_t iov_len) {
- ssize_t wlen = writev(wsgi_req->poll.fd, iovec, iov_len);
- if (wlen < 0) {
- uwsgi_req_error("writev()");
- return 0;
- }
- return wlen;
+ struct sctp_sndrcvinfo sinfo;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ sinfo.sinfo_stream = wsgi_req->stream_id;
+
+ ssize_t wlen = sctp_sendv(wsgi_req->poll.fd, iovec, iov_len, &sinfo, 0);
+ if (wlen < 0) {
+ uwsgi_req_error("writev()");
+ return 0;
+ }
+ return wlen;
}
ssize_t uwsgi_proto_sctp_writev(struct wsgi_request * wsgi_req, struct iovec * iovec, size_t iov_len) {
- ssize_t wlen = writev(wsgi_req->poll.fd, iovec, iov_len);
- if (wlen < 0) {
- uwsgi_req_error("writev()");
- return 0;
- }
- return wlen;
+ struct sctp_sndrcvinfo sinfo;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ sinfo.sinfo_stream = wsgi_req->stream_id;
+
+ ssize_t wlen = sctp_sendv(wsgi_req->poll.fd, iovec, iov_len, &sinfo, 0);
+ if (wlen < 0) {
+ uwsgi_req_error("writev()");
+ return 0;
+ }
+ return wlen;
}
ssize_t uwsgi_proto_sctp_write(struct wsgi_request * wsgi_req, char *buf, size_t len) {
- ssize_t wlen = write(wsgi_req->poll.fd, buf, len);
+ struct sctp_sndrcvinfo sinfo;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ sinfo.sinfo_stream = wsgi_req->stream_id;
+
+ ssize_t wlen = sctp_send(wsgi_req->poll.fd, buf, len, &sinfo, 0);
if (wlen < 0) {
uwsgi_req_error("write()");
return 0;
@@ -83,7 +137,11 @@ ssize_t uwsgi_proto_sctp_write(struct wsgi_request * wsgi_req, char *buf, size_t
}
ssize_t uwsgi_proto_sctp_write_header(struct wsgi_request * wsgi_req, char *buf, size_t len) {
- ssize_t wlen = write(wsgi_req->poll.fd, buf, len);
+ struct sctp_sndrcvinfo sinfo;
+ memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
+ sinfo.sinfo_stream = wsgi_req->stream_id;
+
+ ssize_t wlen = sctp_send(wsgi_req->poll.fd, buf, len, &sinfo, 0);
if (wlen < 0) {
uwsgi_req_error("write()");
return 0;
@@ -98,15 +156,57 @@ int uwsgi_proto_sctp_accept(struct wsgi_request *wsgi_req, int fd) {
void uwsgi_proto_sctp_close(struct wsgi_request *wsgi_req) {
+ // this function could be called in uwsgi_destroy_request too
+ if (wsgi_req->fd_closed) return;
+
+ struct uwsgi_header uh;
+ uh.modifier1 = 200;
+ uh.pktsize = 0;
+ uh.modifier2 = 0;
+
struct sctp_sndrcvinfo sinfo;
memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo));
- // stream 2 is used for closing requests
- sinfo.sinfo_stream = 2;
- memcpy(&sinfo.sinfo_ppid, &wsgi_req->uh, sizeof(uint32_t));
+ // ppid->modifier1 200 is used for closing requests
+ memcpy(&sinfo.sinfo_ppid, &uh, sizeof(uint32_t));
+ sinfo.sinfo_stream = wsgi_req->stream_id;
if (wsgi_req->async_post) {
fclose(wsgi_req->async_post);
}
- sctp_send(wsgi_req->poll.fd, &wsgi_req->uh , sizeof(uint32_t), &sinfo, 0);
+
+ sctp_send(wsgi_req->poll.fd, &uh, sizeof(uh), &sinfo, 0);
+}
+
+ssize_t uwsgi_proto_sctp_sendfile(struct wsgi_request * wsgi_req) {
+
+ ssize_t len;
+ char buf[65536];
+ size_t remains = wsgi_req->sendfile_fd_size - wsgi_req->sendfile_fd_pos;
+
+ wsgi_req->sendfile_fd_chunk = 65536;
+
+ if (uwsgi.async > 1) {
+ len = read(wsgi_req->sendfile_fd, buf, UMIN(remains, wsgi_req->sendfile_fd_chunk));
+ if (len != (int) UMIN(remains, wsgi_req->sendfile_fd_chunk)) {
+ uwsgi_error("read()");
+ return -1;
+ }
+ wsgi_req->sendfile_fd_pos += len;
+ return uwsgi_proto_sctp_write(wsgi_req, buf, len);
+ }
+
+ while (remains) {
+ len = read(wsgi_req->sendfile_fd, buf, UMIN(remains, wsgi_req->sendfile_fd_chunk));
+ if (len != (int) UMIN(remains, wsgi_req->sendfile_fd_chunk)) {
+ uwsgi_error("read()");
+ return -1;
+ }
+ wsgi_req->sendfile_fd_pos += len;
+ len = uwsgi_proto_sctp_write(wsgi_req, buf, len);
+ remains = wsgi_req->sendfile_fd_size - wsgi_req->sendfile_fd_pos;
+ }
+
+ return wsgi_req->sendfile_fd_pos;
+
}
diff --git a/socket.c b/socket.c
index 465d0dd8..cf6815ca 100644
--- a/socket.c
+++ b/socket.c
@@ -184,30 +184,40 @@ int connect_to_sctp(char *socket_names, int queue) {
struct sctp_initmsg initmsg;
memset(&initmsg, 0, sizeof(initmsg));
- initmsg.sinit_max_instreams = 17;
- initmsg.sinit_num_ostreams = 17;
+ initmsg.sinit_max_instreams = 0xffff;
+ initmsg.sinit_num_ostreams = 0xffff;
if (setsockopt(serverfd, IPPROTO_SCTP,
SCTP_INITMSG, &initmsg, sizeof(initmsg))) {
uwsgi_error("setsockopt()");
+ close(serverfd);
goto clear;
}
memset( (void *)&events, 0, sizeof(events) );
events.sctp_data_io_event = 1;
-/*
+ /*
events.sctp_peer_error_event = 1;
events.sctp_shutdown_event = 1;
-*/
+ */
if (setsockopt( serverfd, SOL_SCTP, SCTP_EVENTS,
(const void *)&events, sizeof(events) )) {
uwsgi_error("setsockopt()");
+ close(serverfd);
+ goto clear;
+ }
+
+ int sctp_nodelay = 1;
+ if (setsockopt( serverfd, SOL_SCTP, SCTP_NODELAY, &sctp_nodelay, sizeof(sctp_nodelay))) {
+ uwsgi_error("setsockopt()");
+ close(serverfd);
goto clear;
}
if (sctp_connectx(serverfd, (struct sockaddr *) sins, addresses, NULL)) {
uwsgi_error("sctp_connectx()");
+ close(serverfd);
goto clear;
}
@@ -279,8 +289,8 @@ int bind_to_sctp(char *socket_names) {
struct sctp_initmsg initmsg;
memset(&initmsg, 0, sizeof(initmsg));
- initmsg.sinit_max_instreams = 17;
- initmsg.sinit_num_ostreams = 17;
+ initmsg.sinit_max_instreams = 0xffff;
+ initmsg.sinit_num_ostreams = 0xffff;
if (setsockopt(serverfd, IPPROTO_SCTP,
SCTP_INITMSG, &initmsg, sizeof(initmsg))) {
@@ -290,10 +300,10 @@ int bind_to_sctp(char *socket_names) {
memset( (void *)&events, 0, sizeof(events) );
events.sctp_data_io_event = 1;
-/*
+ /*
events.sctp_peer_error_event = 1;
events.sctp_shutdown_event = 1;
-*/
+ */
if (setsockopt( serverfd, SOL_SCTP, SCTP_EVENTS,
(const void *)&events, sizeof(events) )) {
diff --git a/uwsgi.c b/uwsgi.c
index 9107c3ce..f147dd72 100644
--- a/uwsgi.c
+++ b/uwsgi.c
@@ -2302,7 +2302,7 @@ skipzero:
uwsgi_sock->proto_writev = uwsgi_proto_sctp_writev;
uwsgi_sock->proto_write_header = uwsgi_proto_sctp_write_header;
uwsgi_sock->proto_writev_header = uwsgi_proto_sctp_writev_header;
- uwsgi_sock->proto_sendfile = NULL;
+ uwsgi_sock->proto_sendfile = uwsgi_proto_sctp_sendfile;
uwsgi_sock->proto_close = uwsgi_proto_sctp_close;
}
#endif
diff --git a/uwsgi.h b/uwsgi.h
index 9dede7c6..7f12a819 100644
--- a/uwsgi.h
+++ b/uwsgi.h
@@ -934,6 +934,8 @@ struct wsgi_request {
int sigwait;
int signal_received;
+ uint16_t stream_id;
+
struct msghdr msg;
union {
struct cmsghdr cmsg;
@@ -2141,6 +2143,8 @@ int uwsgi_parse_array(char *, uint16_t, char **, uint16_t *, uint8_t *);
struct uwsgi_gateway *register_gateway(char *, void (*)(int));
void gateway_respawn(int);
+void uwsgi_gateway_go_cheap(char *, int, int *);
+
char *uwsgi_open_and_read(char *, int *, int, char *[]);
char *uwsgi_get_last_char(char *, char);
@@ -2341,6 +2345,7 @@ ssize_t uwsgi_proto_sctp_write(struct wsgi_request *, char *, size_t);
ssize_t uwsgi_proto_sctp_write_header(struct wsgi_request *, char *, size_t);
int uwsgi_proto_sctp_accept(struct wsgi_request *, int);
void uwsgi_proto_sctp_close(struct wsgi_request *);
+ssize_t uwsgi_proto_sctp_sendfile(struct wsgi_request *);
#endif
int uwsgi_proto_http_parser(struct wsgi_request *);