From 7ebbac0dd877b241d7a022f396ec723937043675 Mon Sep 17 00:00:00 2001 From: Roberto De Ioris Date: Tue, 4 Dec 2012 12:35:17 +0100 Subject: [PATCH] completed legion management and logic --- core/init.c | 1 + core/legion.c | 191 +++++++++++++++++++++++++++++++++++++++++++++----- core/socket.c | 3 +- core/uwsgi.c | 4 +- uwsgi.h | 9 +++ 5 files changed, 187 insertions(+), 21 deletions(-) diff --git a/core/init.c b/core/init.c index 6d1c034e..b1389f0b 100644 --- a/core/init.c +++ b/core/init.c @@ -87,6 +87,7 @@ void uwsgi_init_default() { #ifdef UWSGI_MULTICAST uwsgi.multicast_ttl = 1; + uwsgi.multicast_loop = 1; #endif } diff --git a/core/legion.c b/core/legion.c index 87c90d2a..5d096521 100644 --- a/core/legion.c +++ b/core/legion.c @@ -49,6 +49,19 @@ struct uwsgi_legion *uwsgi_legion_get_by_socket(int fd) { return NULL; } +struct uwsgi_legion *uwsgi_legion_get_by_name(char *name) { + struct uwsgi_legion *ul = uwsgi.legions; + while(ul) { + if (!strcmp(name, ul->legion)) { + return ul; + } + ul = ul->next; + } + + return NULL; +} + + void uwsgi_parse_legion(char *key, uint16_t keylen, char *value, uint16_t vallen, void *data) { struct uwsgi_legion *ul = (struct uwsgi_legion *) data; @@ -59,8 +72,41 @@ void uwsgi_parse_legion(char *key, uint16_t keylen, char *value, uint16_t vallen else if (!uwsgi_strncmp(key, keylen, "valor", 5)) { ul->valor = uwsgi_str_num(value, vallen); } + else if (!uwsgi_strncmp(key, keylen, "name", 4)) { + ul->name = value; + ul->name_len = vallen; + } + else if (!uwsgi_strncmp(key, keylen, "pid", 3)) { + ul->pid = uwsgi_str_num(value, vallen); + } } +static void legions_check_lord() { + + struct uwsgi_legion *legion = uwsgi.legions; + while(legion) { + time_t now = uwsgi_now(); + + if (!legion->last_seen_lord) { + legion->last_seen_lord = now; + } + + if (legion->lord) { + goto next; + } + + if (now - legion->last_seen_lord > uwsgi.legion_tolerance) { + uwsgi_log("[uwsgi-legion] i am now the Lord of the Legion %s\n", legion->legion); + legion->lord = now; + // triggering lord hooks + } + +next: + legion = legion->next; + } +} + + static void *legion_loop(void *foobar) { time_t last_round = uwsgi_now(); @@ -70,20 +116,35 @@ static void *legion_loop(void *foobar) { struct uwsgi_legion legion_msg; + if (!uwsgi.legion_freq) uwsgi.legion_freq = 3; + if (!uwsgi.legion_tolerance) uwsgi.legion_tolerance = 15; + for(;;) { + int timeout = uwsgi.legion_freq; time_t now = uwsgi_now(); - int timeout = 0; if (now > last_round) { - timeout = now - last_round; + timeout -= (now - last_round); + if (timeout < 0) { + timeout = 0; + } } + last_round = now; // wait for event int interesting_fd = -1; - int rlen = event_queue_wait(uwsgi.listen_queue, timeout, &interesting_fd); + int rlen = event_queue_wait(uwsgi.legion_queue, timeout, &interesting_fd); - // check elapsed time - if (1) { + now = uwsgi_now(); + if (timeout == 0 || rlen == 0 || (now - last_round) >= timeout) { + struct uwsgi_legion *legions = uwsgi.legions; + while(legions) { + uwsgi_legion_announce(legions); + legions = legions->next; + } + last_round = now; } + legions_check_lord(); + if (rlen > 0) { struct uwsgi_legion *ul = uwsgi_legion_get_by_socket(interesting_fd); if (!ul) continue; @@ -99,6 +160,12 @@ static void *legion_loop(void *foobar) { } struct uwsgi_header *uh = (struct uwsgi_header *) crypted_buf; + + if (uh->modifier1 != 109) { + uwsgi_log("[uwsgi-legion] invalid modifier1"); + continue; + } + int d_len = 0; int d2_len = 0; // decrypt packet using the secret @@ -106,19 +173,22 @@ static void *legion_loop(void *foobar) { uwsgi_error("[uwsgi-legion] EVP_DecryptInit_ex()"); continue; } - if (EVP_DecryptUpdate(ul->decrypt_ctx, clear_buf, &d_len, crypted_buf+4, len) <= 0) { + + if (EVP_DecryptUpdate(ul->decrypt_ctx, clear_buf, &d_len, crypted_buf+4, len-4) <= 0) { uwsgi_error("[uwsgi-legion] EVP_DecryptUpdate()"); continue; } + if (EVP_DecryptFinal_ex(ul->decrypt_ctx, clear_buf + d_len, &d2_len) <= 0) { - uwsgi_error("[uwsgi-legion] EVP_DecryptFinal_ex()"); + ERR_print_errors_fp(stderr); + uwsgi_log("[uwsgi-legion] EVP_DecryptFinal_ex()\n"); continue; } d_len += d2_len; if (d_len != uh->pktsize) { - uwsgi_log("[uwsgi-legion] invalid packet"); + uwsgi_log("[uwsgi-legion] invalid packet size\n"); continue; } @@ -134,12 +204,32 @@ static void *legion_loop(void *foobar) { continue; } + // check for loop packets... (expecially when in multicast mode) + if (!uwsgi_strncmp(uwsgi.hostname, uwsgi.hostname_len, legion_msg.name, legion_msg.name_len)) { + if (legion_msg.pid == ul->pid) { + if (legion_msg.valor == ul->valor) { + continue; + } + } + } + if (ul->lord > 0) { if (legion_msg.valor > ul->valor) { - uwsgi_log("[uwsgi-legion] a new Lord raised for legion %s...\n", ul->legion); - // no more lord, trigger unlord events + uwsgi_log("[uwsgi-legion] a new Lord raised for Legion %s...\n", ul->legion); + // no more lord, trigger unlord hooks + ul->last_seen_lord = uwsgi_now(); + ul->lord = 0; + continue; } } + + if (legion_msg.valor > ul->valor) { + // a lord + ul->last_seen_lord = uwsgi_now(); + } + else if (legion_msg.valor == ul->valor) { + uwsgi_log("[uwsgi-legion] a node with the same valor announced itself !!!\n"); + } } } @@ -163,6 +253,7 @@ void uwsgi_start_legions() { exit(1); } uwsgi_socket_nb(legion->socket); + legion->pid = uwsgi.mypid; legion = legion->next; } @@ -194,29 +285,84 @@ void uwsgi_legion_add(struct uwsgi_legion *ul) { int uwsgi_legion_announce(struct uwsgi_legion *ul) { struct uwsgi_buffer *ub = uwsgi_buffer_new(4096); - if (uwsgi_buffer_append(ub, "\0\0\0\0", 4)) goto err; - if (uwsgi_buffer_append_keyval(ub, "legion", 6, ul->legion, ul->legion_len)) goto err; if (uwsgi_buffer_append_keynum(ub, "valor", 5, ul->valor)) goto err; if (uwsgi_buffer_append_keynum(ub, "unix", 4, uwsgi_now())) goto err; if (uwsgi_buffer_append_keynum(ub, "lord", 4, ul->lord ? ul->lord : 0)) goto err; if (uwsgi_buffer_append_keyval(ub, "name", 4, uwsgi.hostname, uwsgi.hostname_len)) goto err; + if (uwsgi_buffer_append_keynum(ub, "pid", 3, ul->pid)) goto err; + + unsigned char *encrypted = uwsgi_malloc(ub->pos + 4 + EVP_MAX_BLOCK_LENGTH); + if (EVP_EncryptInit_ex(ul->encrypt_ctx, NULL, NULL, NULL, NULL) <= 0) { + uwsgi_error("[uwsgi-legion] EVP_EncryptInit_ex()"); + goto err; + } + + int e_len = 0; + + if (EVP_EncryptUpdate(ul->encrypt_ctx, encrypted+4, &e_len, (unsigned char *)ub->buf, ub->pos) <= 0) { + uwsgi_error("[uwsgi-legion] EVP_EncryptUpdate()"); + goto err; + } + + int tmplen = 0; + if (EVP_EncryptFinal_ex(ul->encrypt_ctx, encrypted+4+e_len, &tmplen) <= 0) { + uwsgi_error("[uwsgi-legion] EVP_EncryptFinal_ex()"); + goto err; + } + + e_len += tmplen; + uint16_t pktsize = ub->pos; + encrypted[0] = 109; + encrypted[1] = (unsigned char) (pktsize & 0xff); + encrypted[2] = (unsigned char) ((pktsize >> 8) & 0xff); + encrypted[3] = 0; struct uwsgi_string_list *usl = ul->nodes; while(usl) { -/* - if (uwsgi_buffer_dgram(ub, ul->socket, usl->custom_ptr)) { - uwsgi_log("[uwsgi-legion] unable to announce presence to legion \"%s\" using addres \"%s\"\n", ul->legion, usl->value); + if (sendto(ul->socket, encrypted, e_len + 4, 0, usl->custom_ptr, usl->custom) != e_len + 4) { + uwsgi_error("[uwsgi-legion] sendto()"); } -*/ usl = usl->next; } + uwsgi_buffer_destroy(ub); + free(encrypted); + return 0; err: uwsgi_buffer_destroy(ub); return -1; } +void uwsgi_opt_legion_node(char *opt, char *value, void *foobar) { + + char *legion = uwsgi_str(value); + + char *space = strchr(legion, ' '); + if (!space) { + uwsgi_log("invalid legion-node syntax, must be \n"); + exit(1); + } + *space = 0; + + struct uwsgi_legion *ul = uwsgi_legion_get_by_name(legion); + if (!ul) { + uwsgi_log("unknown legion: %s\n", legion); + exit(1); + } + + struct uwsgi_string_list *usl = uwsgi_string_new_list(&ul->nodes, space+1); + char *port = strchr(usl->value, ':'); + if (!port) { + uwsgi_log("[uwsgi-legion] invalid udp address: %s\n", usl->value); + exit(1); + } + // no need to zero the memory, socket_to_in_addr will do that + struct sockaddr_in *sin = uwsgi_malloc(sizeof(struct sockaddr_in)); + usl->custom = socket_to_in_addr(usl->value, port, 0, sin); + usl->custom_ptr = sin; +} + void uwsgi_opt_legion(char *opt, char *value, void *foobar) { // legion addr valor algo:secret @@ -266,13 +412,22 @@ void uwsgi_opt_legion(char *opt, char *value, void *foobar) { exit(1); } + int cipher_len = EVP_CIPHER_key_length(cipher); + size_t s_len = strlen(secret); + if ((unsigned int)cipher_len > s_len) { + char *secret_tmp = uwsgi_malloc(cipher_len); + memcpy(secret_tmp, secret, s_len); + memset(secret_tmp + s_len, 0, cipher_len - s_len); + secret = secret_tmp; + } + char *iv = uwsgi_ssl_rand(strlen(secret)); if (!iv) { uwsgi_log("[uwsgi-legion] unable to generate iv for legion %s\n", legion); exit(1); } - if (EVP_EncryptInit_ex(ctx, cipher, NULL, (const unsigned char *)secret, (const unsigned char *) iv) <= 0) { + if (EVP_EncryptInit_ex(ctx, cipher, NULL, (const unsigned char *)secret, (const unsigned char *) "12345678") <= 0) {// (const unsigned char *) iv) <= 0) { uwsgi_error("EVP_EncryptInit_ex()"); exit(1); } @@ -280,7 +435,7 @@ void uwsgi_opt_legion(char *opt, char *value, void *foobar) { EVP_CIPHER_CTX *ctx2 = uwsgi_malloc(sizeof(EVP_CIPHER_CTX)); EVP_CIPHER_CTX_init(ctx2); - if (EVP_DecryptInit_ex(ctx2, cipher, NULL, (const unsigned char *)secret, NULL) <= 0) { + if (EVP_DecryptInit_ex(ctx2, cipher, NULL, (const unsigned char *)secret, (const unsigned char *) "12345678") <= 0) { uwsgi_error("EVP_DecryptInit_ex()"); exit(1); } diff --git a/core/socket.c b/core/socket.c index de5c287d..88568aa7 100644 --- a/core/socket.c +++ b/core/socket.c @@ -178,7 +178,6 @@ int bind_to_udp(char *socket_name, int multicast, int broadcast) { #ifdef UWSGI_MULTICAST struct ip_mreq mc; - uint8_t loop = 1; #endif udp_port = strchr(socket_name, ':'); @@ -262,7 +261,7 @@ int bind_to_udp(char *socket_name, int multicast, int broadcast) { #ifdef UWSGI_MULTICAST 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))) { + if (setsockopt(serverfd, IPPROTO_IP, IP_MULTICAST_LOOP, &uwsgi.multicast_loop, sizeof(uwsgi.multicast_loop))) { uwsgi_error("setsockopt()"); } diff --git a/core/uwsgi.c b/core/uwsgi.c index 26b9599f..36ff04a1 100644 --- a/core/uwsgi.c +++ b/core/uwsgi.c @@ -322,14 +322,16 @@ static struct uwsgi_option uwsgi_base_options[] = { #ifdef UWSGI_MULTICAST {"multicast", required_argument, 0, "subscribe to specified multicast group", uwsgi_opt_set_str, &uwsgi.multicast_group, UWSGI_OPT_MASTER}, {"multicast-ttl", required_argument, 0, "set multicast ttl", uwsgi_opt_set_int, &uwsgi.multicast_ttl, 0}, + {"multicast-loop", required_argument, 0, "set multicast loop (default 1)", uwsgi_opt_set_int, &uwsgi.multicast_loop, 0}, {"cluster", required_argument, 0, "join specified uWSGI cluster", uwsgi_opt_set_str, &uwsgi.cluster, UWSGI_OPT_MASTER}, {"cluster-nodes", required_argument, 0, "get nodes list from the specified cluster", uwsgi_opt_true, &uwsgi.cluster_nodes, UWSGI_OPT_MASTER | UWSGI_OPT_CLUSTER}, {"cluster-reload", required_argument, 0, "send a reload message to the cluster", uwsgi_opt_cluster_reload, NULL, UWSGI_OPT_IMMEDIATE}, {"cluster-log", required_argument, 0, "send a log line to the cluster", uwsgi_opt_cluster_log, NULL, UWSGI_OPT_IMMEDIATE}, #endif - {"legion", required_argument, 0, "became a member of a legion", uwsgi_opt_legion, NULL, UWSGI_OPT_MASTER}, #ifdef UWSGI_SSL + {"legion", required_argument, 0, "became a member of a legion", uwsgi_opt_legion, NULL, UWSGI_OPT_MASTER}, + {"legion-node", required_argument, 0, "add a node to a legion", uwsgi_opt_legion_node, NULL, UWSGI_OPT_MASTER}, {"subscriptions-sign-check", required_argument, 0, "set digest algorithm and certificate directory for secured subscription system", uwsgi_opt_scd, NULL, UWSGI_OPT_MASTER}, {"subscriptions-sign-check-tolerance", required_argument, 0, "set the maximum tolerance (in seconds) of clock skew for secured subscription system", uwsgi_opt_set_int, &uwsgi.subscriptions_sign_check_tolerance, UWSGI_OPT_MASTER}, #endif diff --git a/uwsgi.h b/uwsgi.h index c52178cf..4f0fc236 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -512,6 +512,10 @@ struct uwsgi_legion { uint64_t valor; char *addr; time_t lord; + time_t last_seen_lord; + char *name; + uint16_t name_len; + pid_t pid; int socket; EVP_CIPHER_CTX *encrypt_ctx; EVP_CIPHER_CTX *decrypt_ctx; @@ -1608,6 +1612,7 @@ struct uwsgi_server { #ifdef UWSGI_MULTICAST int multicast_ttl; + int multicast_loop; char *multicast_group; #endif @@ -1921,6 +1926,8 @@ struct uwsgi_server { #ifdef UWSGI_SSL struct uwsgi_legion *legions; int legion_queue; + int legion_freq; + int legion_tolerance; #endif #ifdef __linux__ @@ -3535,9 +3542,11 @@ char *uwsgi_strip(char *); #ifdef UWSGI_SSL void uwsgi_opt_legion(char *, char *, void *); +void uwsgi_opt_legion_node(char *, char *, void *); void uwsgi_legion_add(struct uwsgi_legion *); char *uwsgi_ssl_rand(size_t); void uwsgi_start_legions(void); +int uwsgi_legion_announce(struct uwsgi_legion *); #endif void uwsgi_check_emperor(void);