From 3dec2821f49eef23d39bf95f98e5b35b9519bf98 Mon Sep 17 00:00:00 2001 From: "roberto@fiorenzo" Date: Wed, 24 Nov 2010 11:42:15 +0100 Subject: [PATCH] various new clustering implementations --- master.c | 102 ++++++++++++++++++++++++++++++------------------ protocol.c | 111 +++++++++++++++++++++++++++++++++++++++++++++++++++++ socket.c | 103 ++++++++++++++++++++++++++++++++++++++----------- utils.c | 17 ++++++++ uwsgi.c | 63 ++++++++++++++++++++++++++++-- uwsgi.h | 15 +++++++- 6 files changed, 345 insertions(+), 66 deletions(-) diff --git a/master.c b/master.c index 84ca00a9..684e2335 100644 --- a/master.c +++ b/master.c @@ -40,12 +40,14 @@ void master_loop(char **argv, char **environ) { int master_has_children = 0; #ifdef UWSGI_UDP - struct pollfd uwsgi_poll; + struct pollfd uwsgi_poll[2]; + int uwsgi_poll_size = 0; struct sockaddr_in udp_client; socklen_t udp_len; char udp_client_addr[16]; int udp_managed = 0; int rlen; + int udp_fd = -1 ; #endif int i,j; @@ -67,15 +69,24 @@ void master_loop(char **argv, char **environ) { uwsgi.wsgi_req->buffer = uwsgi.async_buf[0]; #ifdef UWSGI_UDP if (uwsgi.udp_socket) { - uwsgi_poll.fd = bind_to_udp(uwsgi.udp_socket); - if (uwsgi_poll.fd < 0) { + udp_fd = bind_to_udp(uwsgi.udp_socket, 0); + uwsgi_poll[uwsgi_poll_size].fd = udp_fd; + if (uwsgi_poll[uwsgi_poll_size].fd < 0) { uwsgi_log( "unable to bind to udp socket. SNMP and cluster management services will be disabled.\n"); } else { uwsgi_log( "UDP server enabled.\n"); - uwsgi_poll.events = POLLIN; + uwsgi_poll[uwsgi_poll_size].events = POLLIN; + uwsgi_poll_size++; } } +#ifdef UWSGI_MULTICAST + if (uwsgi.cluster) { + uwsgi_poll[uwsgi_poll_size].fd = uwsgi.cluster_fd; + uwsgi_poll[uwsgi_poll_size].events = POLLIN; + uwsgi_poll_size++; + } +#endif #endif #ifdef UWSGI_SNMP @@ -200,39 +211,49 @@ void master_loop(char **argv, char **environ) { check_interval.tv_sec = 1; #ifdef UWSGI_UDP - if (uwsgi.udp_socket && uwsgi_poll.fd >= 0) { - rlen = poll(&uwsgi_poll, 1, check_interval.tv_sec * 1000); +#ifdef UWSGI_MULTICAST + if ( (uwsgi.udp_socket && udp_fd >= 0) || (uwsgi.cluster && uwsgi.cluster_fd >= 0)) { +#else + if ((uwsgi.udp_socket && udp_fd >= 0)) { +#endif + rlen = poll(uwsgi_poll, uwsgi_poll_size, check_interval.tv_sec * 1000); if (rlen < 0) { uwsgi_error("poll()"); } else if (rlen > 0) { - udp_len = sizeof(udp_client); - rlen = recvfrom(uwsgi_poll.fd, uwsgi.wsgi_req->buffer, uwsgi.buffer_size, 0, (struct sockaddr *) &udp_client, &udp_len); - if (rlen < 0) { - uwsgi_error("recvfrom()"); - } - else if (rlen > 0) { - memset(udp_client_addr, 0, 16); - if (inet_ntop(AF_INET, &udp_client.sin_addr.s_addr, udp_client_addr, 16)) { - if (uwsgi.wsgi_req->buffer[0] == UWSGI_MODIFIER_MULTICAST_ANNOUNCE) { - } -#ifdef UWSGI_SNMP - else if (uwsgi.wsgi_req->buffer[0] == 0x30 && uwsgi.snmp) { - manage_snmp(uwsgi_poll.fd, (uint8_t *) uwsgi.wsgi_req->buffer, rlen, &udp_client); - } -#endif - else { - // loop the various udp manager until one returns true - udp_managed = 0; - for(i=0;i<0xFF;i++) { - if (uwsgi.p[i]->manage_udp) { - if (uwsgi.p[i]->manage_udp(udp_client_addr, udp_client.sin_port, uwsgi.wsgi_req->buffer, rlen)) { - udp_managed = 1; - break; - } - } + for(i=0;ibuffer, uwsgi.buffer_size, 0, (struct sockaddr *) &udp_client, &udp_len); + if (rlen < 0) { + uwsgi_error("recvfrom()"); } + else if (rlen > 0) { + memset(udp_client_addr, 0, 16); + if (inet_ntop(AF_INET, &udp_client.sin_addr.s_addr, udp_client_addr, 16)) { + if (uwsgi.wsgi_req->buffer[0] == UWSGI_MODIFIER_MULTICAST_ANNOUNCE) { + } +#ifdef UWSGI_SNMP + else if (uwsgi.wsgi_req->buffer[0] == 0x30 && uwsgi.snmp) { + manage_snmp(udp_fd, (uint8_t *) uwsgi.wsgi_req->buffer, rlen, &udp_client); + } +#endif + else { + + // loop the various udp manager until one returns true + udp_managed = 0; + for(i=0;i<0xFF;i++) { + if (uwsgi.p[i]->manage_udp) { + if (uwsgi.p[i]->manage_udp(udp_client_addr, udp_client.sin_port, uwsgi.wsgi_req->buffer, rlen)) { + udp_managed = 1; + break; + } + } + } /* if (udp_callable && udp_callable_args) { UWSGI_GET_GIL @@ -251,15 +272,20 @@ void master_loop(char **argv, char **environ) { else { // a simple udp logger */ - - if (!udp_managed) { - uwsgi_log( "[udp:%s:%d] %.*s", udp_client_addr, ntohs(udp_client.sin_port), rlen, uwsgi.wsgi_req->buffer); + if (!udp_managed) { + uwsgi_log( "[udp:%s:%d] %.*s", udp_client_addr, ntohs(udp_client.sin_port), rlen, uwsgi.wsgi_req->buffer); + } + } + } + else { + uwsgi_error("inet_ntop()"); + } } - //} } - } - else { - uwsgi_error("inet_ntop()"); + + if (uwsgi_poll[i].fd == uwsgi.cluster_fd) { + uwsgi_log("received a cluster message\n"); + } } } } diff --git a/protocol.c b/protocol.c index 2a4527e5..caeb5e4a 100644 --- a/protocol.c +++ b/protocol.c @@ -558,3 +558,114 @@ int uwsgi_ping_node(int node, struct wsgi_request *wsgi_req) { return 0; } + +ssize_t uwsgi_send_empty_pkt(int fd, char *socket_name, uint8_t modifier1, uint8_t modifier2) { + + struct uwsgi_header uh; + char *port; + uint16_t s_port; + struct sockaddr_in uaddr; + int ret; + + uh.modifier1 = modifier1; + uh.pktsize = 0; + uh.modifier2 = modifier2; + + if (socket_name) { + port = strchr(socket_name, ':'); + if (!port) return -1; + s_port = atoi(port+1); + port[0] = 0; + memset(&uaddr, 0, sizeof(struct sockaddr_in)); + uaddr.sin_family = AF_INET; + uaddr.sin_addr.s_addr = inet_addr(socket_name); + uaddr.sin_port = htons(s_port); + + port[0] = ':'; + + ret = sendto(fd, &uh, 4, 0, (struct sockaddr *) &uaddr, sizeof(struct sockaddr_in)); + } + else { + ret = send(fd, &uh, 4, 0); + } + + if (ret < 0) { + uwsgi_error("sendto()"); + } + + return ret; +} + +int uwsgi_hooked_parse_dict_dgram(int fd, char *buffer, size_t len, uint8_t modifier1, uint8_t modifier2, void (*hook)()) { + + struct uwsgi_header *uh; + ssize_t rlen; + + char *ptrbuf, *bufferend; + uint16_t keysize = 0, valsize = 0; + char *key; + + ptrbuf = buffer; + + rlen = read(fd, buffer, len); + + if (rlen < 0) { + uwsgi_error("read()"); + return -1; + } + + // check for valid dict 4(header) 2(non-zero key)+1 2(value) + if (rlen < (4+2+1+2)) { + uwsgi_log("invalid uwsgi dictionary\n"); + return -1; + } + + uh = (struct uwsgi_header *) buffer; + + if (uh->modifier1 != modifier1 || uh->modifier2 != modifier2) { + uwsgi_log("invalid uwsgi dictionary received, modifier1: %d modifier2: %d\n", uh->modifier1, uh->modifier2); + return -1; + } + + if (uh->pktsize > len) { + uwsgi_log("* WARNING * the uwsgi dictionary received is too big, data will be truncated\n"); + bufferend = ptrbuf + len; + } + else { + bufferend = ptrbuf + uh->pktsize; + } + + + while (ptrbuf < bufferend) { + if (ptrbuf + 2 >= bufferend) return -1; + memcpy(&keysize, ptrbuf, 2); +#ifdef __BIG_ENDIAN__ + keysize = uwsgi_swap16(keysize); +#endif + /* key cannot be null */ + if (!keysize) return -1; + + ptrbuf += 2; + if (ptrbuf + keysize > bufferend) return -1; + + // key + key = ptrbuf; + ptrbuf += keysize; + // value can be null + if (ptrbuf + 2 > bufferend) return -1; + + memcpy(&valsize, ptrbuf, 2); +#ifdef __BIG_ENDIAN__ + valsize = uwsgi_swap16(valsize); +#endif + ptrbuf += 2; + if (ptrbuf + valsize > bufferend) return -1; + + // now call the hook + hook(key, keysize, ptrbuf, valsize); + ptrbuf += valsize; + } + + return 0; + +} diff --git a/socket.c b/socket.c index 0b66fd6e..9046efdf 100644 --- a/socket.c +++ b/socket.c @@ -132,13 +132,14 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst #endif #ifdef UWSGI_UDP - int bind_to_udp(char *socket_name) { + int bind_to_udp(char *socket_name, int multicast) { int serverfd; struct sockaddr_in uws_addr; char *udp_port; #ifdef UWSGI_MULTICAST struct ip_mreq mc; + uint8_t loop = 0; #endif udp_port = strchr(socket_name, ':'); @@ -147,6 +148,11 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst } udp_port[0] = 0; + + if (socket_name[0] == 0 && multicast) { + uwsgi_log("invalid multicast address\n"); + return -1; + } memset(&uws_addr, 0, sizeof(struct sockaddr_in)); uws_addr.sin_family = AF_INET; uws_addr.sin_port = htons(atoi(udp_port + 1)); @@ -167,21 +173,14 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst } #ifdef UWSGI_MULTICAST - if (uwsgi.multicast_group) { - uws_addr.sin_addr.s_addr = INADDR_ANY; + if (multicast) { // if multicast is enabled remember to bind to INADDR_ANY - mc.imr_multiaddr.s_addr = inet_addr(uwsgi.multicast_group); - if (socket_name[0] == 0) { - mc.imr_interface.s_addr = INADDR_ANY; - } - else { - mc.imr_interface.s_addr = inet_addr(socket_name); - } + uws_addr.sin_addr.s_addr = INADDR_ANY; + mc.imr_multiaddr.s_addr = inet_addr(socket_name); + mc.imr_interface.s_addr = INADDR_ANY; } #endif - uwsgi_log( "binding on UDP port: %d\n", ntohs(uws_addr.sin_port)); - if (bind(serverfd, (struct sockaddr *) &uws_addr, sizeof(uws_addr)) != 0) { uwsgi_error("bind()"); close(serverfd); @@ -189,14 +188,19 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst } #ifdef UWSGI_MULTICAST - if (uwsgi.multicast_group) { - uwsgi_log( "joining uWSGI multicast group: %s:%d\n", uwsgi.multicast_group, ntohs(uws_addr.sin_port)); + if (multicast) { + uwsgi_log( "[uWSGI] joining multicast group: %s:%d\n", socket_name, ntohs(uws_addr.sin_port)); + if (setsockopt(serverfd, IPPROTO_IP, IP_MULTICAST_LOOP, &loop, sizeof(loop))) { + uwsgi_error("setsockopt()"); + } + if (setsockopt(serverfd, IPPROTO_IP, IP_ADD_MEMBERSHIP, &mc, sizeof(mc))) { uwsgi_error("setsockopt()"); } } #endif + udp_port[0] = ':'; return serverfd; } @@ -280,11 +284,12 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst } - int bind_to_tcp(char *socket_name, int listen_queue, char *tcp_port) { + int bind_to_tcp(char **socket_name, int listen_queue, char *tcp_port) { int serverfd; struct sockaddr_in uws_addr; int reuse = 1; + int i; tcp_port[0] = 0; memset(&uws_addr, 0, sizeof(struct sockaddr_in)); @@ -292,20 +297,72 @@ int bind_to_unix(char *socket_name, int listen_queue, int chmod_socket, int abst uws_addr.sin_family = AF_INET; uws_addr.sin_port = htons(atoi(tcp_port + 1)); - - if (socket_name[0] == 0) { - uws_addr.sin_addr.s_addr = INADDR_ANY; - } - else { - uws_addr.sin_addr.s_addr = inet_addr(socket_name); - } - serverfd = socket(AF_INET, SOCK_STREAM, 0); if (serverfd < 0) { uwsgi_error("socket()"); exit(1); } + if (*socket_name[0] == 0) { + uws_addr.sin_addr.s_addr = INADDR_ANY; + } + else { + char *asterisk = strchr(*socket_name, '*'); + if (asterisk) { + struct ifconf ifc; + struct ifreq *ifr, *aifr; + // get all the AF_INET addresses available + + ifc.ifc_len = 0; + ifc.ifc_req = NULL; + if (ioctl(serverfd, SIOCGIFCONF, &ifc)) { + uwsgi_error("ioctl()"); + exit(1); + } + + if (ifc.ifc_len <= 0) { + uwsgi_log("unable to get ip address list\n"); + exit(1); + } + + ifr = malloc(ifc.ifc_len); + if (!ifr) { + uwsgi_error("malloc()"); + } + memset(ifr, 0, ifc.ifc_len); + + ifc.ifc_req = ifr; + + if (ioctl(serverfd, SIOCGIFCONF, &ifc)) { + uwsgi_error("ioctl()"); + exit(1); + } + + // here socket_name will be truncated + asterisk[0] = 0; + + char new_addr[16]; + struct sockaddr_in *sin; + for(i=0;i< (int)(ifc.ifc_len/sizeof(struct ifreq));i++) { + aifr = ifr + i ; + memset(new_addr, 0, 16); + sin = (struct sockaddr_in *) &aifr->ifr_addr; + if (inet_ntop(AF_INET, (void *) &sin->sin_addr.s_addr, new_addr, 16)) { + if (!strncmp( *socket_name, new_addr, strlen(*socket_name)) ) { + asterisk[0] = '*'; + uwsgi_log("found %s for %s on interface %s\n", new_addr, *socket_name, aifr->ifr_name); + *socket_name = uwsgi_concat3(new_addr, ":", tcp_port+1); + break; + } + } + } + } + else { + uws_addr.sin_addr.s_addr = inet_addr(*socket_name); + } + } + + if (setsockopt(serverfd, SOL_SOCKET, SO_REUSEADDR, (const void *) &reuse, sizeof(int)) < 0) { uwsgi_error("setsockopt()"); exit(1); diff --git a/utils.c b/utils.c index 1d73090c..b60d7a21 100644 --- a/utils.c +++ b/utils.c @@ -1110,3 +1110,20 @@ void add_exported_option(int i, char *value) { uwsgi.exported_opts_cnt++; } + +int uwsgi_waitfd(int fd, int timeout) { + + int ret; + struct pollfd upoll[1]; + + upoll[0].fd = fd; + upoll[0].events = POLLIN; + + ret = poll(upoll, 1, timeout*1000); + + if (ret < 0) { + uwsgi_error("poll()"); + } + + return ret; +} diff --git a/uwsgi.c b/uwsgi.c index 2af70dad..11408245 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -95,6 +95,7 @@ static struct option long_base_options[] = { #endif #ifdef UWSGI_MULTICAST {"multicast", required_argument, 0, LONG_ARGS_MULTICAST}, + {"cluster", required_argument, 0, LONG_ARGS_CLUSTER}, #endif #ifdef UWSGI_SNMP {"snmp", no_argument, 0, LONG_ARGS_SNMP}, @@ -568,6 +569,34 @@ int main(int argc, char *argv[], char *envp[]) //parse environ parse_sys_envs(environ); + // get cluster configuration + if (uwsgi.cluster != NULL) { + // get multicast socket + uwsgi.cluster_fd = bind_to_udp(uwsgi.cluster, 1); + + // ask for cluster options only if bot pre-existent options are set + if (uwsgi.exported_opts_cnt == 1) { + // now wait max 60 seconds and resend multicast request every 10 seconds + for(i=0;i<6;i++) { + uwsgi_log("asking \"%s\" uWSGI cluster for configuration data:\n", uwsgi.cluster); + if (uwsgi_send_empty_pkt(uwsgi.cluster_fd, uwsgi.cluster, 99, 0) < 0) { + uwsgi_log("unable to send multicast message to %s\n", uwsgi.cluster); + continue; + } + rlen = uwsgi_waitfd(uwsgi.cluster_fd, 10); + if (rlen < 0) { + break; + } + else if (rlen > 0) { + // receive the packet + char clusterbuf[4096]; + uwsgi_hooked_parse_dict_dgram(uwsgi.cluster_fd, clusterbuf, 4096, 99, 1, manage_string_opt); + break; + } + } + } + } + //call after_opt hooks if (uwsgi.binary_path == argv[0]) { @@ -663,7 +692,7 @@ int main(int argc, char *argv[], char *envp[]) char *tcp_port = strchr(uwsgi.http, ':'); if (tcp_port) { uwsgi.http_server_port = tcp_port + 1; - uwsgi.http_fd = bind_to_tcp(uwsgi.http, uwsgi.listen_queue, tcp_port); + uwsgi.http_fd = bind_to_tcp(&uwsgi.http, uwsgi.listen_queue, tcp_port); #ifdef UWSGI_DEBUG uwsgi_debug("HTTP FD: %d\n", uwsgi.http_fd); #endif @@ -952,7 +981,7 @@ int main(int argc, char *argv[], char *envp[]) uwsgi.sockets[i].family = AF_UNIX; uwsgi_log("uwsgi socket %d bound to UNIX address %s fd %d\n", i, uwsgi.sockets[i].name, uwsgi.sockets[i].fd); } else { - uwsgi.sockets[i].fd = bind_to_tcp(uwsgi.sockets[i].name, uwsgi.listen_queue, tcp_port); + uwsgi.sockets[i].fd = bind_to_tcp(&uwsgi.sockets[i].name, uwsgi.listen_queue, tcp_port); uwsgi.sockets[i].family = AF_INET; uwsgi_log("uwsgi socket %d bound to TCP address %s fd %d\n", i, uwsgi.sockets[i].name, uwsgi.sockets[i].fd); } @@ -1422,7 +1451,7 @@ end: if (tcp_port == NULL) { uwsgi.proxyfd = bind_to_unix(uwsgi.proxy_socket_name, UWSGI_LISTEN_QUEUE, uwsgi.chmod_socket, uwsgi.abstract_socket); } else { - uwsgi.proxyfd = bind_to_tcp(uwsgi.proxy_socket_name, UWSGI_LISTEN_QUEUE, tcp_port); + uwsgi.proxyfd = bind_to_tcp(&uwsgi.proxy_socket_name, UWSGI_LISTEN_QUEUE, tcp_port); tcp_port[0] = ':'; } @@ -1553,6 +1582,9 @@ end: uwsgi.multicast_group = optarg; uwsgi.master_process = 1; return 1; + case LONG_ARGS_CLUSTER: + uwsgi.cluster = optarg; + return 1; #endif case LONG_ARGS_CHROOT: uwsgi.chroot = optarg; @@ -2127,3 +2159,28 @@ void build_options() { uwsgi.long_options[opt_count].flag = 0; uwsgi.long_options[opt_count].val = 0; } + + +void manage_string_opt(char *key, int keylen, char *val, int vallen) { + + struct option *lopt, *aopt; + + // never free this value + char *key2 = uwsgi_concat2(key, ""); + char *val2 = uwsgi_concat2(val, ""); + + lopt = uwsgi.long_options; + while ((aopt = lopt)) { + if (!aopt->name) break; + if (!strcmp(key2, aopt->name)) { + if (aopt->flag) { + *aopt->flag = aopt->val; + add_exported_option(0, key2); + } + else { + manage_opt(aopt->val, val2); + } + } + lopt++; + } +} diff --git a/uwsgi.h b/uwsgi.h index 405924f1..542003fe 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -66,6 +66,7 @@ #include #include +#include #include @@ -240,6 +241,7 @@ struct uwsgi_opt { #define LONG_ARGS_VHOSTHOST 17059 #define LONG_ARGS_UPLOAD_PROGRESS 17060 #define LONG_ARGS_REMAP_MODIFIER 17061 +#define LONG_ARGS_CLUSTER 17062 @@ -797,6 +799,8 @@ struct uwsgi_server { char *upload_progress; + char *cluster; + int cluster_fd; }; struct uwsgi_cluster_node { @@ -909,8 +913,8 @@ void grace_them_all(void); void reload_me(void); void end_me(void); int bind_to_unix(char *, int, int, int); -int bind_to_tcp(char *, int, char *); -int bind_to_udp(char *); +int bind_to_tcp(char **, int, char *); +int bind_to_udp(char *, int); int timed_connect(struct pollfd *, const struct sockaddr *, int, int); int uwsgi_connect(char *, int); int connect_to_tcp(char *, int, int); @@ -1144,3 +1148,10 @@ void uwsgi_register_loop(char *, void *); void *uwsgi_get_loop(char *); void add_exported_option(int, char *); + +ssize_t uwsgi_send_empty_pkt(int , char *, uint8_t , uint8_t); + +int uwsgi_waitfd(int, int); + +int uwsgi_hooked_parse_dict_dgram(int, char *, size_t, uint8_t, uint8_t, void (*)()); +void manage_string_opt(char *, int, char*, int);