@cryptotaxi247 / netdata-1 / commits / 0bda31e21

moved [stream] configuration to stream.conf

Costa Tsaousis (ktsaou) committed Feb 25, 2017 at 12:46 UTC 0bda31e2154e0557d293fada5cbe5804dd702cdf
12 files changed +213 -96
conf.d/Makefile.am
+1 -1
@@ -4,7 +4,6 @@
4 MAINTAINERCLEANFILES= $(srcdir)/Makefile.in
5
6 dist_config_DATA = \
7 - aggregated_hosts.conf \
7 apps_groups.conf \
8 charts.d.conf \
9 fping.conf \
@@ -12,6 +11,7 @@ dist_config_DATA = \
11 python.d.conf \
12 health_alarm_notify.conf \
13 health_email_recipients.conf \
14 + stream.conf \
15 $(NULL)
16
17 nodeconfigdir=$(configdir)/node.d
conf.d/stream.conf renamed
+56 -35
@@ -8,27 +8,31 @@
8 #
9 # -----------------------------------------------------------------------------
10 # 1. SLAVE NETDATA - THE ONE THAT WILL BE SENDING METRICS
11 -#
12 -# In /etc/netdata/netdata.conf you have a configuration like this:
13 -#
14 -# [stream]
15 -#
16 -# # enable or disable sending metrics.
17 -# # default: no
18 -# enabled = yes
19 -#
20 -# # one or more netdata hosts to send metrics to.
21 -# # only one will get them, the first available.
22 -# destination = HOST1:PORT1 HOST2:PORT2 ...
23 -#
24 -# # the API key to use, to authorize ourselves
25 -# api key = API_KEY
26 -#
11 +
12 +[stream]
13 + enabled = no
14 +
15 + # where to send metrics to?
16 + # A space separated list of IP:PORT is accepted. The first available will
17 + # get the metrics.
18 + # IPv6 addresses should be [IP]:PORT
19 + destination =
20 +
21 + # The API_KEY to use (as the sender)
22 + api key =
23 +
24 + # other options (uncomment and set)
25 +
26 + # timeout seconds = 60
27 + # default port = 19999
28 + # buffer size bytes = 1048576
29 + # reconnect delay seconds = 5
30 + # initial clock resync iterations = 60
31 +
32 +
33 # -----------------------------------------------------------------------------
34 # 2. MASTER NETDATA - THE ONE THAT WILL BE RECEIVING METRICS
35 #
30 -# Use this file to define the API keys is accepts.
31 -#
36 # You can have one API key per slave, or the same API key for all slaves.
37 #
38 # All options below are used in this order:
@@ -38,8 +42,6 @@
42 # c) this netdata defaults (as in netdata.conf)
43 #
44 # You can combine the above (the more specific setting will be used).
41 -#
42 -# -----------------------------------------------------------------------------
45
46 # API key authentication
47 # If the key is not listed here, it will not be able to connect.
@@ -49,29 +51,41 @@
51
52 # You can disable the API key, by setting this to: no
53 # The default (for unknown API keys) is also: no
52 -# enabled = yes
54 + enabled = no
55
56 # The default history in entries, for all hosts using this API key.
57 # You can also set it per host below.
58 # If you don't set it here, the history size of the central netdata
59 # will be used
58 -# default history = 3600
60 + default history = 3600
61
62 # The default memory mode to be used for all hosts using this API key.
63 # You can also set it per host below.
62 - # If you don't set it here, the memory mode of the central netdata
63 - # will be used.
64 - # Valid modes: save (load/save db), map (like swap), ram (no disk)
65 -# default memory mode = save
64 + # If you don't set it here, the memory mode of netdata.conf will be used.
65 + # Valid modes:
66 + # save save on exit, load on start
67 + # map like swap (continuously syncing to disks)
68 + # ram keep it in RAM, don't touch the disk
69 + # none no database (passing through this netdata)
70 + default memory mode = ram
71
72 # Shall we enable health monitoring for the hosts using this API key?
73 # 3 values:
74 # yes enable alarms
75 # no do not enable alarms
71 - # auto enable alarms, only when the host is streaming metrics
76 + # auto enable alarms, only when the sending netdata is connected
77 # You can also set it per host, below.
73 - # The default is the same as the central netdata
74 -# health enabled by default = auto
78 + # The default is the same as to netdata.conf
79 + health enabled by default = auto
80 +
81 + # postpone alarms for a short period after the sender is connected
82 + default postpone alarms on connect seconds = 60
83 +
84 + # need to route metrics differently? set these.
85 + # the defaults are the ones at the [stream] section
86 + #default proxy enabled = yes | no
87 + #default proxy destination = IP:PORT IP:PORT ...
88 + #default proxy api key = API_KEY
89
90
91 # -----------------------------------------------------------------------------
@@ -81,17 +95,24 @@
95 # you can give settings for each specific host here.
96
97 [MACHINE_GUID]
84 - # This can be used to stop receiving data
98 + # enable this host: yes | no
99 # THIS IS NOT A SECURITY MECHANISM - AN ATTACKER CAN SET ANY OTHER GUID.
100 # Use only the API key for security.
87 -# enabled = yes
101 + enabled = no
102
103 # The number of entries in the database
90 -# history = 3600
104 + history = 3600
105
92 - # The memory mode of the database
93 -# memory mode = save
106 + # The memory mode of the database: save | map | ram | none
107 + memory mode = save
108
109 # Health / alarms control: yes | no | auto
96 -# health enabled = yes
110 + health enabled = yes
111 +
112 + # postpone alarms when the sender connects
113 + postpone alarms on connect seconds = 60
114
115 + # need to route metrics differently?
116 + #proxy enabled = yes | no
117 + #proxy destination = IP:PORT IP:PORT ...
118 + #proxy api key = API_KEY
src/main.c
+12 -9
@@ -705,6 +705,8 @@ int main(int argc, char **argv) {
705 if(!config_loaded)
706 config_load(NULL, 0);
707
708 + // ------------------------------------------------------------------------
709 + // initialize netdata
710 {
711 char *pmax = config_get(CONFIG_SECTION_GLOBAL, "glibc malloc arena max for plugins", "1");
712 if(pmax && *pmax)
@@ -757,6 +759,16 @@ int main(int argc, char **argv) {
759 log_init();
760 error_log_limit_unlimited();
761
762 +
763 + // --------------------------------------------------------------------
764 + // load stream.conf
765 + {
766 + char filename[FILENAME_MAX + 1];
767 + snprintfz(filename, FILENAME_MAX, "%s/stream.conf", netdata_configured_config_dir);
768 + appconfig_load(&stream_config, filename, 0);
769 + }
770 +
771 +
772 // --------------------------------------------------------------------
773 // setup process signals
774
@@ -851,15 +863,6 @@ int main(int argc, char **argv) {
863
864 if(web_server_mode != WEB_SERVER_MODE_NONE)
865 create_listen_sockets();
854 -
855 -
856 - // --------------------------------------------------------------------
857 - // load the aggregated host configuration file
858 - {
859 - char filename[FILENAME_MAX + 1];
860 - snprintfz(filename, FILENAME_MAX, "%s/aggregated_hosts.conf", netdata_configured_config_dir);
861 - appconfig_load(&stream_config, filename, 0);
862 - }
866 }
867
868 // initialize the log files
src/registry.h
+1
@@ -70,5 +70,6 @@ extern int registry_request_hello_json(RRDHOST *host, struct web_client *w);
70 extern void registry_statistics(void);
71
72 extern char *registry_get_this_machine_guid(void);
73 +extern int regenerate_guid(const char *guid, char *result);
74
75 #endif /* NETDATA_REGISTRY_H */
src/registry_internals.c
+6 -6
@@ -7,7 +7,7 @@ struct registry registry;
7
8 // parse a GUID and re-generated to be always lower case
9 // this is used as a protection against the variations of GUIDs
10 -int registry_regenerate_guid(const char *guid, char *result) {
10 +int regenerate_guid(const char *guid, char *result) {
11 uuid_t uuid;
12 if(unlikely(uuid_parse(guid, uuid) == -1)) {
13 info("Registry: GUID '%s' is not a valid GUID.", guid);
@@ -18,7 +18,7 @@ int registry_regenerate_guid(const char *guid, char *result) {
18
19 #ifdef NETDATA_INTERNAL_CHECKS
20 if(strcmp(guid, result))
21 - info("Registry: source GUID '%s' and re-generated GUID '%s' differ!", guid, result);
21 + info("GUID '%s' and re-generated GUID '%s' differ!", guid, result);
22 #endif /* NETDATA_INTERNAL_CHECKS */
23 }
24
@@ -96,14 +96,14 @@ REGISTRY_PERSON_URL *registry_verify_request(char *person_guid, char *machine_gu
96 url = registry_fix_url(url, NULL);
97
98 // make sure the person GUID is valid
99 - if(registry_regenerate_guid(person_guid, pbuf) == -1) {
99 + if(regenerate_guid(person_guid, pbuf) == -1) {
100 info("Registry Request Verification: invalid person GUID, person: '%s', machine '%s', url '%s'", person_guid, machine_guid, url);
101 return NULL;
102 }
103 person_guid = pbuf;
104
105 // make sure the machine GUID is valid
106 - if(registry_regenerate_guid(machine_guid, mbuf) == -1) {
106 + if(regenerate_guid(machine_guid, mbuf) == -1) {
107 info("Registry Request Verification: invalid machine GUID, person: '%s', machine '%s', url '%s'", person_guid, machine_guid, url);
108 return NULL;
109 }
@@ -226,7 +226,7 @@ REGISTRY_MACHINE *registry_request_machine(char *person_guid, char *machine_guid
226 if(!pu || !p || !m) return NULL;
227
228 // make sure the machine GUID is valid
229 - if(registry_regenerate_guid(request_machine, mbuf) == -1) {
229 + if(regenerate_guid(request_machine, mbuf) == -1) {
230 info("Registry Machine URLs request: invalid machine GUID, person: '%s', machine '%s', url '%s', request machine '%s'", p->guid, m->guid, pu->url->url, request_machine);
231 return NULL;
232 }
@@ -288,7 +288,7 @@ char *registry_get_this_machine_guid(void) {
288 error("Failed to read machine GUID from '%s'", registry.machine_guid_filename);
289 else {
290 buf[GUID_LEN] = '\0';
291 - if(registry_regenerate_guid(buf, guid) == -1) {
291 + if(regenerate_guid(buf, guid) == -1) {
292 error("Failed to validate machine GUID '%s' from '%s'. Ignoring it - this might mean this netdata will appear as duplicate in the registry.",
293 buf, registry.machine_guid_filename);
294
src/registry_internals.h
+1 -1
@@ -59,7 +59,7 @@ struct registry {
59 pthread_mutex_t lock;
60 };
61
62 -extern int registry_regenerate_guid(const char *guid, char *result);
62 +extern int regenerate_guid(const char *guid, char *result);
63
64 #include "registry_url.h"
65 #include "registry_machine.h"
src/registry_machine.c
+1 -1
@@ -58,7 +58,7 @@ REGISTRY_MACHINE *registry_machine_get(const char *machine_guid, time_t when) {
58 if(likely(machine_guid && *machine_guid)) {
59 // validate it is a GUID
60 char buf[GUID_LEN + 1];
61 - if(unlikely(registry_regenerate_guid(machine_guid, buf) == -1))
61 + if(unlikely(regenerate_guid(machine_guid, buf) == -1))
62 info("Registry: machine guid '%s' is not a valid guid. Ignoring it.", machine_guid);
63 else {
64 machine_guid = buf;
src/registry_person.c
+1 -1
@@ -183,7 +183,7 @@ REGISTRY_PERSON *registry_person_get(const char *person_guid, time_t when) {
183 if(person_guid && *person_guid) {
184 char buf[GUID_LEN + 1];
185 // validate it is a GUID
186 - if(unlikely(registry_regenerate_guid(person_guid, buf) == -1))
186 + if(unlikely(regenerate_guid(person_guid, buf) == -1))
187 info("Registry: person guid '%s' is not a valid guid. Ignoring it.", person_guid);
188 else {
189 person_guid = buf;
src/rrd.h
+14 -1
@@ -342,6 +342,8 @@ struct rrdhost {
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 + char *rrdpush_destination; // where to send metrics to
346 + char *rrdpush_api_key; // the api key at the receiving netdata
347 volatile int rrdpush_connected; // 1 when the sender is ready to push metrics
348 volatile int rrdpush_spawn; // 1 when the sender thread has been spawn
349 volatile int rrdpush_error_shown; // 1 when we have logged a communication error
@@ -423,7 +425,18 @@ extern pthread_rwlock_t rrd_rwlock;
425 extern void rrd_init(char *hostname);
426
427 extern RRDHOST *rrdhost_find(const char *guid, uint32_t hash);
426 -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);
428 +extern RRDHOST *rrdhost_find_or_create(
429 + const char *hostname
430 + , const char *guid
431 + , const char *os
432 + , int update_every
433 + , int history
434 + , RRD_MEMORY_MODE mode
435 + , int health_enabled
436 + , int rrdpush_enabled
437 + , char *rrdpush_destination
438 + , char *rrdpush_api_key
439 +);
440
441 #ifdef NETDATA_INTERNAL_CHECKS
442 extern void rrdhost_check_wrlock_int(RRDHOST *host, const char *file, const char *function, const unsigned long line);
src/rrdhost.c
+49 -12
@@ -65,6 +65,9 @@ RRDHOST *rrdhost_create(const char *hostname,
65 int entries,
66 RRD_MEMORY_MODE memory_mode,
67 int health_enabled,
68 + int rrdpush_enabled,
69 + char *rrdpush_destination,
70 + char *rrdpush_api_key,
71 int is_localhost
72 ) {
73
@@ -76,11 +79,13 @@ RRDHOST *rrdhost_create(const char *hostname,
79 host->rrd_history_entries = entries;
80 host->rrd_memory_mode = memory_mode;
81 host->health_enabled = (memory_mode == RRD_MEMORY_MODE_NONE)? 0 : health_enabled;
79 - host->rrdpush_enabled = default_rrdpush_enabled;
82 + host->rrdpush_enabled = (rrdpush_enabled && rrdpush_destination && *rrdpush_destination && rrdpush_api_key && *rrdpush_api_key);
83 + host->rrdpush_destination = (host->rrdpush_enabled)?strdupz(rrdpush_destination):NULL;
84 + host->rrdpush_api_key = (host->rrdpush_enabled)?strdupz(rrdpush_api_key):NULL;
85
86 host->rrdpush_pipe[0] = -1;
87 host->rrdpush_pipe[1] = -1;
83 - host->rrdpush_socket = -1;
88 + host->rrdpush_socket = -1;
89
90 pthread_mutex_init(&host->rrdpush_mutex, NULL);
91 pthread_rwlock_init(&host->rrdhost_rwlock, NULL);
@@ -201,6 +206,7 @@ RRDHOST *rrdhost_create(const char *hostname,
206 ", memory mode: %s"
207 ", history entries: %d"
208 ", streaming: %s"
209 + " to: '%s' (api key: '%s')"
210 ", health: %s"
211 ", cache_dir: '%s'"
212 ", varlib_dir: '%s'"
@@ -214,6 +220,8 @@ RRDHOST *rrdhost_create(const char *hostname,
220 , rrd_memory_mode_name(host->rrd_memory_mode)
221 , host->rrd_history_entries
222 , host->rrdpush_enabled?"enabled":"disabled"
223 + , host->rrdpush_destination
224 + , host->rrdpush_api_key
225 , host->health_enabled?"enabled":"disabled"
226 , host->cache_dir
227 , host->varlib_dir
@@ -228,12 +236,35 @@ RRDHOST *rrdhost_create(const char *hostname,
236 return host;
237 }
238
231 -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) {
239 +RRDHOST *rrdhost_find_or_create(
240 + const char *hostname
241 + , const char *guid
242 + , const char *os
243 + , int update_every
244 + , int history
245 + , RRD_MEMORY_MODE mode
246 + , int health_enabled
247 + , int rrdpush_enabled
248 + , char *rrdpush_destination
249 + , char *rrdpush_api_key
250 +) {
251 debug(D_RRDHOST, "Searching for host '%s' with guid '%s'", hostname, guid);
252
253 RRDHOST *host = rrdhost_find(guid, 0);
254 if(!host) {
236 - host = rrdhost_create(hostname, guid, os, update_every, history, mode, health_enabled, 0);
255 + host = rrdhost_create(
256 + hostname
257 + , guid
258 + , os
259 + , update_every
260 + , history
261 + , mode
262 + , health_enabled
263 + , rrdpush_enabled
264 + , rrdpush_destination
265 + , rrdpush_api_key
266 + , 0
267 + );
268 }
269 else {
270 host->health_enabled = health_enabled;
@@ -267,14 +298,18 @@ void rrd_init(char *hostname) {
298 rrdpush_init();
299
300 debug(D_RRDHOST, "Initializing localhost with hostname '%s'", hostname);
270 - localhost = rrdhost_create(hostname,
271 - registry_get_this_machine_guid(),
272 - os_type,
273 - default_rrd_update_every,
274 - default_rrd_history_entries,
275 - default_rrd_memory_mode,
276 - default_health_enabled,
277 - 1
301 + localhost = rrdhost_create(
302 + hostname
303 + , registry_get_this_machine_guid()
304 + , os_type
305 + , default_rrd_update_every
306 + , default_rrd_history_entries
307 + , default_rrd_memory_mode
308 + , default_health_enabled
309 + , default_rrdpush_enabled
310 + , default_rrdpush_destination
311 + , default_rrdpush_api_key
312 + , 1
313 );
314 }
315
@@ -369,6 +404,8 @@ void rrdhost_free(RRDHOST *host) {
404 freez(host->os);
405 freez(host->cache_dir);
406 freez(host->varlib_dir);
407 + freez(host->rrdpush_api_key);
408 + freez(host->rrdpush_destination);
409 freez(host->health_default_exec);
410 freez(host->health_default_recipient);
411 freez(host->health_log_filename);
src/rrdpush.c
+69 -29
@@ -23,18 +23,22 @@
23 *
24 */
25
26 +#define START_STREAMING_PROMPT "Hit me baby, push them over..."
27 +
28 int default_rrdpush_enabled = 0;
27 -static char *rrdpush_destination = NULL;
28 -static char *rrdpush_api_key = NULL;
29 +char *default_rrdpush_destination = NULL;
30 +char *default_rrdpush_api_key = NULL;
31
32 int rrdpush_init() {
31 - default_rrdpush_enabled = config_get_boolean(CONFIG_SECTION_STREAM, "enabled", default_rrdpush_enabled);
32 - rrdpush_destination = config_get(CONFIG_SECTION_STREAM, "destination", "");
33 - rrdpush_api_key = config_get(CONFIG_SECTION_STREAM, "api key", "");
33 + default_rrdpush_enabled = appconfig_get_boolean(&stream_config, CONFIG_SECTION_STREAM, "enabled", default_rrdpush_enabled);
34 + default_rrdpush_destination = appconfig_get(&stream_config, CONFIG_SECTION_STREAM, "destination", "");
35 + default_rrdpush_api_key = appconfig_get(&stream_config, CONFIG_SECTION_STREAM, "api key", "");
36
35 - if(default_rrdpush_enabled && (!rrdpush_destination || !*rrdpush_destination || !rrdpush_api_key || !*rrdpush_api_key)) {
37 + if(default_rrdpush_enabled && (!default_rrdpush_destination || !*default_rrdpush_destination || !default_rrdpush_api_key || !*default_rrdpush_api_key)) {
38 error("STREAM [send]: cannot enable sending thread - information is missing.");
39 default_rrdpush_enabled = 0;
40 + default_rrdpush_api_key = NULL;
41 + default_rrdpush_destination = NULL;
42 }
43
44 return default_rrdpush_enabled;
@@ -237,14 +241,14 @@ void *rrdpush_sender_thread(void *ptr) {
241 if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
242 error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
243
240 - int timeout = (int)config_get_number(CONFIG_SECTION_STREAM, "timeout seconds", 60);
241 - int default_port = (int)config_get_number(CONFIG_SECTION_STREAM, "default port", 19999);
242 - size_t max_size = (size_t)config_get_number(CONFIG_SECTION_STREAM, "buffer size bytes", 1024 * 1024);
243 - unsigned int reconnect_delay = (unsigned int)config_get_number(CONFIG_SECTION_STREAM, "reconnect delay seconds", 5);
244 - remote_clock_resync_iterations = (unsigned int)config_get_number(CONFIG_SECTION_STREAM, "initial clock resync iterations", remote_clock_resync_iterations);
244 + int timeout = (int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "timeout seconds", 60);
245 + int default_port = (int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "default port", 19999);
246 + size_t max_size = (size_t)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "buffer size bytes", 1024 * 1024);
247 + unsigned int reconnect_delay = (unsigned int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "reconnect delay seconds", 5);
248 + remote_clock_resync_iterations = (unsigned int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "initial clock resync iterations", remote_clock_resync_iterations);
249 char connected_to[CONNECTED_TO_SIZE + 1] = "";
250
247 - if(!host->rrdpush_enabled || !rrdpush_destination || !*rrdpush_destination || !rrdpush_api_key || !*rrdpush_api_key)
251 + if(!host->rrdpush_enabled || !host->rrdpush_destination || !*host->rrdpush_destination || !host->rrdpush_api_key || !*host->rrdpush_api_key)
252 goto cleanup;
253
254 // initialize rrdpush globals
@@ -276,11 +280,11 @@ void *rrdpush_sender_thread(void *ptr) {
280 // they will be lost, so there is no point to do it
281 host->rrdpush_connected = 0;
282
279 - info("STREAM %s [send to %s]: connecting...", host->hostname, rrdpush_destination);
280 - host->rrdpush_socket = connect_to_one_of(rrdpush_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
283 + info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_destination);
284 + host->rrdpush_socket = connect_to_one_of(host->rrdpush_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
285
286 if(unlikely(host->rrdpush_socket == -1)) {
283 - error("STREAM %s [send to %s]: failed to connect", host->hostname, rrdpush_destination);
287 + error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_destination);
288 sleep(reconnect_delay);
289 continue;
290 }
@@ -292,7 +296,7 @@ void *rrdpush_sender_thread(void *ptr) {
296 "STREAM key=%s&hostname=%s&machine_guid=%s&os=%s&update_every=%d HTTP/1.1\r\n"
297 "User-Agent: netdata-push-service/%s\r\n"
298 "Accept: */*\r\n\r\n"
295 - , rrdpush_api_key
299 + , host->rrdpush_api_key
300 , host->hostname
301 , host->machine_guid
302 , host->os
@@ -318,7 +322,7 @@ void *rrdpush_sender_thread(void *ptr) {
322 continue;
323 }
324
321 - if(strncmp(http, "STREAM", 6)) {
325 + if(strncmp(http, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT))) {
326 close(host->rrdpush_socket);
327 host->rrdpush_socket = -1;
328 error("STREAM %s [send to %s]: server is not replying properly.", host->hostname, connected_to);
@@ -428,6 +432,9 @@ int rrdpush_receive(int fd, const char *key, const char *hostname, const char *m
432 int history = default_rrd_history_entries;
433 RRD_MEMORY_MODE mode = default_rrd_memory_mode;
434 int health_enabled = default_health_enabled;
435 + int rrdpush_enabled = default_rrdpush_enabled;
436 + char *rrdpush_destination = default_rrdpush_destination;
437 + char *rrdpush_api_key = default_rrdpush_api_key;
438 time_t alarms_delay = 60;
439
440 update_every = (int)appconfig_get_number(&stream_config, machine_guid, "update every", update_every);
@@ -446,10 +453,30 @@ int rrdpush_receive(int fd, const char *key, const char *hostname, const char *m
453 alarms_delay = appconfig_get_number(&stream_config, key, "default postpone alarms on connect seconds", alarms_delay);
454 alarms_delay = appconfig_get_number(&stream_config, machine_guid, "postpone alarms on connect seconds", alarms_delay);
455
456 + rrdpush_enabled = appconfig_get_boolean(&stream_config, key, "default proxy enabled", rrdpush_enabled);
457 + rrdpush_enabled = appconfig_get_boolean(&stream_config, machine_guid, "proxy enabled", rrdpush_enabled);
458 +
459 + rrdpush_destination = appconfig_get(&stream_config, key, "default proxy destination", rrdpush_destination);
460 + rrdpush_destination = appconfig_get(&stream_config, machine_guid, "proxy destination", rrdpush_destination);
461 +
462 + rrdpush_api_key = appconfig_get(&stream_config, key, "default proxy api key", rrdpush_api_key);
463 + rrdpush_api_key = appconfig_get(&stream_config, machine_guid, "proxy api key", rrdpush_api_key);
464 +
465 if(!strcmp(machine_guid, "localhost"))
466 host = localhost;
467 else
452 - host = rrdhost_find_or_create(hostname, machine_guid, os, update_every, history, mode, (health_enabled == CONFIG_BOOLEAN_NO)?0:1);
468 + host = rrdhost_find_or_create(
469 + hostname
470 + , machine_guid
471 + , os
472 + , update_every
473 + , history
474 + , mode
475 + , (health_enabled != CONFIG_BOOLEAN_NO)
476 + , (rrdpush_enabled && rrdpush_destination && *rrdpush_destination && rrdpush_api_key && *rrdpush_api_key)
477 + , rrdpush_destination
478 + , rrdpush_api_key
479 + );
480
481 if(!host) {
482 close(fd);
@@ -457,7 +484,8 @@ int rrdpush_receive(int fd, const char *key, const char *hostname, const char *m
484 return 1;
485 }
486
460 - info("STREAM %s [receive from [%s]:%s]: metrics for host '%s' with machine_guid '%s': update every = %d, history = %d, memory mode = %s, health %s"
487 +#ifdef NETDATA_INTERNAL_CHECKS
488 + info("STREAM %s [receive from [%s]:%s]: client willing to stream metrics for host '%s' with machine_guid '%s': update every = %d, history = %d, memory mode = %s, health %s"
489 , hostname
490 , client_ip
491 , client_port
@@ -468,6 +496,7 @@ int rrdpush_receive(int fd, const char *key, const char *hostname, const char *m
496 , rrd_memory_mode_name(host->rrd_memory_mode)
497 , (health_enabled == CONFIG_BOOLEAN_NO)?"disabled":((health_enabled == CONFIG_BOOLEAN_YES)?"enabled":"auto")
498 );
499 +#endif // NETDATA_INTERNAL_CHECKS
500
501 struct plugind cd = {
502 .enabled = 1,
@@ -487,8 +516,8 @@ int rrdpush_receive(int fd, const char *key, const char *hostname, const char *m
516 snprintfz(cd.cmd, PLUGINSD_CMD_MAX, "%s:%s", client_ip, client_port);
517
518 info("STREAM %s [receive from [%s]:%s]: initializing communication...", host->hostname, client_ip, client_port);
490 - if(send_timeout(fd, "STREAM", 6, 0, 60) != 6) {
491 - error("STREAM %s [receive from [%s]:%s]: cannot send STREAM command.", host->hostname, client_ip, client_port);
519 + if(send_timeout(fd, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT), 0, 60) != strlen(START_STREAMING_PROMPT)) {
520 + error("STREAM %s [receive from [%s]:%s]: cannot send ready command.", host->hostname, client_ip, client_port);
521 return 0;
522 }
523
@@ -510,9 +539,9 @@ int rrdpush_receive(int fd, const char *key, const char *hostname, const char *m
539 rrdhost_unlock(host);
540
541 // call the plugins.d processor to receive the metrics
513 - info("STREAM %s [receive from [%s]:%s]: receiving metrics... (host '%s', machine GUID '%s').", host->hostname, client_ip, client_port, host->hostname, host->machine_guid);
542 + info("STREAM %s [receive from [%s]:%s]: receiving metrics...", host->hostname, client_ip, client_port);
543 size_t count = pluginsd_process(host, &cd, fp, 1);
515 - error("STREAM %s [receive from [%s]:%s]: disconnected (host '%s', machine GUID '%s', completed updates %zu).", host->hostname, client_ip, client_port, host->hostname, host->machine_guid, count);
544 + error("STREAM %s [receive from [%s]:%s]: disconnected (completed updates %zu).", host->hostname, client_ip, client_port, count);
545
546 rrdhost_wrlock(host);
547 host->use_counter--;
@@ -564,10 +593,6 @@ void *rrdpush_receiver_thread(void *ptr) {
593 return NULL;
594 }
595
567 -static inline int rrdpush_receive_validate_api_key(const char *key) {
568 - return appconfig_get_boolean(&stream_config, key, "enabled", 0);
569 -}
570 -
596 void rrdpush_sender_thread_spawn(RRDHOST *host) {
597 if(pthread_create(&host->rrdpush_thread, NULL, rrdpush_sender_thread, (void *)host))
598 error("STREAM %s [send]: failed to create new thread for client.", host->hostname);
@@ -585,6 +610,7 @@ int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url
610
611 char *key = NULL, *hostname = NULL, *machine_guid = NULL, *os = NULL;
612 int update_every = default_rrd_update_every;
613 + char buf[GUID_LEN + 1];
614
615 while(url) {
616 char *value = mystrsep(&url, "?&");
@@ -627,8 +653,22 @@ int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url
653 return 400;
654 }
655
630 - if(!rrdpush_receive_validate_api_key(key)) {
631 - error("STREAM [receive from [%s]:%s]: API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
656 + if(regenerate_guid(key, buf) == -1) {
657 + error("STREAM [receive from [%s]:%s]: API key '%s' is not valid GUID. Forbidding access.", w->client_ip, w->client_port, key);
658 + buffer_flush(w->response.data);
659 + buffer_sprintf(w->response.data, "Your API key is invalid.");
660 + return 401;
661 + }
662 +
663 + if(regenerate_guid(machine_guid, buf) == -1) {
664 + error("STREAM [receive from [%s]:%s]: machine GUID '%s' is not GUID. Forbidding access.", w->client_ip, w->client_port, key);
665 + buffer_flush(w->response.data);
666 + buffer_sprintf(w->response.data, "Your machine GUID is invalid.");
667 + return 404;
668 + }
669 +
670 + if(!appconfig_get_boolean(&stream_config, key, "enabled", 1)) {
671 + error("STREAM [receive from [%s]:%s]: API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, machine_guid);
672 buffer_flush(w->response.data);
673 buffer_sprintf(w->response.data, "Your API key is not permitted access.");
674 return 401;
src/rrdpush.h
+2
@@ -2,6 +2,8 @@
2 #define NETDATA_RRDPUSH_H
3
4 extern int default_rrdpush_enabled;
5 +extern char *default_rrdpush_destination;
6 +extern char *default_rrdpush_api_key;
7
8 extern int rrdpush_init();
9 extern void rrdset_done_push(RRDSET *st);