@cryptotaxi247 / netdata-1 / commits / 20190a026

added /stream/ URL that redirects socket to plugins.d for processing - api key validation missing

Costa Tsaousis (ktsaou) committed Feb 19, 2017 at 21:07 UTC 20190a026beeec38dee18ba38ef008a4b9f275da
3 files changed +340 -269
src/plugins_d.c
+235 -262
@@ -83,17 +83,18 @@ static int pluginsd_split_words(char *str, char **words, int max_words) {
83 return i;
84 }
85
86 +inline size_t pluginsd_process(struct plugind *cd, FILE *fp, int trust_durations) {
87 + int enabled = cd->enabled;
88
87 -void *pluginsd_worker_thread(void *arg) {
88 - struct plugind *cd = (struct plugind *)arg;
89 - cd->obsolete = 0;
89 + if(!fp || !enabled) {
90 + cd->enabled = 0;
91 + return 0;
92 + }
93
91 - char line[PLUGINSD_LINE_MAX + 1];
94 + size_t count = 0;
95 + RRDHOST *host = localhost;
96
93 -#ifdef DETACH_PLUGINS_FROM_NETDATA
94 - usec_t usec = 0, susec = 0;
95 - struct timeval last = {0, 0} , now = {0, 0};
96 -#endif
97 + char line[PLUGINSD_LINE_MAX + 1];
98
99 char *words[MAX_WORDS] = { NULL };
100 uint32_t HOST_HASH = simple_hash("HOST");
@@ -103,305 +104,277 @@ void *pluginsd_worker_thread(void *arg) {
104 uint32_t CHART_HASH = simple_hash("CHART");
105 uint32_t DIMENSION_HASH = simple_hash("DIMENSION");
106 uint32_t DISABLE_HASH = simple_hash("DISABLE");
106 -#ifdef DETACH_PLUGINS_FROM_NETDATA
107 - uint32_t MYPID_HASH = simple_hash("MYPID");
108 - uint32_t STOPPING_WAKE_ME_UP_PLEASE_HASH = simple_hash("STOPPING_WAKE_ME_UP_PLEASE");
109 -#endif
107
111 - size_t count = 0;
112 - RRDHOST *host = localhost;
108 + RRDSET *st = NULL;
109 + uint32_t hash;
110
114 - for(;;) {
111 + while(likely(fgets(line, PLUGINSD_LINE_MAX, fp) != NULL)) {
112 if(unlikely(netdata_exit)) break;
113
117 - FILE *fp = mypopen(cd->cmd, &cd->pid);
118 - if(unlikely(!fp)) {
119 - error("Cannot popen(\"%s\", \"r\").", cd->cmd);
120 - break;
121 - }
122 -
123 - info("PLUGINSD: '%s' running on pid %d", cd->fullfilename, cd->pid);
114 + line[PLUGINSD_LINE_MAX] = '\0';
115
125 - RRDSET *st = NULL;
126 - uint32_t hash;
116 + // debug(D_PLUGINSD, "PLUGINSD: %s: %s", cd->filename, line);
117
128 - while(likely(fgets(line, PLUGINSD_LINE_MAX, fp) != NULL)) {
129 - if(unlikely(netdata_exit)) break;
118 + int w = pluginsd_split_words(line, words, MAX_WORDS);
119 + char *s = words[0];
120 + if(unlikely(!s || !*s || !w)) {
121 + // debug(D_PLUGINSD, "PLUGINSD: empty line");
122 + continue;
123 + }
124
131 - line[PLUGINSD_LINE_MAX] = '\0';
125 + // debug(D_PLUGINSD, "PLUGINSD: words 0='%s' 1='%s' 2='%s' 3='%s' 4='%s' 5='%s' 6='%s' 7='%s' 8='%s' 9='%s'", words[0], words[1], words[2], words[3], words[4], words[5], words[6], words[7], words[8], words[9]);
126
133 - // debug(D_PLUGINSD, "PLUGINSD: %s: %s", cd->filename, line);
127 + if(likely(!simple_hash_strcmp(s, "SET", &hash))) {
128 + char *dimension = words[1];
129 + char *value = words[2];
130
135 - int w = pluginsd_split_words(line, words, MAX_WORDS);
136 - char *s = words[0];
137 - if(unlikely(!s || !*s || !w)) {
138 - // debug(D_PLUGINSD, "PLUGINSD: empty line");
139 - continue;
131 + if(unlikely(!dimension || !*dimension)) {
132 + error("PLUGINSD: '%s' is requesting a SET on chart '%s', without a dimension. Disabling it.", cd->fullfilename, st->id);
133 + enabled = 0;
134 + break;
135 }
136
142 - // debug(D_PLUGINSD, "PLUGINSD: words 0='%s' 1='%s' 2='%s' 3='%s' 4='%s' 5='%s' 6='%s' 7='%s' 8='%s' 9='%s'", words[0], words[1], words[2], words[3], words[4], words[5], words[6], words[7], words[8], words[9]);
143 -
144 - if(likely(!simple_hash_strcmp(s, "SET", &hash))) {
145 - char *dimension = words[1];
146 - char *value = words[2];
147 -
148 - if(unlikely(!dimension || !*dimension)) {
149 - error("PLUGINSD: '%s' is requesting a SET on chart '%s', without a dimension. Disabling it.", cd->fullfilename, st->id);
150 - cd->enabled = 0;
151 - killpid(cd->pid, SIGTERM);
152 - break;
153 - }
137 + if(unlikely(!value || !*value)) value = NULL;
138
155 - if(unlikely(!value || !*value)) value = NULL;
139 + if(unlikely(!st)) {
140 + error("PLUGINSD: '%s' is requesting a SET on dimension %s with value %s, without a BEGIN. Disabling it.", cd->fullfilename, dimension, value?value:"<nothing>");
141 + enabled = 0;
142 + break;
143 + }
144
157 - if(unlikely(!st)) {
158 - error("PLUGINSD: '%s' is requesting a SET on dimension %s with value %s, without a BEGIN. Disabling it.", cd->fullfilename, dimension, value?value:"<nothing>");
159 - cd->enabled = 0;
160 - killpid(cd->pid, SIGTERM);
161 - break;
162 - }
145 + if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG))) debug(D_PLUGINSD, "PLUGINSD: '%s' is setting dimension %s/%s to %s", cd->fullfilename, st->id, dimension, value?value:"<nothing>");
146
164 - if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG))) debug(D_PLUGINSD, "PLUGINSD: '%s' is setting dimension %s/%s to %s", cd->fullfilename, st->id, dimension, value?value:"<nothing>");
147 + if(value) rrddim_set(st, dimension, strtoll(value, NULL, 0));
148 + }
149 + else if(likely(hash == BEGIN_HASH && !strcmp(s, "BEGIN"))) {
150 + char *id = words[1];
151 + char *microseconds_txt = words[2];
152
166 - if(value) rrddim_set(st, dimension, strtoll(value, NULL, 0));
153 + if(unlikely(!id)) {
154 + error("PLUGINSD: '%s' is requesting a BEGIN without a chart id. Disabling it.", cd->fullfilename);
155 + enabled = 0;
156 + break;
157 }
168 - else if(likely(hash == BEGIN_HASH && !strcmp(s, "BEGIN"))) {
169 - char *id = words[1];
170 - char *microseconds_txt = words[2];
158
172 - if(unlikely(!id)) {
173 - error("PLUGINSD: '%s' is requesting a BEGIN without a chart id. Disabling it.", cd->fullfilename);
174 - cd->enabled = 0;
175 - killpid(cd->pid, SIGTERM);
176 - break;
177 - }
159 + st = rrdset_find(host, id);
160 + if(unlikely(!st)) {
161 + error("PLUGINSD: '%s' is requesting a BEGIN on chart '%s', which does not exist. Disabling it.", cd->fullfilename, id);
162 + enabled = 0;
163 + break;
164 + }
165
179 - st = rrdset_find(host, id);
180 - if(unlikely(!st)) {
181 - error("PLUGINSD: '%s' is requesting a BEGIN on chart '%s', which does not exist. Disabling it.", cd->fullfilename, id);
182 - cd->enabled = 0;
183 - killpid(cd->pid, SIGTERM);
184 - break;
185 - }
166 + if(likely(st->counter_done)) {
167 + usec_t microseconds = 0;
168 + if(microseconds_txt && *microseconds_txt) microseconds = str2ull(microseconds_txt);
169
187 - if(likely(st->counter_done)) {
188 - usec_t microseconds = 0;
189 - if(microseconds_txt && *microseconds_txt) microseconds = str2ull(microseconds_txt);
190 - if(microseconds) rrdset_next_usec(st, microseconds);
191 - else rrdset_next(st);
170 + if(likely(microseconds)) {
171 + if(trust_durations)
172 + rrdset_next_usec_unfiltered(st, microseconds);
173 + else
174 + rrdset_next_usec(st, microseconds);
175 }
176 + else rrdset_next(st);
177 + }
178 + }
179 + else if(likely(hash == END_HASH && !strcmp(s, "END"))) {
180 + if(unlikely(!st)) {
181 + error("PLUGINSD: '%s' is requesting an END, without a BEGIN. Disabling it.", cd->fullfilename);
182 + enabled = 0;
183 + break;
184 }
194 - else if(likely(hash == END_HASH && !strcmp(s, "END"))) {
195 - if(unlikely(!st)) {
196 - error("PLUGINSD: '%s' is requesting an END, without a BEGIN. Disabling it.", cd->fullfilename);
197 - cd->enabled = 0;
198 - killpid(cd->pid, SIGTERM);
199 - break;
200 - }
185
202 - if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG))) debug(D_PLUGINSD, "PLUGINSD: '%s' is requesting an END on chart %s", cd->fullfilename, st->id);
186 + if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG))) debug(D_PLUGINSD, "PLUGINSD: '%s' is requesting an END on chart %s", cd->fullfilename, st->id);
187
204 - rrdset_done(st);
205 - st = NULL;
188 + rrdset_done(st);
189 + st = NULL;
190
207 - count++;
208 - }
209 - else if(likely(hash == HOST_HASH && !strcmp(s, "HOST"))) {
210 - char *guid = words[1];
211 - char *hostname = words[2];
212 -
213 - if(unlikely(!guid || !*guid)) {
214 - error("PLUGINSD: '%s' is requesting a HOST, without a guid. Disabling it.", cd->fullfilename);
215 - cd->enabled = 0;
216 - killpid(cd->pid, SIGTERM);
217 - break;
218 - }
219 - if(unlikely(!hostname || !*hostname)) {
220 - error("PLUGINSD: '%s' is requesting a HOST, without a hostname. Disabling it.", cd->fullfilename);
221 - cd->enabled = 0;
222 - killpid(cd->pid, SIGTERM);
223 - break;
224 - }
191 + count++;
192 + }
193 + else if(likely(hash == HOST_HASH && !strcmp(s, "HOST"))) {
194 + char *guid = words[1];
195 + char *hostname = words[2];
196
226 - host = rrdhost_find_or_create(hostname, guid);
197 + if(unlikely(!guid || !*guid)) {
198 + error("PLUGINSD: '%s' is requesting a HOST, without a guid. Disabling it.", cd->fullfilename);
199 + enabled = 0;
200 + break;
201 }
228 - else if(likely(hash == FLUSH_HASH && !strcmp(s, "FLUSH"))) {
229 - debug(D_PLUGINSD, "PLUGINSD: '%s' is requesting a FLUSH", cd->fullfilename);
230 - st = NULL;
202 + if(unlikely(!hostname || !*hostname)) {
203 + error("PLUGINSD: '%s' is requesting a HOST, without a hostname. Disabling it.", cd->fullfilename);
204 + enabled = 0;
205 + break;
206 }
232 - else if(likely(hash == CHART_HASH && !strcmp(s, "CHART"))) {
233 - int noname = 0;
234 - st = NULL;
235 -
236 - if((words[1]) != NULL && (words[2]) != NULL && strcmp(words[1], words[2]) == 0)
237 - noname = 1;
238 -
239 - char *type = words[1];
240 - char *id = NULL;
241 - if(likely(type)) {
242 - id = strchr(type, '.');
243 - if(likely(id)) { *id = '\0'; id++; }
244 - }
245 - char *name = words[2];
246 - char *title = words[3];
247 - char *units = words[4];
248 - char *family = words[5];
249 - char *context = words[6];
250 - char *chart = words[7];
251 - char *priority_s = words[8];
252 - char *update_every_s = words[9];
253 -
254 - if(unlikely(!type || !*type || !id || !*id)) {
255 - error("PLUGINSD: '%s' is requesting a CHART, without a type.id. Disabling it.", cd->fullfilename);
256 - cd->enabled = 0;
257 - killpid(cd->pid, SIGTERM);
258 - break;
259 - }
207
261 - int priority = 1000;
262 - if(likely(priority_s)) priority = str2i(priority_s);
263 -
264 - int update_every = cd->update_every;
265 - if(likely(update_every_s)) update_every = str2i(update_every_s);
266 - if(unlikely(!update_every)) update_every = cd->update_every;
267 -
268 - RRDSET_TYPE chart_type = RRDSET_TYPE_LINE;
269 - if(unlikely(chart)) chart_type = rrdset_type_id(chart);
270 -
271 - if(unlikely(noname || !name || !*name || strcasecmp(name, "NULL") == 0 || strcasecmp(name, "(NULL)") == 0)) name = NULL;
272 - if(unlikely(!family || !*family)) family = NULL;
273 - if(unlikely(!context || !*context)) context = NULL;
274 -
275 - st = rrdset_find_bytype(host, type, id);
276 - if(unlikely(!st)) {
277 - debug(D_PLUGINSD, "PLUGINSD: Creating chart type='%s', id='%s', name='%s', family='%s', context='%s', chart='%s', priority=%d, update_every=%d"
278 - , type, id
279 - , name?name:""
280 - , family?family:""
281 - , context?context:""
282 - , rrdset_type_name(chart_type)
283 - , priority
284 - , update_every
285 - );
286 -
287 - st = rrdset_create(host, type, id, name, family, context, title, units, priority, update_every
288 - , chart_type);
289 - cd->update_every = update_every;
290 - }
291 - else debug(D_PLUGINSD, "PLUGINSD: Chart '%s' already exists. Not adding it again.", st->id);
208 + host = rrdhost_find_or_create(hostname, guid);
209 + }
210 + else if(likely(hash == FLUSH_HASH && !strcmp(s, "FLUSH"))) {
211 + debug(D_PLUGINSD, "PLUGINSD: '%s' is requesting a FLUSH", cd->fullfilename);
212 + st = NULL;
213 + }
214 + else if(likely(hash == CHART_HASH && !strcmp(s, "CHART"))) {
215 + int noname = 0;
216 + st = NULL;
217 +
218 + if((words[1]) != NULL && (words[2]) != NULL && strcmp(words[1], words[2]) == 0)
219 + noname = 1;
220 +
221 + char *type = words[1];
222 + char *id = NULL;
223 + if(likely(type)) {
224 + id = strchr(type, '.');
225 + if(likely(id)) { *id = '\0'; id++; }
226 + }
227 + char *name = words[2];
228 + char *title = words[3];
229 + char *units = words[4];
230 + char *family = words[5];
231 + char *context = words[6];
232 + char *chart = words[7];
233 + char *priority_s = words[8];
234 + char *update_every_s = words[9];
235 +
236 + if(unlikely(!type || !*type || !id || !*id)) {
237 + error("PLUGINSD: '%s' is requesting a CHART, without a type.id. Disabling it.", cd->fullfilename);
238 + enabled = 0;
239 + break;
240 }
293 - else if(likely(hash == DIMENSION_HASH && !strcmp(s, "DIMENSION"))) {
294 - char *id = words[1];
295 - char *name = words[2];
296 - char *algorithm = words[3];
297 - char *multiplier_s = words[4];
298 - char *divisor_s = words[5];
299 - char *options = words[6];
300 -
301 - if(unlikely(!id || !*id)) {
302 - error("PLUGINSD: '%s' is requesting a DIMENSION, without an id. Disabling it.", cd->fullfilename);
303 - cd->enabled = 0;
304 - killpid(cd->pid, SIGTERM);
305 - break;
306 - }
307 -
308 - if(unlikely(!st)) {
309 - error("PLUGINSD: '%s' is requesting a DIMENSION, without a CHART. Disabling it.", cd->fullfilename);
310 - cd->enabled = 0;
311 - killpid(cd->pid, SIGTERM);
312 - break;
313 - }
241
315 - long multiplier = 1;
316 - if(multiplier_s && *multiplier_s) multiplier = strtol(multiplier_s, NULL, 0);
317 - if(unlikely(!multiplier)) multiplier = 1;
318 -
319 - long divisor = 1;
320 - if(likely(divisor_s && *divisor_s)) divisor = strtol(divisor_s, NULL, 0);
321 - if(unlikely(!divisor)) divisor = 1;
322 -
323 - if(unlikely(!algorithm || !*algorithm)) algorithm = "absolute";
324 -
325 - if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG)))
326 - debug(D_PLUGINSD, "PLUGINSD: Creating dimension in chart %s, id='%s', name='%s', algorithm='%s', multiplier=%ld, divisor=%ld, hidden='%s'"
327 - , st->id
328 - , id
329 - , name?name:""
330 - , rrd_algorithm_name(rrd_algorithm_id(algorithm))
331 - , multiplier
332 - , divisor
333 - , options?options:""
334 - );
335 -
336 - RRDDIM *rd = rrddim_find(st, id);
337 - if(unlikely(!rd)) {
338 - rd = rrddim_add(st, id, name, multiplier, divisor, rrd_algorithm_id(algorithm));
339 - rd->flags = 0x00000000;
340 - if(options && *options) {
341 - if(strstr(options, "hidden") != NULL) rrddim_flag_set(rd, RRDDIM_FLAG_HIDDEN);
342 - if(strstr(options, "noreset") != NULL) rrddim_flag_set(rd, RRDDIM_FLAG_DONT_DETECT_RESETS_OR_OVERFLOWS);
343 - if(strstr(options, "nooverflow") != NULL) rrddim_flag_set(rd, RRDDIM_FLAG_DONT_DETECT_RESETS_OR_OVERFLOWS);
344 - }
345 - }
346 - else if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG)))
347 - debug(D_PLUGINSD, "PLUGINSD: dimension %s/%s already exists. Not adding it again.", st->id, id);
242 + int priority = 1000;
243 + if(likely(priority_s)) priority = str2i(priority_s);
244 +
245 + int update_every = cd->update_every;
246 + if(likely(update_every_s)) update_every = str2i(update_every_s);
247 + if(unlikely(!update_every)) update_every = cd->update_every;
248 +
249 + RRDSET_TYPE chart_type = RRDSET_TYPE_LINE;
250 + if(unlikely(chart)) chart_type = rrdset_type_id(chart);
251 +
252 + if(unlikely(noname || !name || !*name || strcasecmp(name, "NULL") == 0 || strcasecmp(name, "(NULL)") == 0)) name = NULL;
253 + if(unlikely(!family || !*family)) family = NULL;
254 + if(unlikely(!context || !*context)) context = NULL;
255 +
256 + st = rrdset_find_bytype(host, type, id);
257 + if(unlikely(!st)) {
258 + debug(D_PLUGINSD, "PLUGINSD: Creating chart type='%s', id='%s', name='%s', family='%s', context='%s', chart='%s', priority=%d, update_every=%d"
259 + , type, id
260 + , name?name:""
261 + , family?family:""
262 + , context?context:""
263 + , rrdset_type_name(chart_type)
264 + , priority
265 + , update_every
266 + );
267 +
268 + st = rrdset_create(host, type, id, name, family, context, title, units, priority, update_every, chart_type);
269 + cd->update_every = update_every;
270 }
349 - else if(unlikely(hash == DISABLE_HASH && !strcmp(s, "DISABLE"))) {
350 - info("PLUGINSD: '%s' called DISABLE. Disabling it.", cd->fullfilename);
351 - cd->enabled = 0;
352 - killpid(cd->pid, SIGTERM);
271 + else debug(D_PLUGINSD, "PLUGINSD: Chart '%s' already exists. Not adding it again.", st->id);
272 + }
273 + else if(likely(hash == DIMENSION_HASH && !strcmp(s, "DIMENSION"))) {
274 + char *id = words[1];
275 + char *name = words[2];
276 + char *algorithm = words[3];
277 + char *multiplier_s = words[4];
278 + char *divisor_s = words[5];
279 + char *options = words[6];
280 +
281 + if(unlikely(!id || !*id)) {
282 + error("PLUGINSD: '%s' is requesting a DIMENSION, without an id. Disabling it.", cd->fullfilename);
283 + enabled = 0;
284 break;
285 }
355 -#ifdef DETACH_PLUGINS_FROM_NETDATA
356 - else if(likely(hash == MYPID_HASH && !strcmp(s, "MYPID"))) {
357 - char *pid_s = words[1];
358 - pid_t pid = strtod(pid_s, NULL, 0);
286
360 - if(likely(pid)) cd->pid = pid;
361 - debug(D_PLUGINSD, "PLUGINSD: %s is on pid %d", cd->id, cd->pid);
287 + if(unlikely(!st)) {
288 + error("PLUGINSD: '%s' is requesting a DIMENSION, without a CHART. Disabling it.", cd->fullfilename);
289 + enabled = 0;
290 + break;
291 }
363 - else if(likely(hash == STOPPING_WAKE_ME_UP_PLEASE_HASH && !strcmp(s, "STOPPING_WAKE_ME_UP_PLEASE"))) {
364 - error("PLUGINSD: '%s' (pid %d) called STOPPING_WAKE_ME_UP_PLEASE.", cd->fullfilename, cd->pid);
292
366 - now_realtime_timeval(&now);
367 - if(unlikely(!usec && !susec)) {
368 - // our first run
369 - susec = cd->rrd_update_every * USEC_PER_SEC;
370 - }
371 - else {
372 - // second+ run
373 - usec = dt_usec(&now, &last) - susec;
374 - error("PLUGINSD: %s last loop took %llu usec (worked for %llu, sleeped for %llu).\n", cd->fullfilename, usec + susec, usec, susec);
375 - if(unlikely(usec < (localhost->rrd_update_every * USEC_PER_SEC / 2ULL))) susec = (localhost->rrd_update_every * USEC_PER_SEC) - usec;
376 - else susec = localhost->rrd_update_every * USEC_PER_SEC / 2ULL;
293 + long multiplier = 1;
294 + if(multiplier_s && *multiplier_s) multiplier = strtol(multiplier_s, NULL, 0);
295 + if(unlikely(!multiplier)) multiplier = 1;
296 +
297 + long divisor = 1;
298 + if(likely(divisor_s && *divisor_s)) divisor = strtol(divisor_s, NULL, 0);
299 + if(unlikely(!divisor)) divisor = 1;
300 +
301 + if(unlikely(!algorithm || !*algorithm)) algorithm = "absolute";
302 +
303 + if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG)))
304 + debug(D_PLUGINSD, "PLUGINSD: Creating dimension in chart %s, id='%s', name='%s', algorithm='%s', multiplier=%ld, divisor=%ld, hidden='%s'"
305 + , st->id
306 + , id
307 + , name?name:""
308 + , rrd_algorithm_name(rrd_algorithm_id(algorithm))
309 + , multiplier
310 + , divisor
311 + , options?options:""
312 + );
313 +
314 + RRDDIM *rd = rrddim_find(st, id);
315 + if(unlikely(!rd)) {
316 + rd = rrddim_add(st, id, name, multiplier, divisor, rrd_algorithm_id(algorithm));
317 + rd->flags = 0x00000000;
318 + if(options && *options) {
319 + if(strstr(options, "hidden") != NULL) rrddim_flag_set(rd, RRDDIM_FLAG_HIDDEN);
320 + if(strstr(options, "noreset") != NULL) rrddim_flag_set(rd, RRDDIM_FLAG_DONT_DETECT_RESETS_OR_OVERFLOWS);
321 + if(strstr(options, "nooverflow") != NULL) rrddim_flag_set(rd, RRDDIM_FLAG_DONT_DETECT_RESETS_OR_OVERFLOWS);
322 }
378 -
379 - error("PLUGINSD: %s sleeping for %llu. Will kill with SIGCONT pid %d to wake it up.\n", cd->fullfilename, susec, cd->pid);
380 - usleep(susec);
381 - killpid(cd->pid, SIGCONT);
382 - memmove(&last, &now, sizeof(struct timeval));
383 - break;
384 - }
385 -#endif
386 - else {
387 - error("PLUGINSD: '%s' is sending command '%s' which is not known by netdata. Disabling it.", cd->fullfilename, s);
388 - cd->enabled = 0;
389 - killpid(cd->pid, SIGTERM);
390 - break;
323 }
324 + else if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_DEBUG)))
325 + debug(D_PLUGINSD, "PLUGINSD: dimension %s/%s already exists. Not adding it again.", st->id, id);
326 + }
327 + else if(unlikely(hash == DISABLE_HASH && !strcmp(s, "DISABLE"))) {
328 + info("PLUGINSD: '%s' called DISABLE. Disabling it.", cd->fullfilename);
329 + enabled = 0;
330 + break;
331 + }
332 + else {
333 + error("PLUGINSD: '%s' is sending command '%s' which is not known by netdata. Disabling it.", cd->fullfilename, s);
334 + enabled = 0;
335 + break;
336 }
393 - if(likely(count)) {
394 - cd->successful_collections += count;
395 - cd->serial_failures = 0;
337 + }
338 +
339 + cd->enabled = enabled;
340 +
341 + if(likely(count)) {
342 + cd->successful_collections += count;
343 + cd->serial_failures = 0;
344 + }
345 + else
346 + cd->serial_failures++;
347 +
348 + return count;
349 +}
350 +
351 +void *pluginsd_worker_thread(void *arg) {
352 + struct plugind *cd = (struct plugind *)arg;
353 + cd->obsolete = 0;
354 +
355 + size_t count = 0;
356 +
357 + for(;;) {
358 + if(unlikely(netdata_exit)) break;
359 +
360 + FILE *fp = mypopen(cd->cmd, &cd->pid);
361 + if(unlikely(!fp)) {
362 + error("Cannot popen(\"%s\", \"r\").", cd->cmd);
363 + break;
364 }
397 - else
398 - cd->serial_failures++;
365 +
366 + info("PLUGINSD: '%s' running on pid %d", cd->fullfilename, cd->pid);
367 +
368 + count = pluginsd_process(cd, fp, 0);
369 + error("PLUGINSD: plugin '%s' disconnected.", cd->fullfilename);
370 +
371 + killpid(cd->pid, SIGTERM);
372
373 info("PLUGINSD: '%s' on pid %d stopped after %zu successful data collections (ENDs).", cd->fullfilename, cd->pid, count);
374
375 // get the return code
376 int code = mypclose(fp, cd->pid);
404 -
377 +
378 if(unlikely(netdata_exit)) break;
379 else if(code != 0) {
380 // the plugin reports failure
src/plugins_d.h
+1
@@ -34,5 +34,6 @@ struct plugind {
34 extern struct plugind *pluginsd_root;
35
36 extern void *pluginsd_main(void *ptr);
37 +extern size_t pluginsd_process(struct plugind *cd, FILE *fp, int trust_durations);
38
39 #endif /* NETDATA_PLUGINS_D_H */
src/web_client.c
+104 -7
@@ -1667,6 +1667,92 @@ int web_client_api_old_data_request(RRDHOST *host, struct web_client *w, char *u
1667 return 200;
1668 }
1669
1670 +int validate_stream_api_key(const char *key) {
1671 + return 0;
1672 +}
1673 +
1674 +int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1675 + (void)host;
1676 +
1677 + info("STREAM request from client '%s:%s'", w->client_ip, w->client_port);
1678 +
1679 + char *key = NULL;
1680 +
1681 + while(url) {
1682 + char *value = mystrsep(&url, "?&");
1683 + if(!value || !*value) continue;
1684 +
1685 + char *name = mystrsep(&value, "=");
1686 + if(!name || !*name) continue;
1687 + if(!value || !*value) continue;
1688 +
1689 + if(!strcmp(name, "key"))
1690 + key = value;
1691 + }
1692 +
1693 + if(!key || !*key) {
1694 + buffer_flush(w->response.data);
1695 + buffer_sprintf(w->response.data, "You need an API key for this request.");
1696 + error("STREAM request from client '%s:%s', without an API key. Forbidding access.", w->client_ip, w->client_port);
1697 + return 401;
1698 + }
1699 +
1700 + if(!validate_stream_api_key(key)) {
1701 + buffer_flush(w->response.data);
1702 + buffer_sprintf(w->response.data, "Your API key is not permitted access.");
1703 + error("STREAM request from client '%s:%s': API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
1704 + return 401;
1705 + }
1706 +
1707 + struct plugind cd = {
1708 + .enabled = 1,
1709 + .update_every = default_localhost_rrd_update_every,
1710 + .pid = 0,
1711 + .serial_failures = 0,
1712 + .successful_collections = 0,
1713 + .obsolete = 0,
1714 + .started_t = now_realtime_sec(),
1715 + .next = NULL,
1716 + };
1717 +
1718 + // put the client IP and port into the buffers used by plugins.d
1719 + snprintfz(cd.id, CONFIG_MAX_NAME, "%s:%s", w->client_ip, w->client_port);
1720 + snprintfz(cd.filename, FILENAME_MAX, "%s:%s", w->client_ip, w->client_port);
1721 + snprintfz(cd.fullfilename, FILENAME_MAX, "%s:%s", w->client_ip, w->client_port);
1722 + snprintfz(cd.cmd, PLUGINSD_CMD_MAX, "%s:%s", w->client_ip, w->client_port);
1723 +
1724 + // remove the non-blocking flag from the socket
1725 + if(fcntl(w->ifd, F_SETFL, fcntl(w->ifd, F_GETFL, 0) & ~O_NONBLOCK) == -1)
1726 + error("STREAM from '%s:%s': cannot remove the non-blocking flag from socket %d", w->client_ip, w->client_port, w->ifd);
1727 +
1728 + // convert the socket to a FILE *
1729 + FILE *fp = fdopen(w->ifd, "r");
1730 + if(!fp) {
1731 + error("STREAM from '%s:%s': failed to get a FILE for FD %d.", w->client_ip, w->client_port, w->ifd);
1732 + buffer_flush(w->response.data);
1733 + buffer_sprintf(w->response.data, "Failed to get a FILE for an FD.");
1734 + return 500;
1735 + }
1736 +
1737 + // call the plugins.d processor to receive the metrics
1738 + size_t count = pluginsd_process(&cd, fp, 1);
1739 + error("STREAM from '%s:%s': client disconnected.", w->client_ip, w->client_port);
1740 +
1741 + // close all sockets, to let the socket worker we are done
1742 + fclose(fp);
1743 + w->ifd = -1;
1744 + if(w->ofd != -1 && w->ofd != w->ifd) {
1745 + close(w->ofd);
1746 + w->ofd = -1;
1747 + }
1748 +
1749 + // this will not send anything
1750 + // the socket is closed
1751 + buffer_flush(w->response.data);
1752 + if(count) return 200;
1753 + return 400;
1754 +}
1755 +
1756 const char *web_content_type_to_string(uint8_t contenttype) {
1757 switch(contenttype) {
1758 case CT_TEXT_HTML:
@@ -2113,7 +2199,7 @@ static inline int web_client_switch_host(RRDHOST *host, struct web_client *w, ch
2199 }
2200
2201 buffer_flush(w->response.data);
2116 - buffer_strcat(w->response.data, "Host is not found: ");
2202 + buffer_strcat(w->response.data, "This netdata does not maintain a database for host: ");
2203 buffer_strcat_htmlescape(w->response.data, tok?tok:"");
2204 return 404;
2205 }
@@ -2127,7 +2213,8 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2213 hash_graph = 0,
2214 hash_list = 0,
2215 hash_all_json = 0,
2130 - hash_host = 0;
2216 + hash_host = 0,
2217 + hash_stream = 0;
2218
2219 #ifdef NETDATA_INTERNAL_CHECKS
2220 static uint32_t hash_exit = 0, hash_debug = 0, hash_mirror = 0;
@@ -2142,6 +2229,7 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2229 hash_list = simple_hash("list");
2230 hash_all_json = simple_hash("all.json");
2231 hash_host = simple_hash("host");
2232 + hash_stream = simple_hash("stream");
2233 #ifdef NETDATA_INTERNAL_CHECKS
2234 hash_exit = simple_hash("exit");
2235 hash_debug = simple_hash("debug");
@@ -2154,14 +2242,18 @@ static inline int web_client_process_url(RRDHOST *host, struct web_client *w, ch
2242 uint32_t hash = simple_hash(tok);
2243 debug(D_WEB_CLIENT, "%llu: Processing command '%s'.", w->id, tok);
2244
2157 - if(unlikely(hash == hash_host && strcmp(tok, "host") == 0)) {
2158 - debug(D_WEB_CLIENT_ACCESS, "%llu: host switch request ...", w->id);
2159 - return web_client_switch_host(host, w, url);
2160 - }
2245 if(unlikely(hash == hash_api && strcmp(tok, "api") == 0)) {
2246 debug(D_WEB_CLIENT_ACCESS, "%llu: API request ...", w->id);
2247 return web_client_api_request(host, w, url);
2248 }
2249 + else if(unlikely(hash == hash_host && strcmp(tok, "host") == 0)) {
2250 + debug(D_WEB_CLIENT_ACCESS, "%llu: host switch request ...", w->id);
2251 + return web_client_switch_host(host, w, url);
2252 + }
2253 + else if(unlikely(hash == hash_stream && strcmp(tok, "stream") == 0)) {
2254 + debug(D_WEB_CLIENT_ACCESS, "%llu: stream request ...", w->id);
2255 + return web_client_stream_request(host, w, url);
2256 + }
2257 else if(unlikely(hash == hash_netdata_conf && strcmp(tok, "netdata.conf") == 0)) {
2258 debug(D_WEB_CLIENT_ACCESS, "%llu: Sending netdata.conf ...", w->id);
2259 w->response.data->contenttype = CT_TEXT_PLAIN;
@@ -2709,7 +2801,8 @@ void *web_client_main(void *ptr)
2801
2802 struct web_client *w = ptr;
2803 struct pollfd fds[2], *ifd, *ofd;
2712 - int retval, fdmax = 0, timeout;
2804 + int retval, timeout;
2805 + nfds_t fdmax = 0;
2806
2807 log_access("%llu: %s port %s connected on thread task id %d", w->id, w->client_ip, w->client_port, gettid());
2808
@@ -2798,6 +2891,10 @@ void *web_client_main(void *ptr)
2891 if(w->mode == WEB_CLIENT_MODE_NORMAL) {
2892 debug(D_WEB_CLIENT, "%llu: Attempting to process received data.", w->id);
2893 web_client_process_request(w);
2894 +
2895 + // if the sockets are closed, may have transferred this client
2896 + // to plugins.d
2897 + if(w->ifd == -1 && w->ofd == -1) break;
2898 }
2899 }
2900