another series of sctp patch

This commit is contained in:
roberto@precise64
2012-03-05 04:29:11 +01:00
parent da9ce52294
commit 629fbf4f07
3 changed files with 102 additions and 58 deletions
+9 -29
View File
@@ -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;
+23 -7
View File
@@ -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;
}
+70 -22
View File
@@ -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) {