@cryptotaxi247 / netdata-1 / commits / f154e189c

complete statsd support, single and multi-threaded

Costa Tsaousis (ktsaou) committed Apr 25, 2017 at 00:27 UTC f154e189cb7375bdf865cd82bc931443eba86d31
1 file changed +468 -277
src/statsd.c
+468 -277
@@ -1,264 +1,396 @@
1 #include "common.h"
2
3 -static struct statsd {
4 - size_t events;
5 - size_t events_gauge;
6 - size_t events_counter;
7 - size_t events_timer;
8 - size_t events_meter;
9 - size_t events_histogram;
10 -
11 - size_t metrics;
12 - size_t metrics_gauge;
13 - size_t metrics_counter;
14 - size_t metrics_timer;
15 - size_t metrics_meter;
16 - size_t metrics_histogram;
17 -} statsd = {
18 - .events = 0,
19 - .metrics = 0
20 -};
21 -
3 // --------------------------------------------------------------------------------------
4
24 -static LISTEN_SOCKETS statsd_sockets = {
25 - .config_section = CONFIG_SECTION_STATSD,
26 - .default_bind_to = "udp:* tcp:*",
27 - .default_port = STATSD_LISTEN_PORT,
28 - .backlog = STATSD_LISTEN_BACKLOG
29 -};
30 -
31 -int statsd_listen_sockets_setup(void) {
32 - return listen_sockets_setup(&statsd_sockets);
33 -}
34 -
35 -// --------------------------------------------------------------------------------------
36 -
37 -typedef enum statsd_metric_type {
38 - STATSD_METRIC_TYPE_GAUGE = 'g',
39 - STATSD_METRIC_TYPE_COUNTER = 'c',
40 - STATSD_METRIC_TYPE_TIMER = 't',
41 - STATSD_METRIC_TYPE_HISTOGRAM = 'h',
42 - STATSD_METRIC_TYPE_METER = 'm'
43 -} STATSD_METRIC_TYPE;
5 +// #define STATSD_MULTITHREADED 1
6 +
7 +#ifdef STATSD_MULTITHREADED
8 +#define STATSD_AVL_TREE avl_tree_lock
9 +#define STATSD_AVL_INSERT avl_insert_lock
10 +#define STATSD_AVL_SEARCH avl_search_lock
11 +#define STATSD_AVL_INDEX_INIT { .avl_tree = { NULL, statsd_metric_compare }, .rwlock = AVL_LOCK_INITIALIZER }
12 +#define STATSD_FIRST_PTR_MUTEX netdata_mutex_t first_mutex
13 +#define STATSD_FIRST_PTR_MUTEX_INIT .first_mutex = NETDATA_MUTEX_INITIALIZER
14 +#define STATSD_FIRST_PTR_MUTEX_LOCK(index) netdata_mutex_lock(&((index)->first_mutex))
15 +#define STATSD_FIRST_PTR_MUTEX_UNLOCK(index) netdata_mutex_unlock(&((index)->first_mutex))
16 +#define STATSD_DICTIONARY_OPTIONS DICTIONARY_FLAG_DEFAULT
17 +#else
18 +#define STATSD_AVL_TREE avl_tree
19 +#define STATSD_AVL_INSERT avl_insert
20 +#define STATSD_AVL_SEARCH avl_search
21 +#define STATSD_AVL_INDEX_INIT { .root = NULL, .compar = statsd_metric_compare }
22 +#define STATSD_FIRST_PTR_MUTEX
23 +#define STATSD_FIRST_PTR_MUTEX_INIT
24 +#define STATSD_FIRST_PTR_MUTEX_LOCK(index)
25 +#define STATSD_FIRST_PTR_MUTEX_UNLOCK(index)
26 +#define STATSD_DICTIONARY_OPTIONS DICTIONARY_FLAG_SINGLE_THREADED
27 +#endif
28 +
29 +typedef struct statsd_metric_gauge {
30 + long double value;
31 +} STATSD_METRIC_GAUGE;
32 +
33 +typedef struct statsd_metric_counter {
34 + long long value;
35 +} STATSD_METRIC_COUNTER;
36 +
37 +typedef struct statsd_metric_histogram {
38 + size_t size;
39 + size_t used;
40 + long double *values;
41 +} STATSD_METRIC_HISTOGRAM;
42 +
43 +typedef struct statsd_metric_set {
44 + DICTIONARY *dict;
45 + unsigned long long unique;
46 +} STATSD_METRIC_SET;
47 +
48 +typedef struct statsd_metric_set_value {
49 + char *value;
50 + uint32_t hash;
51 +} STATSD_METRIC_SET_VALUE;
52
53 typedef struct statsd_metric {
54 avl avl; // indexing
55 + struct statsd_metric *next;
56
57 const char *name;
58 uint32_t hash; // hash of the name
59
51 - STATSD_METRIC_TYPE type;
60 + size_t count; // the number of times this metrics has been collected
61 + calculated_number sampling; // the sampling rate of this metric
62
53 - usec_t last_collected_ut; // the last time this metric was updated
54 - usec_t last_exposed_ut; // the last time this metric was sent to netdata
55 - size_t events; // the number of times this metrics has been collected
63 + char reset; // set to 1 to reset this metric to zero
64
57 - calculated_number count; // number of events since the last exposure to netdata
58 - calculated_number last; // the last value collected
59 - calculated_number value; // the value of the metric
60 - calculated_number min; // the min value collected since the last exposure to netdata
61 - calculated_number max; // the max value collected since the last exposure to netdata
65 + union {
66 + STATSD_METRIC_GAUGE gauge;
67 + STATSD_METRIC_COUNTER counter;
68 + STATSD_METRIC_HISTOGRAM histogram;
69 + STATSD_METRIC_SET set;
70 + };
71 +} STATSD_METRIC;
72
63 - RRDSET *st;
64 - RRDDIM *rd_min;
65 - RRDDIM *rd_max;
66 - RRDDIM *rd_avg;
73 +typedef struct statsd_index {
74 + char *name;
75 + size_t events;
76 + size_t metrics;
77 + STATSD_METRIC *first;
78 + STATSD_AVL_TREE index;
79 + STATSD_FIRST_PTR_MUTEX;
80 +} STATSD_INDEX;
81
68 - struct statsd_metric *next;
69 -} STATSD_METRIC;
82 +static int statsd_metric_compare(void* a, void* b);
83 +
84 +static struct statsd {
85 + STATSD_INDEX gauges;
86 + STATSD_INDEX counters;
87 + STATSD_INDEX timers;
88 + STATSD_INDEX histograms;
89 + STATSD_INDEX meters;
90 + STATSD_INDEX sets;
91 +
92 + size_t histogram_increase_step;
93 + int threads;
94 + LISTEN_SOCKETS sockets;
95 +} statsd = {
96 + .gauges = {
97 + .name = "gauge",
98 + .events = 0,
99 + .metrics = 0,
100 + .first = NULL,
101 + .index = STATSD_AVL_INDEX_INIT,
102 + STATSD_FIRST_PTR_MUTEX_INIT
103 + },
104 + .counters = {
105 + .name = "counter",
106 + .events = 0,
107 + .metrics = 0,
108 + .first = NULL,
109 + .index = STATSD_AVL_INDEX_INIT,
110 + STATSD_FIRST_PTR_MUTEX_INIT
111 + },
112 + .timers = {
113 + .name = "timer",
114 + .events = 0,
115 + .metrics = 0,
116 + .first = NULL,
117 + .index = STATSD_AVL_INDEX_INIT,
118 + STATSD_FIRST_PTR_MUTEX_INIT
119 + },
120 + .histograms = {
121 + .name = "histogram",
122 + .events = 0,
123 + .metrics = 0,
124 + .first = NULL,
125 + .index = STATSD_AVL_INDEX_INIT,
126 + STATSD_FIRST_PTR_MUTEX_INIT
127 + },
128 + .meters = {
129 + .name = "meter",
130 + .events = 0,
131 + .metrics = 0,
132 + .first = NULL,
133 + .index = STATSD_AVL_INDEX_INIT,
134 + STATSD_FIRST_PTR_MUTEX_INIT
135 + },
136 + .sets = {
137 + .name = "set",
138 + .events = 0,
139 + .metrics = 0,
140 + .first = NULL,
141 + .index = STATSD_AVL_INDEX_INIT,
142 + STATSD_FIRST_PTR_MUTEX_INIT
143 + },
144 +
145 + .histogram_increase_step = 10,
146 + .threads = 0,
147 + .sockets = {
148 + .config_section = CONFIG_SECTION_STATSD,
149 + .default_bind_to = "udp:* tcp:*",
150 + .default_port = STATSD_LISTEN_PORT,
151 + .backlog = STATSD_LISTEN_BACKLOG
152 + },
153 +};
154
155 +// --------------------------------------------------------------------------------------
156 +
157 +int statsd_listen_sockets_setup(void) {
158 + return listen_sockets_setup(&statsd.sockets);
159 +}
160
161 // --------------------------------------------------------------------------------------------------------------------
162 // statsd index
163
75 -int statsd_compare(void* a, void* b) {
164 +static int statsd_metric_compare(void* a, void* b) {
165 if(((STATSD_METRIC *)a)->hash < ((STATSD_METRIC *)b)->hash) return -1;
166 else if(((STATSD_METRIC *)a)->hash > ((STATSD_METRIC *)b)->hash) return 1;
167 else return strcmp(((STATSD_METRIC *)a)->name, ((STATSD_METRIC *)b)->name);
168 }
169
81 -static avl_tree statsd_index = {
82 - .compar = statsd_compare,
83 - .root = NULL
84 -};
170 +static int statsd_metric_set_compare(void* a, void* b) {
171 + if(((STATSD_METRIC_SET_VALUE *)a)->hash < ((STATSD_METRIC_SET_VALUE *)b)->hash) return -1;
172 + else if(((STATSD_METRIC_SET_VALUE *)a)->hash > ((STATSD_METRIC_SET_VALUE *)b)->hash) return 1;
173 + else return strcmp(((STATSD_METRIC_SET_VALUE *)a)->value, ((STATSD_METRIC_SET_VALUE *)b)->value);
174 +}
175
86 -static inline STATSD_METRIC *stasd_metric_index_find(const char *name, uint32_t hash) {
176 +static inline STATSD_METRIC *stasd_metric_index_find(STATSD_INDEX *index, const char *name, uint32_t hash) {
177 STATSD_METRIC tmp;
178 tmp.name = name;
179 tmp.hash = (hash)?hash:simple_hash(tmp.name);
180
91 - return (STATSD_METRIC *)avl_search(&statsd_index, (avl *)&tmp);
181 + return (STATSD_METRIC *)STATSD_AVL_SEARCH(&index->index, (avl *)&tmp);
182 +}
183 +
184 +static inline STATSD_METRIC *statsd_find_or_add_metric(STATSD_INDEX *index, char *metric) {
185 + debug(D_STATSD, "finding or adding metric '%s' under '%s'", metric, index->name);
186 +
187 + uint32_t hash = simple_hash(metric);
188 +
189 + STATSD_METRIC *m = stasd_metric_index_find(index, metric, hash);
190 + if(unlikely(!m)) {
191 + debug(D_STATSD, "Creating new %s metric '%s'", index->name, metric);
192 +
193 + m = (STATSD_METRIC *)callocz(sizeof(STATSD_METRIC), 1);
194 + m->name = strdupz(metric);
195 + m->hash = hash;
196 + STATSD_METRIC *n = (STATSD_METRIC *)STATSD_AVL_INSERT(&index->index, (avl *)m);
197 + if(unlikely(n != m)) {
198 + freez((void *)m->name);
199 + freez((void *)m);
200 + m = n;
201 + }
202 + else {
203 + STATSD_FIRST_PTR_MUTEX_LOCK(index);
204 + index->metrics++;
205 + m->next = index->first;
206 + index->first = m;
207 + STATSD_FIRST_PTR_MUTEX_UNLOCK(index);
208 + }
209 + }
210 +
211 + index->events++;
212 + return m;
213 +}
214 +
215 +
216 +// --------------------------------------------------------------------------------------------------------------------
217 +// statsd parsing numbers
218 +
219 +static inline long double statsd_parse_float(const char *v, long double def) {
220 + long double value = def;
221 +
222 + if(likely(v && *v)) {
223 + char *e = NULL;
224 + value = strtold(v, &e);
225 + if(e && *e)
226 + error("STATSD: excess data '%s' after value '%s'", e, v);
227 + }
228 +
229 + return value;
230 +}
231 +
232 +static inline long long statsd_parse_int(const char *v, long long def) {
233 + long long value = def;
234 +
235 + if(likely(v && *v)) {
236 + char *e = NULL;
237 + value = strtoll(v, &e, 10);
238 + if(e && *e)
239 + error("STATSD: excess data '%s' after value '%s'", e, v);
240 + }
241 +
242 + return value;
243 }
244
245
246 // --------------------------------------------------------------------------------------------------------------------
96 -// statsd data collection
247 +// statsd processors per metric type
248
98 -#define STATSD_GAUGE_COLLECTION_RELATIVE 0x00000001
249 +static inline void statsd_reset_metric(STATSD_METRIC *m) {
250 + m->reset = 0;
251 + m->count = 0;
252 + m->sampling = 0.0;
253 +}
254
100 -static inline void statsd_collected_value(STATSD_METRIC *m, calculated_number value, calculated_number sample_rate, uint32_t options) {
101 - debug(D_STATSD, "Updating metric '%s'", m->name);
255 +static inline void statsd_process_gauge(STATSD_METRIC *m, char *v, char *r) {
256 + if(unlikely(!v || !*v)) {
257 + error("STATSD: metric '%s' of type gauge, with empty value is ignored.", m->name);
258 + return;
259 + }
260
103 - statsd.events++;
261 + if(unlikely(m->reset)) statsd_reset_metric(m);
262 + m->count++;
263 + m->sampling += statsd_parse_float(r, 1.0);
264
105 - m->last_collected_ut = now_realtime_usec();
106 - m->events++;
107 - m->count += sample_rate;
265 + if(*v == '+' || *v == '-')
266 + m->gauge.value += statsd_parse_float(v, 1.0);
267 + else
268 + m->gauge.value = statsd_parse_float(v, 1.0);
269 +}
270
109 - m->last = value;
110 - if(value < m->min)
111 - m->min = value;
112 - if(value > m->max)
113 - m->max = value;
271 +static inline void statsd_process_counter(STATSD_METRIC *m, char *v, char *r) {
272 + // we accept empty values for counters
273
115 - switch(m->type) {
116 - case STATSD_METRIC_TYPE_HISTOGRAM:
117 - statsd.events_histogram++;
118 - // FIXME: not implemented yet
119 - m->value += value;
120 - break;
274 + if(unlikely(m->reset)) statsd_reset_metric(m);
275 + m->count++;
276 + m->sampling += statsd_parse_float(r, 1.0);
277
122 - case STATSD_METRIC_TYPE_METER:
123 - statsd.events_meter++;
124 - // we add to this metric
125 - m->value += value;
126 - break;
278 + m->counter.value += statsd_parse_int(v, 1);
279 +}
280
128 - case STATSD_METRIC_TYPE_TIMER:
129 - statsd.events_timer++;
130 - // we add time to this metric
131 - m->value += value;
132 - break;
281 +static inline void statsd_process_meter(STATSD_METRIC *m, char *v, char *r) {
282 + // this is the same with the counter
283 + statsd_process_counter(m, v, r);
284 +}
285
134 - case STATSD_METRIC_TYPE_GAUGE:
135 - statsd.events_gauge++;
136 - if(unlikely(options & STATSD_GAUGE_COLLECTION_RELATIVE))
137 - // we add the collected value
138 - m->value += value;
139 - else
140 - // we replace the value of the metric
141 - m->value = value;
142 - break;
286 +static inline void statsd_process_histogram(STATSD_METRIC *m, char *v, char *r) {
287 + if(unlikely(!v || !*v)) {
288 + error("STATSD: metric '%s' of type histogram, with empty value is ignored.", m->name);
289 + return;
290 + }
291
144 - case STATSD_METRIC_TYPE_COUNTER:
145 - statsd.events_counter++;
146 - // we add the collected value
147 - m->value += value;
148 - break;
292 + if(unlikely(m->reset)) {
293 + m->histogram.used = 0;
294 + statsd_reset_metric(m);
295 }
296
151 - debug(D_STATSD, "Updated metric '%s', type '%c', events %zu, count %0.5Lf, last_collected %llu, last value %0.5Lf, value %0.5Lf, min %0.5Lf, max %0.5Lf"
152 - , m->name
153 - , (char)m->type
154 - , m->events
155 - , m->count
156 - , m->last_collected_ut
157 - , m->last
158 - , m->value
159 - , m->min
160 - , m->max
161 - );
297 + m->count++;
298 + m->sampling += statsd_parse_float(r, 1.0);
299 +
300 + if(m->histogram.used == m->histogram.size) {
301 + m->histogram.size += statsd.histogram_increase_step;
302 + m->histogram.values = reallocz(m->histogram.values, sizeof(long double) * m->histogram.size);
303 + }
304 +
305 + m->histogram.values[m->histogram.used++] = statsd_parse_float(v, 1.0);
306 }
307
164 -static inline void statsd_process_metric(const char *metric, STATSD_METRIC_TYPE type, calculated_number value, calculated_number sample_rate, uint32_t options) {
165 - debug(D_STATSD, "processing metric '%s', type '%c', value %0.5Lf, sample_rate %0.5Lf, options '0x%08x'", metric, (char)type, value, sample_rate, options);
308 +static inline void statsd_process_timer(STATSD_METRIC *m, char *v, char *r) {
309 + if(unlikely(!v || !*v)) {
310 + error("STATSD: metric of type set, with empty value is ignored.");
311 + return;
312 + }
313
167 - uint32_t hash = simple_hash(metric);
314 + // timers are a use case of histogram
315 + statsd_process_histogram(m, v, r);
316 +}
317
169 - STATSD_METRIC *m = stasd_metric_index_find(metric, hash);
170 - if(unlikely(!m)) {
171 - debug(D_STATSD, "Creating new metric '%s'", metric);
318 +static inline void statsd_process_set(STATSD_METRIC *m, char *v, char *r) {
319 + if(unlikely(!v || !*v)) {
320 + error("STATSD: metric of type set, with empty value is ignored.");
321 + return;
322 + }
323
173 - m = (STATSD_METRIC *)callocz(sizeof(STATSD_METRIC), 1);
174 - m->name = strdupz(metric);
175 - m->hash = hash;
176 - m->type = type;
177 - m = (STATSD_METRIC *)avl_insert(&statsd_index, (avl *)m);
178 -
179 - statsd.metrics++;
180 - switch(type) {
181 - case STATSD_METRIC_TYPE_COUNTER:
182 - statsd.metrics_counter++;
183 - break;
184 -
185 - case STATSD_METRIC_TYPE_GAUGE:
186 - statsd.metrics_gauge++;
187 - break;
188 -
189 - case STATSD_METRIC_TYPE_METER:
190 - statsd.metrics_meter++;
191 - break;
192 -
193 - case STATSD_METRIC_TYPE_HISTOGRAM:
194 - statsd.metrics_histogram++;
195 - break;
196 -
197 - case STATSD_METRIC_TYPE_TIMER:
198 - statsd.metrics_timer++;
199 - break;
324 + if(unlikely(m->reset)) {
325 + if(likely(m->set.dict)) {
326 + dictionary_destroy(m->set.dict);
327 + m->set.dict = NULL;
328 }
329 + statsd_reset_metric(m);
330 }
331
203 - statsd_collected_value(m, value, sample_rate, options);
332 + if(unlikely(!m->set.dict)) {
333 + m->set.dict = dictionary_create(STATSD_DICTIONARY_OPTIONS|DICTIONARY_FLAG_VALUE_LINK_DONT_CLONE);
334 + m->set.unique = 0;
335 + }
336 +
337 + void *t = dictionary_get(m->set.dict, v);
338 + if(unlikely(!t)) {
339 + dictionary_set(m->set.dict, v, v, 1);
340 + m->set.unique++;
341 + }
342 }
343
344
345 // --------------------------------------------------------------------------------------------------------------------
346 // statsd parsing
347
210 -static void statsd_process_metric_raw(char *m, char *v, char *t, char *r) {
348 +static void statsd_process_metric(char *m, char *v, char *t, char *r) {
349 debug(D_STATSD, "STATSD: raw metric '%s', value '%s', type '%s', rate '%s'", m, v, t, r);
350
351 if(unlikely(!m || !*m)) return;
352 + if(unlikely(!t || !*t)) t = "m";
353
215 - STATSD_METRIC_TYPE type = STATSD_METRIC_TYPE_METER;
216 - calculated_number value = 1.0;
217 - calculated_number sample_rate = 1.0;
218 - uint32_t options = 0;
219 - char *e;
354 + switch (*t) {
355 + case 'g':
356 + statsd_process_gauge(
357 + statsd_find_or_add_metric(&statsd.gauges, m),
358 + v, r);
359 + break;
360
221 - // collect the value
222 - if(likely(v && *v)) {
223 - e = NULL;
224 - value = strtold(v, &e);
225 - if(e && *e)
226 - error("STATSD: excess data '%s' after value, metric '%s', value '%s'", e, m, v);
227 - }
361 + case 'c':
362 + statsd_process_counter(
363 + statsd_find_or_add_metric(&statsd.counters, m),
364 + v, r);
365 + break;
366
229 - if(likely(r && *r)) {
230 - e = NULL;
231 - sample_rate = strtold(r, &e);
232 - if(e && *e)
233 - error("STATSD: excess data '%s' after sampling rate, metric '%s', value '%s', sampling rate '%s'", e, m, v, r);
234 - }
367 + case 'm':
368 + if (t[1] == 's')
369 + statsd_process_timer(
370 + statsd_find_or_add_metric(&statsd.timers, m),
371 + v, r);
372 + else
373 + statsd_process_meter(
374 + statsd_find_or_add_metric(&statsd.meters, m),
375 + v, r);
376 + break;
377
236 - // we have the metric name, value and type
237 - if(likely(t && *t)) {
238 - switch (*t) {
239 - case 'g':
240 - type = STATSD_METRIC_TYPE_GAUGE;
241 - if (unlikely(*v == '-' || *v == '+'))
242 - options |= STATSD_GAUGE_COLLECTION_RELATIVE;
243 - break;
244 -
245 - default:
246 - case 'c':
247 - type = STATSD_METRIC_TYPE_COUNTER;
248 - break;
249 -
250 - case 'm':
251 - if (t[1] == 's') type = STATSD_METRIC_TYPE_TIMER;
252 - else type = STATSD_METRIC_TYPE_METER;
253 - break;
254 -
255 - case 'h':
256 - type = STATSD_METRIC_TYPE_HISTOGRAM;
257 - break;
258 - }
259 - }
378 + case 'h':
379 + statsd_process_histogram(
380 + statsd_find_or_add_metric(&statsd.histograms, m),
381 + v, r);
382 + break;
383
261 - statsd_process_metric(m, type, value, sample_rate, options);
384 + case 's':
385 + statsd_process_set(
386 + statsd_find_or_add_metric(&statsd.sets, m),
387 + v, r);
388 + break;
389 +
390 + default:
391 + error("STATSD: metric '%s' with value '%s' specifies an unknown type '%s'.", m, v?v:"<unset>", t);
392 + break;
393 + }
394 }
395
396 static inline int is_statsd_space(char c) {
@@ -289,7 +421,7 @@ static inline char *skip_to_next_separator(char *s) {
421 return s;
422 }
423
292 -static char *statsd_get_field(char *s, char **field, int *line_break) {
424 +static inline char *statsd_get_field(char *s, char **field, int *line_break) {
425 *line_break = 0;
426
427 while(is_statsd_separator(*s) || is_statsd_space(*s))
@@ -320,7 +452,10 @@ static char *statsd_get_field(char *s, char **field, int *line_break) {
452 return s;
453 }
454
323 -static void statsd_process(char *buffer, size_t size) {
455 +static inline size_t statsd_process(char *buffer, size_t size, int require_newlines) {
456 + (void)require_newlines;
457 + // FIXME: respect require_newlines to support metrics split in multiple TCP packets
458 +
459 buffer[size] = '\0';
460 debug(D_STATSD, "RECEIVED: '%s'", buffer);
461
@@ -331,25 +466,25 @@ static void statsd_process(char *buffer, size_t size) {
466
467 s = statsd_get_field(s, &m, &line_break);
468 if(unlikely(line_break)) {
334 - statsd_process_metric_raw(m, v, t, r);
469 + statsd_process_metric(m, v, t, r);
470 continue;
471 }
472
473 s = statsd_get_field(s, &v, &line_break);
474 if(unlikely(line_break)) {
340 - statsd_process_metric_raw(m, v, t, r);
475 + statsd_process_metric(m, v, t, r);
476 continue;
477 }
478
479 s = statsd_get_field(s, &t, &line_break);
480 if(likely(line_break)) {
346 - statsd_process_metric_raw(m, v, t, r);
481 + statsd_process_metric(m, v, t, r);
482 continue;
483 }
484
485 s = statsd_get_field(s, &r, &line_break);
486 if(likely(line_break)) {
352 - statsd_process_metric_raw(m, v, t, r);
487 + statsd_process_metric(m, v, t, r);
488 continue;
489 }
490
@@ -361,29 +496,41 @@ static void statsd_process(char *buffer, size_t size) {
496 s++;
497 }
498
364 - statsd_process_metric_raw(m, v, t, r);
499 + statsd_process_metric(m, v, t, r);
500 }
501 +
502 + return 0;
503 }
504
505
506 // --------------------------------------------------------------------------------------------------------------------
507 // statsd pollfd interface
508
372 -static char statsd_read_buffer[65536];
509 +#define STATSD_TCP_BUFFER_SIZE 16384 // minimize reads
510 +#define STATSD_UDP_BUFFER_SIZE 9000 // this should be up to MTU
511 +
512 +struct statsd_tcp {
513 + size_t size;
514 + size_t len;
515 + char buffer[];
516 +};
517
518 // new TCP client connected
519 static void *statsd_add_callback(int fd, short int *events) {
520 (void)fd;
521 *events = POLLIN;
522
379 - return NULL;
380 -}
523 + struct statsd_tcp *data = (struct statsd_tcp *)callocz(sizeof(struct statsd_tcp) + STATSD_TCP_BUFFER_SIZE, 1);
524 + data->size = STATSD_TCP_BUFFER_SIZE - 1;
525
526 + return data;
527 +}
528
529 // TCP client disconnected
530 static void statsd_del_callback(int fd, void *data) {
531 (void)fd;
386 - (void)data;
532 +
533 + freez(data);
534
535 return;
536 }
@@ -394,32 +541,48 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
541
542 switch(socktype) {
543 case SOCK_STREAM: {
544 + struct statsd_tcp *d = (struct statsd_tcp *)data;
545 + if(unlikely(!d)) {
546 + error("STATSD: internal error - tcp receive buffer is null");
547 + return -1;
548 + }
549 +
550 + int ret = 0;
551 ssize_t rc;
552 do {
399 - rc = recv(fd, statsd_read_buffer, sizeof(statsd_read_buffer), MSG_DONTWAIT);
553 + rc = recv(fd, &d->buffer[d->len], d->size - d->len, MSG_DONTWAIT);
554 if (rc < 0) {
555 // read failed
402 - if (errno != EWOULDBLOCK && errno != EAGAIN) {
403 - error("STATSD: recv() failed.");
404 - return -1;
405 - }
406 - } else if (!rc) {
556 + if (errno != EWOULDBLOCK && errno != EAGAIN)
557 + ret = -1;
558 + }
559 + else if (!rc) {
560 // connection closed
561 error("STATSD: client disconnected.");
409 - return -1;
410 - } else {
562 + ret = -1;
563 + }
564 + else {
565 // data received
412 - statsd_process(statsd_read_buffer, (size_t) rc);
566 + d->len += rc;
567 }
568 +
569 + if(likely(d->len > 0))
570 + d->len = statsd_process(d->buffer, d->len, 1);
571 +
572 + if(unlikely(ret == -1))
573 + return -1;
574 +
575 } while (rc != -1);
576 break;
577 }
578
579 case SOCK_DGRAM: {
580 + char buffer[STATSD_UDP_BUFFER_SIZE + 1];
581 +
582 ssize_t rc;
583 do {
584 // FIXME: collect sender information
422 - rc = recvfrom(fd, statsd_read_buffer, sizeof(statsd_read_buffer), MSG_DONTWAIT, NULL, NULL);
585 + rc = recvfrom(fd, buffer, STATSD_UDP_BUFFER_SIZE, MSG_DONTWAIT, NULL, NULL);
586 if (rc < 0) {
587 // read failed
588 if (errno != EWOULDBLOCK && errno != EAGAIN) {
@@ -428,7 +591,7 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
591 }
592 } else if (rc) {
593 // data received
431 - statsd_process(statsd_read_buffer, (size_t) rc);
594 + statsd_process(buffer, (size_t) rc, 0);
595 }
596 } while (rc != -1);
597 break;
@@ -448,8 +611,10 @@ static int statsd_rcv_callback(int fd, int socktype, void *data, short int *even
611 // --------------------------------------------------------------------------------------------------------------------
612 // statsd child thread to update netdata
613
451 -void *statsd_child_thread(void *ptr) {
452 - info("STATSD thread created with task id %d", gettid());
614 +void *statsd_collector_thread(void *ptr) {
615 + int id = *((int *)ptr);
616 +
617 + info("STATSD collector thread No %d created with task id %d", id + 1, gettid());
618
619 if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
620 error("Cannot set pthread cancel type to DEFERRED.");
@@ -457,6 +622,61 @@ void *statsd_child_thread(void *ptr) {
622 if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
623 error("Cannot set pthread cancel state to ENABLE.");
624
625 + poll_events(&statsd.sockets
626 + , statsd_add_callback
627 + , statsd_del_callback
628 + , statsd_rcv_callback
629 + );
630 +
631 + debug(D_WEB_CLIENT, "STATSD: exit!");
632 + listen_sockets_close(&statsd.sockets);
633 +
634 + pthread_exit(NULL);
635 + return NULL;
636 +}
637 +
638 +
639 +// --------------------------------------------------------------------------------------------------------------------
640 +// statsd main thread
641 +
642 +void *statsd_main(void *ptr) {
643 + (void)ptr;
644 +
645 + info("STATSD main thread created with task id %d", gettid());
646 +
647 + if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
648 + error("Cannot set pthread cancel type to DEFERRED.");
649 +
650 + if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
651 + error("Cannot set pthread cancel state to ENABLE.");
652 +
653 + statsd_listen_sockets_setup();
654 + if(!statsd.sockets.opened) {
655 + error("STATSD: No statsd sockets to listen to.");
656 + pthread_exit(NULL);
657 + }
658 +
659 +#ifdef STATSD_MULTITHREADED
660 + statsd.threads = (int)config_get_number(CONFIG_SECTION_STATSD, "collector threads", processors);
661 + if(statsd.threads < 1) {
662 + error("STATSD: Invalid number of threads %d, using %d", statsd.threads, processors);
663 + statsd.threads = processors;
664 + config_set_number(CONFIG_SECTION_STATSD, "collector threads", statsd.threads);
665 + }
666 +#else
667 + statsd.threads = 1;
668 +#endif
669 +
670 + pthread_t threads[statsd.threads];
671 + int i;
672 +
673 + for(i = 0; i < statsd.threads ;i++) {
674 + if(pthread_create(&threads[i], NULL, statsd_collector_thread, &i))
675 + error("STATSD: failed to create child thread.");
676 +
677 + else if(pthread_detach(threads[i]))
678 + error("STATSD: cannot request detach of child thread.");
679 + }
680
681 RRDSET *st_metrics = rrdset_create_localhost(
682 "netdata"
@@ -470,11 +690,12 @@ void *statsd_child_thread(void *ptr) {
690 , localhost->rrd_update_every
691 , RRDSET_TYPE_STACKED
692 );
473 - RRDDIM *rd_metrics_gauge = rrddim_add(st_metrics, "gauge", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
474 - RRDDIM *rd_metrics_counter = rrddim_add(st_metrics, "counter", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
475 - RRDDIM *rd_metrics_timer = rrddim_add(st_metrics, "timer", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
476 - RRDDIM *rd_metrics_meter = rrddim_add(st_metrics, "meter", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
477 - RRDDIM *rd_metrics_histogram = rrddim_add(st_metrics, "histogram", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
693 + RRDDIM *rd_metrics_gauge = rrddim_add(st_metrics, "gauges", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
694 + RRDDIM *rd_metrics_counter = rrddim_add(st_metrics, "counters", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
695 + RRDDIM *rd_metrics_timer = rrddim_add(st_metrics, "timers", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
696 + RRDDIM *rd_metrics_meter = rrddim_add(st_metrics, "meters", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
697 + RRDDIM *rd_metrics_histogram = rrddim_add(st_metrics, "histograms", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
698 + RRDDIM *rd_metrics_set = rrddim_add(st_metrics, "sets", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
699
700 RRDSET *st_events = rrdset_create_localhost(
701 "netdata"
@@ -488,11 +709,12 @@ void *statsd_child_thread(void *ptr) {
709 , localhost->rrd_update_every
710 , RRDSET_TYPE_STACKED
711 );
491 - RRDDIM *rd_events_gauge = rrddim_add(st_events, "gauge", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
492 - RRDDIM *rd_events_counter = rrddim_add(st_events, "counter", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
493 - RRDDIM *rd_events_timer = rrddim_add(st_events, "timer", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
494 - RRDDIM *rd_events_meter = rrddim_add(st_events, "meter", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
495 - RRDDIM *rd_events_histogram = rrddim_add(st_events, "histogram", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
712 + RRDDIM *rd_events_gauge = rrddim_add(st_events, "gauges", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
713 + RRDDIM *rd_events_counter = rrddim_add(st_events, "counters", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
714 + RRDDIM *rd_events_timer = rrddim_add(st_events, "timers", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
715 + RRDDIM *rd_events_meter = rrddim_add(st_events, "meters", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
716 + RRDDIM *rd_events_histogram = rrddim_add(st_events, "histograms", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
717 + RRDDIM *rd_events_set = rrddim_add(st_events, "sets", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
718
719 usec_t step = localhost->rrd_update_every * USEC_PER_SEC;
720 heartbeat_t hb;
@@ -508,64 +730,33 @@ void *statsd_child_thread(void *ptr) {
730 rrdset_next(st_events);
731 }
732
511 - rrddim_set_by_pointer(st_metrics, rd_metrics_gauge, (collected_number)statsd.metrics_gauge);
512 - rrddim_set_by_pointer(st_metrics, rd_metrics_counter, (collected_number)statsd.metrics_counter);
513 - rrddim_set_by_pointer(st_metrics, rd_metrics_timer, (collected_number)statsd.metrics_timer);
514 - rrddim_set_by_pointer(st_metrics, rd_metrics_meter, (collected_number)statsd.metrics_meter);
515 - rrddim_set_by_pointer(st_metrics, rd_metrics_histogram, (collected_number)statsd.metrics_histogram);
733 + rrddim_set_by_pointer(st_metrics, rd_metrics_gauge, (collected_number)statsd.gauges.metrics);
734 + rrddim_set_by_pointer(st_metrics, rd_metrics_counter, (collected_number)statsd.counters.metrics);
735 + rrddim_set_by_pointer(st_metrics, rd_metrics_timer, (collected_number)statsd.timers.metrics);
736 + rrddim_set_by_pointer(st_metrics, rd_metrics_meter, (collected_number)statsd.meters.metrics);
737 + rrddim_set_by_pointer(st_metrics, rd_metrics_histogram, (collected_number)statsd.histograms.metrics);
738 + rrddim_set_by_pointer(st_metrics, rd_metrics_set, (collected_number)statsd.sets.metrics);
739
517 - rrddim_set_by_pointer(st_events, rd_events_gauge, (collected_number)statsd.events_gauge);
518 - rrddim_set_by_pointer(st_events, rd_events_counter, (collected_number)statsd.events_counter);
519 - rrddim_set_by_pointer(st_events, rd_events_timer, (collected_number)statsd.events_timer);
520 - rrddim_set_by_pointer(st_events, rd_events_meter, (collected_number)statsd.events_meter);
521 - rrddim_set_by_pointer(st_events, rd_events_histogram, (collected_number)statsd.events_histogram);
740 + rrddim_set_by_pointer(st_events, rd_events_gauge, (collected_number)statsd.gauges.events);
741 + rrddim_set_by_pointer(st_events, rd_events_counter, (collected_number)statsd.counters.events);
742 + rrddim_set_by_pointer(st_events, rd_events_timer, (collected_number)statsd.timers.events);
743 + rrddim_set_by_pointer(st_events, rd_events_meter, (collected_number)statsd.meters.events);
744 + rrddim_set_by_pointer(st_events, rd_events_histogram, (collected_number)statsd.histograms.events);
745 + rrddim_set_by_pointer(st_events, rd_events_set, (collected_number)statsd.sets.events);
746 +
747 + if(unlikely(netdata_exit))
748 + break;
749
750 rrdset_done(st_metrics);
751 rrdset_done(st_events);
525 - }
526 -
527 - pthread_exit(NULL);
528 - return NULL;
529 -}
530 -
752
532 -// --------------------------------------------------------------------------------------------------------------------
533 -// statsd main thread
534 -
535 -void *statsd_main(void *ptr) {
536 - info("STATSD thread created with task id %d", gettid());
537 -
538 - if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
539 - error("Cannot set pthread cancel type to DEFERRED.");
540 -
541 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
542 - error("Cannot set pthread cancel state to ENABLE.");
753 + if(unlikely(netdata_exit))
754 + break;
755
544 - statsd_listen_sockets_setup();
545 - if(!statsd_sockets.opened) {
546 - error("STATSD: No statsd sockets to listen to.");
547 - goto cleanup;
756 }
757
550 - pthread_t thread;
551 -
552 - if(pthread_create(&thread, NULL, statsd_child_thread, (void *)NULL))
553 - error("STATSD: failed to create child thread.");
554 -
555 - else if(pthread_detach(thread))
556 - error("STATSD: cannot request detach of child thread.");
557 -
558 - poll_events(&statsd_sockets
559 - , statsd_add_callback
560 - , statsd_del_callback
561 - , statsd_rcv_callback
562 - );
563 -
564 -cleanup:
565 - pthread_cancel(thread);
566 -
567 - debug(D_WEB_CLIENT, "STATSD: exit!");
568 - listen_sockets_close(&statsd_sockets);
758 + for(i = 0; i < statsd.threads ;i++)
759 + pthread_cancel(threads[i]);
760
761 pthread_exit(NULL);
762 return NULL;