From 629fbf4f07f7562ee03cd32c53a5bf08ce772399 Mon Sep 17 00:00:00 2001 From: "roberto@precise64" Date: Mon, 5 Mar 2012 04:29:11 +0100 Subject: [PATCH] another series of sctp patch --- plugins/fastrouter/fastrouter.c | 38 ++++---------- plugins/fastrouter/fr_events.c | 30 ++++++++--- plugins/fastrouter/fr_sctp.c | 92 +++++++++++++++++++++++++-------- 3 files changed, 102 insertions(+), 58 deletions(-) diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index 19a1e8d7..d059c73d 100644 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -18,7 +18,8 @@ extern struct uwsgi_server uwsgi; #ifdef UWSGI_SCTP -extern struct uwsgi_fr_sctp_node *uwsgi_fastrouter_sctp_nodes; +extern struct uwsgi_fr_sctp_node **uwsgi_fastrouter_sctp_nodes; +extern struct uwsgi_fr_sctp_node **uwsgi_fastrouter_sctp_nodes_current; #endif void fastrouter_send_stats(int); @@ -148,7 +149,7 @@ struct uwsgi_option fastrouter_options[] = { {"fastrouter-subscription-use-regexp", no_argument, 0, "enable regexp for subscription system", uwsgi_opt_true, &ufr.subscription_regexp, 0}, #ifdef UWSGI_SCTP - {"fastrouter-sctp", required_argument, 0, "run the fastrouter SCTP server on the spcified address", uwsgi_opt_fastrouter_sctp, NULL, 0}, + {"fastrouter-sctp", required_argument, 0, "run the fastrouter SCTP server on the specified address", uwsgi_opt_fastrouter_sctp, NULL, 0}, #endif {"fastrouter-timeout", required_argument, 0, "set fastrouter timeout", uwsgi_opt_set_int, &ufr.socket_timeout, 0}, @@ -479,6 +480,11 @@ void fastrouter_loop(int id) { init_magic_table(ufr.magic_table); } +#ifdef UWSGI_SCTP + uwsgi_fastrouter_sctp_nodes = uwsgi_calloc(sizeof(struct uwsgi_fastrouter_sctp_nodes*)); + uwsgi_fastrouter_sctp_nodes_current = uwsgi_calloc(sizeof(struct uwsgi_fastrouter_sctp_nodes*)); +#endif + struct sockaddr_un fr_addr; socklen_t fr_addr_len = sizeof(struct sockaddr_un); @@ -537,26 +543,8 @@ void fastrouter_loop(int id) { ufr.fr_table[new_connection] = alloc_fr_session(); ufr.fr_table[new_connection]->fd = new_connection; - //ufr.fr_table[new_connection]->modifier1 = 0; ufr.fr_table[new_connection]->instance_fd = -1; ufr.fr_table[new_connection]->status = FASTROUTER_STATUS_RECV_HDR; - /* - ufr.fr_table[new_connection]->h_pos = 0; - ufr.fr_table[new_connection]->pos = 0; - ufr.fr_table[new_connection]->un = NULL; - ufr.fr_table[new_connection]->static_node = NULL; - ufr.fr_table[new_connection]->buf_file = NULL; - ufr.fr_table[new_connection]->buf_file_name = NULL; - ufr.fr_table[new_connection]->instance_failed = 0; - ufr.fr_table[new_connection]->instance_address_len = 0; - ufr.fr_table[new_connection]->hostname_len = 0; - ufr.fr_table[new_connection]->hostname = NULL; - ufr.fr_table[new_connection]->fallback = NULL; - ufr.fr_table[new_connection]->soopt = 0; - ufr.fr_table[new_connection]->timed_out = 0; - ufr.fr_table[new_connection]->do_not_close = 0; - ufr.fr_table[new_connection]->tmp_socket_name = NULL; - */ ufr.fr_table[new_connection]->timeout = add_timeout(ufr.fr_table[new_connection]); @@ -574,15 +562,7 @@ void fastrouter_loop(int id) { break; } uwsgi_fr_sctp_add_node(new_connection); - uwsgi_log("new SCTP peer:\n"); - struct uwsgi_fr_sctp_node *ufsn = uwsgi_fastrouter_sctp_nodes; - while(ufsn) { - uwsgi_log("\tfd = %d\n", ufsn->fd); - if (ufsn->next == uwsgi_fastrouter_sctp_nodes) { - break; - } - ufsn = ufsn->next; - } + uwsgi_log("new SCTP peer: %s:%d\n", inet_ntoa(((struct sockaddr_in *)&fr_addr)->sin_addr), ntohs(((struct sockaddr_in *) &fr_addr)->sin_port)); ufr.fr_table[new_connection] = alloc_fr_session(); ufr.fr_table[new_connection]->instance_fd = new_connection; diff --git a/plugins/fastrouter/fr_events.c b/plugins/fastrouter/fr_events.c index 4b74fa0a..fd3ec507 100644 --- a/plugins/fastrouter/fr_events.c +++ b/plugins/fastrouter/fr_events.c @@ -6,7 +6,8 @@ extern struct uwsgi_server uwsgi; extern struct uwsgi_fastrouter ufr; #ifdef UWSGI_SCTP -extern struct uwsgi_fr_sctp_node *uwsgi_fastrouter_sctp_nodes; +extern struct uwsgi_fr_sctp_node **uwsgi_fastrouter_sctp_nodes; +extern struct uwsgi_fr_sctp_node **uwsgi_fastrouter_sctp_nodes_current; #endif void uwsgi_fastrouter_switch_events(struct fastrouter_session *fr_session, int interesting_fd, char **magic_table) { @@ -26,6 +27,7 @@ void uwsgi_fastrouter_switch_events(struct fastrouter_session *fr_session, int i int tmp_socket_name_len; + switch (fr_session->status) { case FASTROUTER_STATUS_RECV_HDR: @@ -174,7 +176,11 @@ choose_node: #ifdef UWSGI_SCTP else if (ufr.has_sctp_sockets > 0) { - struct uwsgi_fr_sctp_node *ufsn = 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) { @@ -182,7 +188,7 @@ choose_node: choosen_fd = ufsn->fd; break; } - if (ufsn->next == uwsgi_fastrouter_sctp_nodes) { + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes_current) { break; } @@ -190,17 +196,26 @@ choose_node: } // no nodes available - if (choosen_fd == -1) { fr_session->retry = 1; break; } + 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 @@ -350,7 +365,7 @@ choose_node: 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_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]); @@ -375,7 +390,7 @@ choose_node: 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_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]); @@ -383,11 +398,12 @@ choose_node: 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 STREAM !!!\n"); + uwsgi_log("INVALID SCTP STREAM !!!\n"); close_session(ufr.fr_table[fr_session->fd]); break; } diff --git a/plugins/fastrouter/fr_sctp.c b/plugins/fastrouter/fr_sctp.c index 0d80453d..9bf41dfa 100644 --- a/plugins/fastrouter/fr_sctp.c +++ b/plugins/fastrouter/fr_sctp.c @@ -6,69 +6,117 @@ extern struct uwsgi_fastrouter ufr; -struct uwsgi_fr_sctp_node *uwsgi_fastrouter_sctp_nodes = NULL; +struct uwsgi_fr_sctp_node **uwsgi_fastrouter_sctp_nodes; +struct uwsgi_fr_sctp_node **uwsgi_fastrouter_sctp_nodes_current; struct uwsgi_fr_sctp_node *uwsgi_fr_sctp_add_node(int fd) { - struct uwsgi_fr_sctp_node *ufsn = uwsgi_fastrouter_sctp_nodes; + struct uwsgi_fr_sctp_node *ufsn = *uwsgi_fastrouter_sctp_nodes; if (!ufsn) { - uwsgi_fastrouter_sctp_nodes = uwsgi_malloc(sizeof(struct uwsgi_fr_sctp_node)); - uwsgi_fastrouter_sctp_nodes->next = uwsgi_fastrouter_sctp_nodes; - uwsgi_fastrouter_sctp_nodes->prev = uwsgi_fastrouter_sctp_nodes; - uwsgi_fastrouter_sctp_nodes->requests = 0; - uwsgi_fastrouter_sctp_nodes->fd = fd; + *uwsgi_fastrouter_sctp_nodes = uwsgi_malloc(sizeof(struct uwsgi_fr_sctp_node)); + (*uwsgi_fastrouter_sctp_nodes)->next = *uwsgi_fastrouter_sctp_nodes; + (*uwsgi_fastrouter_sctp_nodes)->prev = *uwsgi_fastrouter_sctp_nodes; + (*uwsgi_fastrouter_sctp_nodes)->requests = 0; + (*uwsgi_fastrouter_sctp_nodes)->fd = fd; + + *uwsgi_fastrouter_sctp_nodes_current = *uwsgi_fastrouter_sctp_nodes; + + ufsn = *uwsgi_fastrouter_sctp_nodes; + while(ufsn) { + uwsgi_log("prev %p fd %d next = %p\n", ufsn->prev, ufsn->fd, ufsn->next); + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes) { + break; + } + ufsn = ufsn->next; + } - return uwsgi_fastrouter_sctp_nodes; } else { while(ufsn) { - if (ufsn->next == uwsgi_fastrouter_sctp_nodes) { + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes) { break; } ufsn = ufsn->next; } ufsn->next = uwsgi_malloc(sizeof(struct uwsgi_fr_sctp_node)); - ufsn->next->next = uwsgi_fastrouter_sctp_nodes; + ufsn->next->next = *uwsgi_fastrouter_sctp_nodes; ufsn->next->prev = ufsn; ufsn->next->requests = 0; ufsn->next->fd = fd; + + *uwsgi_fastrouter_sctp_nodes_current = ufsn->next; + + ufsn = *uwsgi_fastrouter_sctp_nodes; + while(ufsn) { + uwsgi_log("prev %p fd %d next = %p\n", ufsn->prev, ufsn->fd, ufsn->next); + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes) { + break; + } + ufsn = ufsn->next; + } } - return ufsn->next; + return *uwsgi_fastrouter_sctp_nodes_current; } void uwsgi_fr_sctp_del_node(int fd) { - struct uwsgi_fr_sctp_node *ufsn = uwsgi_fastrouter_sctp_nodes; + 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; + struct uwsgi_fr_sctp_node *next = ufsn->next; + struct uwsgi_fr_sctp_node *prev = ufsn->prev; - // check: am i the only node ? - if ( ufsn == prev || ufsn == next ) { - free(uwsgi_fastrouter_sctp_nodes); - uwsgi_fastrouter_sctp_nodes = NULL; - break; + // am i the only node ? + if (ufsn == ufsn->next) { + *uwsgi_fastrouter_sctp_nodes = NULL; + free(ufsn); + goto end; + } + + // am i the first node ? + if (ufsn == *uwsgi_fastrouter_sctp_nodes) { + uwsgi_log("--- first node ---\n"); + *uwsgi_fastrouter_sctp_nodes = ufsn->next; + (*uwsgi_fastrouter_sctp_nodes)->prev = prev; + prev->next = *uwsgi_fastrouter_sctp_nodes; + } + else { + uwsgi_log("--- normal node ---\n"); + prev->next = ufsn->next; + next->prev = ufsn->prev; + } + + if (ufsn == *uwsgi_fastrouter_sctp_nodes_current) { + *uwsgi_fastrouter_sctp_nodes_current = *uwsgi_fastrouter_sctp_nodes; } free(ufsn); break; } - if (ufsn->next == uwsgi_fastrouter_sctp_nodes) { + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes) { break; } ufsn = ufsn->next; } + +end: + + ufsn = *uwsgi_fastrouter_sctp_nodes; + while(ufsn) { + uwsgi_log("prev %p fd %d node %p next = %p\n", ufsn->prev, ufsn->fd, ufsn, ufsn->next); + if (ufsn->next == *uwsgi_fastrouter_sctp_nodes) { + break; + } + ufsn = ufsn->next; + } } void uwsgi_opt_fastrouter_sctp(char *opt, char *value, void *foobar) {