X-Git-Url: https://git.octo.it/?a=blobdiff_plain;f=src%2Fplugin.c;h=23ea9017cce27aa71370c32b41d7496b0944c837;hb=58ecdb7ff9a4d90e87daecca05b9cb4a43070fd7;hp=2daeea9fd67e42e8ef36c5b6eee4e185542ed8e1;hpb=5962b08efed2a0e5d438e3829ce4abcf555ffaae;p=collectd.git diff --git a/src/plugin.c b/src/plugin.c index 2daeea9f..23ea9017 100644 --- a/src/plugin.c +++ b/src/plugin.c @@ -23,10 +23,14 @@ #include +#if HAVE_PTHREAD_H +# include +#endif + +#include "common.h" #include "plugin.h" #include "configfile.h" #include "utils_llist.h" -#include "utils_debug.h" /* * Private structures @@ -36,6 +40,7 @@ struct read_func_s int wait_time; int wait_left; int (*callback) (void); + enum { DONE = 0, TODO = 1, ACTIVE = 2 } needs_read; }; typedef struct read_func_s read_func_t; @@ -47,10 +52,15 @@ static llist_t *list_read; static llist_t *list_write; static llist_t *list_shutdown; static llist_t *list_data_set; +static llist_t *list_log; static char *plugindir = NULL; -char hostname[DATA_MAX_NAME_LEN] = "localhost"; +static int read_loop = 1; +static pthread_mutex_t read_lock = PTHREAD_MUTEX_INITIALIZER; +static pthread_cond_t read_cond = PTHREAD_COND_INITIALIZER; +static pthread_t *read_threads = NULL; +static int read_threads_num = 0; /* * Static functions @@ -113,7 +123,7 @@ static int plugin_load_file (char *file) lt_dlhandle dlh; void (*reg_handle) (void); - DBG ("file = %s", file); + DEBUG ("file = %s", file); lt_dlinit (); lt_dlerror (); /* clear errors */ @@ -122,14 +132,14 @@ static int plugin_load_file (char *file) { const char *error = lt_dlerror (); - syslog (LOG_ERR, "lt_dlopen failed: %s", error); - DBG ("lt_dlopen failed: %s", error); + ERROR ("lt_dlopen failed: %s", error); + fprintf (stderr, "lt_dlopen failed: %s\n", error); return (1); } if ((reg_handle = (void (*) (void)) lt_dlsym (dlh, "module_register")) == NULL) { - syslog (LOG_WARNING, "Couldn't find symbol ``module_register'' in ``%s'': %s\n", + WARNING ("Couldn't find symbol ``module_register'' in ``%s'': %s\n", file, lt_dlerror ()); lt_dlclose (dlh); return (-1); @@ -140,6 +150,130 @@ static int plugin_load_file (char *file) return (0); } +static void *plugin_read_thread (void *args) +{ + llentry_t *le; + read_func_t *rf; + int status; + int done; + + pthread_mutex_lock (&read_lock); + + while (read_loop != 0) + { + le = llist_head (list_read); + done = 0; + + while ((read_loop != 0) && (le != NULL)) + { + rf = (read_func_t *) le->value; + + if (rf->needs_read != TODO) + { + le = le->next; + continue; + } + + /* We will do this read function */ + rf->needs_read = ACTIVE; + + DEBUG ("[thread #%5lu] plugin: plugin_read_thread: Handling %s", + (unsigned long int) pthread_self (), le->key); + pthread_mutex_unlock (&read_lock); + + status = rf->callback (); + done++; + + if (status != 0) + { + if (rf->wait_time < interval_g) + rf->wait_time = interval_g; + rf->wait_left = rf->wait_time; + rf->wait_time = rf->wait_time * 2; + if (rf->wait_time > 86400) + rf->wait_time = 86400; + + NOTICE ("read-function of plugin `%s' " + "failed. Will suspend it for %i " + "seconds.", le->key, rf->wait_left); + } + else + { + rf->wait_left = 0; + rf->wait_time = interval_g; + } + + pthread_mutex_lock (&read_lock); + + rf->needs_read = DONE; + le = le->next; + } /* while (le != NULL) */ + + if ((read_loop != 0) && (done == 0)) + { + DEBUG ("[thread #%5lu] plugin: plugin_read_thread: Waiting on read_cond.", + (unsigned long int) pthread_self ()); + pthread_cond_wait (&read_cond, &read_lock); + } + } /* while (read_loop) */ + + pthread_mutex_unlock (&read_lock); + + pthread_exit (NULL); +} /* void *plugin_read_thread */ + +static void start_threads (int num) +{ + int i; + + if (read_threads != NULL) + return; + + read_threads = (pthread_t *) calloc (num, sizeof (pthread_t)); + if (read_threads == NULL) + { + ERROR ("plugin: start_threads: calloc failed."); + return; + } + + read_threads_num = 0; + for (i = 0; i < num; i++) + { + if (pthread_create (read_threads + read_threads_num, NULL, + plugin_read_thread, NULL) == 0) + { + read_threads_num++; + } + else + { + ERROR ("plugin: start_threads: pthread_create failed."); + return; + } + } /* for (i) */ +} /* void start_threads */ + +static void stop_threads (void) +{ + int i; + + pthread_mutex_lock (&read_lock); + read_loop = 0; + DEBUG ("plugin: stop_threads: Signalling `read_cond'"); + pthread_cond_broadcast (&read_cond); + pthread_mutex_unlock (&read_lock); + + for (i = 0; i < read_threads_num; i++) + { + if (pthread_join (read_threads[i], NULL) != 0) + { + ERROR ("plugin: stop_threads: pthread_join failed."); + } + read_threads[i] = (pthread_t) 0; + } + sfree (read_threads); + read_threads_num = 0; +} /* void stop_threads */ + /* * Public functions */ @@ -151,7 +285,11 @@ void plugin_set_dir (const char *dir) if (dir == NULL) plugindir = NULL; else if ((plugindir = strdup (dir)) == NULL) - syslog (LOG_ERR, "strdup failed: %s", strerror (errno)); + { + char errbuf[1024]; + ERROR ("strdup failed: %s", + sstrerror (errno, errbuf, sizeof (errbuf))); + } } #define BUFSIZE 512 @@ -166,7 +304,7 @@ int plugin_load (const char *type) struct stat statbuf; struct dirent *de; - DBG ("type = %s", type); + DEBUG ("type = %s", type); dir = plugin_get_dir (); ret = 1; @@ -175,14 +313,16 @@ int plugin_load (const char *type) * type when matching the filename */ if (snprintf (typename, BUFSIZE, "%s.so", type) >= BUFSIZE) { - syslog (LOG_WARNING, "snprintf: truncated: `%s.so'", type); + WARNING ("snprintf: truncated: `%s.so'", type); return (-1); } typename_len = strlen (typename); if ((dh = opendir (dir)) == NULL) { - syslog (LOG_ERR, "opendir (%s): %s", dir, strerror (errno)); + char errbuf[1024]; + ERROR ("opendir (%s): %s", dir, + sstrerror (errno, errbuf, sizeof (errbuf))); return (-1); } @@ -193,13 +333,15 @@ int plugin_load (const char *type) if (snprintf (filename, BUFSIZE, "%s/%s", dir, de->d_name) >= BUFSIZE) { - syslog (LOG_WARNING, "snprintf: truncated: `%s/%s'", dir, de->d_name); + WARNING ("snprintf: truncated: `%s/%s'", dir, de->d_name); continue; } if (lstat (filename, &statbuf) == -1) { - syslog (LOG_WARNING, "stat %s: %s", filename, strerror (errno)); + char errbuf[1024]; + WARNING ("stat %s: %s", filename, + sstrerror (errno, errbuf, sizeof (errbuf))); continue; } else if (!S_ISREG (statbuf.st_mode)) @@ -214,6 +356,10 @@ int plugin_load (const char *type) ret = 0; break; } + else + { + fprintf (stderr, "Unable to load plugin %s.\n", type); + } } closedir (dh); @@ -232,6 +378,12 @@ int plugin_register_config (const char *name, return (0); } /* int plugin_register_config */ +int plugin_register_complex_config (const char *type, + int (*callback) (oconfig_item_t *)) +{ + return (cf_register_complex (type, callback)); +} /* int plugin_register_complex_config */ + int plugin_register_init (const char *name, int (*callback) (void)) { @@ -246,15 +398,17 @@ int plugin_register_read (const char *name, rf = (read_func_t *) malloc (sizeof (read_func_t)); if (rf == NULL) { - syslog (LOG_ERR, "plugin_register_read: malloc failed: %s", - strerror (errno)); + char errbuf[1024]; + ERROR ("plugin_register_read: malloc failed: %s", + sstrerror (errno, errbuf, sizeof (errbuf))); return (-1); } memset (rf, '\0', sizeof (read_func_t)); - rf->wait_time = atoi (COLLECTD_STEP); + rf->wait_time = interval_g; rf->wait_left = 0; rf->callback = callback; + rf->needs_read = DONE; return (register_callback (&list_read, name, (void *) rf)); } /* int plugin_register_read */ @@ -273,9 +427,53 @@ int plugin_register_shutdown (char *name, int plugin_register_data_set (const data_set_t *ds) { - return (register_callback (&list_data_set, ds->type, (void *) ds)); + data_set_t *ds_copy; + int i; + + if ((list_data_set != NULL) + && (llist_search (list_data_set, ds->type) != NULL)) + { + NOTICE ("Replacing DS `%s' with another version.", ds->type); + plugin_unregister_data_set (ds->type); + } + + ds_copy = (data_set_t *) malloc (sizeof (data_set_t)); + if (ds_copy == NULL) + return (-1); + memcpy(ds_copy, ds, sizeof (data_set_t)); + + ds_copy->ds = (data_source_t *) malloc (sizeof (data_source_t) + * ds->ds_num); + if (ds_copy->ds == NULL) + { + free (ds_copy); + return (-1); + } + + for (i = 0; i < ds->ds_num; i++) + memcpy (ds_copy->ds + i, ds->ds + i, sizeof (data_source_t)); + + return (register_callback (&list_data_set, ds->type, (void *) ds_copy)); } /* int plugin_register_data_set */ +int plugin_register_log (char *name, + void (*callback) (int priority, const char *msg)) +{ + return (register_callback (&list_log, name, (void *) callback)); +} /* int plugin_register_log */ + +int plugin_unregister_config (const char *name) +{ + cf_unregister (name); + return (0); +} /* int plugin_unregister_config */ + +int plugin_unregister_complex_config (const char *name) +{ + cf_unregister_complex (name); + return (0); +} /* int plugin_unregister_complex_config */ + int plugin_unregister_init (const char *name) { return (plugin_unregister (list_init, name)); @@ -283,7 +481,6 @@ int plugin_unregister_init (const char *name) int plugin_unregister_read (const char *name) { - return (plugin_unregister (list_read, name)); llentry_t *e; e = llist_search (list_read, name); @@ -310,15 +507,47 @@ int plugin_unregister_shutdown (const char *name) int plugin_unregister_data_set (const char *name) { - return (plugin_unregister (list_data_set, name)); + llentry_t *e; + data_set_t *ds; + + if (list_data_set == NULL) + return (-1); + + e = llist_search (list_data_set, name); + + if (e == NULL) + return (-1); + + llist_remove (list_data_set, e); + ds = (data_set_t *) e->value; + llentry_destroy (e); + + sfree (ds->ds); + sfree (ds); + + return (0); +} /* int plugin_unregister_data_set */ + +int plugin_unregister_log (const char *name) +{ + return (plugin_unregister (list_log, name)); } void plugin_init_all (void) { int (*callback) (void); llentry_t *le; + int status; - gethostname (hostname, sizeof (hostname)); + /* Start read-threads */ + if (list_read != NULL) + { + const char *rt; + int num; + rt = global_option_get ("ReadThreads"); + num = atoi (rt); + start_threads ((num > 0) ? num : 5); + } if (list_init == NULL) return; @@ -326,8 +555,18 @@ void plugin_init_all (void) le = llist_head (list_init); while (le != NULL) { - callback = le->value; - (*callback) (); + callback = (int (*) (void)) le->value; + status = (*callback) (); + + if (status != 0) + { + ERROR ("Initialization of plugin `%s' " + "failed with status %i. " + "Plugin will be unloaded.", + le->key, status); + /* FIXME: Unload _all_ functions */ + plugin_unregister_read (le->key); + } le = le->next; } @@ -337,47 +576,37 @@ void plugin_read_all (const int *loop) { llentry_t *le; read_func_t *rf; - int status; - int step; if (list_read == NULL) return; - step = atoi (COLLECTD_STEP); + pthread_mutex_lock (&read_lock); le = llist_head (list_read); - while ((*loop == 0) && (le != NULL)) + while (le != NULL) { rf = (read_func_t *) le->value; - if (rf->wait_left > 0) - rf->wait_left -= step; - if (rf->wait_left > 0) + if (rf->needs_read != DONE) { le = le->next; continue; } - status = rf->callback (); - if (status != 0) - { - rf->wait_left = rf->wait_time; - rf->wait_time = rf->wait_time * 2; - if (rf->wait_time > 86400) - rf->wait_time = 86400; - - syslog (LOG_NOTICE, "read-function of plugin `%s' " - "failed. Will syspend it for %i " - "seconds.", le->key, rf->wait_left); - } - else + if (rf->wait_left > 0) + rf->wait_left -= interval_g; + + if (rf->wait_left <= 0) { - rf->wait_left = 0; - rf->wait_time = step; + rf->needs_read = TODO; } le = le->next; - } /* while ((*loop == 0) && (le != NULL)) */ + } + + DEBUG ("plugin: plugin_read_all: Signalling `read_cond'"); + pthread_cond_broadcast (&read_cond); + pthread_mutex_unlock (&read_lock); } /* void plugin_read_all */ void plugin_shutdown_all (void) @@ -385,61 +614,117 @@ void plugin_shutdown_all (void) int (*callback) (void); llentry_t *le; + stop_threads (); + if (list_shutdown == NULL) return; le = llist_head (list_shutdown); while (le != NULL) { - callback = le->value; - (*callback) (); + callback = (int (*) (void)) le->value; + /* Advance the pointer before calling the callback allows + * shutdown functions to unregister themselves. If done the + * other way around the memory `le' points to will be freed + * after callback returns. */ le = le->next; + + (*callback) (); } } /* void plugin_shutdown_all */ -int plugin_dispatch_values (const char *name, const value_list_t *vl) +int plugin_dispatch_values (const char *name, value_list_t *vl) { int (*callback) (const data_set_t *, const value_list_t *); data_set_t *ds; llentry_t *le; - if (list_write == NULL) + if ((list_write == NULL) || (list_data_set == NULL)) return (-1); le = llist_search (list_data_set, name); if (le == NULL) { - DBG ("No such dataset registered: %s", name); + DEBUG ("No such dataset registered: %s", name); return (-1); } ds = (data_set_t *) le->value; - DBG ("time = %u; host = %s; " + DEBUG ("plugin: plugin_dispatch_values: time = %u; interval = %i; " + "host = %s; " "plugin = %s; plugin_instance = %s; " "type = %s; type_instance = %s;", - (unsigned int) vl->time, vl->host, + (unsigned int) vl->time, vl->interval, + vl->host, vl->plugin, vl->plugin_instance, ds->type, vl->type_instance); +#if COLLECT_DEBUG + assert (ds->ds_num == vl->values_len); +#else + if (ds->ds_num != vl->values_len) + { + ERROR ("plugin: ds->type = %s: (ds->ds_num = %i) != " + "(vl->values_len = %i)", + ds->type, ds->ds_num, vl->values_len); + return (-1); + } +#endif + + escape_slashes (vl->host, sizeof (vl->host)); + escape_slashes (vl->plugin, sizeof (vl->plugin)); + escape_slashes (vl->plugin_instance, sizeof (vl->plugin_instance)); + escape_slashes (vl->type_instance, sizeof (vl->type_instance)); + le = llist_head (list_write); while (le != NULL) { - callback = le->value; + callback = (int (*) (const data_set_t *, const value_list_t *)) le->value; (*callback) (ds, vl); le = le->next; } return (0); -} +} /* int plugin_dispatch_values */ + +void plugin_log (int level, const char *format, ...) +{ + char msg[512]; + va_list ap; + + void (*callback) (int, const char *); + llentry_t *le; + + if (list_log == NULL) + return; + +#if !COLLECT_DEBUG + if (level >= LOG_DEBUG) + return; +#endif + + va_start (ap, format); + vsnprintf (msg, 512, format, ap); + msg[511] = '\0'; + va_end (ap); + + le = llist_head (list_log); + while (le != NULL) + { + callback = (void (*) (int, const char *)) le->value; + (*callback) (level, msg); + + le = le->next; + } +} /* void plugin_log */ void plugin_complain (int level, complain_t *c, const char *format, ...) { char message[512]; va_list ap; - int step; if (c->delay > 0) { @@ -447,25 +732,22 @@ void plugin_complain (int level, complain_t *c, const char *format, ...) return; } - step = atoi (COLLECTD_STEP); - assert (step > 0); - - if (c->interval < step) - c->interval = step; + if (c->interval < interval_g) + c->interval = interval_g; else c->interval *= 2; if (c->interval > 86400) c->interval = 86400; - c->delay = c->interval / step; + c->delay = c->interval / interval_g; va_start (ap, format); vsnprintf (message, 512, format, ap); message[511] = '\0'; va_end (ap); - syslog (level, message); + plugin_log (level, message); } void plugin_relief (int level, complain_t *c, const char *format, ...) @@ -483,5 +765,22 @@ void plugin_relief (int level, complain_t *c, const char *format, ...) message[511] = '\0'; va_end (ap); - syslog (level, message); + plugin_log (level, message); } + +const data_set_t *plugin_get_ds (const char *name) +{ + data_set_t *ds; + llentry_t *le; + + le = llist_search (list_data_set, name); + if (le == NULL) + { + DEBUG ("No such dataset registered: %s", name); + return (NULL); + } + + ds = (data_set_t *) le->value; + + return (ds); +} /* data_set_t *plugin_get_ds */