diff --git a/master.c b/master.c index 875df45a..f2a9027a 100644 --- a/master.c +++ b/master.c @@ -100,7 +100,7 @@ void uwsgi_subscribe(char *subscription, uint8_t cmd) { modifier1_len = strlen(modifier1); keysize = strlen(key); } - uwsgi_send_subscription(udp_address, key, keysize, modifier1, modifier1_len, cmd); + uwsgi_send_subscription(udp_address, key, keysize, uwsgi_str_num(modifier1, modifier1_len), 0, cmd); modifier1 = NULL; modifier1_len = 0; } @@ -118,7 +118,7 @@ void uwsgi_subscribe(char *subscription, uint8_t cmd) { modifier1_len = strlen(modifier1); keysize = strlen(key); } - uwsgi_send_subscription(udp_address, key, keysize, modifier1, modifier1_len, cmd); + uwsgi_send_subscription(udp_address, key, keysize, uwsgi_str_num(modifier1, modifier1_len), 0, cmd); modifier1 = NULL; modifier1_len = 0; lines[i] = '\n'; @@ -142,7 +142,7 @@ void uwsgi_subscribe(char *subscription, uint8_t cmd) { modifier1_len = strlen(modifier1); } - uwsgi_send_subscription(udp_address, subscription_key+1, strlen(subscription_key+1), modifier1, modifier1_len, cmd); + uwsgi_send_subscription(udp_address, subscription_key+1, strlen(subscription_key+1), uwsgi_str_num(modifier1, modifier1_len), 0, cmd); if (modifier1) modifier1[-1] = ','; } @@ -163,6 +163,8 @@ void get_linux_tcp_info(int fd) { return; } + uwsgi.shared->load = uwsgi.shared->ti.tcpi_unacked; + uwsgi.shared->options[UWSGI_OPTION_BACKLOG_STATUS] = uwsgi.shared->ti.tcpi_unacked; if (uwsgi.vassal_sos_backlog > 0 && uwsgi.has_emperor) { if ((int)uwsgi.shared->ti.tcpi_unacked >= uwsgi.vassal_sos_backlog) { diff --git a/master_utils.c b/master_utils.c index a1d72f77..089cdd75 100644 --- a/master_utils.c +++ b/master_utils.c @@ -320,6 +320,8 @@ void uwsgi_send_stats(int fd) { stats_send_llu("\"listen_queue\": %llu,\n", uwsgi.shared->ti.tcpi_unacked); #endif + stats_send_llu("\"load\": %llu,\n", uwsgi.shared->load); + fprintf(output,"\"uid\": %d,\n", (int)(getuid())); fprintf(output,"\"gid\": %d,\n", (int)(getgid())); diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index 94cea37d..25a03a0a 100644 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -151,9 +151,29 @@ void fastrouter_manage_subscription(char *key, uint16_t keylen, char *val, uint1 usr->address = val; usr->address_len = vallen; } - else if (!uwsgi_strncmp("modifier1", 9, key, keylen)) { - usr->modifier1 = uwsgi_str_num(val, vallen); + if (vallen == 1) { + usr->modifier1 = *val; + } + else { + usr->modifier1 = uwsgi_str_num(val, vallen); + } + } + else if (!uwsgi_strncmp("cores", 5, key, keylen)) { + if (vallen == 8) { + memcpy(&usr->cores, val, vallen); + } + else { + usr->cores = uwsgi_str_num(val, vallen); + } + } + else if (!uwsgi_strncmp("load", 4, key, keylen)) { + if (vallen == 8) { + memcpy(&usr->cores, val, vallen); + } + else { + usr->cores = uwsgi_str_num(val, vallen); + } } } @@ -898,8 +918,9 @@ void fastrouter_send_stats(int fd) { fprintf(output, "\t\t\"nodes\": [\n"); struct uwsgi_subscribe_node *s_node = s_slot->nodes; while(s_node) { - fprintf(output, "\t\t\t{\"name\": \"%.*s\", \"modifier1\": %d, \"modifier2\": %d, \"last_check\": %llu, \"requests\": %llu, \"tx\": %llu, \"ref\": %llu, \"death_mark\": %d}", s_node->len, s_node->name, s_node->modifier1, s_node->modifier2, - (unsigned long long) s_node->last_check, (unsigned long long) s_node->requests, (unsigned long long) s_node->transferred, (unsigned long long) s_node->reference, s_node->death_mark); + fprintf(output, "\t\t\t{\"name\": \"%.*s\", \"modifier1\": %d, \"modifier2\": %d, \"last_check\": %llu, \"requests\": %llu, \"tx\": %llu, \"cores\": %llu, \"load\": %llu, \"ref\": %llu, \"death_mark\": %d}", s_node->len, s_node->name, s_node->modifier1, s_node->modifier2, + (unsigned long long) s_node->last_check, (unsigned long long) s_node->requests, (unsigned long long) s_node->transferred, + (unsigned long long) s_node->cores, (unsigned long long) s_node->load, (unsigned long long) s_node->reference, s_node->death_mark); if (s_node->next) { fprintf(output, ",\n"); } diff --git a/subscription.c b/subscription.c index 46ba1e92..bf96a328 100644 --- a/subscription.c +++ b/subscription.c @@ -236,6 +236,8 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo node->transferred = 0; node->reference = 0; node->death_mark = 0; + node->cores = usr->cores; + node->load = usr->load; node->last_check = time(NULL); node->slot = current_slot; memcpy(node->name, usr->address, usr->address_len); @@ -275,6 +277,8 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo current_slot->nodes->death_mark = 0; current_slot->nodes->modifier1 = usr->modifier1; current_slot->nodes->modifier2 = usr->modifier2; + current_slot->nodes->cores = usr->cores; + current_slot->nodes->load = usr->load; memcpy(current_slot->nodes->name, usr->address, usr->address_len); current_slot->nodes->last_check = time(NULL); @@ -350,15 +354,13 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo } -void uwsgi_send_subscription(char *udp_address, char *key, size_t keysize, char *modifier1, size_t modifier1_len, uint8_t cmd) { +void uwsgi_send_subscription(char *udp_address, char *key, size_t keysize, uint8_t modifier1, uint8_t modifier2, uint8_t cmd) { + + uint64_t value64 = 0; if (!uwsgi.sockets) return; - size_t ssb_size = 4 + (2 + 3) + (2 + keysize) + (2 + 7) + (2 + strlen(uwsgi.sockets->name)); - - if (modifier1) { - ssb_size += (2 + 9) + (2 + modifier1_len); - } + size_t ssb_size = 4 + (2 + 3) + (2 + keysize) + (2 + 7) + (2 + strlen(uwsgi.sockets->name)) + (2+9 + 2+1) + (2+9 + 2+1) + (2+5 + 2+8) + (2+4 + 2+8); char *subscrbuf = uwsgi_malloc(ssb_size); // leave space for uwsgi header @@ -391,19 +393,57 @@ void uwsgi_send_subscription(char *udp_address, char *key, size_t keysize, char ssb+=ustrlen; // modifier1 = "modifier1" - if (modifier1) { - ustrlen = 9; - *ssb++ = (uint8_t) (ustrlen & 0xff); - *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); - memcpy(ssb, "modifier1", ustrlen); - ssb+=ustrlen; + ustrlen = 9; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + memcpy(ssb, "modifier1", ustrlen); + ssb+=ustrlen; - ustrlen = modifier1_len; - *ssb++ = (uint8_t) (ustrlen & 0xff); - *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); - memcpy(ssb, modifier1, ustrlen); - ssb+=ustrlen; - } + ustrlen = 1; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + *ssb++ = modifier1; + + // modifier2 = "modifier2" + ustrlen = 9; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + memcpy(ssb, "modifier2", ustrlen); + ssb+=ustrlen; + + ustrlen = 1; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + *ssb++ = modifier2; + + // cores = uwsgi.numproc * uwsgi.cores + ustrlen = 5; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + memcpy(ssb, "cores", ustrlen); + ssb+=ustrlen; + + value64 = uwsgi.numproc * uwsgi.cores; + ustrlen = 8; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + memcpy(ssb, &value64, 8); + ssb+=ustrlen; + + // load + ustrlen = 4; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + memcpy(ssb, "load", ustrlen); + ssb+=ustrlen; + + value64 = uwsgi.shared->load; + ustrlen = 8; + *ssb++ = (uint8_t) (ustrlen & 0xff); + *ssb++ = (uint8_t) ((ustrlen >>8) & 0xff); + memcpy(ssb, &value64, 8); + ssb+=ustrlen; + send_udp_message(224, cmd, udp_address, subscrbuf, ssb_size-4); free(subscrbuf); diff --git a/uwsgi.h b/uwsgi.h index e68901bd..d31199d2 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -1708,7 +1708,7 @@ struct uwsgi_shared { #ifdef __linux__ struct tcp_info ti; #endif - + uint64_t load; struct uwsgi_cron cron[MAX_CRONS]; int cron_cnt; }; @@ -2218,6 +2218,9 @@ struct uwsgi_subscribe_req { uint8_t modifier1; uint8_t modifier2; + + uint64_t cores; + uint64_t load; }; #ifndef _NO_UWSGI_RB @@ -2461,6 +2464,8 @@ struct uwsgi_subscribe_node { int death_mark; uint64_t reference; + uint64_t cores; + uint64_t load; struct uwsgi_subscribe_slot *slot; @@ -2511,7 +2516,7 @@ void manage_cluster_announce(char *, uint16_t, char *, uint16_t, void *); int uwsgi_read_response(int, struct uwsgi_header *, int, char **); char *uwsgi_simple_file_read(char *); -void uwsgi_send_subscription(char *, char *, size_t , char *, size_t, uint8_t); +void uwsgi_send_subscription(char *, char *, size_t , uint8_t, uint8_t , uint8_t); void uwsgi_subscribe(char *, uint8_t);