updated carbon plugin

This commit is contained in:
Unbit
2013-02-27 01:15:56 +01:00
parent 61a1ca4ca4
commit 3c810e5ca8
7 changed files with 76 additions and 54 deletions
-1
View File
@@ -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);
}
}
+1 -1
View File
@@ -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);
+17 -8
View File
@@ -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));
+44 -33
View File
@@ -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,
};
+1 -1
View File
@@ -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);
@@ -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) {
+12 -9
View File
@@ -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 <zlib.h>
int uwsgi_deflate_init(z_stream *, char *, size_t);