subscription optimizations

This commit is contained in:
roberto@oneiric64
2011-12-04 13:45:18 +01:00
parent c7bde15f39
commit 0b010124ef
5 changed files with 97 additions and 27 deletions
+5 -3
View File
@@ -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) {
+2
View File
@@ -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()));
+25 -4
View File
@@ -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");
}
+58 -18
View File
@@ -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);
+7 -2
View File
@@ -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);