@cryptotaxi247 / netdata-1 / commits / f1fed1161

allow metrics streaming to work in parallel with local database; propagate O/S type to central netdata; allow machine_guid=localhost to use the same host

Costa Tsaousis (ktsaou) committed Feb 23, 2017 at 00:09 UTC f1fed11619c1a4afeeb7ddda81f5ee6526a37bd4
17 files changed +220 -117
src/appconfig.c
+7 -5
@@ -503,11 +503,13 @@ void appconfig_generate(struct config *root, BUFFER *wb, int only_changed)
503
504 appconfig_wrlock(root);
505 for(co = root->sections; co ; co = co->next) {
506 - if(!strcmp(co->name, "global") ||
507 - !strcmp(co->name, "plugins") ||
508 - !strcmp(co->name, "registry") ||
509 - !strcmp(co->name, "health") ||
510 - !strcmp(co->name, "backend"))
506 + if(!strcmp(co->name, "global")
507 + || !strcmp(co->name, "plugins")
508 + || !strcmp(co->name, "registry")
509 + || !strcmp(co->name, "health")
510 + || !strcmp(co->name, "backend")
511 + || !strcmp(co->name, "stream")
512 + )
513 pri = 0;
514 else if(!strncmp(co->name, "plugin:", 7)) pri = 1;
515 else pri = 2;
src/backends.c
+2 -2
@@ -147,7 +147,7 @@ void *backends_main(void *ptr) {
147 // ------------------------------------------------------------------------
148 // collect configuration options
149
150 - if(central_netdata_to_push_data) {
150 + if(rrdpush_exclusive) {
151 info("Backend is disabled - use the central netdata");
152 goto cleanup;
153 }
@@ -399,7 +399,7 @@ void *backends_main(void *ptr) {
399 usec_t start_ut = now_monotonic_usec();
400 size_t reconnects = 0;
401
402 - sock = connect_to_one_of(destination, default_port, &timeout, &reconnects);
402 + sock = connect_to_one_of(destination, default_port, &timeout, &reconnects, NULL, 0);
403
404 chart_backend_reconnects += reconnects;
405 chart_backend_latency += now_monotonic_usec() - start_ut;
src/common.h
-1
@@ -226,7 +226,6 @@ extern char *netdata_configured_cache_dir;
226 extern char *netdata_configured_varlib_dir;
227 extern char *netdata_configured_home_dir;
228 extern char *netdata_configured_host_prefix;
229 -extern char *central_netdata_to_push_data;
229
230 extern void netdata_fix_chart_id(char *s);
231 extern void netdata_fix_chart_name(char *s);
src/health.c
+1 -1
@@ -15,7 +15,7 @@ inline char *health_config_dir(void) {
15 void health_init(void) {
16 debug(D_HEALTH, "Health configuration initializing");
17
18 - if(!central_netdata_to_push_data) {
18 + if(!rrdpush_exclusive) {
19 if(!(default_health_enabled = config_get_boolean("health", "enabled", 1))) {
20 debug(D_HEALTH, "Health is disabled.");
21 return;
src/main.c
+16 -16
@@ -2,8 +2,6 @@
2
3 extern void *cgroups_main(void *ptr);
4
5 -char *central_netdata_to_push_data = NULL;
6 -
5 void netdata_cleanup_and_exit(int ret) {
6 netdata_exit = 1;
7
@@ -64,15 +62,16 @@ struct netdata_static_thread static_threads[] = {
62 void web_server_threading_selection(void) {
63 int multi_threaded = 0;
64 int single_threaded = 0;
67 - int central_thread = 0;
65 + int rrdpush_thread = 1;
66 + int backends_thread = 1;
67
69 - if(!central_netdata_to_push_data) {
70 - multi_threaded = config_get_boolean("global", "multi threaded web server", 1);
71 - single_threaded = !multi_threaded;
68 + if(rrdpush_exclusive) {
69 + backends_thread = 0;
70 + info("Web servers and backends thread are disabled - use the central netdata.");
71 }
72 else {
74 - central_thread = 1;
75 - info("Web servers are disabled - use the central netdata.");
73 + multi_threaded = config_get_boolean("global", "multi threaded web server", 1);
74 + single_threaded = !multi_threaded;
75 }
76
77 int i;
@@ -84,10 +83,13 @@ void web_server_threading_selection(void) {
83 static_threads[i].enabled = single_threaded;
84
85 if(static_threads[i].start_routine == central_netdata_push_thread)
87 - static_threads[i].enabled = central_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;
90 }
91
90 - if(central_netdata_to_push_data)
92 + if(rrdpush_exclusive)
93 return;
94
95 web_client_timeout = (int) config_get_number("global", "disconnect idle web clients after seconds", DEFAULT_DISCONNECT_IDLE_WEB_CLIENTS_AFTER_SECONDS);
@@ -694,15 +696,13 @@ int main(int argc, char **argv) {
696 // --------------------------------------------------------------------
697 // find we need to send data to a central netdata
698
697 - central_netdata_to_push_data = config_get("global", "central netdata to send all data", "");
698 - if(central_netdata_to_push_data && !*central_netdata_to_push_data)
699 - central_netdata_to_push_data = NULL;
699 + rrdpush_init();
700
701
702 // --------------------------------------------------------------------
703 // get default memory mode for the database
704
705 - if(central_netdata_to_push_data) {
705 + if(rrdpush_exclusive) {
706 default_rrd_memory_mode = RRD_MEMORY_MODE_RAM;
707 config_set("global", "memory mode", rrd_memory_mode_name(default_rrd_memory_mode));
708 }
@@ -713,7 +713,7 @@ int main(int argc, char **argv) {
713 // --------------------------------------------------------------------
714 // get default database size
715
716 - if(central_netdata_to_push_data) {
716 + if(rrdpush_exclusive) {
717 default_rrd_history_entries = 10;
718 config_set_number("global", "history", default_rrd_history_entries);
719 }
@@ -844,7 +844,7 @@ int main(int argc, char **argv) {
844 // --------------------------------------------------------------------
845 // create the listening sockets
846
847 - if(!check_config && !central_netdata_to_push_data) {
847 + if(!check_config && !rrdpush_exclusive) {
848 char filename[FILENAME_MAX + 1];
849 snprintfz(filename, FILENAME_MAX, "%s/aggregated_hosts.conf", netdata_configured_config_dir);
850 appconfig_load(&stream_config, filename, 0);
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(!central_netdata_to_push_data) {
7 + if(!rrdpush_exclusive) {
8 registry.enabled = config_get_boolean("registry", "enabled", 0);
9 }
10 else {
src/rrd.h
+5 -1
@@ -367,6 +367,10 @@ struct rrdhost {
367 // are created or renamed, that match them
368 RRDCALCTEMPLATE *templates;
369
370 + char *os; // the O/S type of the host
371 + volatile size_t use_counter; // when remote hosts are streaming to this
372 + // host, this is the counter of connected clients
373 +
374 // health / alarm settings
375 char *health_default_exec;
376 char *health_default_recipient;
@@ -406,7 +410,7 @@ extern pthread_rwlock_t rrd_rwlock;
410 extern void rrd_init(char *hostname);
411
412 extern RRDHOST *rrdhost_find(const char *guid, uint32_t hash);
409 -extern RRDHOST *rrdhost_find_or_create(const char *hostname, const char *guid, int update_every, int history, RRD_MEMORY_MODE mode, int health_enabled);
413 +extern RRDHOST *rrdhost_find_or_create(const char *hostname, const char *guid, const char *os, int update_every, int history, RRD_MEMORY_MODE mode, int health_enabled);
414
415 #ifdef NETDATA_INTERNAL_CHECKS
416 extern void rrdhost_check_wrlock_int(RRDHOST *host, const char *file, const char *function, const unsigned long line);
src/rrd2json.c
+1 -1
@@ -92,7 +92,7 @@ void rrd_stats_api_v1_charts(RRDHOST *host, BUFFER *wb)
92 ",\n\t\"charts\": {"
93 , host->hostname
94 , program_version
95 - , os_type
95 + , host->os
96 , host->rrd_update_every
97 , host->rrd_history_entries
98 );
src/rrdhost.c
+11 -2
@@ -43,6 +43,11 @@ static inline void rrdhost_init_hostname(RRDHOST *host, const char *hostname) {
43 host->hash_hostname = simple_hash(host->hostname);
44 }
45
46 +static inline void rrdhost_init_os(RRDHOST *host, const char *os) {
47 + freez(host->os);
48 + host->os = strdupz(os?os:"unknown");
49 +}
50 +
51 static inline void rrdhost_init_machine_guid(RRDHOST *host, const char *machine_guid) {
52 strncpy(host->machine_guid, machine_guid, GUID_LEN);
53 host->machine_guid[GUID_LEN] = '\0';
@@ -55,6 +60,7 @@ static inline void rrdhost_init_machine_guid(RRDHOST *host, const char *machine_
60
61 RRDHOST *rrdhost_create(const char *hostname,
62 const char *guid,
63 + const char *os,
64 int update_every,
65 int entries,
66 RRD_MEMORY_MODE memory_mode,
@@ -73,6 +79,7 @@ RRDHOST *rrdhost_create(const char *hostname,
79
80 rrdhost_init_hostname(host, hostname);
81 rrdhost_init_machine_guid(host, guid);
82 + rrdhost_init_os(host, os);
83
84 avl_init_lock(&(host->rrdset_root_index), rrdset_compare);
85 avl_init_lock(&(host->rrdset_root_index_name), rrdset_compare_name);
@@ -176,12 +183,12 @@ RRDHOST *rrdhost_create(const char *hostname,
183 return host;
184 }
185
179 -RRDHOST *rrdhost_find_or_create(const char *hostname, const char *guid, int update_every, int history, RRD_MEMORY_MODE mode, int health_enabled) {
186 +RRDHOST *rrdhost_find_or_create(const char *hostname, const char *guid, const char *os, int update_every, int history, RRD_MEMORY_MODE mode, int health_enabled) {
187 debug(D_RRDHOST, "Searching for host '%s' with guid '%s'", hostname, guid);
188
189 RRDHOST *host = rrdhost_find(guid, 0);
190 if(!host) {
184 - host = rrdhost_create(hostname, guid, update_every, history, mode, health_enabled);
191 + host = rrdhost_create(hostname, guid, os, update_every, history, mode, health_enabled);
192 }
193 else {
194 host->health_enabled = health_enabled;
@@ -214,6 +221,7 @@ void rrd_init(char *hostname) {
221
222 localhost = rrdhost_create(hostname,
223 registry_get_this_machine_guid(),
224 + os_type,
225 default_rrd_update_every,
226 default_rrd_history_entries,
227 default_rrd_memory_mode,
@@ -304,6 +312,7 @@ void rrdhost_free(RRDHOST *host) {
312 // ------------------------------------------------------------------------
313 // free it
314
315 + freez(host->os);
316 freez(host->cache_dir);
317 freez(host->varlib_dir);
318 freez(host->health_default_exec);
src/rrdpush.c
+106 -47
@@ -1,36 +1,55 @@
1 #include "common.h"
2
3 +int rrdpush_enabled = 0;
4 +int rrdpush_exclusive = 1;
5 +
6 +static char *central_netdata = NULL;
7 +static char *api_key = NULL;
8 +
9 +#define CONNECTED_TO_SIZE 100
10 +
11 +// data collection happens from multiple threads
12 +// each of these threads calls rrdset_done()
13 +// which in turn calls rrdset_done_push()
14 +// which uses this pipe to notify the streaming thread
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
6 -int rrdpush_pipe[2];
7 -
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 central netdata
29 +// this is set to 1, otherwise 0.
30 static volatile int rrdpush_connected = 0;
31
12 -static inline void rrdpush_lock() {
13 - pthread_mutex_lock(&rrdpush_mutex);
14 -}
32 +// to have the remote netdata re-sync the charts
33 +// to its current clock, we send for this many
34 +// iterations a BEGIN line without microseconds
35 +// this is for the first iterations of each chart
36 +static unsigned int remote_clock_resync_iterations = 60;
37
16 -static inline void rrdpush_unlock() {
17 - pthread_mutex_unlock(&rrdpush_mutex);
18 -}
38 +#define rrdpush_lock() pthread_mutex_lock(&rrdpush_mutex)
39 +#define rrdpush_unlock() pthread_mutex_unlock(&rrdpush_mutex)
40
41 +// checks if the current chart definition has been sent
42 static inline int need_to_send_chart_definition(RRDSET *st) {
43 RRDDIM *rd;
44 rrddim_foreach_read(rd, st)
45 if(!rrddim_flag_check(rd, RRDDIM_FLAG_EXPOSED))
46 return 1;
47
26 -
27 - // fprintf(stderr, "NOT Sending CHART '%s' '%s'\n", st->id, st->name);
48 return 0;
49 }
50
51 +// sends the current chart definition
52 static inline void send_chart_definition(RRDSET *st) {
32 - // fprintf(stderr, "Sending CHART '%s' '%s'\n", st->id, st->name);
33 -
53 buffer_sprintf(rrdpush_buffer, "CHART '%s' '%s' '%s' '%s' '%s' '%s' '%s' %ld %d\n"
54 , st->id
55 , st->name
@@ -58,8 +77,9 @@ static inline void send_chart_definition(RRDSET *st) {
77 }
78 }
79
80 +// sends the current chart dimensions
81 static inline void send_chart_metrics(RRDSET *st) {
62 - buffer_sprintf(rrdpush_buffer, "BEGIN %s %llu\n", st->id, (st->counter_done > 60)?st->usec_since_last_update:0);
82 + buffer_sprintf(rrdpush_buffer, "BEGIN %s %llu\n", st->id, (st->counter_done > remote_clock_resync_iterations)?st->usec_since_last_update:0);
83
84 RRDDIM *rd;
85 rrddim_foreach_read(rd, st) {
@@ -73,6 +93,8 @@ static inline void send_chart_metrics(RRDSET *st) {
93 buffer_strcat(rrdpush_buffer, "END\n");
94 }
95
96 +// resets all the chart, so that their definitions
97 +// will be resent to the central netdata
98 static void reset_all_charts(void) {
99 rrd_rdlock();
100
@@ -106,7 +128,6 @@ void rrdset_done_push(RRDSET *st) {
128 if(unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ENABLED)))
129 return;
130
109 -
131 rrdpush_lock();
132
133 if(unlikely(!rrdpush_buffer || !rrdpush_connected)) {
@@ -144,6 +165,20 @@ static inline void rrdpush_flush(void) {
165 rrdpush_unlock();
166 }
167
168 +int rrdpush_init() {
169 + rrdpush_enabled = config_get_boolean("stream", "enabled", rrdpush_enabled);
170 + rrdpush_exclusive = config_get_boolean("stream", "exclusive", rrdpush_exclusive);
171 + central_netdata = config_get("stream", "stream metrics to", "");
172 + api_key = config_get("stream", "api key", "");
173 +
174 + if(!rrdpush_enabled || !central_netdata || !*central_netdata || !api_key || !*api_key) {
175 + rrdpush_enabled = 0;
176 + rrdpush_exclusive = 0;
177 + }
178 +
179 + return rrdpush_enabled;
180 +}
181 +
182 void *central_netdata_push_thread(void *ptr) {
183 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
184
@@ -155,24 +190,31 @@ void *central_netdata_push_thread(void *ptr) {
190 if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
191 error("STREAM: cannot set pthread cancel state to ENABLE.");
192
193 + int timeout = (int)config_get_number("stream", "timeout seconds", 60);
194 + int default_port = (int)config_get_number("stream", "default port", 19999);
195 + size_t max_size = (size_t)config_get_number("stream", "buffer size bytes", 1024 * 1024);
196 + unsigned int reconnect_delay = (unsigned int)config_get_number("stream", "reconnect delay seconds", 5);
197 + remote_clock_resync_iterations = (unsigned int)config_get_number("stream", "initial clock resync iterations", remote_clock_resync_iterations);
198 + int sock = -1;
199
159 - rrdpush_buffer = buffer_create(1);
160 -
161 - if(pipe(rrdpush_pipe) == -1)
162 - fatal("STREAM: cannot create required pipe.");
163 -
164 - struct timeval tv = {
165 - .tv_sec = 60,
166 - .tv_usec = 0
167 - };
200 + if(!rrdpush_enabled || !central_netdata || !*central_netdata || !api_key || !*api_key)
201 + goto cleanup;
202
203 + // initialize rrdpush globals
204 + rrdpush_buffer = buffer_create(1);
205 rrdpush_connected = 0;
206 + if(pipe(rrdpush_pipe) == -1) fatal("STREAM: cannot create required pipe.");
207 +
208 + // initialize local variables
209 size_t begin = 0;
171 - size_t max_size = 1024 * 1024;
210 size_t reconnects_counter = 0;
211 size_t sent_bytes = 0;
212 size_t sent_connection = 0;
175 - int sock = -1;
213 +
214 + struct timeval tv = {
215 + .tv_sec = timeout,
216 + .tv_usec = 0
217 + };
218
219 struct pollfd fds[2], *ifd, *ofd;
220 nfds_t fdmax;
@@ -180,6 +222,8 @@ void *central_netdata_push_thread(void *ptr) {
222 ifd = &fds[0];
223 ofd = &fds[1];
224
225 + char connected_to[CONNECTED_TO_SIZE + 1];
226 +
227 for(;;) {
228 if(netdata_exit) break;
229
@@ -188,55 +232,57 @@ void *central_netdata_push_thread(void *ptr) {
232 // they will be lost, so there is no point to do it
233 rrdpush_connected = 0;
234
191 - info("STREAM: connecting to central netdata at: %s", central_netdata_to_push_data);
192 - sock = connect_to_one_of(central_netdata_to_push_data, 19999, &tv, &reconnects_counter);
235 + info("STREAM: connecting to central netdata at: %s", central_netdata);
236 + sock = connect_to_one_of(central_netdata, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
237
238 if(unlikely(sock == -1)) {
195 - error("STREAM: failed to connect to central netdata at: %s", central_netdata_to_push_data);
196 - sleep(5);
239 + error("STREAM: failed to connect to central netdata at: %s", central_netdata);
240 + sleep(reconnect_delay);
241 continue;
242 }
243
200 - info("STREAM: initializing communication to central netdata at: %s", central_netdata_to_push_data);
244 + info("STREAM: initializing communication to central netdata at: %s", connected_to);
245
246 char http[1000 + 1];
203 - snprintfz(http, 1000, "GET /stream?key=%s&hostname=%s&machine_guid=%s&update_every=%d HTTP/1.1\r\n"
247 + snprintfz(http, 1000,
248 + "STREAM key=%s&hostname=%s&machine_guid=%s&os=%s&update_every=%d HTTP/1.1\r\n"
249 "User-Agent: netdata-push-service/%s\r\n"
250 "Accept: */*\r\n\r\n"
206 - , config_get("global", "central netdata api key", "")
251 + , api_key
252 , localhost->hostname
253 , localhost->machine_guid
254 + , localhost->os
255 , default_rrd_update_every
256 , program_version
257 );
258
213 - if(send_timeout(sock, http, strlen(http), 0, 60) == -1) {
259 + if(send_timeout(sock, http, strlen(http), 0, timeout) == -1) {
260 close(sock);
261 sock = -1;
216 - error("STREAM: failed to send http header to netdata at: %s", central_netdata_to_push_data);
217 - sleep(5);
262 + error("STREAM: failed to send http header to netdata at: %s", connected_to);
263 + sleep(reconnect_delay);
264 continue;
265 }
266
221 - info("STREAM: Waiting for STREAM from central netdata at: %s", central_netdata_to_push_data);
267 + info("STREAM: Waiting for STREAM from central netdata at: %s", connected_to);
268
223 - if(recv_timeout(sock, http, 1000, 0, 60) == -1) {
269 + if(recv_timeout(sock, http, 1000, 0, timeout) == -1) {
270 close(sock);
271 sock = -1;
226 - error("STREAM: failed to receive STREAM from netdata at: %s", central_netdata_to_push_data);
227 - sleep(5);
272 + error("STREAM: failed to receive STREAM from netdata at: %s", connected_to);
273 + sleep(reconnect_delay);
274 continue;
275 }
276
277 if(strncmp(http, "STREAM", 6)) {
278 close(sock);
279 sock = -1;
234 - error("STREAM: netdata servers at %s, did not send STREAM", central_netdata_to_push_data);
235 - sleep(5);
280 + error("STREAM: server at %s, did not send STREAM", connected_to);
281 + sleep(reconnect_delay);
282 continue;
283 }
284
239 - info("STREAM: Established STREAM with central netdata at: %s - sending metrics...", central_netdata_to_push_data);
285 + info("STREAM: Established communication with central netdata at: %s - sending metrics...", connected_to);
286
287 if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
288 error("STREAM: cannot set non-blocking mode for socket.");
@@ -264,7 +310,7 @@ void *central_netdata_push_thread(void *ptr) {
310 }
311
312 if(netdata_exit) break;
267 - int retval = poll(fds, fdmax, 60 * 1000);
313 + int retval = poll(fds, fdmax, timeout * 1000);
314 if(netdata_exit) break;
315
316 if(unlikely(retval == -1)) {
@@ -298,7 +344,7 @@ void *central_netdata_push_thread(void *ptr) {
344 ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], buffer_strlen(rrdpush_buffer) - begin, MSG_DONTWAIT);
345 if(ret == -1) {
346 if(errno != EAGAIN && errno != EINTR) {
301 - error("STREAM: failed to send metrics to central netdata at %s. We have sent %zu bytes on this connection.", central_netdata_to_push_data, sent_connection);
347 + error("STREAM: failed to send metrics to central netdata at %s. We have sent %zu bytes on this connection.", connected_to, sent_connection);
348 close(sock);
349 sock = -1;
350 }
@@ -326,10 +372,23 @@ void *central_netdata_push_thread(void *ptr) {
372 }
373 }
374
375 +cleanup:
376 debug(D_WEB_CLIENT, "STREAM: central netdata push thread exits.");
330 - if(sock != -1) {
331 - close(sock);
332 - }
377 +
378 + // make sure the data collection threads do not write data
379 + rrdpush_connected = 0;
380 +
381 + // close the pipe
382 + if(rrdpush_pipe[PIPE_READ] != -1) close(rrdpush_pipe[PIPE_READ]);
383 + if(rrdpush_pipe[PIPE_WRITE] != -1) close(rrdpush_pipe[PIPE_WRITE]);
384 +
385 + // close the socket
386 + if(sock != -1) close(sock);
387 +
388 + rrdpush_lock();
389 + buffer_free(rrdpush_buffer);
390 + rrdpush_buffer = NULL;
391 + rrdpush_unlock();
392
393 static_thread->enabled = 0;
394 pthread_exit(NULL);
src/rrdpush.h
+4
@@ -1,6 +1,10 @@
1 #ifndef NETDATA_RRDPUSH_H
2 #define NETDATA_RRDPUSH_H
3
4 +extern int rrdpush_enabled;
5 +extern int rrdpush_exclusive;
6 +
7 +extern int rrdpush_init();
8 extern void rrdset_done_push(RRDSET *st);
9 extern void *central_netdata_push_thread(void *ptr);
10
src/rrdset.c
+7 -4
@@ -196,7 +196,7 @@ void rrdset_reset(RRDSET *st) {
196 // RRDSET - helpers for rrdset_create()
197
198 inline long align_entries_to_pagesize(long entries) {
199 - if(central_netdata_to_push_data)
199 + if(rrdpush_exclusive)
200 return entries;
201
202 if(entries < 5) entries = 5;
@@ -609,16 +609,16 @@ static inline void rrdset_done_push_int(RRDSET *st) {
609 rrdset_update_last_collected_time(st);
610 }
611
612 - rrdset_done_push(st);
613 -
612 st->counter++;
613 st->counter_done++;
614 +
615 + rrdset_done_push(st);
616 }
617
618 void rrdset_done(RRDSET *st) {
619 if(unlikely(netdata_exit)) return;
620
621 - if(unlikely(central_netdata_to_push_data)) {
621 + if(unlikely(rrdpush_exclusive)) {
622 rrdset_done_push_int(st);
623 return;
624 }
@@ -1207,5 +1207,8 @@ void rrdset_done(RRDSET *st) {
1207
1208 if(unlikely(pthread_setcancelstate(pthreadoldcancelstate, NULL) != 0))
1209 error("Cannot set pthread cancel state to RESTORE (%d).", pthreadoldcancelstate);
1210 +
1211 + if(unlikely(rrdpush_enabled))
1212 + rrdset_done_push_int(st);
1213 }
1214
src/socket.c
+8 -2
@@ -180,7 +180,7 @@ int connect_to(const char *definition, int default_port, struct timeval *timeout
180 return fd;
181 }
182
183 -int connect_to_one_of(const char *destination, int default_port, struct timeval *timeout, size_t *reconnects_counter) {
183 +int connect_to_one_of(const char *destination, int default_port, struct timeval *timeout, size_t *reconnects_counter, char *connected_to, size_t connected_to_size) {
184 int sock = -1;
185
186 const char *s = destination;
@@ -200,7 +200,13 @@ int connect_to_one_of(const char *destination, int default_port, struct timeval
200 strncpyz(buf, s, e - s);
201 if(reconnects_counter) *reconnects_counter += 1;
202 sock = connect_to(buf, default_port, timeout);
203 - if(sock != -1) break;
203 + if(sock != -1) {
204 + if(connected_to && connected_to_size) {
205 + strncpy(connected_to, buf, connected_to_size);
206 + connected_to[connected_to_size - 1] = '\0';
207 + }
208 + break;
209 + }
210 s = e;
211 }
212
src/socket.h
+1 -1
@@ -6,7 +6,7 @@
6 #define NETDATA_SOCKET_H
7
8 extern int connect_to(const char *definition, int default_port, struct timeval *timeout);
9 -extern int connect_to_one_of(const char *destination, int default_port, struct timeval *timeout, size_t *reconnects_counter);
9 +extern int connect_to_one_of(const char *destination, int default_port, struct timeval *timeout, size_t *reconnects_counter, char *connected_to, size_t connected_to_size);
10
11 extern ssize_t recv_timeout(int sockfd, void *buf, size_t len, int flags, int timeout);
12 extern ssize_t send_timeout(int sockfd, void *buf, size_t len, int flags, int timeout);
src/web_buffer.c
+3 -2
@@ -359,8 +359,9 @@ BUFFER *buffer_create(size_t size)
359 return(b);
360 }
361
362 -void buffer_free(BUFFER *b)
363 -{
362 +void buffer_free(BUFFER *b) {
363 + if(unlikely(!b)) return;
364 +
365 buffer_overflow_check(b);
366
367 debug(D_WEB_BUFFER, "Freeing web buffer of size %zu.", b->size);
src/web_client.c
+39 -26
@@ -45,8 +45,7 @@ static inline int web_client_uncrock_socket(struct web_client *w) {
45 return 0;
46 }
47
48 -struct web_client *web_client_create(int listener)
49 -{
48 +struct web_client *web_client_create(int listener) {
49 struct web_client *w;
50
51 w = callocz(1, sizeof(struct web_client));
@@ -223,9 +222,9 @@ struct web_client *web_client_free(struct web_client *w) {
222
223 if(w->prev) w->prev->next = w->next;
224 if(w->next) w->next->prev = w->prev;
226 - if(w->response.header_output) buffer_free(w->response.header_output);
227 - if(w->response.header) buffer_free(w->response.header);
228 - if(w->response.data) buffer_free(w->response.data);
225 + buffer_free(w->response.header_output);
226 + buffer_free(w->response.header);
227 + buffer_free(w->response.data);
228 if(w->ifd != -1) close(w->ifd);
229 if(w->ofd != -1 && w->ofd != w->ifd) close(w->ofd);
230 freez(w);
@@ -1066,8 +1065,7 @@ int web_client_api_request_v1_badge(RRDHOST *host, struct web_client *w, char *u
1065 }
1066
1067 cleanup:
1069 - if(dimensions)
1070 - buffer_free(dimensions);
1068 + buffer_free(dimensions);
1069 return ret;
1070 }
1071
@@ -1242,7 +1240,7 @@ int web_client_api_request_v1_data(RRDHOST *host, struct web_client *w, char *ur
1240 buffer_strcat(w->response.data, ");");
1241
1242 cleanup:
1245 - if(dimensions) buffer_free(dimensions);
1243 + buffer_free(dimensions);
1244 return ret;
1245 }
1246
@@ -1675,7 +1673,7 @@ int validate_stream_api_key(const char *key) {
1673 }
1674
1675 int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1678 - char *key = NULL, *hostname = NULL, *machine_guid = NULL;
1676 + char *key = NULL, *hostname = NULL, *machine_guid = NULL, *os = NULL;
1677 int update_every = default_rrd_update_every;
1678 int history = default_rrd_history_entries;
1679 RRD_MEMORY_MODE mode = default_rrd_memory_mode;
@@ -1697,6 +1695,8 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1695 machine_guid = value;
1696 else if(!strcmp(name, "update_every"))
1697 update_every = (int)strtoul(value, NULL, 0);
1698 + else if(!strcmp(name, "os"))
1699 + os = value;
1700 }
1701
1702 if(!key || !*key) {
@@ -1734,7 +1734,6 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1734 return 404;
1735 }
1736
1737 - // update_every = (int)appconfig_get_number(&stream_config, key, "default update every", update_every);
1737 update_every = (int)appconfig_get_number(&stream_config, machine_guid, "update every", update_every);
1738 if(update_every < 0) update_every = 1;
1739
@@ -1748,7 +1747,10 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1747 health_enabled = appconfig_get_boolean_ondemand(&stream_config, key, "health enabled by default", health_enabled);
1748 health_enabled = appconfig_get_boolean_ondemand(&stream_config, machine_guid, "health enabled", health_enabled);
1749
1751 - host = rrdhost_find_or_create(hostname, machine_guid, update_every, history, mode, health_enabled?1:0);
1750 + if(strcmp(machine_guid, "localhost"))
1751 + host = localhost;
1752 + else
1753 + host = rrdhost_find_or_create(hostname, machine_guid, os, update_every, history, mode, health_enabled?1:0);
1754
1755 info("STREAM request from client '%s:%s' for host '%s' with machine_guid '%s': update every = %d, history = %d, memory mode = %s, health %s",
1756 w->client_ip, w->client_port,
@@ -1806,20 +1808,24 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1808 return 500;
1809 }
1810
1811 + rrdhost_wrlock(host);
1812 + host->use_counter++;
1813 + rrdhost_unlock(host);
1814 +
1815 // call the plugins.d processor to receive the metrics
1816 info("STREAM [%s]:%s: connecting client to plugins.d on host '%s' with machine GUID '%s'.", w->client_ip, w->client_port, host->hostname, host->machine_guid);
1817 size_t count = pluginsd_process(host, &cd, fp, 1);
1818 error("STREAM [%s]:%s: client disconnected (host '%s', machine GUID '%s').", w->client_ip, w->client_port, host->hostname, host->machine_guid);
1813 - if(health_enabled == CONFIG_BOOLEAN_AUTO)
1819 +
1820 + rrdhost_wrlock(host);
1821 + host->use_counter--;
1822 + if(!host->use_counter && health_enabled == CONFIG_BOOLEAN_AUTO)
1823 host->health_enabled = 0;
1824 + rrdhost_unlock(host);
1825
1816 - // close all sockets, to let the socket worker we are done
1826 + // cleanup
1827 fclose(fp);
1828 w->ifd = -1;
1819 - if(w->ofd != -1 && w->ofd != w->ifd) {
1820 - close(w->ofd);
1821 - w->ofd = -1;
1822 - }
1829
1830 // this will not send anything
1831 // the socket is closed
@@ -2032,6 +2038,10 @@ static inline HTTP_VALIDATION http_request_validate(struct web_client *w) {
2038 encoded_url = s = &s[8];
2039 w->mode = WEB_CLIENT_MODE_OPTIONS;
2040 }
2041 + else if(!strncmp(s, "STREAM ", 8)) {
2042 + encoded_url = s = &s[8];
2043 + w->mode = WEB_CLIENT_MODE_STREAM;
2044 + }
2045 else {
2046 w->wait_receive = 0;
2047 return HTTP_VALIDATION_NOT_SUPPORTED;
@@ -2292,8 +2302,7 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2302 hash_graph = 0,
2303 hash_list = 0,
2304 hash_all_json = 0,
2295 - hash_host = 0,
2296 - hash_stream = 0;
2305 + hash_host = 0;
2306
2307 #ifdef NETDATA_INTERNAL_CHECKS
2308 static uint32_t hash_exit = 0, hash_debug = 0, hash_mirror = 0;
@@ -2308,7 +2317,6 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2317 hash_list = simple_hash("list");
2318 hash_all_json = simple_hash("all.json");
2319 hash_host = simple_hash("host");
2311 - hash_stream = simple_hash("stream");
2320 #ifdef NETDATA_INTERNAL_CHECKS
2321 hash_exit = simple_hash("exit");
2322 hash_debug = simple_hash("debug");
@@ -2329,10 +2337,6 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2337 debug(D_WEB_CLIENT_ACCESS, "%llu: host switch request ...", w->id);
2338 return web_client_switch_host(host, w, url);
2339 }
2332 - else if(unlikely(hash == hash_stream && strcmp(tok, "stream") == 0)) {
2333 - debug(D_WEB_CLIENT_ACCESS, "%llu: stream request ...", w->id);
2334 - return web_client_stream_request(host, w, url);
2335 - }
2340 else if(unlikely(hash == hash_netdata_conf && strcmp(tok, "netdata.conf") == 0)) {
2341 debug(D_WEB_CLIENT_ACCESS, "%llu: Sending netdata.conf ...", w->id);
2342 w->response.data->contenttype = CT_TEXT_PLAIN;
@@ -2483,6 +2487,10 @@ void web_client_process_request(struct web_client *w) {
2487 buffer_strcat(w->response.data, "OK");
2488 w->response.code = 200;
2489 }
2490 + else if(unlikely(w->mode == WEB_CLIENT_MODE_STREAM)) {
2491 + w->response.code = web_client_stream_request(localhost, w, w->decoded_url);
2492 + return;
2493 + }
2494 else
2495 w->response.code = web_client_process_url(localhost, w, w->decoded_url);
2496 break;
@@ -2528,6 +2536,10 @@ void web_client_process_request(struct web_client *w) {
2536 else w->wait_send = 0;
2537
2538 switch(w->mode) {
2539 + case WEB_CLIENT_MODE_STREAM:
2540 + debug(D_WEB_CLIENT, "%llu: STREAM done.", w->id);
2541 + break;
2542 +
2543 case WEB_CLIENT_MODE_OPTIONS:
2544 debug(D_WEB_CLIENT, "%llu: Done preparing the OPTIONS response. Sending data (%zu bytes) to client.", w->id, w->response.data->len);
2545 break;
@@ -2559,7 +2571,7 @@ void web_client_process_request(struct web_client *w) {
2571 break;
2572
2573 default:
2562 - fatal("%llu: Unknown client mode %d.", w->id, w->mode);
2574 + fatal("%llu: Unknown client mode %u.", w->id, w->mode);
2575 break;
2576 }
2577 }
@@ -2982,7 +2994,8 @@ void *web_client_main(void *ptr)
2994
2995 // if the sockets are closed, may have transferred this client
2996 // to plugins.d
2985 - if(w->ifd == -1 && w->ofd == -1) break;
2997 + if(unlikely(w->mode == WEB_CLIENT_MODE_STREAM))
2998 + break;
2999 }
3000 }
3001
src/web_client.h
+8 -5
@@ -13,9 +13,12 @@ extern int web_enable_gzip,
13 extern int respect_web_browser_do_not_track_policy;
14 extern char *web_x_frame_options;
15
16 -#define WEB_CLIENT_MODE_NORMAL 0
17 -#define WEB_CLIENT_MODE_FILECOPY 1
18 -#define WEB_CLIENT_MODE_OPTIONS 2
16 +typedef enum web_client_mode {
17 + WEB_CLIENT_MODE_NORMAL = 0,
18 + WEB_CLIENT_MODE_FILECOPY = 1,
19 + WEB_CLIENT_MODE_OPTIONS = 2,
20 + WEB_CLIENT_MODE_STREAM = 3
21 +} WEB_CLIENT_MODE;
22
23 #define URL_MAX 8192
24 #define ZLIB_CHUNK 16384
@@ -55,14 +58,14 @@ struct web_client {
58
59 uint8_t keepalive:1; // if set to 1, the web client will be re-used
60
58 - uint8_t mode:3; // the operational mode of the client
59 -
61 uint8_t wait_receive:1; // 1 = we are waiting more input data
62 uint8_t wait_send:1; // 1 = we have data to send to the client
63
64 uint8_t donottrack:1; // 1 = we should not set cookies on this client
65 uint8_t tracking_required:1; // 1 = if the request requires cookies
66
67 + WEB_CLIENT_MODE mode; // the operational mode of the client
68 +
69 int tcp_cork; // 1 = we have a cork on the socket
70
71 int ifd;