completed legion management and logic

This commit is contained in:
Roberto De Ioris
2012-12-04 12:35:17 +01:00
parent 1662f58a09
commit 7ebbac0dd8
5 changed files with 187 additions and 21 deletions
+1
View File
@@ -87,6 +87,7 @@ void uwsgi_init_default() {
#ifdef UWSGI_MULTICAST
uwsgi.multicast_ttl = 1;
uwsgi.multicast_loop = 1;
#endif
}
+173 -18
View File
@@ -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 <legion> <addr>\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);
}
+1 -2
View File
@@ -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()");
}
+3 -1
View File
@@ -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
+9
View File
@@ -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);