246
static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
247
}
248
249
+/**
250
+ * Set Kinesis variables
251
+ *
252
+ * Set the variables necessaries to work with this specific backend.
253
+ *
254
+ * @param default_port the default port of the backend
255
+ * @param brc function called to check the result.
256
+ * @param brf function called to format the msessage to the backend
257
+ * @param type the backend string selector.
258
+ */
259
+void backend_set_kinesis_variables(int *default_port,
260
+ backend_response_checker_t brc,
261
+ backend_request_formatter_t brf)
262
+{
263
+ (void)default_port;
264
+#ifndef HAVE_KINESIS
265
+ (void)brc;
266
+ (void)brf;
267
+#endif
268
+
269
+#if HAVE_KINESIS
270
+ *brc = process_json_response;
271
+ if (BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
272
+ *brf = format_dimension_collected_json_plaintext;
273
+ else
274
+ *brf = format_dimension_stored_json_plaintext;
275
+#endif
276
+}
277
+
278
+/**
279
+ * Set Prometheus variables
280
+ *
281
+ * Set the variables necessaries to work with this specific backend.
282
+ *
283
+ * @param default_port the default port of the backend
284
+ * @param brc function called to check the result.
285
+ * @param brf function called to format the msessage to the backend
286
+ * @param type the backend string selector.
287
+ */
288
+void backend_set_prometheus_variables(int *default_port,
289
+ backend_response_checker_t brc,
290
+ backend_request_formatter_t brf)
291
+{
292
+ (void)default_port;
293
+ (void)brf;
294
+#ifndef ENABLE_PROMETHEUS_REMOTE_WRITE
295
+ (void)brc;
296
+#endif
297
+
298
+#if ENABLE_PROMETHEUS_REMOTE_WRITE
299
+ *brc = process_prometheus_remote_write_response;
300
+#endif /* ENABLE_PROMETHEUS_REMOTE_WRITE */
301
+}
302
+
303
+/**
304
+ * Set JSON variables
305
+ *
306
+ * Set the variables necessaries to work with this specific backend.
307
+ *
308
+ * @param default_port the default port of the backend
309
+ * @param brc function called to check the result.
310
+ * @param brf function called to format the msessage to the backend
311
+ * @param type the backend string selector.
312
+ */
313
+void backend_set_json_variables(int *default_port,
314
+ backend_response_checker_t brc,
315
+ backend_request_formatter_t brf)
316
+{
317
+ *default_port = 5448;
318
+ *brc = process_json_response;
319
+
320
+ if (BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
321
+ *brf = format_dimension_collected_json_plaintext;
322
+ else
323
+ *brf = format_dimension_stored_json_plaintext;
324
+}
325
+
326
+/**
327
+ * Set OpenTSDB HTTP variables
328
+ *
329
+ * Set the variables necessaries to work with this specific backend.
330
+ *
331
+ * @param default_port the default port of the backend
332
+ * @param brc function called to check the result.
333
+ * @param brf function called to format the msessage to the backend
334
+ * @param type the backend string selector.
335
+ */
336
+void backend_set_opentsdb_http_variables(int *default_port,
337
+ backend_response_checker_t brc,
338
+ backend_request_formatter_t brf)
339
+{
340
+ *default_port = 4242;
341
+ *brc = process_opentsdb_response;
342
+
343
+ if(BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
344
+ *brf = format_dimension_collected_opentsdb_http;
345
+ else
346
+ *brf = format_dimension_stored_opentsdb_http;
347
+
348
+}
349
+
350
+/**
351
+ * Set OpenTSDB Telnet variables
352
+ *
353
+ * Set the variables necessaries to work with this specific backend.
354
+ *
355
+ * @param default_port the default port of the backend
356
+ * @param brc function called to check the result.
357
+ * @param brf function called to format the msessage to the backend
358
+ * @param type the backend string selector.
359
+ */
360
+void backend_set_opentsdb_telnet_variables(int *default_port,
361
+ backend_response_checker_t brc,
362
+ backend_request_formatter_t brf)
363
+{
364
+ *default_port = 4242;
365
+ *brc = process_opentsdb_response;
366
+
367
+ if(BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
368
+ *brf = format_dimension_collected_opentsdb_telnet;
369
+ else
370
+ *brf = format_dimension_stored_opentsdb_telnet;
371
+}
372
+
373
+/**
374
+ * Set Graphite variables
375
+ *
376
+ * Set the variables necessaries to work with this specific backend.
377
+ *
378
+ * @param default_port the default port of the backend
379
+ * @param brc function called to check the result.
380
+ * @param brf function called to format the msessage to the backend
381
+ * @param type the backend string selector.
382
+ */
383
+void backend_set_graphite_variables(int *default_port,
384
+ backend_response_checker_t brc,
385
+ backend_request_formatter_t brf)
386
+{
387
+ *default_port = 2003;
388
+ *brc = process_graphite_response;
389
+
390
+ if(BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
391
+ *brf = format_dimension_collected_graphite_plaintext;
392
+ else
393
+ *brf = format_dimension_stored_graphite_plaintext;
394
+}
395
+
396
+/**
397
+ * Select Type
398
+ *
399
+ * Select the backedn type based in the user input
400
+ *
401
+ * @param type is the string that defines the backend type
402
+ *
403
+ * @return It returns the backend id.
404
+ */
405
+BACKEND_TYPE backend_select_type(const char *type) {
406
+ if(!strcmp(type, "graphite") || !strcmp(type, "graphite:plaintext")) {
407
+ return BACKEND_TYPE_GRAPHITE;
408
+ }
409
+ else if(!strcmp(type, "opentsdb") || !strcmp(type, "opentsdb:telnet")) {
410
+ return BACKEND_TYPE_OPENTSDB_USING_TELNET;
411
+ }
412
+ else if(!strcmp(type, "opentsdb:http") || !strcmp(type, "opentsdb:https")) {
413
+ return BACKEND_TYPE_OPENTSDB_USING_HTTP;
414
+ }
415
+ else if (!strcmp(type, "json") || !strcmp(type, "json:plaintext")) {
416
+ return BACKEND_TYPE_JSON;
417
+ }
418
+ else if (!strcmp(type, "prometheus_remote_write")) {
419
+ return BACKEND_TYPE_PROMETEUS;
420
+ }
421
+ else if (!strcmp(type, "kinesis") || !strcmp(type, "kinesis:plaintext")) {
422
+ return BACKEND_TYPE_KINESIS;
423
+ }
424
+
425
+ return BACKEND_TYPE_UNKNOWN;
426
+}
427
+
428
+/**
429
+ * Backend main
430
+ *
431
+ * The main thread used to control the backedns.
432
+ *
433
+ * @param ptr a pointer to netdata_static_structure.
434
+ *
435
+ * @return It always return NULL.
436
+ */
437
void *backends_main(void *ptr) {
438
netdata_thread_cleanup_push(backends_main_cleanup, ptr);
439
453
BUFFER *http_request_header = buffer_create(1);
454
#endif
455
456
+#ifdef ENABLE_HTTPS
457
+ struct netdata_ssl opentsdb_ssl = {NULL , NETDATA_SSL_START};
458
+#endif
459
460
// ------------------------------------------------------------------------
461
// collect configuration options
504
505
// ------------------------------------------------------------------------
506
// select the backend type
316
-
317
- if(!strcmp(type, "graphite") || !strcmp(type, "graphite:plaintext")) {
318
-
319
- default_port = 2003;
320
- backend_response_checker = process_graphite_response;
321
-
322
- if(BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
323
- backend_request_formatter = format_dimension_collected_graphite_plaintext;
324
- else
325
- backend_request_formatter = format_dimension_stored_graphite_plaintext;
326
-
327
- }
328
- else if(!strcmp(type, "opentsdb") || !strcmp(type, "opentsdb:telnet")) {
329
-
330
- default_port = 4242;
331
- backend_response_checker = process_opentsdb_response;
332
-
333
- if(BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
334
- backend_request_formatter = format_dimension_collected_opentsdb_telnet;
335
- else
336
- backend_request_formatter = format_dimension_stored_opentsdb_telnet;
337
-
507
+ BACKEND_TYPE work_type = backend_select_type(type);
508
+ if (work_type == BACKEND_TYPE_UNKNOWN) {
509
+ error("BACKEND: Unknown backend type '%s'", type);
510
+ goto cleanup;
511
}
339
- else if (!strcmp(type, "json") || !strcmp(type, "json:plaintext")) {
340
-
341
- default_port = 5448;
342
- backend_response_checker = process_json_response;
512
344
- if (BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
345
- backend_request_formatter = format_dimension_collected_json_plaintext;
346
- else
347
- backend_request_formatter = format_dimension_stored_json_plaintext;
348
-
349
- }
350
- else if (!strcmp(type, "kinesis") || !strcmp(type, "kinesis:plaintext")) {
351
-#if HAVE_KINESIS
352
- do_kinesis = 1;
353
-
354
- if(unlikely(read_kinesis_conf(netdata_configured_user_config_dir, &kinesis_auth_key_id, &kinesis_secure_key, &kinesis_stream_name))) {
355
- error("BACKEND: kinesis backend type is set but cannot read its configuration from %s/aws_kinesis.conf", netdata_configured_user_config_dir);
356
- goto cleanup;
513
+ switch (work_type) {
514
+ case BACKEND_TYPE_OPENTSDB_USING_HTTP: {
515
+#ifdef ENABLE_HTTPS
516
+ if (!strcmp(type, "opentsdb:https")) {
517
+ security_start_ssl(NETDATA_SSL_CONTEXT_OPENTSDB);
518
+ }
519
+#endif
520
+ backend_set_opentsdb_http_variables(&default_port,&backend_response_checker,&backend_request_formatter);
521
+ break;
522
}
523
+ case BACKEND_TYPE_PROMETEUS: {
524
+#if ENABLE_PROMETHEUS_REMOTE_WRITE
525
+ do_prometheus_remote_write = 1;
526
359
- kinesis_init(destination, kinesis_auth_key_id, kinesis_secure_key, timeout.tv_sec * 1000 + timeout.tv_usec / 1000);
360
-
361
- backend_response_checker = process_json_response;
362
- if (BACKEND_OPTIONS_DATA_SOURCE(global_backend_options) == BACKEND_SOURCE_DATA_AS_COLLECTED)
363
- backend_request_formatter = format_dimension_collected_json_plaintext;
364
- else
365
- backend_request_formatter = format_dimension_stored_json_plaintext;
527
+ init_write_request();
528
#else
367
- error("AWS Kinesis support isn't compiled");
368
-#endif /* HAVE_KINESIS */
369
- }
370
- else if (!strcmp(type, "prometheus_remote_write")) {
371
-#if ENABLE_PROMETHEUS_REMOTE_WRITE
372
- do_prometheus_remote_write = 1;
529
+ error("BACKEND: Prometheus remote write support isn't compiled");
530
+#endif // ENABLE_PROMETHEUS_REMOTE_WRITE
531
+ backend_set_prometheus_variables(&default_port,&backend_response_checker,&backend_request_formatter);
532
+ break;
533
+ }
534
+ case BACKEND_TYPE_KINESIS: {
535
+#if HAVE_KINESIS
536
+ do_kinesis = 1;
537
374
- backend_response_checker = process_prometheus_remote_write_response;
538
+ if(unlikely(read_kinesis_conf(netdata_configured_user_config_dir, &kinesis_auth_key_id, &kinesis_secure_key, &kinesis_stream_name))) {
539
+ error("BACKEND: kinesis backend type is set but cannot read its configuration from %s/aws_kinesis.conf", netdata_configured_user_config_dir);
540
+ goto cleanup;
541
+ }
542
376
- init_write_request();
543
+ kinesis_init(destination, kinesis_auth_key_id, kinesis_secure_key, timeout.tv_sec * 1000 + timeout.tv_usec / 1000);
544
#else
378
- error("Prometheus remote write support isn't compiled");
379
-#endif /* ENABLE_PROMETHEUS_REMOTE_WRITE */
380
- }
381
- else {
382
- error("BACKEND: Unknown backend type '%s'", type);
383
- goto cleanup;
545
+ error("BACKEND: AWS Kinesis support isn't compiled");
546
+#endif // HAVE_KINESIS
547
+ backend_set_kinesis_variables(&default_port,&backend_response_checker,&backend_request_formatter);
548
+ break;
549
+ }
550
+ case BACKEND_TYPE_GRAPHITE: {
551
+ backend_set_graphite_variables(&default_port,&backend_response_checker,&backend_request_formatter);
552
+ break;
553
+ }
554
+ case BACKEND_TYPE_OPENTSDB_USING_TELNET: {
555
+ backend_set_opentsdb_telnet_variables(&default_port,&backend_response_checker,&backend_request_formatter);
556
+ break;
557
+ }
558
+ case BACKEND_TYPE_JSON: {
559
+ backend_set_json_variables(&default_port,&backend_response_checker,&backend_request_formatter);
560
+ break;
561
+ }
562
+ case BACKEND_TYPE_UNKNOWN: {
563
+ break;
564
+ }
565
}
566
567
#if ENABLE_PROMETHEUS_REMOTE_WRITE
574
}
575
576
396
- // ------------------------------------------------------------------------
397
- // prepare the charts for monitoring the backend operation
577
+// ------------------------------------------------------------------------
578
+// prepare the charts for monitoring the backend operation
579
580
struct rusage thread;
581
582
collected_number
402
- chart_buffered_metrics = 0,
403
- chart_lost_metrics = 0,
404
- chart_sent_metrics = 0,
405
- chart_buffered_bytes = 0,
406
- chart_received_bytes = 0,
407
- chart_sent_bytes = 0,
408
- chart_receptions = 0,
409
- chart_transmission_successes = 0,
410
- chart_transmission_failures = 0,
411
- chart_data_lost_events = 0,
412
- chart_lost_bytes = 0,
413
- chart_backend_reconnects = 0;
414
- // chart_backend_latency = 0;
583
+ chart_buffered_metrics = 0,
584
+ chart_lost_metrics = 0,
585
+ chart_sent_metrics = 0,
586
+ chart_buffered_bytes = 0,
587
+ chart_received_bytes = 0,
588
+ chart_sent_bytes = 0,
589
+ chart_receptions = 0,
590
+ chart_transmission_successes = 0,
591
+ chart_transmission_failures = 0,
592
+ chart_data_lost_events = 0,
593
+ chart_lost_bytes = 0,
594
+ chart_backend_reconnects = 0;
595
+ // chart_backend_latency = 0;
596
597
RRDSET *chart_metrics = rrdset_create_localhost("netdata", "backend_metrics", NULL, "backend", NULL, "Netdata Buffered Metrics", "metrics", "backends", NULL, 130600, global_backend_update_every, RRDSET_TYPE_LINE);
598
rrddim_add(chart_metrics, "buffered", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
613
rrddim_add(chart_ops, "read", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
614
615
/*
435
- * this is misleading - we can only measure the time we need to send data
436
- * this time is not related to the time required for the data to travel to
437
- * the backend database and the time that server needed to process them
438
- *
439
- * issue #1432 and https://www.softlab.ntua.gr/facilities/documentation/unix/unix-socket-faq/unix-socket-faq-2.html
440
- *
616
+ * this is misleading - we can only measure the time we need to send data
617
+ * this time is not related to the time required for the data to travel to
618
+ * the backend database and the time that server needed to process them
619
+ *
620
+ * issue #1432 and https://www.softlab.ntua.gr/facilities/documentation/unix/unix-socket-faq/unix-socket-faq-2.html
621
+ *
622
RRDSET *chart_latency = rrdset_create_localhost("netdata", "backend_latency", NULL, "backend", NULL, "Netdata Backend Latency", "ms", "backends", NULL, 130620, global_backend_update_every, RRDSET_TYPE_AREA);
623
rrddim_add(chart_latency, "latency", NULL, 1, 1000, RRD_ALGORITHM_ABSOLUTE);
624
*/
849
while(sock != -1 && errno != EWOULDBLOCK) {
850
buffer_need_bytes(response, 4096);
851
671
- ssize_t r = recv(sock, &response->buffer[response->len], response->size - response->len, MSG_DONTWAIT);
852
+ ssize_t r;
853
+#ifdef ENABLE_HTTPS
854
+ if(opentsdb_ssl.conn && !opentsdb_ssl.flags) {
855
+ r = SSL_read(opentsdb_ssl.conn, &response->buffer[response->len], response->size - response->len);
856
+ } else {
857
+ r = recv(sock, &response->buffer[response->len], response->size - response->len, MSG_DONTWAIT);
858
+ }
859
+#else
860
+ r = recv(sock, &response->buffer[response->len], response->size - response->len, MSG_DONTWAIT);
861
+#endif
862
if(likely(r > 0)) {
863
// we received some data
864
response->len += r;
891
size_t reconnects = 0;
892
893
sock = connect_to_one_of(destination, default_port, &timeout, &reconnects, NULL, 0);
894
+#ifdef ENABLE_HTTPS
895
+ if(sock != -1) {
896
+ if(netdata_opentsdb_ctx) {
897
+ if(!opentsdb_ssl.conn) {
898
+ opentsdb_ssl.conn = SSL_new(netdata_opentsdb_ctx);
899
+ if(!opentsdb_ssl.conn) {
900
+ error("Failed to allocate SSL structure %d.", sock);
901
+ opentsdb_ssl.flags = NETDATA_SSL_NO_HANDSHAKE;
902
+ }
903
+ } else {
904
+ SSL_clear(opentsdb_ssl.conn);
905
+ }
906
+ }
907
908
+ if(opentsdb_ssl.conn) {
909
+ if(SSL_set_fd(opentsdb_ssl.conn, sock) != 1) {
910
+ error("Failed to set the socket to the SSL on socket fd %d.", host->rrdpush_sender_socket);
911
+ opentsdb_ssl.flags = NETDATA_SSL_NO_HANDSHAKE;
912
+ } else {
913
+ opentsdb_ssl.flags = NETDATA_SSL_HANDSHAKE_COMPLETE;
914
+ SSL_set_connect_state(opentsdb_ssl.conn);
915
+ int err = SSL_connect(opentsdb_ssl.conn);
916
+ if (err != 1) {
917
+ err = SSL_get_error(opentsdb_ssl.conn, err);
918
+ error("SSL cannot connect with the server: %s ", ERR_error_string((long)SSL_get_error(opentsdb_ssl.conn, err), NULL));
919
+ opentsdb_ssl.flags = NETDATA_SSL_NO_HANDSHAKE;
920
+ } //TODO: check certificate here
921
+ }
922
+ }
923
+ }
924
+#endif
925
chart_backend_reconnects += reconnects;
926
// chart_backend_latency += now_monotonic_usec() - start_ut;
927
}
976
}
977
#endif
978
759
- ssize_t written = send(sock, buffer_tostring(b), len, flags);
979
+ ssize_t written;
980
+#ifdef ENABLE_HTTPS
981
+ if(opentsdb_ssl.conn && !opentsdb_ssl.flags) {
982
+ written = SSL_write(opentsdb_ssl.conn, buffer_tostring(b), len);
983
+ } else {
984
+ written = send(sock, buffer_tostring(b), len, flags);
985
+ }
986
+#else
987
+ written = send(sock, buffer_tostring(b), len, flags);
988
+#endif
989
+
990
// chart_backend_latency += now_monotonic_usec() - start_ut;
991
if(written != -1 && (size_t)written == len) {
992
// we sent the data successfully
1113
buffer_free(b);
1114
buffer_free(response);
1115
1116
+#ifdef ENABLE_HTTPS
1117
+ if(netdata_opentsdb_ctx) {
1118
+ if(opentsdb_ssl.conn) {
1119
+ SSL_free(opentsdb_ssl.conn);
1120
+ }
1121
+ }
1122
+#endif
1123
+
1124
netdata_thread_cleanup_pop(1);
1125
return NULL;
1126
}