little refactoring for the fastrouter

This commit is contained in:
roberto@precise64
2012-03-04 16:10:35 +01:00
parent 205605a3c4
commit da9ce52294
12 changed files with 845 additions and 625 deletions
+13
View File
@@ -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;
}
+1
View File
@@ -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 "<h1>Hello World</h1>"
+2 -15
View File
@@ -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);
}
}
}
+28 -576
View File
@@ -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;
}
}
}
}
+74
View File
@@ -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 *);
+547
View File
@@ -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;
}
}
+31
View File
@@ -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");
+1 -1
View File
@@ -4,4 +4,4 @@ CFLAGS = []
LDFLAGS = []
LIBS = []
GCC_LIST = ['fr_sctp', 'fastrouter']
GCC_LIST = ['fr_sctp', 'fastrouter', 'fr_events']
+124 -24
View File
@@ -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;
}
+18 -8
View File
@@ -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) )) {
+1 -1
View File
@@ -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
+5
View File
@@ -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 *);