From 4c854a83c95a7c2ca1b06e78fc9366913ccaa034 Mon Sep 17 00:00:00 2001 From: "roberto@centos6" Date: Wed, 7 Mar 2012 01:40:57 +0100 Subject: [PATCH] SCTP improvements, rack optimizations --- buildconf/default.ini | 2 +- plugins/fastrouter/fastrouter.c | 1 + plugins/fastrouter/fr_events.c | 1025 ++++++++++++++++--------------- plugins/rack/rack_plugin.c | 29 +- proto/sctp.c | 47 +- utils.c | 27 +- uwsgi.c | 4 +- uwsgi.h | 2 + 8 files changed, 592 insertions(+), 545 deletions(-) diff --git a/buildconf/default.ini b/buildconf/default.ini index 43a040c5..229fbb40 100755 --- a/buildconf/default.ini +++ b/buildconf/default.ini @@ -6,7 +6,7 @@ json = auto sqlite3 = auto zeromq = auto snmp = true -sctp = true +sctp = false spooler = true embedded = true udp = true diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index d059c73d..88d749b4 100755 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -566,6 +566,7 @@ void fastrouter_loop(int id) { ufr.fr_table[new_connection] = alloc_fr_session(); ufr.fr_table[new_connection]->instance_fd = new_connection; + ufr.fr_table[new_connection]->fd = -1; ufr.fr_table[new_connection]->persistent = 1; ufr.fr_table[new_connection]->status = FASTROUTER_STATUS_SCTP_NODE_FREE; diff --git a/plugins/fastrouter/fr_events.c b/plugins/fastrouter/fr_events.c index fd3ec507..05e117c7 100755 --- a/plugins/fastrouter/fr_events.c +++ b/plugins/fastrouter/fr_events.c @@ -16,548 +16,563 @@ void uwsgi_fastrouter_switch_events(struct fastrouter_session *fr_session, int i struct iovec iov[2]; struct msghdr msg; - union { - struct cmsghdr cmsg; - char control[CMSG_SPACE(sizeof(int))]; - } msg_control; - struct cmsghdr *cmsg; + union { + struct cmsghdr cmsg; + char control[CMSG_SPACE(sizeof(int))]; + } msg_control; + struct cmsghdr *cmsg; ssize_t len; - char *post_tmp_buf[0xffff]; + char *post_tmp_buf[0xffff]; int tmp_socket_name_len; + switch (fr_session->status) { - 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) { + 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); + 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; + fr_session->status = FASTROUTER_STATUS_RECV_VARS; + } + break; - case FASTROUTER_STATUS_RECV_VARS: + case FASTROUTER_STATUS_RECV_VARS: - if (interesting_fd == -1) goto choose_node; + 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; - } + 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; - } + 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); + //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); - } + 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.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.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; } - 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; + if (fr_session->static_node) { - // is it a dead node ? - if (fr_session->static_node->custom > 0) { + 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; + } - // 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) { + else if (ufr.has_sctp_sockets > 0) { - if (!*uwsgi_fastrouter_sctp_nodes_current) - *uwsgi_fastrouter_sctp_nodes_current = *uwsgi_fastrouter_sctp_nodes; + if (!*uwsgi_fastrouter_sctp_nodes_current) + *uwsgi_fastrouter_sctp_nodes_current = *uwsgi_fastrouter_sctp_nodes; - struct uwsgi_fr_sctp_node *ufsn = *uwsgi_fastrouter_sctp_nodes_current; - 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_current) { - break; - } - - ufsn = ufsn->next; - } - - // no nodes available - if (choosen_fd == -1) { - fr_session->retry = 1; - del_timeout(fr_session); - fr_session->timeout = add_fake_timeout(fr_session); - 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; - - // round robin - *uwsgi_fastrouter_sctp_nodes_current = (*uwsgi_fastrouter_sctp_nodes_current)->next; - 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); - } + struct uwsgi_fr_sctp_node *ufsn = *uwsgi_fastrouter_sctp_nodes_current; + 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; } - 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("[0] 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("[1] 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); - uwsgi_log("DONE\n"); - break; - } - - if (sinfo.sinfo_stream != fr_session->fd) { - uwsgi_log("INVALID SCTP 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); + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes_current) { break; } - if (fwrite(post_tmp_buf, len, 1, fr_session->buf_file) != 1) { - uwsgi_error("fwrite()"); - close_session(fr_session); - break; - } + ufsn = ufsn->next; + } - 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); - } + // no nodes available + if (choosen_fd == -1) { + fr_session->retry = 1; + del_timeout(fr_session); + fr_session->timeout = add_fake_timeout(fr_session); 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; + // round robin + *uwsgi_fastrouter_sctp_nodes_current = (*uwsgi_fastrouter_sctp_nodes_current)->next; + break; + } +#endif - // fallback to destroy !!! - default: - uwsgi_log("unknown event: closing session\n"); + // 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("[0] 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("[1] 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 (fr_session->fd != -1 && sinfo.sinfo_stream != fr_session->fd) { + if (fr_session->fd != -1) { + uwsgi_log("INVALID SCTP 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; + if (fr_session->fd != -1) { + close_session(ufr.fr_table[fr_session->fd]); + } + break; + } + + if (fr_session->fd == -1) { + 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) { + + len = recv(fr_session->fd, fr_session->buffer, 0xffff, 0); + if (len <= 0) { + if (len < 0) + uwsgi_error("recv()"); + // mark session as broken + ufr.fr_table[fr_session->instance_fd]->fd = -1; + close_session(fr_session); + break; + } + + struct sctp_sndrcvinfo sinfo; + memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo)); + // map the stream id to the file descriptor + struct uwsgi_header uh; + uh.modifier1 = 199; + uh.pktsize = 0; + uh.modifier2 = 0; + memcpy(&sinfo.sinfo_ppid, &uh, sizeof(uint32_t)); + 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/rack/rack_plugin.c b/plugins/rack/rack_plugin.c index b63d7551..909ef27a 100755 --- a/plugins/rack/rack_plugin.c +++ b/plugins/rack/rack_plugin.c @@ -546,9 +546,7 @@ VALUE send_header(VALUE obj, VALUE headers) { struct wsgi_request *wsgi_req = current_wsgi_req(); - size_t len; VALUE hkey, hval; - //uwsgi_log("HEADERS %d\n", TYPE(obj)); if (TYPE(obj) == T_ARRAY) { @@ -581,18 +579,16 @@ VALUE send_header(VALUE obj, VALUE headers) { size_t header_value_len = RSTRING_LEN(hval); size_t i,cnt=0; char *this_header = header_value; + struct iovec iov[4]; for(i=0;isocket->proto_write_header( wsgi_req, RSTRING_PTR(hkey), RSTRING_LEN(hkey)); - wsgi_req->headers_size += len; - len = wsgi_req->socket->proto_write_header( wsgi_req, (char *)": ", 2); - wsgi_req->headers_size += len; - len = wsgi_req->socket->proto_write_header( wsgi_req, this_header, cnt); - wsgi_req->headers_size += len; - len = wsgi_req->socket->proto_write_header( wsgi_req, (char *)"\r\n", 2); - wsgi_req->headers_size += len; + iov[0].iov_base = RSTRING_PTR(hkey); iov[0].iov_len = RSTRING_LEN(hkey); + iov[1].iov_base = (char *)": "; iov[1].iov_len = 2; + iov[2].iov_base = this_header; iov[2].iov_len = cnt; + iov[3].iov_base = (char *)"\r\n"; iov[3].iov_len = 2; + wsgi_req->headers_size += wsgi_req->socket->proto_writev_header( wsgi_req, iov, 4); //uwsgi_log("(multi) --%.*s: %.*s--\n", RSTRING_LEN(hkey), RSTRING_PTR(hkey), cnt, this_header); @@ -606,14 +602,11 @@ VALUE send_header(VALUE obj, VALUE headers) { } if (cnt > 0) { - len = wsgi_req->socket->proto_write_header( wsgi_req, RSTRING_PTR(hkey), RSTRING_LEN(hkey)); - wsgi_req->headers_size += len; - len = wsgi_req->socket->proto_write_header( wsgi_req, (char *)": ", 2); - wsgi_req->headers_size += len; - len = wsgi_req->socket->proto_write_header( wsgi_req, this_header, cnt); - wsgi_req->headers_size += len; - len = wsgi_req->socket->proto_write_header( wsgi_req, (char *)"\r\n", 2); - wsgi_req->headers_size += len; + iov[0].iov_base = RSTRING_PTR(hkey); iov[0].iov_len = RSTRING_LEN(hkey); + iov[1].iov_base = (char *)": "; iov[1].iov_len = 2; + iov[2].iov_base = this_header; iov[2].iov_len = cnt; + iov[3].iov_base = (char *)"\r\n"; iov[3].iov_len = 2; + wsgi_req->headers_size += wsgi_req->socket->proto_writev_header( wsgi_req, iov, 4); wsgi_req->header_cnt++; //uwsgi_log("--%.*s: %.*s--\n", RSTRING_LEN(hkey), RSTRING_PTR(hkey), cnt, this_header); } diff --git a/proto/sctp.c b/proto/sctp.c index 5a4a835d..79ecec7b 100755 --- a/proto/sctp.c +++ b/proto/sctp.c @@ -46,25 +46,17 @@ int uwsgi_proto_sctp_parser(struct wsgi_request *wsgi_req) { ssize_t len = sctp_recvmsg(wsgi_req->socket->fd, wsgi_req->buffer, uwsgi.buffer_size, NULL, NULL, &sinfo, &msg_flags); - if (len < 0) { - uwsgi_error("sctp_recvmsg()"); - if (msg_flags == 0) { - // 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) { + if (len <= 0) { + if (len < 0) + uwsgi_error("sctp_recvmsg()"); 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; + // no special message needed + return -3; } // get the uwsgi 4 bytes header from ppid @@ -160,13 +152,13 @@ void uwsgi_proto_sctp_close(struct wsgi_request *wsgi_req) { if (wsgi_req->fd_closed) return; struct uwsgi_header uh; + // ppid->modifier1 200 is used for closing requests uh.modifier1 = 200; uh.pktsize = 0; uh.modifier2 = 0; struct sctp_sndrcvinfo sinfo; memset(&sinfo, 0, sizeof(struct sctp_sndrcvinfo)); - // ppid->modifier1 200 is used for closing requests memcpy(&sinfo.sinfo_ppid, &uh, sizeof(uint32_t)); sinfo.sinfo_stream = wsgi_req->stream_id; @@ -210,3 +202,30 @@ ssize_t uwsgi_proto_sctp_sendfile(struct wsgi_request * wsgi_req) { } +ssize_t uwsgi_proto_sctp_read_body(struct wsgi_request * wsgi_req, char *buf, size_t len) { + + struct sctp_sndrcvinfo sinfo; + memset(&sinfo, 0, sizeof(sinfo)); + int msg_flags = 0; + struct uwsgi_header *uh; + + ssize_t slen = sctp_recvmsg(wsgi_req->socket->fd, buf, len, NULL, NULL, &sinfo, &msg_flags); + + if (slen <= 0) { + return -1; + } + + if (wsgi_req->stream_id != sinfo.sinfo_stream) { + return -1; + } + + uh = (struct uwsgi_header *) &sinfo.sinfo_ppid; + + if (uh->modifier1 != 199) { + return -1; + } + + return slen; + +} + diff --git a/utils.c b/utils.c index 99bdf309..b23dc49b 100755 --- a/utils.c +++ b/utils.c @@ -1558,7 +1558,12 @@ int uwsgi_read_whole_body_in_mem(struct wsgi_request *wsgi_req, char *buf) { return 0; } - len = read(wsgi_req->poll.fd, ptr, post_remains); + if (wsgi_req->socket->proto_read_body) { + len = wsgi_req->socket->proto_read_body(wsgi_req, ptr, post_remains); + } + else { + len = read(wsgi_req->poll.fd, ptr, post_remains); + } if (len <= 0) { uwsgi_error("read()"); @@ -1582,7 +1587,6 @@ int uwsgi_read_whole_body(struct wsgi_request *wsgi_req, char *buf, size_t len) const char *x_progress_id = "X-Progress-ID="; char *xpi_ptr = (char *) x_progress_id; - wsgi_req->async_post = tmpfile(); if (!wsgi_req->async_post) { uwsgi_error("tmpfile()"); @@ -1668,10 +1672,20 @@ int uwsgi_read_whole_body(struct wsgi_request *wsgi_req, char *buf, size_t len) } if (post_remains > len) { - post_chunk = read(wsgi_req->poll.fd, buf, len); + if (wsgi_req->socket->proto_read_body) { + post_chunk = wsgi_req->socket->proto_read_body(wsgi_req, buf, len); + } + else { + post_chunk = read(wsgi_req->poll.fd, buf, len); + } } else { - post_chunk = read(wsgi_req->poll.fd, buf, post_remains); + if (wsgi_req->socket->proto_read_body) { + post_chunk = wsgi_req->socket->proto_read_body(wsgi_req, buf, len); + } + else { + post_chunk = read(wsgi_req->poll.fd, buf, post_remains); + } } if (post_chunk < 0) { @@ -4034,6 +4048,7 @@ int uwsgi_file_to_string_list(char *filename, struct uwsgi_string_list **list) { } void uwsgi_setup_post_buffering(void) { + int i; uwsgi.async_post_buf = uwsgi_malloc(sizeof(char *) * uwsgi.cores); if (!uwsgi.post_buffering_bufsize) uwsgi.post_buffering_bufsize = 8192; @@ -4042,6 +4057,10 @@ void uwsgi_setup_post_buffering(void) { uwsgi_log("setting request body buffering size to %d bytes\n", uwsgi.post_buffering_bufsize); } + for(i=0;i 0) { - uwsgi.async_post_buf[i] = uwsgi_malloc(uwsgi.post_buffering_bufsize); - } } #ifdef UWSGI_DEBUG @@ -2304,6 +2301,7 @@ skipzero: uwsgi_sock->proto_writev_header = uwsgi_proto_sctp_writev_header; uwsgi_sock->proto_sendfile = uwsgi_proto_sctp_sendfile; uwsgi_sock->proto_close = uwsgi_proto_sctp_close; + uwsgi_sock->proto_read_body = uwsgi_proto_sctp_read_body; } #endif else { diff --git a/uwsgi.h b/uwsgi.h index a31d8d81..6d61cf45 100755 --- a/uwsgi.h +++ b/uwsgi.h @@ -579,6 +579,7 @@ struct uwsgi_socket { ssize_t(*proto_write_header) (struct wsgi_request *, char *, size_t); ssize_t(*proto_writev_header) (struct wsgi_request *, struct iovec *, size_t); ssize_t(*proto_sendfile) (struct wsgi_request *); + ssize_t(*proto_read_body) (struct wsgi_request *, char *, size_t); void (*proto_close) (struct wsgi_request *); int edge_trigger; @@ -2349,6 +2350,7 @@ 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 *); +ssize_t uwsgi_proto_sctp_read_body(struct wsgi_request *, char *, size_t); #endif int uwsgi_proto_http_parser(struct wsgi_request *);