@cryptotaxi247 / netdata-1 / commits / eb13ee695

properly lock all sensitive linked lists

Costa Tsaousis (ktsaou) committed Feb 20, 2017 at 22:45 UTC eb13ee695866a42e1504799775c3bff5a8dab471
10 files changed +214 -101
src/backends.c
+5 -11
@@ -309,31 +309,25 @@ void *backends_main(void *ptr) {
309 if(unlikely(pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &pthreadoldcancelstate) != 0))
310 error("Cannot set pthread cancel state to DISABLE.");
311
312 + rrd_rdlock();
313 RRDHOST *host;
313 - for(host = localhost; host ; host = host->next) {
314 - // for each host
315 -
314 + rrdhost_foreach_read(host) {
315 rrdhost_rdlock(host);
316
317 RRDSET *st;
319 - for(st = host->rrdset_root; st; st = st->next) {
320 - // for each chart
321 -
318 + rrdset_foreach_read(st, host) {
319 rrdset_rdlock(st);
320
321 RRDDIM *rd;
325 - for(rd = st->dimensions; rd; rd = rd->next) {
326 - // for each dimension
327 -
322 + rrddim_foreach_read(rd, st) {
323 if(rd->last_collected_time.tv_sec >= after)
324 chart_buffered_metrics += backend_request_formatter(b, prefix, host, (host == localhost)?hostname:host->hostname, st, rd, after, before, options);
325 }
331 -
326 rrdset_unlock(st);
327 }
334 -
328 rrdhost_unlock(host);
329 }
330 + rrd_unlock();
331
332 if(unlikely(pthread_setcancelstate(pthreadoldcancelstate, NULL) != 0))
333 error("Cannot set pthread cancel state to RESTORE (%d).", pthreadoldcancelstate);
src/health.c
+17 -10
@@ -57,34 +57,37 @@ void health_reload_host(RRDHOST *host) {
57 t->flags |= HEALTH_ENTRY_FLAG_UPDATED;
58 }
59
60 + rrdhost_rdlock(host);
61 // reset all thresholds to all charts
62 RRDSET *st;
62 - for(st = host->rrdset_root; st ; st = st->next) {
63 + rrdset_foreach_read(st, host) {
64 st->green = NAN;
65 st->red = NAN;
66 }
67 + rrdhost_unlock(host);
68
69 // load the new alarms
70 rrdhost_wrlock(host);
71 health_readdir(host, path);
70 - rrdhost_unlock(host);
72
73 // link the loaded alarms to their charts
73 - for(st = host->rrdset_root; st ; st = st->next) {
74 - rrdhost_wrlock(host);
75 -
74 + rrdset_foreach_write(st, host) {
75 rrdsetcalc_link_matching(st);
76 rrdcalctemplate_link_matching(st);
78 -
79 - rrdhost_unlock(host);
77 }
78 +
79 + rrdhost_unlock(host);
80 }
81
82 void health_reload(void) {
84 - RRDHOST *host;
83
86 - for(host = localhost; host ; host = host->next)
84 + rrd_rdlock();
85 +
86 + RRDHOST *host;
87 + rrdhost_foreach_read(host)
88 health_reload_host(host);
89 +
90 + rrd_unlock();
91 }
92
93 // ----------------------------------------------------------------------------
@@ -348,8 +351,10 @@ void *health_main(void *ptr) {
351 if(unlikely(pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, &oldstate) != 0))
352 error("Cannot set pthread cancel state to DISABLE.");
353
354 + rrd_rdlock();
355 +
356 RRDHOST *host;
352 - for(host = localhost; host ; host = host->next) {
357 + rrdhost_foreach_read(host) {
358 if(unlikely(!host->health_enabled)) continue;
359
360 rrdhost_rdlock(host);
@@ -596,6 +601,8 @@ void *health_main(void *ptr) {
601
602 } /* host loop */
603
604 + rrd_unlock();
605 +
606 if(unlikely(pthread_setcancelstate(oldstate, NULL) != 0))
607 error("Cannot set pthread cancel state to RESTORE (%d).", oldstate);
608
src/plugins_d.c
+18 -1
@@ -107,7 +107,23 @@ inline size_t pluginsd_process(RRDHOST *host, struct plugind *cd, FILE *fp, int
107 RRDSET *st = NULL;
108 uint32_t hash;
109
110 - while(likely(fgets(line, PLUGINSD_LINE_MAX, fp) != NULL)) {
110 + errno = 0;
111 + clearerr(fp);
112 +
113 + if(unlikely(fileno(fp) == -1)) {
114 + error("PLUGINSD: %s: file is not a valid stream.", cd->fullfilename);
115 + goto cleanup;
116 + }
117 +
118 + while(!ferror(fp)) {
119 + if(unlikely(netdata_exit)) break;
120 +
121 + char *r = fgets(line, PLUGINSD_LINE_MAX, fp);
122 + if(unlikely(!r)) {
123 + error("PLUGINSD: %s : read failed.", cd->fullfilename);
124 + break;
125 + }
126 +
127 if(unlikely(netdata_exit)) break;
128
129 line[PLUGINSD_LINE_MAX] = '\0';
@@ -335,6 +351,7 @@ inline size_t pluginsd_process(RRDHOST *host, struct plugind *cd, FILE *fp, int
351 }
352 }
353
354 +cleanup:
355 cd->enabled = enabled;
356
357 if(likely(count)) {
src/rrd.h
+53 -14
@@ -185,6 +185,15 @@ struct rrddim {
185 };
186 typedef struct rrddim RRDDIM;
187
188 +// ----------------------------------------------------------------------------
189 +// these loop macros make sure the linked list is accessed with the right lock
190 +
191 +#define rrddim_foreach_read(rd, st) \
192 + for(rd = st->dimensions, rrdset_check_rdlock(st); rd ; rd = rd->next)
193 +
194 +#define rrddim_foreach_write(rd, st) \
195 + for(rd = st->dimensions, rrdset_check_wrlock(st); rd ; rd = rd->next)
196 +
197
198 // ----------------------------------------------------------------------------
199 // RRDSET - this is a chart
@@ -252,7 +261,7 @@ struct rrdset {
261 char *cache_dir; // the directory to store dimensions
262 char cache_filename[FILENAME_MAX+1]; // the filename to store this set
263
255 - pthread_rwlock_t rrdset_rwlock;
264 + pthread_rwlock_t rrdset_rwlock; // protects dimensions linked list
265
266 unsigned long counter; // the number of times we added values to this rrd
267 unsigned long counter_done; // the number of times we added values to this rrd
@@ -305,6 +314,16 @@ typedef struct rrdset RRDSET;
314 #define rrdset_wrlock(st) pthread_rwlock_wrlock(&((st)->rrdset_rwlock))
315 #define rrdset_unlock(st) pthread_rwlock_unlock(&((st)->rrdset_rwlock))
316
317 +// ----------------------------------------------------------------------------
318 +// these loop macros make sure the linked list is accessed with the right lock
319 +
320 +#define rrdset_foreach_read(st, host) \
321 + for(st = host->rrdset_root, rrdhost_check_rdlock(host); st ; st = st->next)
322 +
323 +#define rrdset_foreach_write(st, host) \
324 + for(st = host->rrdset_root, rrdhost_check_wrlock(host); st ; st = st->next)
325 +
326 +
327 // ----------------------------------------------------------------------------
328 // RRD HOST
329
@@ -325,7 +344,7 @@ struct rrdhost {
344
345 RRDSET *rrdset_root; // the host charts
346
328 - pthread_rwlock_t rrdhost_rwlock; // lock for this RRDHOST
347 + pthread_rwlock_t rrdhost_rwlock; // lock for this RRDHOST (protects rrdset_root linked list)
348
349 avl_tree_lock rrdset_root_index; // the host's charts index (by id)
350 avl_tree_lock rrdset_root_index_name; // the host's charts index (by name)
@@ -364,35 +383,55 @@ extern RRDHOST *localhost;
383 #define rrdhost_wrlock(h) pthread_rwlock_wrlock(&((h)->rrdhost_rwlock))
384 #define rrdhost_unlock(h) pthread_rwlock_unlock(&((h)->rrdhost_rwlock))
385
386 +// ----------------------------------------------------------------------------
387 +// these loop macros make sure the linked list is accessed with the right lock
388 +
389 +#define rrdhost_foreach_read(var) \
390 + for(var = localhost, rrd_check_rdlock(); var ; var = var->next)
391 +
392 +#define rrdhost_foreach_write(var) \
393 + for(var = localhost, rrd_check_wrlock(); var ; var = var->next)
394 +
395 +
396 +// ----------------------------------------------------------------------------
397 +// global lock for all RRDHOSTs
398 +
399 +extern pthread_rwlock_t rrd_rwlock;
400 +#define rrd_rdlock() pthread_rwlock_rdlock(&rrd_rwlock)
401 +#define rrd_wrlock() pthread_rwlock_wrlock(&rrd_rwlock)
402 +#define rrd_unlock() pthread_rwlock_unlock(&rrd_rwlock)
403 +
404 +// ----------------------------------------------------------------------------
405 +
406 extern void rrd_init(char *hostname);
407
408 extern RRDHOST *rrdhost_find(const char *guid, uint32_t hash);
409 extern RRDHOST *rrdhost_find_or_create(const char *hostname, const char *guid);
410
411 #ifdef NETDATA_INTERNAL_CHECKS
412 +extern void rrdhost_check_wrlock_int(RRDHOST *host, const char *file, const char *function, const unsigned long line);
413 +extern void rrdhost_check_rdlock_int(RRDHOST *host, const char *file, const char *function, const unsigned long line);
414 +extern void rrdset_check_rdlock_int(RRDSET *st, const char *file, const char *function, const unsigned long line);
415 +extern void rrdset_check_wrlock_int(RRDSET *st, const char *file, const char *function, const unsigned long line);
416 +extern void rrd_check_rdlock_int(const char *file, const char *function, const unsigned long line);
417 +extern void rrd_check_wrlock_int(const char *file, const char *function, const unsigned long line);
418 +
419 #define rrdhost_check_rdlock(host) rrdhost_check_rdlock_int(host, __FILE__, __FUNCTION__, __LINE__)
420 #define rrdhost_check_wrlock(host) rrdhost_check_wrlock_int(host, __FILE__, __FUNCTION__, __LINE__)
421 +#define rrdset_check_rdlock(st) rrdset_check_rdlock_int(st, __FILE__, __FUNCTION__, __LINE__)
422 +#define rrdset_check_wrlock(st) rrdset_check_wrlock_int(st, __FILE__, __FUNCTION__, __LINE__)
423 #define rrd_check_rdlock() rrd_check_rdlock_int(__FILE__, __FUNCTION__, __LINE__)
424 #define rrd_check_wrlock() rrd_check_wrlock_int(__FILE__, __FUNCTION__, __LINE__)
425 +
426 #else
427 #define rrdhost_check_rdlock(host) (void)0
428 #define rrdhost_check_wrlock(host) (void)0
429 +#define rrdset_check_rdlock(host) (void)0
430 +#define rrdset_check_wrlock(host) (void)0
431 #define rrd_check_rdlock() (void)0
432 #define rrd_check_wrlock() (void)0
433 #endif
434
384 -extern void rrdhost_check_wrlock_int(RRDHOST *host, const char *file, const char *function, const unsigned long line);
385 -extern void rrdhost_check_rdlock_int(RRDHOST *host, const char *file, const char *function, const unsigned long line);
386 -
387 -// ----------------------------------------------------------------------------
388 -// global lock for all RRDHOSTs
389 -
390 -extern pthread_rwlock_t rrd_rwlock;
391 -#define rrd_rdlock() pthread_rwlock_rdlock(&rrd_rwlock)
392 -#define rrd_wrlock() pthread_rwlock_wrlock(&rrd_rwlock)
393 -#define rrd_unlock() pthread_rwlock_unlock(&rrd_rwlock)
394 -
395 -
435 // ----------------------------------------------------------------------------
436 // RRDSET functions
437
src/rrd2json.c
+26 -14
@@ -43,7 +43,7 @@ void rrd_stats_api_v1_chart_with_data(RRDSET *st, BUFFER *wb, size_t *dimensions
43
44 size_t dimensions = 0;
45 RRDDIM *rd;
46 - for(rd = st->dimensions; rd ; rd = rd->next) {
46 + rrddim_foreach_read(rd, st) {
47 if(rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN)) continue;
48
49 memory += rd->memsize;
@@ -97,8 +97,9 @@ void rrd_stats_api_v1_charts(RRDHOST *host, BUFFER *wb)
97 , host->rrd_history_entries
98 );
99
100 + c = 0;
101 rrdhost_rdlock(host);
101 - for(st = host->rrdset_root, c = 0; st ; st = st->next) {
102 + rrdset_foreach_read(st, host) {
103 if(rrdset_flag_check(st, RRDSET_FLAG_ENABLED) && st->dimensions) {
104 if(c) buffer_strcat(wb, ",");
105 buffer_strcat(wb, "\n\t\t\"");
@@ -157,7 +158,7 @@ void rrd_stats_api_v1_charts_allmetrics_prometheus(RRDHOST *host, BUFFER *wb) {
158
159 // for each chart
160 RRDSET *st;
160 - for(st = host->rrdset_root; st ; st = st->next) {
161 + rrdset_foreach_read(st, host) {
162 char chart[PROMETHEUS_ELEMENT_MAX + 1];
163 prometheus_name_copy(chart, st->id, PROMETHEUS_ELEMENT_MAX);
164
@@ -167,7 +168,7 @@ void rrd_stats_api_v1_charts_allmetrics_prometheus(RRDHOST *host, BUFFER *wb) {
168
169 // for each dimension
170 RRDDIM *rd;
170 - for(rd = st->dimensions; rd ; rd = rd->next) {
171 + rrddim_foreach_read(rd, st) {
172 if(rd->counter) {
173 char dimension[PROMETHEUS_ELEMENT_MAX + 1];
174 prometheus_name_copy(dimension, rd->id, PROMETHEUS_ELEMENT_MAX);
@@ -228,7 +229,7 @@ void rrd_stats_api_v1_charts_allmetrics_shell(RRDHOST *host, BUFFER *wb) {
229
230 // for each chart
231 RRDSET *st;
231 - for(st = host->rrdset_root; st ; st = st->next) {
232 + rrdset_foreach_read(st, host) {
233 calculated_number total = 0.0;
234 char chart[SHELL_ELEMENT_MAX + 1];
235 shell_name_copy(chart, st->id, SHELL_ELEMENT_MAX);
@@ -239,7 +240,7 @@ void rrd_stats_api_v1_charts_allmetrics_shell(RRDHOST *host, BUFFER *wb) {
240
241 // for each dimension
242 RRDDIM *rd;
242 - for(rd = st->dimensions; rd ; rd = rd->next) {
243 + rrddim_foreach_read(rd, st) {
244 if(rd->counter) {
245 char dimension[SHELL_ELEMENT_MAX + 1];
246 shell_name_copy(dimension, rd->id, SHELL_ELEMENT_MAX);
@@ -350,7 +351,7 @@ unsigned long rrd_stats_one_json(RRDSET *st, char *options, BUFFER *wb)
351 unsigned long memory = st->memsize;
352
353 RRDDIM *rd;
353 - for(rd = st->dimensions; rd ; rd = rd->next) {
354 + rrddim_foreach_read(rd, st) {
355
356 memory += rd->memsize;
357
@@ -411,20 +412,20 @@ void rrd_stats_graph_json(RRDSET *st, char *options, BUFFER *wb)
412 void rrd_stats_all_json(RRDHOST *host, BUFFER *wb)
413 {
414 unsigned long memory = 0;
414 - long c;
415 + long c = 0;
416 RRDSET *st;
417
418 buffer_strcat(wb, RRD_GRAPH_JSON_HEADER);
419
420 rrdhost_rdlock(host);
420 -
421 - for(st = host->rrdset_root, c = 0; st ; st = st->next) {
421 + rrdset_foreach_read(st, host) {
422 if(rrdset_flag_check(st, RRDSET_FLAG_ENABLED) && st->dimensions) {
423 if(c) buffer_strcat(wb, ",\n");
424 memory += rrd_stats_one_json(st, NULL, wb);
425 c++;
426 }
427 }
428 + rrdhost_unlock(host);
429
430 buffer_sprintf(wb, "\n\t],\n"
431 "\t\"hostname\": \"%s\",\n"
@@ -437,8 +438,6 @@ void rrd_stats_all_json(RRDHOST *host, BUFFER *wb)
438 , host->rrd_history_entries
439 , memory
440 );
440 -
441 - rrdhost_unlock(host);
441 }
442
443
@@ -544,6 +543,8 @@ static void rrdr_dump(RRDR *r)
543 */
544
545 void rrdr_disable_not_selected_dimensions(RRDR *r, uint32_t options, const char *dims) {
546 + rrdset_check_rdlock(r->st);
547 +
548 if(unlikely(!dims || !*dims)) return;
549
550 char b[strlen(dims) + 1];
@@ -649,6 +650,8 @@ void rrdr_buffer_print_format(BUFFER *wb, uint32_t format)
650
651 uint32_t rrdr_check_options(RRDR *r, uint32_t options, const char *dims)
652 {
653 + rrdset_check_rdlock(r->st);
654 +
655 (void)dims;
656
657 if(options & RRDR_OPTION_NONZERO) {
@@ -684,6 +687,8 @@ uint32_t rrdr_check_options(RRDR *r, uint32_t options, const char *dims)
687
688 void rrdr_json_wrapper_begin(RRDR *r, BUFFER *wb, uint32_t format, uint32_t options, int string_value)
689 {
690 + rrdset_check_rdlock(r->st);
691 +
692 long rows = rrdr_rows(r);
693 long c, i;
694 RRDDIM *rd;
@@ -865,6 +870,8 @@ void rrdr_json_wrapper_end(RRDR *r, BUFFER *wb, uint32_t format, uint32_t option
870
871 static void rrdr2json(RRDR *r, BUFFER *wb, uint32_t options, int datatable)
872 {
873 + rrdset_check_rdlock(r->st);
874 +
875 //info("RRD2JSON(): %s: BEGIN", r->st->id);
876 int row_annotations = 0, dates, dates_with_new = 0;
877 char kq[2] = "", // key quote
@@ -1091,6 +1098,8 @@ static void rrdr2json(RRDR *r, BUFFER *wb, uint32_t options, int datatable)
1098
1099 static void rrdr2csv(RRDR *r, BUFFER *wb, uint32_t options, const char *startline, const char *separator, const char *endline, const char *betweenlines)
1100 {
1101 + rrdset_check_rdlock(r->st);
1102 +
1103 //info("RRD2CSV(): %s: BEGIN", r->st->id);
1104 long c, i;
1105 RRDDIM *d;
@@ -1196,6 +1205,8 @@ static void rrdr2csv(RRDR *r, BUFFER *wb, uint32_t options, const char *startlin
1205 }
1206
1207 inline static calculated_number rrdr2value(RRDR *r, long i, uint32_t options, int *all_values_are_null) {
1208 + rrdset_check_rdlock(r->st);
1209 +
1210 long c;
1211 RRDDIM *d;
1212
@@ -1396,7 +1407,7 @@ static RRDR *rrdr_create(RRDSET *st, long n)
1407 rrdr_lock_rrdset(r);
1408
1409 RRDDIM *rd;
1399 - for(rd = st->dimensions ; rd ; rd = rd->next) r->d++;
1410 + rrddim_foreach_read(rd, st) r->d++;
1411
1412 r->n = n;
1413
@@ -1599,6 +1610,7 @@ RRDR *rrd2rrdr(RRDSET *st, long points, long long after, long long before, int g
1610 // initialize them
1611 RRDDIM *rd;
1612 long c;
1613 + rrdset_check_rdlock(st);
1614 for( rd = st->dimensions, c = 0 ; rd && c < dimensions ; rd = rd->next, c++) {
1615 last_values[c] = 0;
1616 group_values[c] = (group_method == GROUP_MAX || group_method == GROUP_MIN)?NAN:0;
@@ -2036,7 +2048,7 @@ time_t rrd_stats_json(int type, RRDSET *st, BUFFER *wb, long points, long group,
2048
2049 int dimensions = 0;
2050 RRDDIM *rd;
2039 - for( rd = st->dimensions ; rd ; rd = rd->next) dimensions++;
2051 + rrddim_foreach_read(rd, st) dimensions++;
2052 if(!dimensions) {
2053 rrdset_unlock(st);
2054 buffer_strcat(wb, "No dimensions yet.");
src/rrdcalc.c
+1 -1
@@ -285,7 +285,7 @@ inline void rrdcalc_create_part2(RRDHOST *host, RRDCALC *rc) {
285
286 // link it to its chart
287 RRDSET *st;
288 - for(st = host->rrdset_root; st ; st = st->next) {
288 + rrdset_foreach_read(st, host) {
289 if(rrdcalc_is_matching_this_rrdset(rc, st)) {
290 rrdsetcalc_link(st, rc);
291 break;
src/rrdhost.c
+3 -3
@@ -318,7 +318,7 @@ void rrdhost_save(RRDHOST *host) {
318 // to ensure only one thread is saving the database
319 rrdhost_wrlock(host);
320
321 - for(st = host->rrdset_root; st ; st = st->next) {
321 + rrdset_foreach_write(st, host) {
322 rrdset_rdlock(st);
323
324 if(st->rrd_memory_mode == RRD_MEMORY_MODE_SAVE) {
@@ -326,7 +326,7 @@ void rrdhost_save(RRDHOST *host) {
326 savememory(st->cache_filename, st, st->memsize);
327 }
328
329 - for(rd = st->dimensions; rd ; rd = rd->next) {
329 + rrddim_foreach_read(rd, st) {
330 if(likely(rd->rrd_memory_mode == RRD_MEMORY_MODE_SAVE)) {
331 debug(D_RRD_STATS, "Saving dimension '%s' to '%s'.", rd->name, rd->cache_filename);
332 savememory(rd->cache_filename, rd, rd->memsize);
@@ -345,7 +345,7 @@ void rrdhost_save_all(void) {
345 rrd_rdlock();
346
347 RRDHOST *host;
348 - for(host = localhost; host ; host = host->next)
348 + rrdhost_foreach_read(host)
349 rrdhost_save(host);
350
351 rrd_unlock();
src/rrdpush.c
+51 -32
@@ -19,7 +19,7 @@ static inline void rrdpush_unlock() {
19
20 static inline int need_to_send_chart_definitions(RRDSET *st) {
21 RRDDIM *rd;
22 - for(rd = st->dimensions; rd ;rd = rd->next)
22 + rrddim_foreach_read(rd, st)
23 if(rrddim_flag_check(rd, RRDDIM_FLAG_UPDATED) && !rrddim_flag_check(rd, RRDDIM_FLAG_EXPOSED))
24 return 1;
25
@@ -40,7 +40,7 @@ static inline void send_chart_definitions(RRDSET *st) {
40 );
41
42 RRDDIM *rd;
43 - for(rd = st->dimensions; rd ;rd = rd->next) {
43 + rrddim_foreach_read(rd, st) {
44 buffer_sprintf(rrdpush_buffer, "DIMENSION '%s' '%s' '%s' " COLLECTED_NUMBER_FORMAT " " COLLECTED_NUMBER_FORMAT " '%s %s'\n"
45 , rd->id
46 , rd->name
@@ -50,6 +50,7 @@ static inline void send_chart_definitions(RRDSET *st) {
50 , rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN)?"hidden":""
51 , rrddim_flag_check(rd, RRDDIM_FLAG_DONT_DETECT_RESETS_OR_OVERFLOWS)?"noreset":""
52 );
53 + rrddim_flag_set(rd, RRDDIM_FLAG_EXPOSED);
54 }
55 }
56
@@ -57,7 +58,7 @@ static inline void send_chart_metrics(RRDSET *st) {
58 buffer_sprintf(rrdpush_buffer, "BEGIN %s %llu\n", st->id, st->usec_since_last_update);
59
60 RRDDIM *rd;
60 - for(rd = st->dimensions; rd ;rd = rd->next) {
61 + rrddim_foreach_read(rd, st) {
62 if(rrddim_flag_check(rd, RRDDIM_FLAG_UPDATED))
63 buffer_sprintf(rrdpush_buffer, "SET %s = " COLLECTED_NUMBER_FORMAT "\n"
64 , rd->id
@@ -71,35 +72,39 @@ static inline void send_chart_metrics(RRDSET *st) {
72 static void reset_all_charts(void) {
73 rrd_rdlock();
74
74 - RRDHOST *h;
75 - for(h = localhost; h ;h = h->next) {
75 + RRDHOST *host;
76 + rrdhost_foreach_read(host) {
77 + rrdhost_rdlock(host);
78 +
79 RRDSET *st;
77 - for(st = h->rrdset_root ; st ; st = st->next) {
80 + rrdset_foreach_read(st, host) {
81 rrdset_rdlock(st);
82
83 RRDDIM *rd;
81 - for(rd = st->dimensions; rd ;rd = rd->next)
84 + rrddim_foreach_read(rd, st)
85 rrddim_flag_clear(rd, RRDDIM_FLAG_EXPOSED);
86
87 rrdset_unlock(st);
88 }
89 + rrdhost_unlock(host);
90 }
91 + rrd_unlock();
92
93 last_host = NULL;
89 -
90 - rrd_unlock();
94 }
95
96 void rrdset_done_push(RRDSET *st) {
97
95 - if(!rrdset_flag_check(st, RRDSET_FLAG_ENABLED))
98 + if(unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ENABLED) || !rrdpush_buffer))
99 return;
100
101 rrdpush_lock();
102 rrdset_rdlock(st);
103
101 - if(st->rrdhost != last_host)
102 - buffer_sprintf(rrdpush_buffer, "HOST '%s' '%s'\n", st->rrdhost->hostname, st->rrdhost->machine_guid);
104 + if(st->rrdhost != last_host) {
105 + buffer_sprintf(rrdpush_buffer, "HOST '%s' '%s'\n", st->rrdhost->machine_guid, st->rrdhost->hostname);
106 + last_host = st->rrdhost;
107 + }
108
109 if(need_to_send_chart_definitions(st))
110 send_chart_definitions(st);
@@ -139,36 +144,58 @@ void *central_netdata_push_thread(void *ptr) {
144 size_t begin = 0;
145 size_t max_size = 1024 * 1024;
146 size_t reconnects_counter = 0;
147 + size_t sent_bytes = 0;
148 + size_t sent_connection = 0;
149 int sock = -1;
150 char buffer[1];
151
152 for(;;) {
153 if(unlikely(sock == -1)) {
154 + info("PUSH: connecting to central netdata at: %s", central_netdata_to_push_data);
155 sock = connect_to_one_of(central_netdata_to_push_data, 19999, &tv, &reconnects_counter);
156
157 if(unlikely(sock != -1)) {
150 - if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
151 - error("Cannot set non-blocking mode for socket.");
158 + info("PUSH: connected to central netdata at: %s", central_netdata_to_push_data);
159
153 - buffer_sprintf(rrdpush_buffer, "GET /stream?key=%s\r\n\r\n", config_get("global", "central netdata api key", ""));
154 - reset_all_charts();
160 + if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
161 + error("PUSH: cannot set non-blocking mode for socket.");
162 }
163 + else
164 + error("PUSH: failed to connect to central netdata at: %s", central_netdata_to_push_data);
165 +
166 + rrdpush_lock();
167 + if(buffer_strlen(rrdpush_buffer))
168 + error("PUSH: discarding %zu bytes of metrics data already in the buffer.", buffer_strlen(rrdpush_buffer));
169 +
170 + buffer_flush(rrdpush_buffer);
171 + buffer_sprintf(rrdpush_buffer, "GET /stream?key=%s HTTP/1.1\r\nUser-Agent: netdata-push-service/%s\r\nAccept: */*\r\n\r\n", config_get("global", "central netdata api key", ""), VERSION);
172 + reset_all_charts();
173 + rrdpush_unlock();
174 + sent_connection = 0;
175 }
176
177 if(read(rrdpush_pipe[PIPE_READ], buffer, 1) == -1) {
159 - error("Cannot read from internal pipe.");
178 + error("PUSH: Cannot read from internal pipe.");
179 sleep(1);
180 }
181
163 - if(likely(sock != -1)) {
182 + if(likely(sock != -1 && begin < rrdpush_buffer->len)) {
183 + // fprintf(stderr, "PUSH BEGIN\n");
184 + // fwrite(&rrdpush_buffer->buffer[begin], 1, rrdpush_buffer->len - begin, stderr);
185 + // fprintf(stderr, "\nPUSH END\n");
186 +
187 rrdpush_lock();
165 - ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], rrdpush_buffer->len, MSG_DONTWAIT);
188 + ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], rrdpush_buffer->len - begin, MSG_DONTWAIT);
189 if(ret == -1) {
167 - error("Failed to send metrics to central netdata at %s", central_netdata_to_push_data);
168 - close(sock);
169 - sock = -1;
190 + if(errno != EAGAIN) {
191 + error("PUSH: failed to send metrics to central netdata at %s. We have sent %zu bytes on this connection.", central_netdata_to_push_data, sent_connection);
192 + close(sock);
193 + sock = -1;
194 + }
195 }
196 else {
197 + sent_connection += ret;
198 + sent_bytes += ret;
199 begin += ret;
200 if(begin == rrdpush_buffer->len) {
201 buffer_flush(rrdpush_buffer);
@@ -180,23 +207,15 @@ void *central_netdata_push_thread(void *ptr) {
207
208 // protection from overflow
209 if(rrdpush_buffer->len > max_size) {
183 - rrdpush_lock();
184 -
185 - error("Discarding %zu bytes of metrics data, because we cannot connect to central netdata at %s"
186 - , buffer_strlen(rrdpush_buffer), central_netdata_to_push_data);
187 -
188 - buffer_flush(rrdpush_buffer);
189 -
210 + errno = 0;
211 + error("PUSH: 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.", rrdpush_buffer->len, rrdpush_buffer->len - begin, sent_bytes, sent_connection);
212 if(sock != -1) {
213 close(sock);
214 sock = -1;
215 }
194 -
195 - rrdpush_unlock();
216 }
217 }
218
199 -cleanup:
219 debug(D_WEB_CLIENT, "Central netdata push thread exits.");
220 if(sock != -1)
221 close(sock);
src/rrdset.c
+26 -7
@@ -3,6 +3,23 @@
3
4 #define RRD_DEFAULT_GAP_INTERPOLATIONS 1
5
6 +void rrdset_check_rdlock_int(RRDSET *st, const char *file, const char *function, const unsigned long line) {
7 + debug(D_RRD_CALLS, "Checking read lock on chart '%s'", st->id);
8 +
9 + int ret = pthread_rwlock_trywrlock(&st->rrdset_rwlock);
10 + if(ret == 0)
11 + fatal("RRDSET '%s' should be read-locked, but it is not, at function %s() at line %lu of file '%s'", st->id, function, line, file);
12 +}
13 +
14 +void rrdset_check_wrlock_int(RRDSET *st, const char *file, const char *function, const unsigned long line) {
15 + debug(D_RRD_CALLS, "Checking write lock on chart '%s'", st->id);
16 +
17 + int ret = pthread_rwlock_tryrdlock(&st->rrdset_rwlock);
18 + if(ret == 0)
19 + fatal("RRDSET '%s' should be write-locked, but it is not, at function %s() at line %lu of file '%s'", st->id, function, line, file);
20 +}
21 +
22 +
23 // ----------------------------------------------------------------------------
24 // RRDSET index
25
@@ -143,7 +160,7 @@ void rrdset_set_name(RRDSET *st, const char *name) {
160
161 rrdset_wrlock(st);
162 RRDDIM *rd;
146 - for(rd = st->dimensions; rd ;rd = rd->next)
163 + rrddim_foreach_write(rd, st)
164 rrddimvar_rename_all(rd);
165 rrdset_unlock(st);
166
@@ -167,7 +184,7 @@ void rrdset_reset(RRDSET *st) {
184 st->counter_done = 0;
185
186 RRDDIM *rd;
170 - for(rd = st->dimensions; rd ; rd = rd->next) {
187 + rrddim_foreach_read(rd, st) {
188 rd->last_collected_time.tv_sec = 0;
189 rd->last_collected_time.tv_usec = 0;
190 rd->counter = 0;
@@ -674,18 +691,20 @@ void rrdset_done(RRDSET *st) {
691 st->counter_done++;
692
693 // calculate totals and count the dimensions
677 - int dimensions;
694 + int dimensions = 0;
695 st->collected_total = 0;
679 - for( rd = st->dimensions, dimensions = 0 ; rd ; rd = rd->next, dimensions++ )
696 + rrddim_foreach_read(rd, st) {
697 + dimensions++;
698 if(likely(rrddim_flag_check(rd, RRDDIM_FLAG_UPDATED)))
699 st->collected_total += rd->collected_value;
700 + }
701
702 uint32_t storage_flags = SN_EXISTS;
703
704 // process all dimensions to calculate their values
705 // based on the collected figures only
706 // at this stage we do not interpolate anything
688 - for( rd = st->dimensions ; rd ; rd = rd->next ) {
707 + rrddim_foreach_read(rd, st) {
708
709 if(unlikely(!rrddim_flag_check(rd, RRDDIM_FLAG_UPDATED))) {
710 rd->calculated_value = 0;
@@ -890,7 +909,7 @@ void rrdset_done(RRDSET *st) {
909 st->last_updated.tv_sec = (time_t) (next_store_ut / USEC_PER_SEC);
910 st->last_updated.tv_usec = 0;
911
893 - for( rd = st->dimensions ; likely(rd) ; rd = rd->next ) {
912 + rrddim_foreach_read(rd, st) {
913 calculated_number new_value;
914
915 switch(rd->algorithm) {
@@ -1035,7 +1054,7 @@ void rrdset_done(RRDSET *st) {
1054
1055 st->last_collected_total = st->collected_total;
1056
1038 - for( rd = st->dimensions; rd ; rd = rd->next ) {
1057 + rrddim_foreach_read(rd, st) {
1058 if(unlikely(!rrddim_flag_check(rd, RRDDIM_FLAG_UPDATED)))
1059 continue;
1060
src/web_client.c
+14 -8
@@ -1672,7 +1672,7 @@ int validate_stream_api_key(const char *key) {
1672 }
1673
1674 int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1675 - info("STREAM request from client '%s:%s' for host '%s'", w->client_ip, w->client_port, host->hostname);
1675 + info("STREAM request from client '%s:%s', starting as host '%s'", w->client_ip, w->client_port, host->hostname);
1676
1677 char *key = NULL;
1678
@@ -1689,16 +1689,16 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1689 }
1690
1691 if(!key || !*key) {
1692 + error("STREAM request from client '%s:%s', without an API key. Forbidding access.", w->client_ip, w->client_port);
1693 buffer_flush(w->response.data);
1694 buffer_sprintf(w->response.data, "You need an API key for this request.");
1694 - error("STREAM request from client '%s:%s', without an API key. Forbidding access.", w->client_ip, w->client_port);
1695 return 401;
1696 }
1697
1698 if(!validate_stream_api_key(key)) {
1699 + error("STREAM request from client '%s:%s': API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
1700 buffer_flush(w->response.data);
1701 buffer_sprintf(w->response.data, "Your API key is not permitted access.");
1701 - error("STREAM request from client '%s:%s': API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
1702 return 401;
1703 }
1704
@@ -1733,6 +1733,7 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1733 }
1734
1735 // call the plugins.d processor to receive the metrics
1736 + info("STREAM connecting client '%s:%s' to plugins.d.", w->client_ip, w->client_port);
1737 size_t count = pluginsd_process(host, &cd, fp, 1);
1738 error("STREAM from '%s:%s': client disconnected.", w->client_ip, w->client_port);
1739
@@ -2188,12 +2189,16 @@ static inline int web_client_switch_host(RRDHOST *host, struct web_client *w, ch
2189 if(unlikely(hash == hash_localhost && !strcmp(tok, "localhost")))
2190 return web_client_process_url(localhost, w, url);
2191
2192 + rrd_rdlock();
2193 RRDHOST *h;
2192 - for(h = localhost; h; h = h->next) {
2194 + rrdhost_foreach_read(h) {
2195 if(unlikely((hash == h->hash_hostname && !strcmp(tok, h->hostname)) ||
2194 - (hash == h->hash_machine_guid && !strcmp(tok, h->machine_guid))))
2196 + (hash == h->hash_machine_guid && !strcmp(tok, h->machine_guid)))) {
2197 + rrd_unlock();
2198 return web_client_process_url(h, w, url);
2199 + }
2200 }
2201 + rrd_unlock();
2202 }
2203
2204 buffer_flush(w->response.data);
@@ -2301,10 +2306,11 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2306 debug(D_WEB_CLIENT_ACCESS, "%llu: Sending list of RRD_STATS...", w->id);
2307
2308 buffer_flush(w->response.data);
2304 - RRDSET *st = host->rrdset_root;
2309 + RRDSET *st;
2310
2306 - for ( ; st ; st = st->next )
2307 - buffer_sprintf(w->response.data, "%s\n", st->name);
2311 + rrdhost_rdlock(host);
2312 + rrdset_foreach_read(st, host) buffer_sprintf(w->response.data, "%s\n", st->name);
2313 + rrdhost_unlock(host);
2314
2315 return 200;
2316 }