@cryptotaxi247 / netdata-1 / commits / c98acefa3

allow each netdata host to have its own thread for streaming metrics

Costa Tsaousis (ktsaou) committed Feb 24, 2017 at 01:47 UTC c98acefa301f23d8c7d2857b8cde5fbb9c2326a9
9 files changed +224 -191
src/backends.c
+3 -5
@@ -147,11 +147,6 @@ void *backends_main(void *ptr) {
147 // ------------------------------------------------------------------------
148 // collect configuration options
149
150 - if(rrdpush_exclusive) {
151 - info("Backend is disabled - use the central netdata");
152 - goto cleanup;
153 - }
154 -
150 struct timeval timeout = {
151 .tv_sec = 0,
152 .tv_usec = 0
@@ -312,6 +307,9 @@ void *backends_main(void *ptr) {
307 rrd_rdlock();
308 RRDHOST *host;
309 rrdhost_foreach_read(host) {
310 + if(host->rrd_memory_mode == RRD_MEMORY_MODE_NONE)
311 + continue;
312 +
313 rrdhost_rdlock(host);
314
315 RRDSET *st;
src/health.c
+9 -10
@@ -15,16 +15,9 @@ inline char *health_config_dir(void) {
15 void health_init(void) {
16 debug(D_HEALTH, "Health configuration initializing");
17
18 - if(!rrdpush_exclusive) {
19 - if(!(default_health_enabled = config_get_boolean("health", "enabled", 1))) {
20 - debug(D_HEALTH, "Health is disabled.");
21 - return;
22 - }
23 - }
24 - else {
25 - info("Health is disabled - setup alarms at the central netdata.");
26 - config_set_boolean("health", "enabled", 0);
27 - default_health_enabled = 0;
18 + if(!(default_health_enabled = config_get_boolean("health", "enabled", 1))) {
19 + debug(D_HEALTH, "Health is disabled.");
20 + return;
21 }
22 }
23
@@ -32,6 +25,9 @@ void health_init(void) {
25 // re-load health configuration
26
27 void health_reload_host(RRDHOST *host) {
28 + if(unlikely(!host->health_enabled))
29 + return;
30 +
31 char *path = health_config_dir();
32
33 // free all running alarms
@@ -363,6 +359,9 @@ void *health_main(void *ptr) {
359
360 RRDHOST *host;
361 rrdhost_foreach_read(host) {
362 + if(unlikely(!host->health_enabled))
363 + continue;
364 +
365 if(unlikely(apply_hibernation_delay)) {
366
367 info("Postponing alarm checks for %ld seconds, on host '%s', due to boottime discrepancy (realtime dt: %ld, boottime dt: %ld)."
src/main.c
+16 -20
@@ -55,19 +55,16 @@ struct netdata_static_thread static_threads[] = {
55 {"plugins.d", NULL, NULL, 1, NULL, NULL, pluginsd_main},
56 {"web", NULL, NULL, 1, NULL, NULL, socket_listen_main_multi_threaded},
57 {"web-single-threaded", NULL, NULL, 0, NULL, NULL, socket_listen_main_single_threaded},
58 - {"central-netdata-push",NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
58 + {"push-metrics", NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
59 {NULL, NULL, NULL, 0, NULL, NULL, NULL}
60 };
61
62 void web_server_threading_selection(void) {
63 int multi_threaded = 0;
64 int single_threaded = 0;
65 - int rrdpush_thread = 1;
66 - int backends_thread = 1;
65
68 - if(rrdpush_exclusive) {
69 - backends_thread = 0;
70 - info("Web servers and backends thread are disabled - use the central netdata.");
66 + if(default_rrdpush_exclusive) {
67 + info("Web server is disabled - use the remote netdata.");
68 }
69 else {
70 multi_threaded = config_get_boolean("global", "multi threaded web server", 1);
@@ -81,15 +78,9 @@ void web_server_threading_selection(void) {
78
79 if(static_threads[i].start_routine == socket_listen_main_single_threaded)
80 static_threads[i].enabled = single_threaded;
84 -
85 - if(static_threads[i].start_routine == rrdpush_sender_thread)
86 - static_threads[i].enabled = rrdpush_thread;
87 -
88 - if(static_threads[i].start_routine == backends_main)
89 - static_threads[i].enabled = backends_thread;
81 }
82
92 - if(rrdpush_exclusive)
83 + if(default_rrdpush_exclusive)
84 return;
85
86 web_client_timeout = (int) config_get_number("global", "disconnect idle web clients after seconds", DEFAULT_DISCONNECT_IDLE_WEB_CLIENTS_AFTER_SECONDS);
@@ -694,7 +685,7 @@ int main(int argc, char **argv) {
685
686
687 // --------------------------------------------------------------------
697 - // find we need to send data to a central netdata
688 + // find we need to send data to another netdata
689
690 rrdpush_init();
691
@@ -702,7 +693,7 @@ int main(int argc, char **argv) {
693 // --------------------------------------------------------------------
694 // get default memory mode for the database
695
705 - if(rrdpush_exclusive) {
696 + if(default_rrdpush_exclusive) {
697 default_rrd_memory_mode = RRD_MEMORY_MODE_NONE;
698 config_set("global", "memory mode", rrd_memory_mode_name(default_rrd_memory_mode));
699 }
@@ -713,14 +704,14 @@ int main(int argc, char **argv) {
704 // --------------------------------------------------------------------
705 // get default database size
706
716 - if(rrdpush_exclusive) {
707 + if(default_rrdpush_exclusive) {
708 default_rrd_history_entries = 10;
709 config_set_number("global", "history", default_rrd_history_entries);
710 }
711 else
721 - default_rrd_history_entries = (int) config_get_number("global", "history", align_entries_to_pagesize(RRD_DEFAULT_HISTORY_ENTRIES));
712 + default_rrd_history_entries = (int) config_get_number("global", "history", align_entries_to_pagesize(default_rrd_memory_mode, RRD_DEFAULT_HISTORY_ENTRIES));
713
723 - long h = align_entries_to_pagesize(default_rrd_history_entries);
714 + long h = align_entries_to_pagesize(default_rrd_memory_mode, default_rrd_history_entries);
715 if(h != default_rrd_history_entries) {
716 config_set_number("global", "history", h);
717 default_rrd_history_entries = (int)h;
@@ -844,11 +835,16 @@ int main(int argc, char **argv) {
835 // --------------------------------------------------------------------
836 // create the listening sockets
837
847 - if(!check_config && !rrdpush_exclusive) {
838 + if(!check_config && !default_rrdpush_exclusive)
839 + create_listen_sockets();
840 +
841 +
842 + // --------------------------------------------------------------------
843 + // load the aggregated host configuration file
844 + {
845 char filename[FILENAME_MAX + 1];
846 snprintfz(filename, FILENAME_MAX, "%s/aggregated_hosts.conf", netdata_configured_config_dir);
847 appconfig_load(&stream_config, filename, 0);
851 - create_listen_sockets();
848 }
849 }
850
src/registry_init.c
+1 -1
@@ -4,7 +4,7 @@ int registry_init(void) {
4 char filename[FILENAME_MAX + 1];
5
6 // registry enabled?
7 - if(!rrdpush_exclusive) {
7 + if(!default_rrdpush_exclusive) {
8 registry.enabled = config_get_boolean("registry", "enabled", 0);
9 }
10 else {
src/rrd.h
+12 -1
@@ -341,6 +341,17 @@ struct rrdhost {
341 int rrd_update_every; // the update frequency of the host
342 int rrd_history_entries; // the number of history entries for the host's charts
343
344 + int rrdpush_enabled; // 1 when this host sends metrics to another netdata
345 + int rrdpush_exclusive; // 1 when this host is exclusively sending metrics without a database
346 + volatile int rrdpush_connected; // 1 when the sender is ready to push metrics
347 + volatile int rrdpush_spawn; // 1 when the sender thread has been spawn
348 + volatile int rrdpush_error_shown; // 1 when we have logged a communication error
349 + int rrdpush_socket; // the fd of the socket to the remote host, or -1
350 + pthread_t rrdpush_thread; // the sender thread
351 + pthread_mutex_t rrdpush_mutex; // exclusive access to rrdpush_buffer
352 + int rrdpush_pipe[2]; // collector to sender thread communication
353 + BUFFER *rrdpush_buffer; // collector fills it, sender sends them
354 +
355 int health_enabled; // 1 when this host has health enabled
356 time_t health_delay_up_to; // a timestamp to delay alarms processing up to
357 RRD_MEMORY_MODE rrd_memory_mode; // the memory more for the charts of this host
@@ -526,7 +537,7 @@ extern int rrddim_unhide(RRDSET *st, const char *id);
537 extern collected_number rrddim_set_by_pointer(RRDSET *st, RRDDIM *rd, collected_number value);
538 extern collected_number rrddim_set(RRDSET *st, const char *id, collected_number value);
539
529 -extern long align_entries_to_pagesize(long entries);
540 +extern long align_entries_to_pagesize(RRD_MEMORY_MODE mode, long entries);
541
542
543 // ----------------------------------------------------------------------------
src/rrdhost.c
+33 -22
@@ -73,18 +73,25 @@ RRDHOST *rrdhost_create(const char *hostname,
73 host->rrd_update_every = update_every;
74 host->rrd_history_entries = entries;
75 host->rrd_memory_mode = memory_mode;
76 - host->health_enabled = health_enabled;
76 + host->health_enabled = (memory_mode == RRD_MEMORY_MODE_NONE)? 0 : health_enabled;
77 + host->rrdpush_enabled = default_rrdpush_enabled;
78 + host->rrdpush_exclusive = default_rrdpush_exclusive;
79
78 - pthread_rwlock_init(&(host->rrdhost_rwlock), NULL);
80 + host->rrdpush_pipe[0] = -1;
81 + host->rrdpush_pipe[1] = -1;
82 + host->rrdpush_socket = -1;
83 +
84 + pthread_mutex_init(&host->rrdpush_mutex, NULL);
85 + pthread_rwlock_init(&host->rrdhost_rwlock, NULL);
86
87 rrdhost_init_hostname(host, hostname);
88 rrdhost_init_machine_guid(host, guid);
89 rrdhost_init_os(host, os);
90
84 - avl_init_lock(&(host->rrdset_root_index), rrdset_compare);
91 + avl_init_lock(&(host->rrdset_root_index), rrdset_compare);
92 avl_init_lock(&(host->rrdset_root_index_name), rrdset_compare_name);
86 - avl_init_lock(&(host->rrdfamily_root_index), rrdfamily_compare);
87 - avl_init_lock(&(host->variables_root_index), rrdvar_compare);
93 + avl_init_lock(&(host->rrdfamily_root_index), rrdfamily_compare);
94 + avl_init_lock(&(host->variables_root_index), rrdvar_compare);
95
96 // ------------------------------------------------------------------------
97 // initialize health variables
@@ -110,12 +117,9 @@ RRDHOST *rrdhost_create(const char *hostname,
117 if(!localhost) {
118 // this is localhost
119
113 - host->cache_dir = strdupz(netdata_configured_cache_dir);
120 + host->cache_dir = strdupz(netdata_configured_cache_dir);
121 host->varlib_dir = strdupz(netdata_configured_varlib_dir);
122
116 - snprintfz(filename, FILENAME_MAX, "%s/health/health-log.db", host->varlib_dir);
117 - host->health_log_filename = strdupz(config_get("health", "health db file", filename));
118 -
123 }
124 else {
125 // this is not localhost - append our GUID to localhost path
@@ -136,18 +140,18 @@ RRDHOST *rrdhost_create(const char *hostname,
140 int r = mkdir(host->varlib_dir, 0775);
141 if(r != 0 && errno != EEXIST)
142 error("Host '%s': cannot create directory '%s'", host->hostname, host->varlib_dir);
139 - }
143
141 - snprintfz(filename, FILENAME_MAX, "%s/health", host->varlib_dir);
142 - int r = mkdir(filename, 0775);
143 - if(r != 0 && errno != EEXIST)
144 - error("Host '%s': cannot create directory '%s'", host->hostname, filename);
145 -
146 - snprintfz(filename, FILENAME_MAX, "%s/health/health-log.db", host->varlib_dir);
147 - host->health_log_filename = strdupz(filename);
144 + snprintfz(filename, FILENAME_MAX, "%s/health", host->varlib_dir);
145 + r = mkdir(filename, 0775);
146 + if(r != 0 && errno != EEXIST)
147 + error("Host '%s': cannot create directory '%s'", host->hostname, filename);
148 + }
149
150 }
151
152 + snprintfz(filename, FILENAME_MAX, "%s/health/health-log.db", host->varlib_dir);
153 + host->health_log_filename = strdupz(config_get("health", "health db file", filename));
154 +
155 snprintfz(filename, FILENAME_MAX, "%s/alarm-notify.sh", netdata_configured_plugins_dir);
156 host->health_default_exec = strdupz(config_get("health", "script to execute on alarm", filename));
157 host->health_default_recipient = strdup("root");
@@ -156,12 +160,14 @@ RRDHOST *rrdhost_create(const char *hostname,
160 // ------------------------------------------------------------------------
161 // load health configuration
162
159 - health_alarm_log_load(host);
160 - health_alarm_log_open(host);
163 + if(host->health_enabled) {
164 + health_alarm_log_load(host);
165 + health_alarm_log_open(host);
166
162 - rrdhost_wrlock(host);
163 - health_readdir(host, health_config_dir());
164 - rrdhost_unlock(host);
167 + rrdhost_wrlock(host);
168 + health_readdir(host, health_config_dir());
169 + rrdhost_unlock(host);
170 + }
171
172
173 // ------------------------------------------------------------------------
@@ -312,6 +318,11 @@ void rrdhost_free(RRDHOST *host) {
318 // ------------------------------------------------------------------------
319 // free it
320
321 + if(host->rrdpush_spawn) {
322 + pthread_cancel(host->rrdpush_thread);
323 + rrdpush_sender_cleanup(host);
324 + }
325 +
326 freez(host->os);
327 freez(host->cache_dir);
328 freez(host->varlib_dir);
src/rrdpush.c
+138 -121
@@ -1,7 +1,7 @@
1 #include "common.h"
2
3 -int rrdpush_enabled = 0;
4 -int rrdpush_exclusive = 1;
3 +int default_rrdpush_enabled = 0;
4 +int default_rrdpush_exclusive = 1;
5
6 static char *remote_netdata_config = NULL;
7 static char *api_key = NULL;
@@ -15,19 +15,6 @@ static char *api_key = NULL;
15 // that there are more data ready to be sent
16 #define PIPE_READ 0
17 #define PIPE_WRITE 1
18 -int rrdpush_pipe[2] = { -1, -1 };
19 -
20 -// a buffer used to store data to be sent.
21 -// the format is the same as external plugins.
22 -static BUFFER *rrdpush_buffer = NULL;
23 -
24 -// locking to get exclusive access to shared resources
25 -// (rrdpush_pipe[PIPE_WRITE], rrdpush_buffer
26 -static pthread_mutex_t rrdpush_mutex = PTHREAD_MUTEX_INITIALIZER;
27 -
28 -// if the streaming thread is connected to a remote netdata
29 -// this is set to 1, otherwise 0.
30 -static volatile int rrdpush_connected = 0;
18
19 // to have the remote netdata re-sync the charts
20 // to its current clock, we send for this many
@@ -35,8 +22,8 @@ static volatile int rrdpush_connected = 0;
22 // this is for the first iterations of each chart
23 static unsigned int remote_clock_resync_iterations = 60;
24
38 -#define rrdpush_lock() pthread_mutex_lock(&rrdpush_mutex)
39 -#define rrdpush_unlock() pthread_mutex_unlock(&rrdpush_mutex)
25 +#define rrdpush_lock(host) pthread_mutex_lock(&((host)->rrdpush_mutex))
26 +#define rrdpush_unlock(host) pthread_mutex_unlock(&((host)->rrdpush_mutex))
27
28 // checks if the current chart definition has been sent
29 static inline int need_to_send_chart_definition(RRDSET *st) {
@@ -50,7 +37,7 @@ static inline int need_to_send_chart_definition(RRDSET *st) {
37
38 // sends the current chart definition
39 static inline void send_chart_definition(RRDSET *st) {
53 - buffer_sprintf(rrdpush_buffer, "CHART '%s' '%s' '%s' '%s' '%s' '%s' '%s' %ld %d\n"
40 + buffer_sprintf(st->rrdhost->rrdpush_buffer, "CHART '%s' '%s' '%s' '%s' '%s' '%s' '%s' %ld %d\n"
41 , st->id
42 , st->name
43 , st->title
@@ -64,7 +51,7 @@ static inline void send_chart_definition(RRDSET *st) {
51
52 RRDDIM *rd;
53 rrddim_foreach_read(rd, st) {
67 - buffer_sprintf(rrdpush_buffer, "DIMENSION '%s' '%s' '%s' " COLLECTED_NUMBER_FORMAT " " COLLECTED_NUMBER_FORMAT " '%s %s'\n"
54 + buffer_sprintf(st->rrdhost->rrdpush_buffer, "DIMENSION '%s' '%s' '%s' " COLLECTED_NUMBER_FORMAT " " COLLECTED_NUMBER_FORMAT " '%s %s'\n"
55 , rd->id
56 , rd->name
57 , rrd_algorithm_name(rd->algorithm)
@@ -79,69 +66,69 @@ static inline void send_chart_definition(RRDSET *st) {
66
67 // sends the current chart dimensions
68 static inline void send_chart_metrics(RRDSET *st) {
82 - buffer_sprintf(rrdpush_buffer, "BEGIN %s %llu\n", st->id, (st->counter_done > remote_clock_resync_iterations)?st->usec_since_last_update:0);
69 + buffer_sprintf(st->rrdhost->rrdpush_buffer, "BEGIN %s %llu\n", st->id, (st->counter_done > remote_clock_resync_iterations)?st->usec_since_last_update:0);
70
71 RRDDIM *rd;
72 rrddim_foreach_read(rd, st) {
73 if(rrddim_flag_check(rd, RRDDIM_FLAG_UPDATED) && rrddim_flag_check(rd, RRDDIM_FLAG_EXPOSED))
87 - buffer_sprintf(rrdpush_buffer, "SET %s = " COLLECTED_NUMBER_FORMAT "\n"
74 + buffer_sprintf(st->rrdhost->rrdpush_buffer, "SET %s = " COLLECTED_NUMBER_FORMAT "\n"
75 , rd->id
76 , rd->collected_value
77 );
78 }
79
93 - buffer_strcat(rrdpush_buffer, "END\n");
80 + buffer_strcat(st->rrdhost->rrdpush_buffer, "END\n");
81 }
82
83 // resets all the chart, so that their definitions
84 // will be resent to the central netdata
98 -static void reset_all_charts(void) {
99 - rrd_rdlock();
100 -
101 - RRDHOST *host;
102 - rrdhost_foreach_read(host) {
103 - rrdhost_rdlock(host);
85 +static void reset_all_charts(RRDHOST *host) {
86 + rrdhost_rdlock(host);
87
105 - RRDSET *st;
106 - rrdset_foreach_read(st, host) {
88 + RRDSET *st;
89 + rrdset_foreach_read(st, host) {
90
108 - // make it re-align the current time
109 - // on the remote host
110 - st->counter_done = 0;
91 + // make it re-align the current time
92 + // on the remote host
93 + st->counter_done = 0;
94
112 - rrdset_rdlock(st);
95 + rrdset_rdlock(st);
96
114 - RRDDIM *rd;
115 - rrddim_foreach_read(rd, st)
116 - rrddim_flag_clear(rd, RRDDIM_FLAG_EXPOSED);
97 + RRDDIM *rd;
98 + rrddim_foreach_read(rd, st)
99 + rrddim_flag_clear(rd, RRDDIM_FLAG_EXPOSED);
100
118 - rrdset_unlock(st);
119 - }
120 - rrdhost_unlock(host);
101 + rrdset_unlock(st);
102 }
122 - rrd_unlock();
103 +
104 + rrdhost_unlock(host);
105 }
106
107 +void rrdpush_sender_thread_spawn(RRDHOST *host);
108 +
109 void rrdset_done_push(RRDSET *st) {
126 - static int error_shown = 0;
110 + RRDHOST *host = st->rrdhost;
111
112 if(unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ENABLED)))
113 return;
114
131 - rrdpush_lock();
115 + rrdpush_lock(host);
116
133 - if(unlikely(!rrdpush_buffer || !rrdpush_connected)) {
134 - if(unlikely(!error_shown))
117 + if(unlikely(host->rrdpush_enabled && !host->rrdpush_spawn))
118 + rrdpush_sender_thread_spawn(host);
119 +
120 + if(unlikely(!host->rrdpush_buffer || !host->rrdpush_connected)) {
121 + if(unlikely(!host->rrdpush_error_shown))
122 error("STREAM [send]: not ready - discarding collected metrics.");
123
137 - error_shown = 1;
124 + host->rrdpush_error_shown = 1;
125
139 - rrdpush_unlock();
126 + rrdpush_unlock(host);
127 return;
128 }
142 - else if(unlikely(error_shown)) {
129 + else if(unlikely(host->rrdpush_error_shown)) {
130 error("STREAM [send]: ready - sending metrics...");
144 - error_shown = 0;
131 + host->rrdpush_error_shown = 0;
132 }
133
134 rrdset_rdlock(st);
@@ -152,38 +139,74 @@ void rrdset_done_push(RRDSET *st) {
139 rrdset_unlock(st);
140
141 // signal the sender there are more data
155 - if(write(rrdpush_pipe[PIPE_WRITE], " ", 1) == -1)
142 + if(write(host->rrdpush_pipe[PIPE_WRITE], " ", 1) == -1)
143 error("STREAM [send]: cannot write to internal pipe");
144
158 - rrdpush_unlock();
145 + rrdpush_unlock(host);
146 }
147
161 -static inline void rrdpush_flush(void) {
162 - rrdpush_lock();
163 - if(buffer_strlen(rrdpush_buffer))
164 - error("STREAM [send]: discarding %zu bytes of metrics already in the buffer.", buffer_strlen(rrdpush_buffer));
148 +static inline void rrdpush_flush(RRDHOST *host) {
149 + rrdpush_lock(host);
150 + if(buffer_strlen(host->rrdpush_buffer))
151 + error("STREAM [send]: discarding %zu bytes of metrics already in the buffer.", buffer_strlen(host->rrdpush_buffer));
152
166 - buffer_flush(rrdpush_buffer);
167 - reset_all_charts();
168 - rrdpush_unlock();
153 + buffer_flush(host->rrdpush_buffer);
154 + reset_all_charts(host);
155 + rrdpush_unlock(host);
156 }
157
158 int rrdpush_init() {
172 - rrdpush_enabled = config_get_boolean("stream", "enabled", rrdpush_enabled);
173 - rrdpush_exclusive = config_get_boolean("stream", "exclusive", rrdpush_exclusive);
159 + default_rrdpush_enabled = config_get_boolean("stream", "enabled", default_rrdpush_enabled);
160 + default_rrdpush_exclusive = config_get_boolean("stream", "exclusive", default_rrdpush_exclusive);
161 remote_netdata_config = config_get("stream", "stream metrics to", "");
162 api_key = config_get("stream", "api key", "");
163
177 - if(!rrdpush_enabled || !remote_netdata_config || !*remote_netdata_config || !api_key || !*api_key) {
178 - rrdpush_enabled = 0;
179 - rrdpush_exclusive = 0;
164 + if(!default_rrdpush_enabled || !remote_netdata_config || !*remote_netdata_config || !api_key || !*api_key) {
165 + default_rrdpush_enabled = 0;
166 + default_rrdpush_exclusive = 0;
167 }
168
182 - return rrdpush_enabled;
169 + return default_rrdpush_enabled;
170 +}
171 +
172 +static inline void rrdpush_sender_lock(RRDHOST *host) {
173 + if(pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL) != 0)
174 + error("STREAM [send]: cannot set pthread cancel state to DISABLE.");
175 +
176 + rrdpush_lock(host);
177 +}
178 +
179 +static inline void rrdpush_sender_unlock(RRDHOST *host) {
180 + if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
181 + error("STREAM [send]: cannot set pthread cancel state to DISABLE.");
182 +
183 + rrdpush_unlock(host);
184 +}
185 +
186 +void rrdpush_sender_cleanup(RRDHOST *host) {
187 + rrdpush_lock(host);
188 +
189 + host->rrdpush_connected = 0;
190 +
191 + if(host->rrdpush_socket != -1) close(host->rrdpush_socket);
192 +
193 + // close the pipe
194 + if(host->rrdpush_pipe[PIPE_READ] != -1) close(host->rrdpush_pipe[PIPE_READ]);
195 + if(host->rrdpush_pipe[PIPE_WRITE] != -1) close(host->rrdpush_pipe[PIPE_WRITE]);
196 + host->rrdpush_pipe[PIPE_READ] = -1;
197 + host->rrdpush_pipe[PIPE_WRITE] = -1;
198 +
199 + buffer_free(host->rrdpush_buffer);
200 + host->rrdpush_buffer = NULL;
201 +
202 + host->rrdpush_spawn = 0;
203 + host->rrdpush_enabled = 0;
204 +
205 + rrdpush_unlock(host);
206 }
207
208 void *rrdpush_sender_thread(void *ptr) {
186 - struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
209 + RRDHOST *host = (RRDHOST *)ptr;
210
211 info("STREAM [send]: thread created (task id %d)", gettid());
212
@@ -198,16 +221,15 @@ void *rrdpush_sender_thread(void *ptr) {
221 size_t max_size = (size_t)config_get_number("stream", "buffer size bytes", 1024 * 1024);
222 unsigned int reconnect_delay = (unsigned int)config_get_number("stream", "reconnect delay seconds", 5);
223 remote_clock_resync_iterations = (unsigned int)config_get_number("stream", "initial clock resync iterations", remote_clock_resync_iterations);
201 - int sock = -1;
224 char connected_to[CONNECTED_TO_SIZE + 1] = "";
225
204 - if(!rrdpush_enabled || !remote_netdata_config || !*remote_netdata_config || !api_key || !*api_key)
226 + if(!host->rrdpush_enabled || !remote_netdata_config || !*remote_netdata_config || !api_key || !*api_key)
227 goto cleanup;
228
229 // initialize rrdpush globals
208 - rrdpush_buffer = buffer_create(1);
209 - rrdpush_connected = 0;
210 - if(pipe(rrdpush_pipe) == -1) fatal("STREAM [send]: cannot create required pipe.");
230 + host->rrdpush_buffer = buffer_create(1);
231 + host->rrdpush_connected = 0;
232 + if(pipe(host->rrdpush_pipe) == -1) fatal("STREAM [send]: cannot create required pipe.");
233
234 // initialize local variables
235 size_t begin = 0;
@@ -226,18 +248,17 @@ void *rrdpush_sender_thread(void *ptr) {
248 ifd = &fds[0];
249 ofd = &fds[1];
250
229 - for(;;) {
230 - if(netdata_exit) break;
251 + for(; host->rrdpush_enabled && !netdata_exit ;) {
252
232 - if(unlikely(sock == -1)) {
253 + if(unlikely(host->rrdpush_socket == -1)) {
254 // stop appending data into rrdpush_buffer
255 // they will be lost, so there is no point to do it
235 - rrdpush_connected = 0;
256 + host->rrdpush_connected = 0;
257
258 info("STREAM [send to %s]: connecting...", remote_netdata_config);
238 - sock = connect_to_one_of(remote_netdata_config, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
259 + host->rrdpush_socket = connect_to_one_of(remote_netdata_config, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
260
240 - if(unlikely(sock == -1)) {
261 + if(unlikely(host->rrdpush_socket == -1)) {
262 error("STREAM [send to %s]: failed to connect", remote_netdata_config);
263 sleep(reconnect_delay);
264 continue;
@@ -258,9 +279,9 @@ void *rrdpush_sender_thread(void *ptr) {
279 , program_version
280 );
281
261 - if(send_timeout(sock, http, strlen(http), 0, timeout) == -1) {
262 - close(sock);
263 - sock = -1;
282 + if(send_timeout(host->rrdpush_socket, http, strlen(http), 0, timeout) == -1) {
283 + close(host->rrdpush_socket);
284 + host->rrdpush_socket = -1;
285 error("STREAM [send to %s]: failed to send http header to netdata", connected_to);
286 sleep(reconnect_delay);
287 continue;
@@ -268,17 +289,17 @@ void *rrdpush_sender_thread(void *ptr) {
289
290 info("STREAM [send to %s]: waiting response from remote netdata...", connected_to);
291
271 - if(recv_timeout(sock, http, 1000, 0, timeout) == -1) {
272 - close(sock);
273 - sock = -1;
292 + if(recv_timeout(host->rrdpush_socket, http, 1000, 0, timeout) == -1) {
293 + close(host->rrdpush_socket);
294 + host->rrdpush_socket = -1;
295 error("STREAM [send to %s]: failed to initialize communication", connected_to);
296 sleep(reconnect_delay);
297 continue;
298 }
299
300 if(strncmp(http, "STREAM", 6)) {
280 - close(sock);
281 - sock = -1;
301 + close(host->rrdpush_socket);
302 + host->rrdpush_socket = -1;
303 error("STREAM [send to %s]: server is not replying properly.", connected_to);
304 sleep(reconnect_delay);
305 continue;
@@ -286,23 +307,23 @@ void *rrdpush_sender_thread(void *ptr) {
307
308 info("STREAM [send to %s]: established communication - sending metrics...", connected_to);
309
289 - if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
310 + if(fcntl(host->rrdpush_socket, F_SETFL, O_NONBLOCK) < 0)
311 error("STREAM [send to %s]: cannot set non-blocking mode for socket.", connected_to);
312
292 - rrdpush_flush();
313 + rrdpush_flush(host);
314 sent_connection = 0;
315
316 // allow appending data into rrdpush_buffer
296 - rrdpush_connected = 1;
317 + host->rrdpush_connected = 1;
318 }
319
299 - ifd->fd = rrdpush_pipe[PIPE_READ];
320 + ifd->fd = host->rrdpush_pipe[PIPE_READ];
321 ifd->events = POLLIN;
322 ifd->revents = 0;
323
303 - ofd->fd = sock;
324 + ofd->fd = host->rrdpush_socket;
325 ofd->revents = 0;
305 - if(begin < buffer_strlen(rrdpush_buffer)) {
326 + if(begin < buffer_strlen(host->rrdpush_buffer)) {
327 ofd->events = POLLOUT;
328 fdmax = 2;
329 }
@@ -320,8 +341,8 @@ void *rrdpush_sender_thread(void *ptr) {
341 continue;
342
343 error("STREAM [send to %s]: failed to poll().", connected_to);
323 - close(sock);
324 - sock = -1;
344 + close(host->rrdpush_socket);
345 + host->rrdpush_socket = -1;
346 break;
347 }
348 else if(unlikely(!retval)) {
@@ -331,39 +352,39 @@ void *rrdpush_sender_thread(void *ptr) {
352
353 if(ifd->revents & POLLIN) {
354 char buffer[1000 + 1];
334 - if(read(rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
355 + if(read(host->rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
356 error("STREAM [send to %s]: cannot read from internal pipe.", connected_to);
357 }
358
338 - if(ofd->revents & POLLOUT && begin < buffer_strlen(rrdpush_buffer)) {
339 - rrdpush_lock();
340 - ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], buffer_strlen(rrdpush_buffer) - begin, MSG_DONTWAIT);
359 + if(ofd->revents & POLLOUT && begin < buffer_strlen(host->rrdpush_buffer)) {
360 + rrdpush_sender_lock(host);
361 + ssize_t ret = send(host->rrdpush_socket, &host->rrdpush_buffer->buffer[begin], buffer_strlen(host->rrdpush_buffer) - begin, MSG_DONTWAIT);
362 if(ret == -1) {
363 if(errno != EAGAIN && errno != EINTR) {
364 error("STREAM [send to %s]: failed to send metrics - closing connection - we have sent %zu bytes on this connection.", connected_to, sent_connection);
344 - close(sock);
345 - sock = -1;
365 + close(host->rrdpush_socket);
366 + host->rrdpush_socket = -1;
367 }
368 }
369 else {
370 sent_connection += ret;
371 sent_bytes += ret;
372 begin += ret;
352 - if(begin == buffer_strlen(rrdpush_buffer)) {
353 - buffer_flush(rrdpush_buffer);
373 + if(begin == buffer_strlen(host->rrdpush_buffer)) {
374 + buffer_flush(host->rrdpush_buffer);
375 begin = 0;
376 }
377 }
357 - rrdpush_unlock();
378 + rrdpush_sender_unlock(host);
379 }
380
381 // protection from overflow
361 - if(rrdpush_buffer->len > max_size) {
382 + if(host->rrdpush_buffer->len > max_size) {
383 errno = 0;
363 - error("STREAM [send to %s]: too many data pending - buffer is %zu bytes long, %zu unsent - we have sent %zu bytes in total, %zu on this connection. Closing connection to flush the data.", connected_to, rrdpush_buffer->len, rrdpush_buffer->len - begin, sent_bytes, sent_connection);
364 - if(sock != -1) {
365 - close(sock);
366 - sock = -1;
384 + error("STREAM [send to %s]: too many data pending - buffer is %zu bytes long, %zu unsent - we have sent %zu bytes in total, %zu on this connection. Closing connection to flush the data.", connected_to, host->rrdpush_buffer->len, host->rrdpush_buffer->len - begin, sent_bytes, sent_connection);
385 + if(host->rrdpush_socket != -1) {
386 + close(host->rrdpush_socket);
387 + host->rrdpush_socket = -1;
388 }
389 }
390 }
@@ -371,22 +392,8 @@ void *rrdpush_sender_thread(void *ptr) {
392 cleanup:
393 debug(D_WEB_CLIENT, "STREAM [send]: sending thread exits.");
394
374 - // make sure the data collection threads do not write data
375 - rrdpush_connected = 0;
395 + rrdpush_sender_cleanup(host);
396
377 - // close the pipe
378 - if(rrdpush_pipe[PIPE_READ] != -1) close(rrdpush_pipe[PIPE_READ]);
379 - if(rrdpush_pipe[PIPE_WRITE] != -1) close(rrdpush_pipe[PIPE_WRITE]);
380 -
381 - // close the socket
382 - if(sock != -1) close(sock);
383 -
384 - rrdpush_lock();
385 - buffer_free(rrdpush_buffer);
386 - rrdpush_buffer = NULL;
387 - rrdpush_unlock();
388 -
389 - static_thread->enabled = 0;
397 pthread_exit(NULL);
398 return NULL;
399 }
@@ -531,6 +538,16 @@ static inline int rrdpush_receive_validate_api_key(const char *key) {
538 return appconfig_get_boolean(&stream_config, key, "enabled", 0);
539 }
540
541 +void rrdpush_sender_thread_spawn(RRDHOST *host) {
542 + if(pthread_create(&host->rrdpush_thread, NULL, rrdpush_sender_thread, (void *)host))
543 + error("STREAM [send for host %s]: failed to create new thread for client.", host->hostname);
544 +
545 + else if(pthread_detach(host->rrdpush_thread))
546 + error("STREAM [send for host %s]: cannot request detach newly created thread.", host->hostname);
547 +
548 + host->rrdpush_spawn = 1;
549 +}
550 +
551 int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url) {
552 (void)host;
553
src/rrdpush.h
+3 -2
@@ -1,13 +1,14 @@
1 #ifndef NETDATA_RRDPUSH_H
2 #define NETDATA_RRDPUSH_H
3
4 -extern int rrdpush_enabled;
5 -extern int rrdpush_exclusive;
4 +extern int default_rrdpush_enabled;
5 +extern int default_rrdpush_exclusive;
6
7 extern int rrdpush_init();
8 extern void rrdset_done_push(RRDSET *st);
9 extern void *rrdpush_sender_thread(void *ptr);
10
11 extern int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url);
12 +extern void rrdpush_sender_cleanup(RRDHOST *host);
13
14 #endif //NETDATA_RRDPUSH_H
src/rrdset.c
+9 -9
@@ -195,16 +195,16 @@ void rrdset_reset(RRDSET *st) {
195 // ----------------------------------------------------------------------------
196 // RRDSET - helpers for rrdset_create()
197
198 -inline long align_entries_to_pagesize(long entries) {
199 - if(rrdpush_exclusive)
200 - return entries;
198 +inline long align_entries_to_pagesize(RRD_MEMORY_MODE mode, long entries) {
199 + if(unlikely(entries < 5)) entries = 5;
200 + if(unlikely(entries > RRD_HISTORY_ENTRIES_MAX)) entries = RRD_HISTORY_ENTRIES_MAX;
201
202 - if(entries < 5) entries = 5;
203 - if(entries > RRD_HISTORY_ENTRIES_MAX) entries = RRD_HISTORY_ENTRIES_MAX;
202 + if(unlikely(mode == RRD_MEMORY_MODE_NONE || mode == RRD_MEMORY_MODE_RAM))
203 + return entries;
204
205 long page = (size_t)sysconf(_SC_PAGESIZE);
206 long size = sizeof(RRDDIM) + entries * sizeof(storage_number);
207 - if(size % page) {
207 + if(unlikely(size % page)) {
208 size -= (size % page);
209 size += page;
210
@@ -322,7 +322,7 @@ RRDSET *rrdset_create(RRDHOST *host, const char *type, const char *id, const cha
322 // get the options from the config, we need to create it
323
324 long rentries = config_get_number(config_section, "history", host->rrd_history_entries);
325 - long entries = align_entries_to_pagesize(rentries);
325 + long entries = align_entries_to_pagesize(host->rrd_memory_mode, rentries);
326 if(entries != rentries) entries = config_set_number(config_section, "history", entries);
327
328 if(host->rrd_memory_mode == RRD_MEMORY_MODE_NONE && entries != rentries)
@@ -628,7 +628,7 @@ static inline void rrdset_done_push_int(RRDSET *st) {
628 void rrdset_done(RRDSET *st) {
629 if(unlikely(netdata_exit)) return;
630
631 - if(unlikely(rrdpush_exclusive)) {
631 + if(unlikely(st->rrdhost->rrdpush_exclusive)) {
632 rrdset_done_push_int(st);
633 return;
634 }
@@ -1218,7 +1218,7 @@ void rrdset_done(RRDSET *st) {
1218 if(unlikely(pthread_setcancelstate(pthreadoldcancelstate, NULL) != 0))
1219 error("Cannot set pthread cancel state to RESTORE (%d).", pthreadoldcancelstate);
1220
1221 - if(unlikely(rrdpush_enabled))
1221 + if(unlikely(st->rrdhost->rrdpush_enabled))
1222 rrdset_done_push_int(st);
1223 }
1224