init_connectors.c 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. // SPDX-License-Identifier: GPL-3.0-or-later
  2. #include "exporting_engine.h"
  3. #include "graphite/graphite.h"
  4. #include "json/json.h"
  5. #include "opentsdb/opentsdb.h"
  6. #ifdef ENABLE_PROMETHEUS_REMOTE_WRITE
  7. #include "prometheus/remote_write/remote_write.h"
  8. #endif
  9. #if HAVE_KINESIS
  10. #include "aws_kinesis/aws_kinesis.h"
  11. #endif
  12. #ifdef ENABLE_EXPORTING_PUBSUB
  13. #include "pubsub/pubsub.h"
  14. #endif
  15. #ifdef HAVE_MONGOC
  16. #include "mongodb/mongodb.h"
  17. #endif
  18. /**
  19. * Initialize connectors
  20. *
  21. * @param engine an engine data structure.
  22. * @return Returns 0 on success, 1 on failure.
  23. */
  24. int init_connectors(struct engine *engine)
  25. {
  26. engine->now = now_realtime_sec();
  27. for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
  28. instance->index = engine->instance_num++;
  29. instance->after = engine->now;
  30. switch (instance->config.type) {
  31. case EXPORTING_CONNECTOR_TYPE_GRAPHITE:
  32. if (init_graphite_instance(instance) != 0)
  33. return 1;
  34. break;
  35. case EXPORTING_CONNECTOR_TYPE_GRAPHITE_HTTP:
  36. if (init_graphite_instance(instance) != 0)
  37. return 1;
  38. break;
  39. case EXPORTING_CONNECTOR_TYPE_JSON:
  40. if (init_json_instance(instance) != 0)
  41. return 1;
  42. break;
  43. case EXPORTING_CONNECTOR_TYPE_JSON_HTTP:
  44. if (init_json_http_instance(instance) != 0)
  45. return 1;
  46. break;
  47. case EXPORTING_CONNECTOR_TYPE_OPENTSDB:
  48. if (init_opentsdb_telnet_instance(instance) != 0)
  49. return 1;
  50. break;
  51. case EXPORTING_CONNECTOR_TYPE_OPENTSDB_HTTP:
  52. if (init_opentsdb_http_instance(instance) != 0)
  53. return 1;
  54. break;
  55. case EXPORTING_CONNECTOR_TYPE_PROMETHEUS_REMOTE_WRITE:
  56. #ifdef ENABLE_PROMETHEUS_REMOTE_WRITE
  57. if (init_prometheus_remote_write_instance(instance) != 0)
  58. return 1;
  59. #endif
  60. break;
  61. case EXPORTING_CONNECTOR_TYPE_KINESIS:
  62. #if HAVE_KINESIS
  63. if (init_aws_kinesis_instance(instance) != 0)
  64. return 1;
  65. #endif
  66. break;
  67. case EXPORTING_CONNECTOR_TYPE_PUBSUB:
  68. #if ENABLE_EXPORTING_PUBSUB
  69. if (init_pubsub_instance(instance) != 0)
  70. return 1;
  71. #endif
  72. break;
  73. case EXPORTING_CONNECTOR_TYPE_MONGODB:
  74. #ifdef HAVE_MONGOC
  75. if (init_mongodb_instance(instance) != 0)
  76. return 1;
  77. #endif
  78. break;
  79. default:
  80. netdata_log_error("EXPORTING: unknown exporting connector type");
  81. return 1;
  82. }
  83. // dispatch the instance worker thread
  84. int error = uv_thread_create(&instance->thread, instance->worker, instance);
  85. if (error) {
  86. netdata_log_error("EXPORTING: cannot create thread worker. uv_thread_create(): %s", uv_strerror(error));
  87. return 1;
  88. }
  89. char threadname[NETDATA_THREAD_NAME_MAX + 1];
  90. snprintfz(threadname, NETDATA_THREAD_NAME_MAX, "EXPORTING-%zu", instance->index);
  91. uv_thread_set_name_np(instance->thread, threadname);
  92. analytics_statistic_t statistic = { "EXPORTING_START", "OK", instance->config.type_name };
  93. analytics_statistic_send(&statistic);
  94. }
  95. return 0;
  96. }
  97. // TODO: use a base64 encoder from a library
  98. static size_t base64_encode(unsigned char *input, size_t input_size, char *output, size_t output_size)
  99. {
  100. uint32_t value;
  101. static char lookup[] = "ABCDEFGHIJKLMNOPQRSTUVWXYZ"
  102. "abcdefghijklmnopqrstuvwxyz"
  103. "0123456789+/";
  104. if ((input_size / 3 + 1) * 4 >= output_size) {
  105. netdata_log_error("Output buffer for encoding size=%zu is not large enough for %zu-bytes input", output_size, input_size);
  106. return 0;
  107. }
  108. size_t count = 0;
  109. while (input_size >= 3) {
  110. value = ((input[0] << 16) + (input[1] << 8) + input[2]) & 0xffffff;
  111. output[0] = lookup[value >> 18];
  112. output[1] = lookup[(value >> 12) & 0x3f];
  113. output[2] = lookup[(value >> 6) & 0x3f];
  114. output[3] = lookup[value & 0x3f];
  115. //netdata_log_error("Base-64 encode (%04x) -> %c %c %c %c\n", value, output[0], output[1], output[2], output[3]);
  116. output += 4;
  117. input += 3;
  118. input_size -= 3;
  119. count += 4;
  120. }
  121. switch (input_size) {
  122. case 2:
  123. value = (input[0] << 10) + (input[1] << 2);
  124. output[0] = lookup[(value >> 12) & 0x3f];
  125. output[1] = lookup[(value >> 6) & 0x3f];
  126. output[2] = lookup[value & 0x3f];
  127. output[3] = '=';
  128. //netdata_log_error("Base-64 encode (%06x) -> %c %c %c %c\n", (value>>2)&0xffff, output[0], output[1], output[2], output[3]);
  129. count += 4;
  130. output[4] = '\0';
  131. break;
  132. case 1:
  133. value = input[0] << 4;
  134. output[0] = lookup[(value >> 6) & 0x3f];
  135. output[1] = lookup[value & 0x3f];
  136. output[2] = '=';
  137. output[3] = '=';
  138. //netdata_log_error("Base-64 encode (%06x) -> %c %c %c %c\n", value, output[0], output[1], output[2], output[3]);
  139. count += 4;
  140. output[4] = '\0';
  141. break;
  142. case 0:
  143. output[0] = '\0';
  144. break;
  145. }
  146. return count;
  147. }
  148. /**
  149. * Initialize a ring buffer and credentials for a simple connector
  150. *
  151. * @param instance an instance data structure.
  152. */
  153. void simple_connector_init(struct instance *instance)
  154. {
  155. struct simple_connector_data *connector_specific_data =
  156. (struct simple_connector_data *)instance->connector_specific_data;
  157. if (connector_specific_data->first_buffer)
  158. return;
  159. connector_specific_data->header = buffer_create(0, &netdata_buffers_statistics.buffers_exporters);
  160. connector_specific_data->buffer = buffer_create(0, &netdata_buffers_statistics.buffers_exporters);
  161. // create a ring buffer
  162. struct simple_connector_buffer *first_buffer = NULL;
  163. if (instance->config.buffer_on_failures < 1)
  164. instance->config.buffer_on_failures = 1;
  165. for (int i = 0; i < instance->config.buffer_on_failures; i++) {
  166. struct simple_connector_buffer *current_buffer = callocz(1, sizeof(struct simple_connector_buffer));
  167. if (!connector_specific_data->first_buffer)
  168. first_buffer = current_buffer;
  169. else
  170. current_buffer->next = connector_specific_data->first_buffer;
  171. connector_specific_data->first_buffer = current_buffer;
  172. }
  173. first_buffer->next = connector_specific_data->first_buffer;
  174. connector_specific_data->last_buffer = connector_specific_data->first_buffer;
  175. if (*instance->config.username || *instance->config.password) {
  176. BUFFER *auth_string = buffer_create(0, &netdata_buffers_statistics.buffers_exporters);
  177. buffer_sprintf(auth_string, "%s:%s", instance->config.username, instance->config.password);
  178. size_t encoded_size = (buffer_strlen(auth_string) / 3 + 1) * 4 + 1;
  179. char *encoded_credentials = callocz(1, encoded_size);
  180. base64_encode((unsigned char*)buffer_tostring(auth_string), buffer_strlen(auth_string), encoded_credentials, encoded_size);
  181. buffer_flush(auth_string);
  182. buffer_sprintf(auth_string, "Authorization: Basic %s\n", encoded_credentials);
  183. freez(encoded_credentials);
  184. connector_specific_data->auth_string = strdupz(buffer_tostring(auth_string));
  185. buffer_free(auth_string);
  186. }
  187. return;
  188. }