various new clustering implementations

This commit is contained in:
roberto@fiorenzo
2010-11-24 11:42:15 +01:00
parent 652a9c5825
commit 3dec2821f4
6 changed files with 345 additions and 66 deletions
+64 -38
View File
@@ -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;i<uwsgi_poll_size;i++) {
if (uwsgi_poll[i].revents & POLLIN) {
if (uwsgi_poll[i].fd == udp_fd) {
udp_len = sizeof(udp_client);
rlen = recvfrom(udp_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(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");
}
}
}
}
+111
View File
@@ -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;
}
+80 -23
View File
@@ -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);
+17
View File
@@ -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;
}
+60 -3
View File
@@ -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++;
}
}
+13 -2
View File
@@ -66,6 +66,7 @@
#include <sys/time.h>
#include <unistd.h>
#include <net/if.h>
#include <dirent.h>
@@ -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);