From 72f693f02a0bbf2b0e4504358b60a73cef162a0e Mon Sep 17 00:00:00 2001 From: "roberto@fiorenzo" Date: Tue, 29 Nov 2011 11:46:50 +0100 Subject: [PATCH] fixed subscription rr and added --fastrouter-stats --- plugins/fastrouter/fastrouter.c | 115 +++++++++++++++++++++++++++++++- subscription.c | 6 +- 2 files changed, 119 insertions(+), 2 deletions(-) diff --git a/plugins/fastrouter/fastrouter.c b/plugins/fastrouter/fastrouter.c index dc96a574..04400b7f 100644 --- a/plugins/fastrouter/fastrouter.c +++ b/plugins/fastrouter/fastrouter.c @@ -21,6 +21,7 @@ #define LONG_ARGS_FASTROUTER_SUBSCRIPTION_SLOT 150007 #define LONG_ARGS_FASTROUTER_USE_CODE_STRING 150008 #define LONG_ARGS_FASTROUTER_TOLERANCE 150009 +#define LONG_ARGS_FASTROUTER_STATS 150010 #define FASTROUTER_STATUS_FREE 0 #define FASTROUTER_STATUS_CONNECTING 1 @@ -33,6 +34,8 @@ #define del_check_timeout(x) rb_erase(&x->rbt, ufr.timeouts); #define del_timeout(x) rb_erase(&x->timeout->rbt, ufr.timeouts); free(x->timeout); +void fastrouter_send_stats(int); + struct uwsgi_fastrouter_socket { char *name; int fd; @@ -56,6 +59,8 @@ struct uwsgi_fastrouter { char *base; int base_len; + char *stats_server; + char *subscription_server; struct uwsgi_subscribe_slot *subscriptions; int subscription_regexp; @@ -126,6 +131,9 @@ struct option fastrouter_options[] = { {"fastrouter-subscription-slot", required_argument, 0, LONG_ARGS_FASTROUTER_SUBSCRIPTION_SLOT}, {"fastrouter-subscription-use-regexp", no_argument, &ufr.subscription_regexp, 1}, {"fastrouter-timeout", required_argument, 0, LONG_ARGS_FASTROUTER_TIMEOUT}, + {"fastrouter-stats", required_argument, 0, LONG_ARGS_FASTROUTER_STATS}, + {"fastrouter-stats-server", required_argument, 0, LONG_ARGS_FASTROUTER_STATS}, + {"fastrouter-ss", required_argument, 0, LONG_ARGS_FASTROUTER_STATS}, {0, 0, 0, 0}, }; @@ -308,6 +316,7 @@ void fastrouter_loop() { socklen_t solen = sizeof(int); int ufr_subserver = -1; + int ufr_stats_server = -1; for(i=0;i<2048;i++) { fr_table[i] = NULL; @@ -364,6 +373,23 @@ void fastrouter_loop() { //ufr.subscriptions_check = add_check_timeout(10); } + if (ufr.stats_server) { + char *tcp_port = strchr(ufr.stats_server, ':'); + if (tcp_port) { + // disable deferred accept for this socket + int current_defer_accept = uwsgi.no_defer_accept; + uwsgi.no_defer_accept = 1; + ufr_stats_server = bind_to_tcp(ufr.stats_server, uwsgi.listen_queue, tcp_port); + uwsgi.no_defer_accept = current_defer_accept; + } + else { + ufr_stats_server = bind_to_unix(ufr.stats_server, uwsgi.listen_queue, uwsgi.chmod_socket, uwsgi.abstract_socket); + } + + event_queue_add_fd_read(ufr.queue, ufr_stats_server); + uwsgi_log("*** FastRouter stats server enabled on %s fd: %d ***\n", ufr.stats_server, ufr_stats_server); + } + if (ufr.pattern) { init_magic_table(magic_table); } @@ -430,7 +456,10 @@ void fastrouter_loop() { continue; } - if (interesting_fd == ufr_subserver) { + if (interesting_fd == ufr_stats_server) { + fastrouter_send_stats(ufr_stats_server); + } + else if (interesting_fd == ufr_subserver) { len = recv(ufr_subserver, bbuf, 4096, 0); #ifdef UWSGI_EVENT_USE_PORT event_queue_add_fd_read(ufr.queue, ufr_subserver); @@ -763,6 +792,9 @@ int fastrouter_opt(int i, char *optarg) { case LONG_ARGS_FASTROUTER_SUBSCRIPTION_SERVER: ufr.subscription_server = optarg; return 1; + case LONG_ARGS_FASTROUTER_STATS: + ufr.stats_server = optarg; + return 1; case LONG_ARGS_FASTROUTER_EVENTS: ufr.nevents = atoi(optarg); return 1; @@ -813,3 +845,84 @@ struct uwsgi_plugin fastrouter_plugin = { .init = fastrouter_init, }; + +#define stats_send_llu(x, y) fprintf(output, x, (long long unsigned int) y) +#define stats_send(x, y) fprintf(output, x, y) + +void fastrouter_send_stats(int fd) { + + struct sockaddr_un client_src; + socklen_t client_src_len = 0; + int client_fd = accept(fd, (struct sockaddr *) &client_src, &client_src_len); + if (client_fd < 0) { + uwsgi_error("accept()"); + return; + } + + FILE *output = fdopen(client_fd, "w"); + if (!output) { + uwsgi_error("fdopen()"); + close(client_fd); + return; + } + + stats_send("{ \"version\": \"%s\",\n", UWSGI_VERSION); + + fprintf(output,"\"uid\": %d,\n", (int)(getuid())); + fprintf(output,"\"gid\": %d,\n", (int)(getgid())); + + char *cwd = uwsgi_get_cwd(); + stats_send("\"cwd\": \"%s\",\n", cwd); + free(cwd); + + fprintf(output, "\"fastrouter\": ["); + struct uwsgi_fastrouter_socket *uwsgi_sock = ufr.sockets; + while(uwsgi_sock) { + if (uwsgi_sock->next) { + stats_send("\"%s\",", uwsgi_sock->name); + } + else { + stats_send("\"%s\"", uwsgi_sock->name); + } + uwsgi_sock = uwsgi_sock->next; + } + fprintf(output, "],\n"); + + if (ufr.subscription_server) { + fprintf(output, "\"subscriptions\": [\n"); + struct uwsgi_subscribe_slot *s_slot = ufr.subscriptions; + while(s_slot) { + fprintf(output, "\t{ \"key\": \"%.*s\",\n", s_slot->keylen, s_slot->key); + fprintf(output, "\t\t\"hits\": %llu,\n", (unsigned long long) s_slot->hits); + 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}", 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); + if (s_node->next) { + fprintf(output, ",\n"); + } + else { + fprintf(output, "\n"); + } + s_node = s_node->next; + } + fprintf(output, "\t\t]\n"); + if (s_slot->next) { + fprintf(output, "\t},\n"); + } + else { + fprintf(output, "\t}\n"); + } + s_slot = s_slot->next; + } + fprintf(output, "],\n"); + } + + fprintf(output,"\"cheap\": %d\n", ufr.i_am_cheap); + + fprintf(output,"}\n"); + fclose(output); + +} + diff --git a/subscription.c b/subscription.c index d435c917..df704c69 100644 --- a/subscription.c +++ b/subscription.c @@ -103,7 +103,7 @@ struct uwsgi_subscribe_node *uwsgi_get_subscribe_node(struct uwsgi_subscribe_slo node = node->next; rr_pos++; } - current_slot->rr = 0; + current_slot->rr = 1; if (current_slot->nodes) { current_slot->nodes->reference++; } @@ -211,6 +211,8 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo node->len = usr->address_len; node->modifier1 = usr->modifier1; node->modifier2 = usr->modifier2; + node->requests = 0; + node->transferred = 0; node->reference = 0; node->death_mark = 0; node->last_check = time(NULL); @@ -247,6 +249,8 @@ struct uwsgi_subscribe_node *uwsgi_add_subscribe_node(struct uwsgi_subscribe_slo current_slot->nodes->slot = current_slot; current_slot->nodes->len = usr->address_len; current_slot->nodes->reference = 0; + current_slot->nodes->requests = 0; + current_slot->nodes->transferred = 0; current_slot->nodes->death_mark = 0; current_slot->nodes->modifier1 = usr->modifier1; current_slot->nodes->modifier2 = usr->modifier2;