From 3c810e5ca80ee4d805a09b3f9d153c40ad4afc87 Mon Sep 17 00:00:00 2001 From: Unbit Date: Wed, 27 Feb 2013 01:15:56 +0100 Subject: [PATCH] updated carbon plugin --- core/exceptions.c | 1 - core/master.c | 2 +- core/stats.c | 25 ++++-- plugins/carbon/carbon.c | 77 +++++++++++-------- plugins/stats_pusher_mongodb/plugin.c | 2 +- .../stats_pusher_mongodb.cc | 2 +- uwsgi.h | 21 ++--- 7 files changed, 76 insertions(+), 54 deletions(-) diff --git a/core/exceptions.c b/core/exceptions.c index b9efe668..fdaf0a04 100644 --- a/core/exceptions.c +++ b/core/exceptions.c @@ -458,7 +458,6 @@ static void uwsgi_exception_handler_thread_loop(struct uwsgi_thread *ut) { } } } - uwsgi_log("STATUS = %d %d %d %d\n", uwsgi.status.gracefully_reloading, uwsgi.status.brutally_reloading, 0, 0); } } diff --git a/core/master.c b/core/master.c index bac54fbb..95bebcc2 100644 --- a/core/master.c +++ b/core/master.c @@ -411,7 +411,7 @@ int master_loop(char **argv, char **environ) { } - if (uwsgi.requested_stats_pushers) { + if (uwsgi.stats_pusher_instances) { if (!uwsgi_thread_new(uwsgi_stats_pusher_loop)) { uwsgi_log("!!! unable to spawn stats pusher thread !!!\n"); exit(1); diff --git a/core/stats.c b/core/stats.c index 98c68eee..5ede5cb2 100644 --- a/core/stats.c +++ b/core/stats.c @@ -374,7 +374,7 @@ struct uwsgi_stats_pusher *uwsgi_stats_pusher_get(char *name) { return usp; } -void uwsgi_stats_pusher_add(struct uwsgi_stats_pusher *pusher, char *arg) { +struct uwsgi_stats_pusher_instance *uwsgi_stats_pusher_add(struct uwsgi_stats_pusher *pusher, char *arg) { struct uwsgi_stats_pusher_instance *old_uspi = NULL, *uspi = uwsgi.stats_pusher_instances; while (uspi) { old_uspi = uspi; @@ -390,6 +390,8 @@ void uwsgi_stats_pusher_add(struct uwsgi_stats_pusher *pusher, char *arg) { else { uwsgi.stats_pusher_instances = uspi; } + + return uspi; } void uwsgi_stats_pusher_loop(struct uwsgi_thread *ut) { @@ -414,12 +416,17 @@ void uwsgi_stats_pusher_loop(struct uwsgi_thread *ut) { while (uspi) { int delta = uspi->freq ? uspi->freq : uwsgi.stats_pusher_default_freq; if ((uspi->last_run + delta) <= now) { - if (!us) { - us = uwsgi_master_generate_stats(); - if (!us) - goto next; + if (uspi->raw) { + uspi->pusher->func(uspi, now, NULL, 0); + } + else { + if (!us) { + us = uwsgi_master_generate_stats(); + if (!us) + goto next; + } + uspi->pusher->func(uspi, now, us->base, us->pos); } - uspi->pusher->func(uspi, us->base, us->pos); uspi->last_run = now; } next: @@ -457,7 +464,7 @@ void uwsgi_stats_pusher_setup() { } } -void uwsgi_register_stats_pusher(char *name, void (*func) (struct uwsgi_stats_pusher_instance *, char *, size_t)) { +struct uwsgi_stats_pusher *uwsgi_register_stats_pusher(char *name, void (*func) (struct uwsgi_stats_pusher_instance *, time_t, char *, size_t)) { struct uwsgi_stats_pusher *pusher = uwsgi.stats_pushers, *old_pusher = NULL; @@ -476,6 +483,8 @@ void uwsgi_register_stats_pusher(char *name, void (*func) (struct uwsgi_stats_pu else { uwsgi.stats_pushers = pusher; } + + return pusher; } struct uwsgi_stats_pusher_file_conf { @@ -484,7 +493,7 @@ struct uwsgi_stats_pusher_file_conf { char *separator; }; -void uwsgi_stats_pusher_file(struct uwsgi_stats_pusher_instance *uspi, char *json, size_t json_len) { +void uwsgi_stats_pusher_file(struct uwsgi_stats_pusher_instance *uspi, time_t now, char *json, size_t json_len) { struct uwsgi_stats_pusher_file_conf *uspic = (struct uwsgi_stats_pusher_file_conf *) uspi->data; if (!uspi->configured) { uspic = uwsgi_calloc(sizeof(struct uwsgi_stats_pusher_file_conf)); diff --git a/plugins/carbon/carbon.c b/plugins/carbon/carbon.c index 5d2a30a3..7324e14f 100644 --- a/plugins/carbon/carbon.c +++ b/plugins/carbon/carbon.c @@ -31,9 +31,10 @@ struct uwsgi_carbon { int max_retries; int retry_delay; char *root_node; + struct uwsgi_stats_pusher *pusher; } u_carbon; -struct uwsgi_option carbon_options[] = { +static struct uwsgi_option carbon_options[] = { {"carbon", required_argument, 0, "push statistics to the specified carbon server", uwsgi_opt_add_string_list, &u_carbon.servers, UWSGI_OPT_MASTER}, {"carbon-timeout", required_argument, 0, "set carbon connection timeout in seconds (default 3)", uwsgi_opt_set_int, &u_carbon.timeout, 0}, {"carbon-freq", required_argument, 0, "set carbon push frequency in seconds (default 60)", uwsgi_opt_set_int, &u_carbon.freq, 0}, @@ -46,8 +47,7 @@ struct uwsgi_option carbon_options[] = { }; - -void carbon_post_init() { +static void carbon_post_init() { int i; struct uwsgi_string_list *usl = u_carbon.servers; @@ -103,9 +103,14 @@ void carbon_post_init() { uwsgi_log("[carbon] carbon plugin started, %is frequency, %is timeout, max retries %i, retry delay %is\n", u_carbon.freq, u_carbon.timeout, u_carbon.max_retries, u_carbon.retry_delay); + + struct uwsgi_stats_pusher_instance *uspi = uwsgi_stats_pusher_add(u_carbon.pusher, NULL); + uspi->freq = u_carbon.freq; + // no need to generate the json + uspi->raw=1; } -int carbon_write(int *fd, char *fmt,...) { +static int carbon_write(int fd, char *fmt,...) { va_list ap; va_start(ap, fmt); @@ -117,15 +122,15 @@ int carbon_write(int *fd, char *fmt,...) { if (rlen < 1) return 0; - if (write(*fd, ptr, rlen) <= 0) { - uwsgi_error("write()"); + if (uwsgi_write_nb(fd, ptr, rlen, u_carbon.timeout)) { + uwsgi_error("carbon_write()"); return 0; } return 1; } -void carbon_push_stats(int retry_cycle) { +static void carbon_push_stats(int retry_cycle, time_t now) { struct carbon_server_list *usl = u_carbon.servers_data; int i; int fd; @@ -181,7 +186,7 @@ void carbon_push_stats(int retry_cycle) { unsigned long long worker_busyness = 0; unsigned long long total_harakiri = 0; - wok = carbon_write(&fd, "%s%s.%s.requests %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) uwsgi.workers[0].requests, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.requests %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) uwsgi.workers[0].requests, (unsigned long long) now); if (!wok) goto clear; for(i=1;i<=uwsgi.numproc;i++) { @@ -220,43 +225,43 @@ void carbon_push_stats(int retry_cycle) { //skip per worker metrics when disabled if (u_carbon.no_workers) continue; - wok = carbon_write(&fd, "%s%s.%s.worker%d.requests %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].requests, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.requests %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].requests, (unsigned long long) now); if (!wok) goto clear; if (uwsgi.shared->options[UWSGI_OPTION_MEMORY_DEBUG] == 1 || uwsgi.force_get_memusage) { - wok = carbon_write(&fd, "%s%s.%s.worker%d.rss_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].rss_size, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.rss_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].rss_size, (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.worker%d.vsz_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].vsz_size, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.vsz_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].vsz_size, (unsigned long long) now); if (!wok) goto clear; } - wok = carbon_write(&fd, "%s%s.%s.worker%d.avg_rt %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) avg_rt, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.avg_rt %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) avg_rt, (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.worker%d.tx %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].tx, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.tx %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].tx, (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.worker%d.busyness %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) worker_busyness, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.busyness %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) worker_busyness, (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.worker%d.harakiri %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].harakiri_count, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.worker%d.harakiri %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, i, (unsigned long long) uwsgi.workers[i].harakiri_count, (unsigned long long) now); if (!wok) goto clear; } if (uwsgi.shared->options[UWSGI_OPTION_MEMORY_DEBUG] == 1 || uwsgi.force_get_memusage) { - wok = carbon_write(&fd, "%s%s.%s.rss_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_rss, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.rss_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_rss, (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.vsz_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_vsz, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.vsz_size %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_vsz, (unsigned long long) now); if (!wok) goto clear; } - wok = carbon_write(&fd, "%s%s.%s.avg_rt %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) (active_workers ? total_avg_rt / active_workers : 0), (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.avg_rt %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) (active_workers ? total_avg_rt / active_workers : 0), (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.tx %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_tx, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.tx %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_tx, (unsigned long long) now); if (!wok) goto clear; if (active_workers > 0) { @@ -265,18 +270,18 @@ void carbon_push_stats(int retry_cycle) { } else { total_avg_busyness = 0; } - wok = carbon_write(&fd, "%s%s.%s.busyness %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_avg_busyness, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.busyness %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_avg_busyness, (unsigned long long) now); if (!wok) goto clear; - wok = carbon_write(&fd, "%s%s.%s.active_workers %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) active_workers, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.active_workers %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) active_workers, (unsigned long long) now); if (!wok) goto clear; if (uwsgi.cheaper) { - wok = carbon_write(&fd, "%s%s.%s.cheaped_workers %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) uwsgi.numproc - active_workers, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.cheaped_workers %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) uwsgi.numproc - active_workers, (unsigned long long) now); if (!wok) goto clear; } - wok = carbon_write(&fd, "%s%s.%s.harakiri %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_harakiri, (unsigned long long) uwsgi.current_time); + wok = carbon_write(fd, "%s%s.%s.harakiri %llu %llu\n", u_carbon.root_node, uwsgi.hostname, u_carbon.id, (unsigned long long) total_harakiri, (unsigned long long) now); if (!wok) goto clear; usl->healthy = 1; @@ -293,28 +298,34 @@ nxt: u_carbon.last_update -= u_carbon.timeout; } -void carbon_master_cycle() { +static void carbon_push(struct uwsgi_stats_pusher_instance *uspi, time_t now, char *json, size_t json_len) { - if (!u_carbon.servers) return; - - if (uwsgi.current_time - u_carbon.last_update >= u_carbon.freq || uwsgi.status.is_cleaning) { + if (u_carbon.need_retry && now >= u_carbon.next_retry) { + carbon_push_stats(1, now); + } + else { // update u_carbon.need_retry = 0; - carbon_push_stats(0); - } else if (u_carbon.need_retry && (uwsgi.current_time >= u_carbon.next_retry)) { - // retry failed servers - carbon_push_stats(1); + carbon_push_stats(0, now); + } } +static void carbon_cleanup() { + carbon_push_stats(0, uwsgi_now()); +} + +static void carbon_register() { + u_carbon.pusher = uwsgi_register_stats_pusher("carbon", carbon_push); +} struct uwsgi_plugin carbon_plugin = { .name = "carbon", - .master_cleanup = carbon_master_cycle, + .master_cleanup = carbon_cleanup, .options = carbon_options, - .master_cycle = carbon_master_cycle, + .on_load = carbon_register, .post_init = carbon_post_init, }; diff --git a/plugins/stats_pusher_mongodb/plugin.c b/plugins/stats_pusher_mongodb/plugin.c index c36c9abb..2a7f3b5c 100644 --- a/plugins/stats_pusher_mongodb/plugin.c +++ b/plugins/stats_pusher_mongodb/plugin.c @@ -1,6 +1,6 @@ #include "../../uwsgi.h" -void stats_pusher_mongodb(struct uwsgi_stats_pusher_instance *, char *, size_t); +void stats_pusher_mongodb(struct uwsgi_stats_pusher_instance *, time_t, char *, size_t); static void stats_pusher_mongodb_init(void) { uwsgi_register_stats_pusher("mongodb", stats_pusher_mongodb); diff --git a/plugins/stats_pusher_mongodb/stats_pusher_mongodb.cc b/plugins/stats_pusher_mongodb/stats_pusher_mongodb.cc index 2d49fd8c..06e4bc74 100644 --- a/plugins/stats_pusher_mongodb/stats_pusher_mongodb.cc +++ b/plugins/stats_pusher_mongodb/stats_pusher_mongodb.cc @@ -12,7 +12,7 @@ struct stats_pusher_mongodb_conf { }; -extern "C" void stats_pusher_mongodb(struct uwsgi_stats_pusher_instance *uspi, char *json, size_t json_len) { +extern "C" void stats_pusher_mongodb(struct uwsgi_stats_pusher_instance *uspi, time_t now, char *json, size_t json_len) { struct stats_pusher_mongodb_conf *spmc = (struct stats_pusher_mongodb_conf *) uspi->data; if (!uspi->configured) { diff --git a/uwsgi.h b/uwsgi.h index d1fb28d4..a95dad57 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -3359,7 +3359,7 @@ void uwsgi_reload(char **); struct uwsgi_stats_pusher { char *name; - void (*func) (struct uwsgi_stats_pusher_instance *, char *, size_t); + void (*func) (struct uwsgi_stats_pusher_instance *, time_t, char *, size_t); struct uwsgi_stats_pusher *next; }; @@ -3367,21 +3367,22 @@ void uwsgi_reload(char **); struct uwsgi_stats_pusher *pusher; char *arg; void *data; + int raw; int configured; int freq; time_t last_run; struct uwsgi_stats_pusher_instance *next; }; - struct uwsgi_thread; - void uwsgi_stats_pusher_loop(struct uwsgi_thread *); - void uwsgi_stats_pusher_file(struct uwsgi_stats_pusher_instance *, char *, size_t); - void uwsgi_stats_pusher_socket(struct uwsgi_stats_pusher_instance *, char *, size_t); +struct uwsgi_thread; +void uwsgi_stats_pusher_loop(struct uwsgi_thread *); +void uwsgi_stats_pusher_file(struct uwsgi_stats_pusher_instance *, time_t, char *, size_t); +void uwsgi_stats_pusher_socket(struct uwsgi_stats_pusher_instance *, time_t, char *, size_t); - void uwsgi_stats_pusher_setup(void); - void uwsgi_send_stats(int, struct uwsgi_stats *(*func) (void)); - struct uwsgi_stats *uwsgi_master_generate_stats(void); - void uwsgi_register_stats_pusher(char *, void (*)(struct uwsgi_stats_pusher_instance *, char *, size_t)); +void uwsgi_stats_pusher_setup(void); +void uwsgi_send_stats(int, struct uwsgi_stats *(*func) (void)); +struct uwsgi_stats *uwsgi_master_generate_stats(void); +struct uwsgi_stats_pusher * uwsgi_register_stats_pusher(char *, void (*)(struct uwsgi_stats_pusher_instance *, time_t, char *, size_t)); struct uwsgi_stats *uwsgi_stats_new(size_t); int uwsgi_stats_symbol(struct uwsgi_stats *, char); @@ -3863,6 +3864,8 @@ void uwsgi_exceptions_handler_thread_start(void); #define uwsgi_response_add_connection_close(x) uwsgi_response_add_header(x, "Connection", 10, "close", 5) #define uwsgi_response_add_content_type(x, y, z) uwsgi_response_add_header(x, "Content-Type", 12, y, z) +struct uwsgi_stats_pusher_instance *uwsgi_stats_pusher_add(struct uwsgi_stats_pusher *, char *); + #ifdef UWSGI_ZLIB #include int uwsgi_deflate_init(z_stream *, char *, size_t);