| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | #ifndef NETDATA_EXPORTING_ENGINE_H |
| 4 | #define NETDATA_EXPORTING_ENGINE_H 1 |
| 5 | |
| 6 | #include "database/rrd.h" |
| 7 | #include <uv.h> |
| 8 | |
| 9 | #define exporter_get(section, name, value) expconfig_get(&exporting_config, section, name, value) |
| 10 | #define exporter_get_number(section, name, value) expconfig_get_number(&exporting_config, section, name, value) |
| 11 | #define exporter_get_boolean(section, name, value) expconfig_get_boolean(&exporting_config, section, name, value) |
| 12 | |
| 13 | extern struct config exporting_config; |
| 14 | |
| 15 | #define EXPORTING_UPDATE_EVERY_OPTION_NAME "update every" |
| 16 | #define EXPORTING_UPDATE_EVERY_DEFAULT 10 |
| 17 | |
| 18 | typedef enum exporting_options { |
| 19 | EXPORTING_OPTION_NON = 0, |
| 20 | |
| 21 | EXPORTING_SOURCE_DATA_AS_COLLECTED = (1 << 0), |
| 22 | EXPORTING_SOURCE_DATA_AVERAGE = (1 << 1), |
| 23 | EXPORTING_SOURCE_DATA_SUM = (1 << 2), |
| 24 | |
| 25 | EXPORTING_OPTION_SEND_CONFIGURED_LABELS = (1 << 3), |
| 26 | EXPORTING_OPTION_SEND_AUTOMATIC_LABELS = (1 << 4), |
| 27 | EXPORTING_OPTION_USE_TLS = (1 << 5), |
| 28 | |
| 29 | EXPORTING_OPTION_SEND_NAMES = (1 << 16), |
| 30 | EXPORTING_OPTION_SEND_VARIABLES = (1 << 17) |
| 31 | } EXPORTING_OPTIONS; |
| 32 | |
| 33 | #define EXPORTING_OPTIONS_SOURCE_BITS \ |
| 34 | (EXPORTING_SOURCE_DATA_AS_COLLECTED | EXPORTING_SOURCE_DATA_AVERAGE | EXPORTING_SOURCE_DATA_SUM) |
| 35 | #define EXPORTING_OPTIONS_DATA_SOURCE(exporting_options) ((exporting_options) & EXPORTING_OPTIONS_SOURCE_BITS) |
| 36 | |
| 37 | extern EXPORTING_OPTIONS global_exporting_options; |
| 38 | extern const char *global_exporting_prefix; |
| 39 | |
| 40 | #define sending_labels_configured(instance) \ |
| 41 | ((instance)->config.options & (EXPORTING_OPTION_SEND_CONFIGURED_LABELS | EXPORTING_OPTION_SEND_AUTOMATIC_LABELS)) |
| 42 | |
| 43 | #define should_send_label(instance, label_source) \ |
| 44 | (((instance)->config.options & EXPORTING_OPTION_SEND_CONFIGURED_LABELS && (label_source)&RRDLABEL_SRC_CONFIG) || \ |
| 45 | ((instance)->config.options & EXPORTING_OPTION_SEND_AUTOMATIC_LABELS && (label_source)&RRDLABEL_SRC_AUTO)) |
| 46 | |
| 47 | #define should_send_variables(instance) ((instance)->config.options & EXPORTING_OPTION_SEND_VARIABLES) |
| 48 | |
| 49 | typedef enum exporting_connector_types { |
| 50 | EXPORTING_CONNECTOR_TYPE_UNKNOWN, // Invalid type |
| 51 | EXPORTING_CONNECTOR_TYPE_GRAPHITE, // Send plain text to Graphite |
| 52 | EXPORTING_CONNECTOR_TYPE_GRAPHITE_HTTP, // Send data to Graphite using HTTP API |
| 53 | EXPORTING_CONNECTOR_TYPE_JSON, // Send data in JSON format |
| 54 | EXPORTING_CONNECTOR_TYPE_JSON_HTTP, // Send data in JSON format using HTTP API |
| 55 | EXPORTING_CONNECTOR_TYPE_OPENTSDB, // Send data to OpenTSDB using telnet API |
| 56 | EXPORTING_CONNECTOR_TYPE_OPENTSDB_HTTP, // Send data to OpenTSDB using HTTP API |
| 57 | EXPORTING_CONNECTOR_TYPE_PROMETHEUS_REMOTE_WRITE, // Send data using Prometheus remote write protocol |
| 58 | EXPORTING_CONNECTOR_TYPE_KINESIS, // Send message to AWS Kinesis |
| 59 | EXPORTING_CONNECTOR_TYPE_PUBSUB, // Send message to Google Cloud Pub/Sub |
| 60 | EXPORTING_CONNECTOR_TYPE_MONGODB, // Send data to MongoDB collection |
| 61 | EXPORTING_CONNECTOR_TYPE_NUM // Number of exporting connector types |
| 62 | } EXPORTING_CONNECTOR_TYPE; |
| 63 | |
| 64 | struct engine; |
| 65 | |
| 66 | struct instance_config { |
| 67 | EXPORTING_CONNECTOR_TYPE type; |
| 68 | const char *type_name; |
| 69 | |
| 70 | const char *name; |
| 71 | const char *destination; |
| 72 | const char *username; |
| 73 | const char *password; |
| 74 | const char *prefix; |
| 75 | const char *label_prefix; |
| 76 | const char *hostname; |
| 77 | const char *thread_tag; |
| 78 | |
| 79 | int update_every; |
| 80 | int buffer_on_failures; |
| 81 | long timeoutms; |
| 82 | |
| 83 | EXPORTING_OPTIONS options; |
| 84 | SIMPLE_PATTERN *charts_pattern; |
| 85 | SIMPLE_PATTERN *hosts_pattern; |
| 86 | |
| 87 | int initialized; |
| 88 | |
| 89 | void *connector_specific_config; |
| 90 | }; |
| 91 | |
| 92 | struct simple_connector_config { |
| 93 | int default_port; |
| 94 | }; |
| 95 | |
| 96 | struct simple_connector_buffer { |
| 97 | BUFFER *header; |
| 98 | BUFFER *buffer; |
| 99 | |
| 100 | size_t buffered_metrics; |
| 101 | size_t buffered_bytes; |
| 102 | |
| 103 | int used; |
| 104 | |
| 105 | struct simple_connector_buffer *next; |
| 106 | }; |
| 107 | |
| 108 | #define CONNECTED_TO_MAX 1024 |
| 109 | |
| 110 | struct simple_connector_data { |
| 111 | void *connector_specific_data; |
| 112 | |
| 113 | char connected_to[CONNECTED_TO_MAX]; |
| 114 | |
| 115 | char *auth_string; |
| 116 | |
| 117 | size_t total_buffered_metrics; |
| 118 | |
| 119 | BUFFER *header; |
| 120 | BUFFER *buffer; |
| 121 | size_t buffered_metrics; |
| 122 | size_t buffered_bytes; |
| 123 | |
| 124 | struct simple_connector_buffer *previous_buffer; |
| 125 | struct simple_connector_buffer *first_buffer; |
| 126 | struct simple_connector_buffer *last_buffer; |
| 127 | |
| 128 | NETDATA_SSL ssl; |
| 129 | }; |
| 130 | |
| 131 | struct prometheus_remote_write_specific_config { |
| 132 | char *remote_write_path; |
| 133 | }; |
| 134 | |
| 135 | struct aws_kinesis_specific_config { |
| 136 | char *stream_name; |
| 137 | char *auth_key_id; |
| 138 | char *secure_key; |
| 139 | }; |
| 140 | |
| 141 | struct pubsub_specific_config { |
| 142 | char *credentials_file; |
| 143 | char *project_id; |
| 144 | char *topic_id; |
| 145 | }; |
| 146 | |
| 147 | struct mongodb_specific_config { |
| 148 | char *database; |
| 149 | char *collection; |
| 150 | }; |
| 151 | |
| 152 | struct engine_config { |
| 153 | const char *hostname; |
| 154 | int update_every; |
| 155 | }; |
| 156 | |
| 157 | struct stats { |
| 158 | collected_number buffered_metrics; |
| 159 | collected_number lost_metrics; |
| 160 | collected_number sent_metrics; |
| 161 | collected_number buffered_bytes; |
| 162 | collected_number lost_bytes; |
| 163 | collected_number sent_bytes; |
| 164 | collected_number received_bytes; |
| 165 | collected_number transmission_successes; |
| 166 | collected_number data_lost_events; |
| 167 | collected_number reconnects; |
| 168 | collected_number transmission_failures; |
| 169 | collected_number receptions; |
| 170 | |
| 171 | int initialized; |
| 172 | |
| 173 | RRDSET *st_metrics; |
| 174 | RRDDIM *rd_buffered_metrics; |
| 175 | RRDDIM *rd_lost_metrics; |
| 176 | RRDDIM *rd_sent_metrics; |
| 177 | |
| 178 | RRDSET *st_bytes; |
| 179 | RRDDIM *rd_buffered_bytes; |
| 180 | RRDDIM *rd_lost_bytes; |
| 181 | RRDDIM *rd_sent_bytes; |
| 182 | RRDDIM *rd_received_bytes; |
| 183 | |
| 184 | RRDSET *st_ops; |
| 185 | RRDDIM *rd_transmission_successes; |
| 186 | RRDDIM *rd_data_lost_events; |
| 187 | RRDDIM *rd_reconnects; |
| 188 | RRDDIM *rd_transmission_failures; |
| 189 | RRDDIM *rd_receptions; |
| 190 | |
| 191 | RRDSET *st_rusage; |
| 192 | RRDDIM *rd_user; |
| 193 | RRDDIM *rd_system; |
| 194 | }; |
| 195 | |
| 196 | struct instance { |
| 197 | struct instance_config config; |
| 198 | void *buffer; |
| 199 | void (*worker)(void *instance_p); |
| 200 | struct stats stats; |
| 201 | |
| 202 | int scheduled; |
| 203 | int disabled; |
| 204 | int skip_host; |
| 205 | int skip_chart; |
| 206 | |
| 207 | BUFFER *labels_buffer; |
| 208 | |
| 209 | time_t after; |
| 210 | time_t before; |
| 211 | |
| 212 | ND_THREAD *thread; |
| 213 | netdata_mutex_t mutex; |
| 214 | netdata_cond_t cond_var; |
| 215 | int data_is_ready; |
| 216 | |
| 217 | int (*start_batch_formatting)(struct instance *instance); |
| 218 | int (*start_host_formatting)(struct instance *instance, RRDHOST *host); |
| 219 | int (*start_chart_formatting)(struct instance *instance, RRDSET *st); |
| 220 | int (*metric_formatting)(struct instance *instance, RRDDIM *rd); |
| 221 | int (*end_chart_formatting)(struct instance *instance, RRDSET *st); |
| 222 | int (*variables_formatting)(struct instance *instance, RRDHOST *host); |
| 223 | int (*end_host_formatting)(struct instance *instance, RRDHOST *host); |
| 224 | int (*end_batch_formatting)(struct instance *instance); |
| 225 | |
| 226 | void (*prepare_header)(struct instance *instance); |
| 227 | int (*check_response)(BUFFER *buffer, struct instance *instance); |
| 228 | |
| 229 | void *connector_specific_data; |
| 230 | |
| 231 | size_t index; |
| 232 | struct instance *next; |
| 233 | struct engine *engine; |
| 234 | |
| 235 | volatile sig_atomic_t exited; |
| 236 | }; |
| 237 | |
| 238 | struct engine { |
| 239 | struct engine_config config; |
| 240 | |
| 241 | size_t instance_num; |
| 242 | time_t now; |
| 243 | |
| 244 | int aws_sdk_initialized; |
| 245 | int protocol_buffers_initialized; |
| 246 | int mongoc_initialized; |
| 247 | |
| 248 | struct instance *instance_root; |
| 249 | |
| 250 | volatile sig_atomic_t exit; |
| 251 | }; |
| 252 | |
| 253 | extern struct instance *prometheus_exporter_instance; |
| 254 | |
| 255 | void exporting_main(void *ptr); |
| 256 | |
| 257 | struct engine *read_exporting_config(); |
| 258 | EXPORTING_CONNECTOR_TYPE exporting_select_type(const char *type); |
| 259 | |
| 260 | int init_connectors(struct engine *engine); |
| 261 | void simple_connector_init(struct instance *instance); |
| 262 | |
| 263 | int mark_scheduled_instances(struct engine *engine); |
| 264 | void prepare_buffers(struct engine *engine); |
| 265 | |
| 266 | size_t exporting_name_copy(char *dst, const char *src, size_t max_len); |
| 267 | |
| 268 | int rrdhost_is_exportable(struct instance *instance, RRDHOST *host); |
| 269 | int rrdset_is_exportable(struct instance *instance, RRDSET *st); |
| 270 | |
| 271 | EXPORTING_OPTIONS exporting_parse_data_source(const char *source, EXPORTING_OPTIONS exporting_options); |
| 272 | |
| 273 | NETDATA_DOUBLE |
| 274 | exporting_calculate_value_from_stored_data( |
| 275 | struct instance *instance, |
| 276 | RRDDIM *rd, |
| 277 | time_t *last_timestamp); |
| 278 | |
| 279 | void start_batch_formatting(struct engine *engine); |
| 280 | void start_host_formatting(struct engine *engine, RRDHOST *host); |
| 281 | void start_chart_formatting(struct engine *engine, RRDSET *st); |
| 282 | void metric_formatting(struct engine *engine, RRDDIM *rd); |
| 283 | void end_chart_formatting(struct engine *engine, RRDSET *st); |
| 284 | void variables_formatting(struct engine *engine, RRDHOST *host); |
| 285 | void end_host_formatting(struct engine *engine, RRDHOST *host); |
| 286 | void end_batch_formatting(struct engine *engine); |
| 287 | int flush_host_labels(struct instance *instance, RRDHOST *host); |
| 288 | int simple_connector_end_batch(struct instance *instance); |
| 289 | |
| 290 | int exporting_discard_response(BUFFER *buffer, struct instance *instance); |
| 291 | void simple_connector_receive_response(int *sock, struct instance *instance); |
| 292 | void simple_connector_send_buffer( |
| 293 | int *sock, int *failures, struct instance *instance, BUFFER *header, BUFFER *buffer, size_t buffered_metrics); |
| 294 | void simple_connector_worker(void *instance_p); |
| 295 | |
| 296 | void create_main_rusage_chart(RRDSET **st_rusage, RRDDIM **rd_user, RRDDIM **rd_system); |
| 297 | void send_main_rusage(RRDSET *st_rusage, RRDDIM *rd_user, RRDDIM *rd_system); |
| 298 | void send_internal_metrics(struct instance *instance); |
| 299 | |
| 300 | void clean_instance(struct instance *ptr); |
| 301 | void simple_connector_cleanup(struct instance *instance); |
| 302 | |
| 303 | /** |
| 304 | * Free exporting configuration |
| 305 | * |
| 306 | * Free all memory associated with the exporting configuration. |
| 307 | * Called during shutdown to prevent memory leaks. |
| 308 | */ |
| 309 | void exporting_config_free(void); |
| 310 | |
| 311 | static inline void disable_instance(struct instance *instance) |
| 312 | { |
| 313 | instance->disabled = 1; |
| 314 | instance->scheduled = 0; |
| 315 | netdata_mutex_unlock(&instance->mutex); |
| 316 | netdata_log_error("EXPORTING: Instance %s disabled", instance->config.name); |
| 317 | } |
| 318 | |
| 319 | #include "exporting/prometheus/prometheus.h" |
| 320 | #include "exporting/opentsdb/opentsdb.h" |
| 321 | #ifdef ENABLE_PROMETHEUS_REMOTE_WRITE |
| 322 | #include "exporting/prometheus/remote_write/remote_write.h" |
| 323 | #endif |
| 324 | |
| 325 | #if HAVE_KINESIS |
| 326 | #include "exporting/aws_kinesis/aws_kinesis.h" |
| 327 | #endif |
| 328 | |
| 329 | #endif /* NETDATA_EXPORTING_ENGINE_H */ |