@cryptotaxi247 / netdata-1 / commits / d79bbbf94

Add an AWS Kinesis connector to the exporting engine (#8145)

* Prepare files for the AWS Kinesis exporting connector * Update the documentation * Rename functions in backends * Include the connector to the Netdata buid * Add initializers and a worker * Add Kinesis specific configuration options * Add a compile time configuration check * Remove the connector data structure * Restore unit tests * Fix the compile-time configuration check * Initialize AWS SDK only once * Don't create an instance for an unknown exporting connector * Separate client and request outcome data for every instance * Fix memory cleanup, document functions * Add unit tests * Update the documentation

Vladimir Kobal committed Feb 25, 2020 at 21:08 UTC d79bbbf943f72495e135eee4afc25723f886592f
29 files changed +952 -390
CMakeLists.txt
+20 -1
@@ -636,6 +636,13 @@ set(EXPORTING_ENGINE_FILES
636 exporting/send_internal_metrics.c
637 )
638
639 +set(KINESIS_EXPORTING_FILES
640 + exporting/aws_kinesis/aws_kinesis.c
641 + exporting/aws_kinesis/aws_kinesis.h
642 + exporting/aws_kinesis/aws_kinesis_put_record.cc
643 + exporting/aws_kinesis/aws_kinesis_put_record.h
644 + )
645 +
646 set(KINESIS_BACKEND_FILES
647 backends/aws_kinesis/aws_kinesis.c
648 backends/aws_kinesis/aws_kinesis.h
@@ -720,7 +727,7 @@ ENDIF()
727
728 IF(ENABLE_BACKEND_KINESIS)
729 message(STATUS "kinesis backend: enabled")
723 - list(APPEND NETDATA_FILES ${KINESIS_BACKEND_FILES})
730 + list(APPEND NETDATA_FILES ${KINESIS_BACKEND_FILES} ${KINESIS_EXPORTING_FILES})
731 list(APPEND NETDATA_COMMON_LIBRARIES aws-cpp-sdk-kinesis aws-cpp-sdk-core ${CRYPTO_LIBRARIES} ${SSL_LIBRARIES} ${CURL_LIBRARIES})
732 list(APPEND NETDATA_COMMON_INCLUDE_DIRS ${CRYPTO_INCLUDE_DIRS} ${SSL_INCLUDE_DIRS} ${CURL_INCLUDE_DIRS})
733 list(APPEND NETDATA_COMMON_CFLAGS ${CRYPTO_CFLAGS_OTHER} ${SSL_CFLAGS_OTHER} ${CURL_CFLAGS_OTHER})
@@ -981,6 +988,17 @@ if(BUILD_TESTING)
988 exporting/tests/system_doubles.c
989 )
990 set(TEST_NAME exporting_engine)
991 + set(KINESIS_LINK_OPTIONS)
992 +if(ENABLE_BACKEND_KINESIS)
993 + list(APPEND EXPORTING_ENGINE_FILES ${KINESIS_EXPORTING_FILES})
994 + list(
995 + APPEND KINESIS_LINK_OPTIONS
996 + -Wl,--wrap=aws_sdk_init
997 + -Wl,--wrap=kinesis_init
998 + -Wl,--wrap=kinesis_put_record
999 + -Wl,--wrap=kinesis_get_result
1000 + )
1001 +endif()
1002 add_executable(${TEST_NAME}_testdriver ${EXPORTING_ENGINE_TEST_FILES} ${EXPORTING_ENGINE_FILES})
1003 target_compile_options(
1004 ${TEST_NAME}_testdriver
@@ -1011,6 +1029,7 @@ if(BUILD_TESTING)
1029 -Wl,--wrap=recv
1030 -Wl,--wrap=send
1031 -Wl,--wrap=connect_to_one_of
1032 + ${KINESIS_LINK_OPTIONS}
1033 )
1034 target_link_libraries(${TEST_NAME}_testdriver libnetdata ${NETDATA_COMMON_LIBRARIES} ${CMOCKA_LIBRARIES})
1035 add_test(NAME test_${TEST_NAME} COMMAND ${TEST_NAME}_testdriver)
Makefile.am
+24
@@ -495,6 +495,13 @@ EXPORTING_ENGINE_FILES = \
495 exporting/send_internal_metrics.c \
496 $(NULL)
497
498 +KINESIS_EXPORTING_FILES = \
499 + exporting/aws_kinesis/aws_kinesis.c \
500 + exporting/aws_kinesis/aws_kinesis.h \
501 + exporting/aws_kinesis/aws_kinesis_put_record.cc \
502 + exporting/aws_kinesis/aws_kinesis_put_record.h \
503 + $(NULL)
504 +
505 KINESIS_BACKEND_FILES = \
506 backends/aws_kinesis/aws_kinesis.c \
507 backends/aws_kinesis/aws_kinesis.h \
@@ -711,6 +718,13 @@ if ENABLE_PLUGIN_SLABINFO
718 $(NULL)
719 endif
720
721 +if ENABLE_EXPORTING
722 +if ENABLE_BACKEND_KINESIS
723 + netdata_SOURCES += $(KINESIS_EXPORTING_FILES)
724 + netdata_LDADD += $(OPTIONAL_KINESIS_LIBS)
725 +endif
726 +endif
727 +
728 if ENABLE_BACKEND_KINESIS
729 netdata_SOURCES += $(KINESIS_BACKEND_FILES)
730 netdata_LDADD += $(OPTIONAL_KINESIS_LIBS)
@@ -847,4 +861,14 @@ if ENABLE_UNITTESTS
861 $(TEST_LDFLAGS) \
862 $(NULL)
863 exporting_tests_exporting_engine_testdriver_LDADD = $(NETDATA_COMMON_LIBS) $(TEST_LIBS)
864 +if ENABLE_BACKEND_KINESIS
865 + exporting_tests_exporting_engine_testdriver_SOURCES += $(KINESIS_EXPORTING_FILES)
866 + exporting_tests_exporting_engine_testdriver_LDADD += $(OPTIONAL_KINESIS_LIBS)
867 + exporting_tests_exporting_engine_testdriver_LDFLAGS += \
868 + -Wl,--wrap=aws_sdk_init \
869 + -Wl,--wrap=kinesis_init \
870 + -Wl,--wrap=kinesis_put_record \
871 + -Wl,--wrap=kinesis_get_result \
872 + $(NULL)
873 +endif
874 endif
backends/aws_kinesis/aws_kinesis_put_record.cc
+7 -7
@@ -10,18 +10,18 @@
10
11 using namespace Aws;
12
13 -SDKOptions options;
13 +static SDKOptions options;
14
15 -Kinesis::KinesisClient *client;
15 +static Kinesis::KinesisClient *client;
16
17 struct request_outcome {
18 Kinesis::Model::PutRecordOutcomeCallable future_outcome;
19 size_t data_len;
20 };
21
22 -Vector<request_outcome> request_outcomes;
22 +static Vector<request_outcome> request_outcomes;
23
24 -void kinesis_init(const char *region, const char *access_key_id, const char *secret_key, const long timeout) {
24 +void backends_kinesis_init(const char *region, const char *access_key_id, const char *secret_key, const long timeout) {
25 InitAPI(options);
26
27 Client::ClientConfiguration config;
@@ -37,13 +37,13 @@ void kinesis_init(const char *region, const char *access_key_id, const char *sec
37 }
38 }
39
40 -void kinesis_shutdown() {
40 +void backends_kinesis_shutdown() {
41 Delete(client);
42
43 ShutdownAPI(options);
44 }
45
46 -int kinesis_put_record(const char *stream_name, const char *partition_key,
46 +int backends_kinesis_put_record(const char *stream_name, const char *partition_key,
47 const char *data, size_t data_len) {
48 Kinesis::Model::PutRecordRequest request;
49
@@ -56,7 +56,7 @@ int kinesis_put_record(const char *stream_name, const char *partition_key,
56 return 0;
57 }
58
59 -int kinesis_get_result(char *error_message, size_t *sent_bytes, size_t *lost_bytes) {
59 +int backends_kinesis_get_result(char *error_message, size_t *sent_bytes, size_t *lost_bytes) {
60 Kinesis::Model::PutRecordOutcome outcome;
61 *sent_bytes = 0;
62 *lost_bytes = 0;
backends/aws_kinesis/aws_kinesis_put_record.h
+4 -4
@@ -9,14 +9,14 @@
9 extern "C" {
10 #endif
11
12 -void kinesis_init(const char *region, const char *access_key_id, const char *secret_key, const long timeout);
12 +void backends_kinesis_init(const char *region, const char *access_key_id, const char *secret_key, const long timeout);
13
14 -void kinesis_shutdown();
14 +void backends_kinesis_shutdown();
15
16 -int kinesis_put_record(const char *stream_name, const char *partition_key,
16 +int backends_kinesis_put_record(const char *stream_name, const char *partition_key,
17 const char *data, size_t data_len);
18
19 -int kinesis_get_result(char *error_message, size_t *sent_bytes, size_t *lost_bytes);
19 +int backends_kinesis_get_result(char *error_message, size_t *sent_bytes, size_t *lost_bytes);
20
21 #ifdef __cplusplus
22 }
backends/backends.c
+5 -5
@@ -578,7 +578,7 @@ void *backends_main(void *ptr) {
578 goto cleanup;
579 }
580
581 - kinesis_init(destination, kinesis_auth_key_id, kinesis_secure_key, timeout.tv_sec * 1000 + timeout.tv_usec / 1000);
581 + backends_kinesis_init(destination, kinesis_auth_key_id, kinesis_secure_key, timeout.tv_sec * 1000 + timeout.tv_usec / 1000);
582 #else
583 error("BACKEND: AWS Kinesis support isn't compiled");
584 #endif // HAVE_KINESIS
@@ -860,18 +860,18 @@ void *backends_main(void *ptr) {
860
861 char error_message[ERROR_LINE_MAX + 1] = "";
862
863 - debug(D_BACKEND, "BACKEND: kinesis_put_record(): dest = %s, id = %s, key = %s, stream = %s, partition_key = %s, \
863 + debug(D_BACKEND, "BACKEND: backends_kinesis_put_record(): dest = %s, id = %s, key = %s, stream = %s, partition_key = %s, \
864 buffer = %zu, record = %zu", destination, kinesis_auth_key_id, kinesis_secure_key, kinesis_stream_name,
865 partition_key, buffer_len, record_len);
866
867 - kinesis_put_record(kinesis_stream_name, partition_key, first_char, record_len);
867 + backends_kinesis_put_record(kinesis_stream_name, partition_key, first_char, record_len);
868
869 sent += record_len;
870 chart_transmission_successes++;
871
872 size_t sent_bytes = 0, lost_bytes = 0;
873
874 - if(unlikely(kinesis_get_result(error_message, &sent_bytes, &lost_bytes))) {
874 + if(unlikely(backends_kinesis_get_result(error_message, &sent_bytes, &lost_bytes))) {
875 // oops! we couldn't send (all or some of the) data
876 error("BACKEND: %s", error_message);
877 error("BACKEND: failed to write data to database backend '%s'. Willing to write %zu bytes, wrote %zu bytes.",
@@ -1199,7 +1199,7 @@ void *backends_main(void *ptr) {
1199 cleanup:
1200 #if HAVE_KINESIS
1201 if(do_kinesis) {
1202 - kinesis_shutdown();
1202 + backends_kinesis_shutdown();
1203 freez(kinesis_auth_key_id);
1204 freez(kinesis_secure_key);
1205 freez(kinesis_stream_name);
configure.ac
+2 -1
@@ -445,7 +445,7 @@ if test "${ACLK}" = "yes"; then
445 AC_MSG_CHECKING([if libmosquitto static lib is present])
446 if test -f "externaldeps/mosquitto/libmosquitto.a"; then
447 HAVE_libmosquitto_a="yes"
448 - else
448 + else
449 HAVE_libmosquitto_a="no"
450 fi
451 AC_MSG_RESULT([${HAVE_libmosquitto_a}])
@@ -1311,6 +1311,7 @@ AC_CONFIG_FILES([
1311 exporting/graphite/Makefile
1312 exporting/json/Makefile
1313 exporting/opentsdb/Makefile
1314 + exporting/aws_kinesis/Makefile
1315 exporting/tests/Makefile
1316 health/Makefile
1317 health/notifications/Makefile
exporting/Makefile.am
+1
@@ -8,6 +8,7 @@ SUBDIRS = \
8 graphite \
9 json \
10 opentsdb \
11 + aws_kinesis \
12 $(NULL)
13
14 dist_noinst_DATA = \
exporting/aws_kinesis/Makefile.am new
+8
@@ -0,0 +1,8 @@
1 +# SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +AUTOMAKE_OPTIONS = subdir-objects
4 +MAINTAINERCLEANFILES = $(srcdir)/Makefile.in
5 +
6 +dist_noinst_DATA = \
7 + README.md \
8 + $(NULL)
exporting/aws_kinesis/README.md new
+48
@@ -0,0 +1,48 @@
1 +# Export metrics to AWS Kinesis Data Streams
2 +
3 +## Prerequisites
4 +
5 +To use AWS Kinesis for metric collecting and processing, you should first
6 +[install](https://docs.aws.amazon.com/en_us/sdk-for-cpp/v1/developer-guide/setup.html) AWS SDK for C++. Netdata
7 +works with the SDK version 1.7.121. Other versions might work correctly as well, but they were not tested with Netdata.
8 +`libcrypto`, `libssl`, and `libcurl` are also required to compile Netdata with Kinesis support enabled. Next, Netdata
9 +should be re-installed from the source. The installer will detect that the required libraries are now available.
10 +
11 +If the AWS SDK for C++ is being installed from source, it is useful to set `-DBUILD_ONLY="kinesis"`. Otherwise, the
12 +building process could take a very long time. Note that the default installation path for the libraries is
13 +`/usr/local/lib64`. Many Linux distributions don't include this path as the default one for a library search, so it is
14 +advisable to use the following options to `cmake` while building the AWS SDK:
15 +
16 +```sh
17 +cmake -DCMAKE_INSTALL_LIBDIR=/usr/lib -DCMAKE_INSTALL_INCLUDEDIR=/usr/include -DBUILD_SHARED_LIBS=OFF -DBUILD_ONLY=kinesis <aws-sdk-cpp sources>
18 +```
19 +
20 +## Configuration
21 +
22 +To enable data sending to the Kinesis service, run `./edit-config exporting.conf` in the Netdata configuration directory
23 +and set the following options:
24 +
25 +```conf
26 +[kinesis:my_instance]
27 + enabled = yes
28 + destination = us-east-1
29 +```
30 +
31 +Set the `destination` option to an AWS region.
32 +
33 +Set AWS credentials and stream name:
34 +
35 +```conf
36 + # AWS credentials
37 + aws_access_key_id = your_access_key_id
38 + aws_secret_access_key = your_secret_access_key
39 + # destination stream
40 + stream name = your_stream_name
41 +```
42 +
43 +Alternatively, you can set AWS credentials for the `netdata` user using AWS SDK for C++ [standard methods](https://docs.aws.amazon.com/sdk-for-cpp/v1/developer-guide/credentials.html).
44 +
45 +Netdata automatically computes a partition key for every record with the purpose to distribute records across
46 +available shards evenly.
47 +
48 +[![analytics](https://www.google-analytics.com/collect?v=1&aip=1&t=pageview&_s=1&ds=github&dr=https%3A%2F%2Fgithub.com%2Fnetdata%2Fnetdata&dl=https%3A%2F%2Fmy-netdata.io%2Fgithub%2Fexporting%2Faws_kinesis%2FREADME&_u=MAC~&cid=5792dfd7-8dc4-476b-af31-da2fdb9f93d2&tid=UA-64295674-3)](<>)
exporting/aws_kinesis/aws_kinesis.c new
+157
@@ -0,0 +1,157 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include "aws_kinesis.h"
4 +
5 +/**
6 + * Initialize AWS Kinesis connector instance
7 + *
8 + * @param instance an instance data structure.
9 + * @return Returns 0 on success, 1 on failure.
10 + */
11 +int init_aws_kinesis_instance(struct instance *instance)
12 +{
13 + instance->worker = aws_kinesis_connector_worker;
14 +
15 + instance->start_batch_formatting = NULL;
16 + instance->start_host_formatting = format_host_labels_json_plaintext;
17 + instance->start_chart_formatting = NULL;
18 +
19 + if (EXPORTING_OPTIONS_DATA_SOURCE(instance->config.options) == EXPORTING_SOURCE_DATA_AS_COLLECTED)
20 + instance->metric_formatting = format_dimension_collected_json_plaintext;
21 + else
22 + instance->metric_formatting = format_dimension_stored_json_plaintext;
23 +
24 + instance->end_chart_formatting = NULL;
25 + instance->end_host_formatting = flush_host_labels;
26 + instance->end_batch_formatting = NULL;
27 +
28 + instance->buffer = (void *)buffer_create(0);
29 + if (!instance->buffer) {
30 + error("EXPORTING: cannot create buffer for AWS Kinesis exporting connector instance %s", instance->config.name);
31 + return 1;
32 + }
33 + uv_mutex_init(&instance->mutex);
34 + uv_cond_init(&instance->cond_var);
35 +
36 + if (!instance->engine->aws_sdk_initialized) {
37 + aws_sdk_init();
38 + instance->engine->aws_sdk_initialized = 1;
39 + }
40 +
41 + struct aws_kinesis_specific_config *connector_specific_config = instance->config.connector_specific_config;
42 + struct aws_kinesis_specific_data *connector_specific_data = callocz(1, sizeof(struct aws_kinesis_specific_data));
43 + instance->connector_specific_data = (void *)connector_specific_data;
44 +
45 + kinesis_init(
46 + (void *)connector_specific_data,
47 + instance->config.destination,
48 + connector_specific_config->auth_key_id,
49 + connector_specific_config->secure_key,
50 + instance->config.timeoutms);
51 +
52 + return 0;
53 +}
54 +
55 +/**
56 + * AWS Kinesis connector worker
57 + *
58 + * Runs in a separate thread for every instance.
59 + *
60 + * @param instance_p an instance data structure.
61 + */
62 +void aws_kinesis_connector_worker(void *instance_p)
63 +{
64 + struct instance *instance = (struct instance *)instance_p;
65 + struct aws_kinesis_specific_config *connector_specific_config = instance->config.connector_specific_config;
66 + struct aws_kinesis_specific_data *connector_specific_data = instance->connector_specific_data;
67 +
68 + while (!netdata_exit) {
69 + unsigned long long partition_key_seq = 0;
70 + struct stats *stats = &instance->stats;
71 +
72 + uv_mutex_lock(&instance->mutex);
73 + uv_cond_wait(&instance->cond_var, &instance->mutex);
74 +
75 + BUFFER *buffer = (BUFFER *)instance->buffer;
76 + size_t buffer_len = buffer_strlen(buffer);
77 +
78 + size_t sent = 0;
79 +
80 + while (sent < buffer_len) {
81 + char partition_key[KINESIS_PARTITION_KEY_MAX + 1];
82 + snprintf(partition_key, KINESIS_PARTITION_KEY_MAX, "netdata_%llu", partition_key_seq++);
83 + size_t partition_key_len = strnlen(partition_key, KINESIS_PARTITION_KEY_MAX);
84 +
85 + const char *first_char = buffer_tostring(buffer) + sent;
86 +
87 + size_t record_len = 0;
88 +
89 + // split buffer into chunks of maximum allowed size
90 + if (buffer_len - sent < KINESIS_RECORD_MAX - partition_key_len) {
91 + record_len = buffer_len - sent;
92 + } else {
93 + record_len = KINESIS_RECORD_MAX - partition_key_len;
94 + while (*(first_char + record_len) != '\n' && record_len)
95 + record_len--;
96 + }
97 + char error_message[ERROR_LINE_MAX + 1] = "";
98 +
99 + debug(
100 + D_BACKEND,
101 + "EXPORTING: kinesis_put_record(): dest = %s, id = %s, key = %s, stream = %s, partition_key = %s, \
102 + buffer = %zu, record = %zu",
103 + instance->config.destination,
104 + connector_specific_config->auth_key_id,
105 + connector_specific_config->secure_key,
106 + connector_specific_config->stream_name,
107 + partition_key,
108 + buffer_len,
109 + record_len);
110 +
111 + kinesis_put_record(
112 + connector_specific_data, connector_specific_config->stream_name, partition_key, first_char, record_len);
113 +
114 + sent += record_len;
115 + stats->chart_transmission_successes++;
116 +
117 + size_t sent_bytes = 0, lost_bytes = 0;
118 +
119 + if (unlikely(kinesis_get_result(
120 + connector_specific_data->request_outcomes, error_message, &sent_bytes, &lost_bytes))) {
121 + // oops! we couldn't send (all or some of the) data
122 + error("EXPORTING: %s", error_message);
123 + error(
124 + "EXPORTING: failed to write data to database backend '%s'. Willing to write %zu bytes, wrote %zu bytes.",
125 + instance->config.destination, sent_bytes, sent_bytes - lost_bytes);
126 +
127 + stats->chart_transmission_failures++;
128 + stats->chart_data_lost_events++;
129 + stats->chart_lost_bytes += lost_bytes;
130 +
131 + // estimate the number of lost metrics
132 + stats->chart_lost_metrics += (collected_number)(
133 + stats->chart_buffered_metrics *
134 + (buffer_len && (lost_bytes > buffer_len) ? (double)lost_bytes / buffer_len : 1));
135 +
136 + break;
137 + } else {
138 + stats->chart_receptions++;
139 + }
140 +
141 + if (unlikely(netdata_exit))
142 + break;
143 + }
144 +
145 + stats->chart_sent_bytes += sent;
146 + if (likely(sent == buffer_len))
147 + stats->chart_sent_metrics = stats->chart_buffered_metrics;
148 +
149 + buffer_flush(buffer);
150 +
151 + uv_mutex_unlock(&instance->mutex);
152 +
153 +#ifdef UNIT_TESTING
154 + break;
155 +#endif
156 + }
157 +}
exporting/aws_kinesis/aws_kinesis.h new
+16
@@ -0,0 +1,16 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_EXPORTING_KINESIS_H
4 +#define NETDATA_EXPORTING_KINESIS_H
5 +
6 +#include "exporting/exporting_engine.h"
7 +#include "exporting/json/json.h"
8 +#include "aws_kinesis_put_record.h"
9 +
10 +#define KINESIS_PARTITION_KEY_MAX 256
11 +#define KINESIS_RECORD_MAX 1024 * 1024
12 +
13 +int init_aws_kinesis_instance(struct instance *instance);
14 +void aws_kinesis_connector_worker(void *instance_p);
15 +
16 +#endif //NETDATA_EXPORTING_KINESIS_H
exporting/aws_kinesis/aws_kinesis_put_record.cc new
+151
@@ -0,0 +1,151 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include <aws/core/Aws.h>
4 +#include <aws/core/client/ClientConfiguration.h>
5 +#include <aws/core/auth/AWSCredentials.h>
6 +#include <aws/core/utils/Outcome.h>
7 +#include <aws/kinesis/KinesisClient.h>
8 +#include <aws/kinesis/model/PutRecordRequest.h>
9 +#include "aws_kinesis_put_record.h"
10 +
11 +using namespace Aws;
12 +
13 +static SDKOptions options;
14 +
15 +struct request_outcome {
16 + Kinesis::Model::PutRecordOutcomeCallable future_outcome;
17 + size_t data_len;
18 +};
19 +
20 +/**
21 + * Initialize AWS SDK API
22 + */
23 +void aws_sdk_init()
24 +{
25 + InitAPI(options);
26 +}
27 +
28 +/**
29 + * Shutdown AWS SDK API
30 + */
31 +void aws_sdk_shutdown()
32 +{
33 + ShutdownAPI(options);
34 +}
35 +
36 +/**
37 + * Initialize a client and a data structure for request outcomes
38 + *
39 + * @param kinesis_specific_data_p a pointer to a structure with client and request outcome information.
40 + * @param region AWS region.
41 + * @param access_key_id AWS account access key ID.
42 + * @param secret_key AWS account secret access key.
43 + * @param timeout communication timeout.
44 + */
45 +void kinesis_init(
46 + void *kinesis_specific_data_p, const char *region, const char *access_key_id, const char *secret_key,
47 + const long timeout)
48 +{
49 + struct aws_kinesis_specific_data *kinesis_specific_data =
50 + (struct aws_kinesis_specific_data *)kinesis_specific_data_p;
51 +
52 + Client::ClientConfiguration config;
53 +
54 + config.region = region;
55 + config.requestTimeoutMs = timeout;
56 + config.connectTimeoutMs = timeout;
57 +
58 + Kinesis::KinesisClient *client;
59 +
60 + if (access_key_id && *access_key_id && secret_key && *secret_key) {
61 + client = New<Kinesis::KinesisClient>("client", Auth::AWSCredentials(access_key_id, secret_key), config);
62 + } else {
63 + client = New<Kinesis::KinesisClient>("client", config);
64 + }
65 + kinesis_specific_data->client = (void *)client;
66 +
67 + Vector<request_outcome> *request_outcomes;
68 +
69 + request_outcomes = new Vector<request_outcome>;
70 + kinesis_specific_data->request_outcomes = (void *)request_outcomes;
71 +}
72 +
73 +/**
74 + * Deallocate Kinesis specific data
75 + *
76 + * @param kinesis_specific_data_p a pointer to a structure with client and request outcome information.
77 + */
78 +void kinesis_shutdown(void *kinesis_specific_data_p)
79 +{
80 + struct aws_kinesis_specific_data *kinesis_specific_data =
81 + (struct aws_kinesis_specific_data *)kinesis_specific_data_p;
82 +
83 + Delete((Kinesis::KinesisClient *)kinesis_specific_data->client);
84 + delete (Vector<request_outcome> *)kinesis_specific_data->request_outcomes;
85 +}
86 +
87 +/**
88 + * Send data to the Kinesis service
89 + *
90 + * @param kinesis_specific_data_p a pointer to a structure with client and request outcome information.
91 + * @param stream_name the name of a stream to send to.
92 + * @param partition_key a partition key which automatically maps data to a specific stream.
93 + * @param data a data buffer to send to the stream.
94 + * @param data_len the length of the data buffer.
95 + */
96 +void kinesis_put_record(
97 + void *kinesis_specific_data_p, const char *stream_name, const char *partition_key, const char *data,
98 + size_t data_len)
99 +{
100 + struct aws_kinesis_specific_data *kinesis_specific_data =
101 + (struct aws_kinesis_specific_data *)kinesis_specific_data_p;
102 + Kinesis::Model::PutRecordRequest request;
103 +
104 + request.SetStreamName(stream_name);
105 + request.SetPartitionKey(partition_key);
106 + request.SetData(Utils::ByteBuffer((unsigned char *)data, data_len));
107 +
108 + ((Vector<request_outcome> *)(kinesis_specific_data->request_outcomes))->push_back(
109 + { ((Kinesis::KinesisClient *)(kinesis_specific_data->client))->PutRecordCallable(request), data_len });
110 +}
111 +
112 +/**
113 + * Get results from service responces
114 + *
115 + * @param request_outcomes_p request outcome information.
116 + * @param error_message report error message to a caller.
117 + * @param sent_bytes report to a caller how many bytes was successfuly sent.
118 + * @param lost_bytes report to a caller how many bytes was lost during transmission.
119 + * @return Returns 0 if all data was sent successfully, 1 when data was lost on transmission
120 + */
121 +int kinesis_get_result(void *request_outcomes_p, char *error_message, size_t *sent_bytes, size_t *lost_bytes)
122 +{
123 + Vector<request_outcome> *request_outcomes = (Vector<request_outcome> *)request_outcomes_p;
124 + Kinesis::Model::PutRecordOutcome outcome;
125 + *sent_bytes = 0;
126 + *lost_bytes = 0;
127 +
128 + for (auto request_outcome = request_outcomes->begin(); request_outcome != request_outcomes->end();) {
129 + std::future_status status = request_outcome->future_outcome.wait_for(std::chrono::microseconds(100));
130 +
131 + if (status == std::future_status::ready || status == std::future_status::deferred) {
132 + outcome = request_outcome->future_outcome.get();
133 + *sent_bytes += request_outcome->data_len;
134 +
135 + if (!outcome.IsSuccess()) {
136 + *lost_bytes += request_outcome->data_len;
137 + outcome.GetError().GetMessage().copy(error_message, ERROR_LINE_MAX);
138 + }
139 +
140 + request_outcomes->erase(request_outcome);
141 + } else {
142 + ++request_outcome;
143 + }
144 + }
145 +
146 + if (*lost_bytes) {
147 + return 1;
148 + }
149 +
150 + return 0;
151 +}
exporting/aws_kinesis/aws_kinesis_put_record.h new
+35
@@ -0,0 +1,35 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_EXPORTING_KINESIS_PUT_RECORD_H
4 +#define NETDATA_EXPORTING_KINESIS_PUT_RECORD_H
5 +
6 +#define ERROR_LINE_MAX 1023
7 +
8 +#ifdef __cplusplus
9 +extern "C" {
10 +#endif
11 +
12 +struct aws_kinesis_specific_data {
13 + void *client;
14 + void *request_outcomes;
15 +};
16 +
17 +void aws_sdk_init();
18 +void aws_sdk_shutdown();
19 +
20 +void kinesis_init(
21 + void *kinesis_specific_data_p, const char *region, const char *access_key_id, const char *secret_key,
22 + const long timeout);
23 +void kinesis_shutdown(void *client);
24 +
25 +void kinesis_put_record(
26 + void *kinesis_specific_data_p, const char *stream_name, const char *partition_key, const char *data,
27 + size_t data_len);
28 +
29 +int kinesis_get_result(void *request_outcomes_p, char *error_message, size_t *sent_bytes, size_t *lost_bytes);
30 +
31 +#ifdef __cplusplus
32 +}
33 +#endif
34 +
35 +#endif //NETDATA_EXPORTING_KINESIS_PUT_RECORD_H
exporting/check_filters.c
+2 -2
@@ -12,7 +12,7 @@
12 int rrdhost_is_exportable(struct instance *instance, RRDHOST *host)
13 {
14 if (host->exporting_flags == NULL)
15 - host->exporting_flags = callocz(instance->connector->engine->instance_num, sizeof(size_t));
15 + host->exporting_flags = callocz(instance->engine->instance_num, sizeof(size_t));
16
17 RRDHOST_FLAGS *flags = &host->exporting_flags[instance->index];
18
@@ -46,7 +46,7 @@ int rrdset_is_exportable(struct instance *instance, RRDSET *st)
46 RRDHOST *host = st->rrdhost;
47
48 if (st->exporting_flags == NULL)
49 - st->exporting_flags = callocz(instance->connector->engine->instance_num, sizeof(size_t));
49 + st->exporting_flags = callocz(instance->engine->instance_num, sizeof(size_t));
50
51 RRDSET_FLAGS *flags = &st->exporting_flags[instance->index];
52
exporting/exporting_engine.h
+18 -14
@@ -43,6 +43,12 @@ extern struct config exporting_config;
43 #define EXPORTER_SEND_NAMES "send names instead of ids"
44 #define EXPORTER_SEND_NAMES_DEFAULT CONFIG_BOOLEAN_YES
45
46 +#define EXPORTER_KINESIS_STREAM_NAME "stream name"
47 +#define EXPORTER_KINESIS_STREAM_NAME_DEFAULT "netdata"
48 +
49 +#define EXPORTER_AWS_ACCESS_KEY_ID "aws_access_key_id"
50 +#define EXPORTER_AWS_SECRET_ACCESS_KEY "aws_secret_access_key"
51 +
52 typedef enum exporting_options {
53 EXPORTING_OPTION_NONE = 0,
54
@@ -72,6 +78,8 @@ typedef enum exporting_options {
78 struct engine;
79
80 struct instance_config {
81 + BACKEND_TYPE type;
82 +
83 const char *name;
84 const char *destination;
85
@@ -90,9 +98,10 @@ struct simple_connector_config {
98 int default_port;
99 };
100
93 -struct connector_config {
94 - BACKEND_TYPE type;
95 - void *connector_specific_config;
101 +struct aws_kinesis_specific_config {
102 + char *stream_name;
103 + char *auth_key_id;
104 + char *secure_key;
105 };
106
107 struct engine_config {
@@ -119,6 +128,7 @@ struct stats {
128 struct instance {
129 struct instance_config config;
130 void *buffer;
131 + void (*worker)(void *instance_p);
132 struct stats stats;
133
134 int scheduled;
@@ -142,18 +152,10 @@ struct instance {
152 int (*end_host_formatting)(struct instance *instance, RRDHOST *host);
153 int (*end_batch_formatting)(struct instance *instance);
154
155 + void *connector_specific_data;
156 +
157 size_t index;
158 struct instance *next;
147 - struct connector *connector;
148 -};
149 -
150 -struct connector {
151 - struct connector_config config;
152 -
153 - void (*worker)(void *instance_p);
154 -
155 - struct instance *instance_root;
156 - struct connector *next;
159 struct engine *engine;
160 };
161
@@ -163,7 +165,9 @@ struct engine {
165 size_t instance_num;
166 time_t now;
167
166 - struct connector *connector_root;
168 + int aws_sdk_initialized;
169 +
170 + struct instance *instance_root;
171 };
172
173 void *exporting_main(void *ptr);
exporting/graphite/graphite.c
+8 -19
@@ -2,23 +2,6 @@
2
3 #include "graphite.h"
4
5 -/**
6 - * Initialize Graphite connector
7 - *
8 - * @param instance a connector data structure.
9 - * @return Always returns 0.
10 - */
11 -int init_graphite_connector(struct connector *connector)
12 -{
13 - connector->worker = simple_connector_worker;
14 -
15 - struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
16 - connector->config.connector_specific_config = (void *)connector_specific_config;
17 - connector_specific_config->default_port = 2003;
18 -
19 - return 0;
20 -}
21 -
5 /**
6 * Initialize Graphite connector instance
7 *
@@ -27,6 +10,12 @@ int init_graphite_connector(struct connector *connector)
10 */
11 int init_graphite_instance(struct instance *instance)
12 {
13 + instance->worker = simple_connector_worker;
14 +
15 + struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
16 + instance->config.connector_specific_config = (void *)connector_specific_config;
17 + connector_specific_config->default_port = 2003;
18 +
19 instance->start_batch_formatting = NULL;
20 instance->start_host_formatting = format_host_labels_graphite_plaintext;
21 instance->start_chart_formatting = NULL;
@@ -115,7 +104,7 @@ int format_host_labels_graphite_plaintext(struct instance *instance, RRDHOST *ho
104 */
105 int format_dimension_collected_graphite_plaintext(struct instance *instance, RRDDIM *rd)
106 {
118 - struct engine *engine = instance->connector->engine;
107 + struct engine *engine = instance->engine;
108 RRDSET *st = rd->rrdset;
109 RRDHOST *host = st->rrdhost;
110
@@ -156,7 +145,7 @@ int format_dimension_collected_graphite_plaintext(struct instance *instance, RRD
145 */
146 int format_dimension_stored_graphite_plaintext(struct instance *instance, RRDDIM *rd)
147 {
159 - struct engine *engine = instance->connector->engine;
148 + struct engine *engine = instance->engine;
149 RRDSET *st = rd->rrdset;
150 RRDHOST *host = st->rrdhost;
151
exporting/graphite/graphite.h
-1
@@ -5,7 +5,6 @@
5
6 #include "exporting/exporting_engine.h"
7
8 -int init_graphite_connector(struct connector *connector);
8 int init_graphite_instance(struct instance *instance);
9
10 void sanitize_graphite_label_value(char *dst, char *src, size_t len);
exporting/init_connectors.c
+24 -40
@@ -4,6 +4,7 @@
4 #include "graphite/graphite.h"
5 #include "json/json.h"
6 #include "opentsdb/opentsdb.h"
7 +#include "aws_kinesis/aws_kinesis.h"
8
9 /**
10 * Initialize connectors
@@ -15,64 +16,47 @@ int init_connectors(struct engine *engine)
16 {
17 engine->now = now_realtime_sec();
18
18 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
19 - switch (connector->config.type) {
19 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
20 + instance->index = engine->instance_num++;
21 + instance->after = engine->now;
22 +
23 + switch (instance->config.type) {
24 case BACKEND_TYPE_GRAPHITE:
21 - if (init_graphite_connector(connector) != 0)
25 + if (init_graphite_instance(instance) != 0)
26 return 1;
27 break;
28 case BACKEND_TYPE_JSON:
25 - if (init_json_connector(connector) != 0)
29 + if (init_json_instance(instance) != 0)
30 return 1;
31 break;
32 case BACKEND_TYPE_OPENTSDB_USING_TELNET:
29 - if (init_opentsdb_connector(connector) != 0)
33 + if (init_opentsdb_telnet_instance(instance) != 0)
34 return 1;
35 break;
36 case BACKEND_TYPE_OPENTSDB_USING_HTTP:
33 - if (init_opentsdb_connector(connector) != 0)
37 + if (init_opentsdb_http_instance(instance) != 0)
38 + return 1;
39 + break;
40 + case BACKEND_TYPE_KINESIS:
41 +#if HAVE_KINESIS
42 + if (init_aws_kinesis_instance(instance) != 0)
43 return 1;
44 +#endif
45 break;
46 default:
47 error("EXPORTING: unknown exporting connector type");
48 return 1;
49 }
40 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
41 - instance->index = engine->instance_num++;
42 - instance->after = engine->now;
50
44 - switch (connector->config.type) {
45 - case BACKEND_TYPE_GRAPHITE:
46 - if (init_graphite_instance(instance) != 0)
47 - return 1;
48 - break;
49 - case BACKEND_TYPE_JSON:
50 - if (init_json_instance(instance) != 0)
51 - return 1;
52 - break;
53 - case BACKEND_TYPE_OPENTSDB_USING_TELNET:
54 - if (init_opentsdb_telnet_instance(instance) != 0)
55 - return 1;
56 - break;
57 - case BACKEND_TYPE_OPENTSDB_USING_HTTP:
58 - if (init_opentsdb_http_instance(instance) != 0)
59 - return 1;
60 - break;
61 - default:
62 - error("EXPORTING: unknown exporting connector type");
63 - return 1;
64 - }
65 -
66 - // dispatch the instance worker thread
67 - int error = uv_thread_create(&instance->thread, connector->worker, instance);
68 - if (error) {
69 - error("EXPORTING: cannot create tread worker. uv_thread_create(): %s", uv_strerror(error));
70 - return 1;
71 - }
72 - char threadname[NETDATA_THREAD_NAME_MAX+1];
73 - snprintfz(threadname, NETDATA_THREAD_NAME_MAX, "EXPORTING-%zu", instance->index);
74 - uv_thread_set_name_np(instance->thread, threadname);
51 + // dispatch the instance worker thread
52 + int error = uv_thread_create(&instance->thread, instance->worker, instance);
53 + if (error) {
54 + error("EXPORTING: cannot create tread worker. uv_thread_create(): %s", uv_strerror(error));
55 + return 1;
56 }
57 + char threadname[NETDATA_THREAD_NAME_MAX+1];
58 + snprintfz(threadname, NETDATA_THREAD_NAME_MAX, "EXPORTING-%zu", instance->index);
59 + uv_thread_set_name_np(instance->thread, threadname);
60 }
61
62 return 0;
exporting/json/json.c
+8 -19
@@ -2,23 +2,6 @@
2
3 #include "json.h"
4
5 -/**
6 - * Initialize JSON connector
7 - *
8 - * @param instance a connector data structure.
9 - * @return Always returns 0.
10 - */
11 -int init_json_connector(struct connector *connector)
12 -{
13 - connector->worker = simple_connector_worker;
14 -
15 - struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
16 - connector->config.connector_specific_config = (void *)connector_specific_config;
17 - connector_specific_config->default_port = 5448;
18 -
19 - return 0;
20 -}
21 -
5 /**
6 * Initialize JSON connector instance
7 *
@@ -27,6 +10,12 @@ int init_json_connector(struct connector *connector)
10 */
11 int init_json_instance(struct instance *instance)
12 {
13 + instance->worker = simple_connector_worker;
14 +
15 + struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
16 + instance->config.connector_specific_config = (void *)connector_specific_config;
17 + connector_specific_config->default_port = 5448;
18 +
19 instance->start_batch_formatting = NULL;
20 instance->start_host_formatting = format_host_labels_json_plaintext;
21 instance->start_chart_formatting = NULL;
@@ -99,7 +88,7 @@ int format_host_labels_json_plaintext(struct instance *instance, RRDHOST *host)
88 */
89 int format_dimension_collected_json_plaintext(struct instance *instance, RRDDIM *rd)
90 {
102 - struct engine *engine = instance->connector->engine;
91 + struct engine *engine = instance->engine;
92 RRDSET *st = rd->rrdset;
93 RRDHOST *host = st->rrdhost;
94
@@ -171,7 +160,7 @@ int format_dimension_collected_json_plaintext(struct instance *instance, RRDDIM
160 */
161 int format_dimension_stored_json_plaintext(struct instance *instance, RRDDIM *rd)
162 {
174 - struct engine *engine = instance->connector->engine;
163 + struct engine *engine = instance->engine;
164 RRDSET *st = rd->rrdset;
165 RRDHOST *host = st->rrdhost;
166
exporting/json/json.h
-1
@@ -5,7 +5,6 @@
5
6 #include "exporting/exporting_engine.h"
7
8 -int init_json_connector(struct connector *connector);
8 int init_json_instance(struct instance *instance);
9
10 int format_host_labels_json_plaintext(struct instance *instance, RRDHOST *host);
exporting/opentsdb/opentsdb.c
+16 -21
@@ -2,23 +2,6 @@
2
3 #include "opentsdb.h"
4
5 -/**
6 - * Initialize OpenTSDB connector
7 - *
8 - * @param instance a connector data structure.
9 - * @return Always returns 0.
10 - */
11 -int init_opentsdb_connector(struct connector *connector)
12 -{
13 - connector->worker = simple_connector_worker;
14 -
15 - struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
16 - connector->config.connector_specific_config = (void *)connector_specific_config;
17 - connector_specific_config->default_port = 4242;
18 -
19 - return 0;
20 -}
21 -
5 /**
6 * Initialize OpenTSDB telnet connector instance
7 *
@@ -27,6 +10,12 @@ int init_opentsdb_connector(struct connector *connector)
10 */
11 int init_opentsdb_telnet_instance(struct instance *instance)
12 {
13 + instance->worker = simple_connector_worker;
14 +
15 + struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
16 + instance->config.connector_specific_config = (void *)connector_specific_config;
17 + connector_specific_config->default_port = 4242;
18 +
19 instance->start_batch_formatting = NULL;
20 instance->start_host_formatting = format_host_labels_opentsdb_telnet;
21 instance->start_chart_formatting = NULL;
@@ -59,6 +48,12 @@ int init_opentsdb_telnet_instance(struct instance *instance)
48 */
49 int init_opentsdb_http_instance(struct instance *instance)
50 {
51 + instance->worker = simple_connector_worker;
52 +
53 + struct simple_connector_config *connector_specific_config = mallocz(sizeof(struct simple_connector_config));
54 + instance->config.connector_specific_config = (void *)connector_specific_config;
55 + connector_specific_config->default_port = 4242;
56 +
57 instance->start_batch_formatting = NULL;
58 instance->start_host_formatting = format_host_labels_opentsdb_http;
59 instance->start_chart_formatting = NULL;
@@ -145,7 +140,7 @@ int format_host_labels_opentsdb_telnet(struct instance *instance, RRDHOST *host)
140 */
141 int format_dimension_collected_opentsdb_telnet(struct instance *instance, RRDDIM *rd)
142 {
148 - struct engine *engine = instance->connector->engine;
143 + struct engine *engine = instance->engine;
144 RRDSET *st = rd->rrdset;
145 RRDHOST *host = st->rrdhost;
146
@@ -186,7 +181,7 @@ int format_dimension_collected_opentsdb_telnet(struct instance *instance, RRDDIM
181 */
182 int format_dimension_stored_opentsdb_telnet(struct instance *instance, RRDDIM *rd)
183 {
189 - struct engine *engine = instance->connector->engine;
184 + struct engine *engine = instance->engine;
185 RRDSET *st = rd->rrdset;
186 RRDHOST *host = st->rrdhost;
187
@@ -293,7 +288,7 @@ int format_host_labels_opentsdb_http(struct instance *instance, RRDHOST *host)
288 */
289 int format_dimension_collected_opentsdb_http(struct instance *instance, RRDDIM *rd)
290 {
296 - struct engine *engine = instance->connector->engine;
291 + struct engine *engine = instance->engine;
292 RRDSET *st = rd->rrdset;
293 RRDHOST *host = st->rrdhost;
294
@@ -347,7 +342,7 @@ int format_dimension_collected_opentsdb_http(struct instance *instance, RRDDIM *
342 */
343 int format_dimension_stored_opentsdb_http(struct instance *instance, RRDDIM *rd)
344 {
350 - struct engine *engine = instance->connector->engine;
345 + struct engine *engine = instance->engine;
346 RRDSET *st = rd->rrdset;
347 RRDHOST *host = st->rrdhost;
348
exporting/opentsdb/opentsdb.h
-1
@@ -5,7 +5,6 @@
5
6 #include "exporting/exporting_engine.h"
7
8 -int init_opentsdb_connector(struct connector *connector);
8 int init_opentsdb_telnet_instance(struct instance *instance);
9 int init_opentsdb_http_instance(struct instance *instance);
10
exporting/process_data.c
+55 -71
@@ -42,13 +42,11 @@ int mark_scheduled_instances(struct engine *engine)
42 {
43 int instances_were_scheduled = 0;
44
45 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
46 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
47 - if (engine->now % instance->config.update_every < localhost->rrd_update_every) {
48 - instance->scheduled = 1;
49 - instances_were_scheduled = 1;
50 - instance->before = engine->now;
51 - }
45 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
46 + if (engine->now % instance->config.update_every < localhost->rrd_update_every) {
47 + instance->scheduled = 1;
48 + instances_were_scheduled = 1;
49 + instance->before = engine->now;
50 }
51 }
52
@@ -166,14 +164,12 @@ calculated_number exporting_calculate_value_from_stored_data(
164 */
165 int start_batch_formatting(struct engine *engine)
166 {
169 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
170 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
171 - if (instance->scheduled) {
172 - uv_mutex_lock(&instance->mutex);
173 - if (instance->start_batch_formatting && instance->start_batch_formatting(instance) != 0) {
174 - error("EXPORTING: cannot start batch formatting for %s", instance->config.name);
175 - return 1;
176 - }
167 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
168 + if (instance->scheduled) {
169 + uv_mutex_lock(&instance->mutex);
170 + if (instance->start_batch_formatting && instance->start_batch_formatting(instance) != 0) {
171 + error("EXPORTING: cannot start batch formatting for %s", instance->config.name);
172 + return 1;
173 }
174 }
175 }
@@ -190,17 +186,15 @@ int start_batch_formatting(struct engine *engine)
186 */
187 int start_host_formatting(struct engine *engine, RRDHOST *host)
188 {
193 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
194 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
195 - if (instance->scheduled) {
196 - if (rrdhost_is_exportable(instance, host)) {
197 - if (instance->start_host_formatting && instance->start_host_formatting(instance, host) != 0) {
198 - error("EXPORTING: cannot start host formatting for %s", instance->config.name);
199 - return 1;
200 - }
201 - } else {
202 - instance->skip_host = 1;
189 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
190 + if (instance->scheduled) {
191 + if (rrdhost_is_exportable(instance, host)) {
192 + if (instance->start_host_formatting && instance->start_host_formatting(instance, host) != 0) {
193 + error("EXPORTING: cannot start host formatting for %s", instance->config.name);
194 + return 1;
195 }
196 + } else {
197 + instance->skip_host = 1;
198 }
199 }
200 }
@@ -217,17 +211,15 @@ int start_host_formatting(struct engine *engine, RRDHOST *host)
211 */
212 int start_chart_formatting(struct engine *engine, RRDSET *st)
213 {
220 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
221 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
222 - if (instance->scheduled && !instance->skip_host) {
223 - if (rrdset_is_exportable(instance, st)) {
224 - if (instance->start_chart_formatting && instance->start_chart_formatting(instance, st) != 0) {
225 - error("EXPORTING: cannot start chart formatting for %s", instance->config.name);
226 - return 1;
227 - }
228 - } else {
229 - instance->skip_chart = 1;
214 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
215 + if (instance->scheduled && !instance->skip_host) {
216 + if (rrdset_is_exportable(instance, st)) {
217 + if (instance->start_chart_formatting && instance->start_chart_formatting(instance, st) != 0) {
218 + error("EXPORTING: cannot start chart formatting for %s", instance->config.name);
219 + return 1;
220 }
221 + } else {
222 + instance->skip_chart = 1;
223 }
224 }
225 }
@@ -244,15 +236,13 @@ int start_chart_formatting(struct engine *engine, RRDSET *st)
236 */
237 int metric_formatting(struct engine *engine, RRDDIM *rd)
238 {
247 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
248 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
249 - if (instance->scheduled && !instance->skip_host && !instance->skip_chart) {
250 - if (instance->metric_formatting && instance->metric_formatting(instance, rd) != 0) {
251 - error("EXPORTING: cannot format metric for %s", instance->config.name);
252 - return 1;
253 - }
254 - instance->stats.chart_buffered_metrics++;
239 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
240 + if (instance->scheduled && !instance->skip_host && !instance->skip_chart) {
241 + if (instance->metric_formatting && instance->metric_formatting(instance, rd) != 0) {
242 + error("EXPORTING: cannot format metric for %s", instance->config.name);
243 + return 1;
244 }
245 + instance->stats.chart_buffered_metrics++;
246 }
247 }
248
@@ -268,16 +258,14 @@ int metric_formatting(struct engine *engine, RRDDIM *rd)
258 */
259 int end_chart_formatting(struct engine *engine, RRDSET *st)
260 {
271 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
272 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
273 - if (instance->scheduled && !instance->skip_host && !instance->skip_chart) {
274 - if (instance->end_chart_formatting && instance->end_chart_formatting(instance, st) != 0) {
275 - error("EXPORTING: cannot end chart formatting for %s", instance->config.name);
276 - return 1;
277 - }
261 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
262 + if (instance->scheduled && !instance->skip_host && !instance->skip_chart) {
263 + if (instance->end_chart_formatting && instance->end_chart_formatting(instance, st) != 0) {
264 + error("EXPORTING: cannot end chart formatting for %s", instance->config.name);
265 + return 1;
266 }
279 - instance->skip_chart = 0;
267 }
268 + instance->skip_chart = 0;
269 }
270
271 return 0;
@@ -292,16 +280,14 @@ int end_chart_formatting(struct engine *engine, RRDSET *st)
280 */
281 int end_host_formatting(struct engine *engine, RRDHOST *host)
282 {
295 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
296 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
297 - if (instance->scheduled && !instance->skip_host) {
298 - if (instance->end_host_formatting && instance->end_host_formatting(instance, host) != 0) {
299 - error("EXPORTING: cannot end host formatting for %s", instance->config.name);
300 - return 1;
301 - }
283 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
284 + if (instance->scheduled && !instance->skip_host) {
285 + if (instance->end_host_formatting && instance->end_host_formatting(instance, host) != 0) {
286 + error("EXPORTING: cannot end host formatting for %s", instance->config.name);
287 + return 1;
288 }
303 - instance->skip_host = 0;
289 }
290 + instance->skip_host = 0;
291 }
292
293 return 0;
@@ -315,19 +301,17 @@ int end_host_formatting(struct engine *engine, RRDHOST *host)
301 */
302 int end_batch_formatting(struct engine *engine)
303 {
318 - for (struct connector *connector = engine->connector_root; connector; connector = connector->next) {
319 - for (struct instance *instance = connector->instance_root; instance; instance = instance->next) {
320 - if (instance->scheduled) {
321 - if (instance->end_batch_formatting && instance->end_batch_formatting(instance) != 0) {
322 - error("EXPORTING: cannot end batch formatting for %s", instance->config.name);
323 - return 1;
324 - }
325 - uv_mutex_unlock(&instance->mutex);
326 - uv_cond_signal(&instance->cond_var);
327 -
328 - instance->scheduled = 0;
329 - instance->after = instance->before;
304 + for (struct instance *instance = engine->instance_root; instance; instance = instance->next) {
305 + if (instance->scheduled) {
306 + if (instance->end_batch_formatting && instance->end_batch_formatting(instance) != 0) {
307 + error("EXPORTING: cannot end batch formatting for %s", instance->config.name);
308 + return 1;
309 }
310 + uv_mutex_unlock(&instance->mutex);
311 + uv_cond_signal(&instance->cond_var);
312 +
313 + instance->scheduled = 0;
314 + instance->after = instance->before;
315 }
316 }
317
exporting/read_config.c
+95 -89
@@ -194,11 +194,12 @@ struct engine *read_exporting_config()
194 static struct engine *engine = NULL;
195 struct connector_instance_list {
196 struct connector_instance local_ci;
197 + BACKEND_TYPE backend_type;
198 +
199 struct connector_instance_list *next;
200 };
201 struct connector_instance local_ci;
200 - struct connector_instance_list *tmp_ci_list, *tmp_ci_list1;
201 - struct connector_instance_list **ci_list;
202 + struct connector_instance_list *tmp_ci_list, *tmp_ci_list1, *tmp_ci_list_prev = NULL;
203
204 if (unlikely(engine))
205 return engine;
@@ -218,17 +219,11 @@ struct engine *read_exporting_config()
219
220 freez(filename);
221
221 - // Will build a list of instances per connector
222 // TODO: change BACKEND to EXPORTING
223 - ci_list = callocz(BACKEND_TYPE_NUM, sizeof(struct connector_instance_list *));
224 -
223 while (get_connector_instance(&local_ci)) {
226 - BACKEND_TYPE backend_type;
227 -
224 info("Processing connector instance (%s)", local_ci.instance_name);
229 - if (exporter_get_boolean(local_ci.instance_name, "enabled", 0)) {
230 - backend_type = exporting_select_type(local_ci.connector_name);
225
226 + if (exporter_get_boolean(local_ci.instance_name, "enabled", 0)) {
227 info(
228 "Instance (%s) on connector (%s) is enabled and scheduled for activation",
229 local_ci.instance_name,
@@ -236,8 +231,9 @@ struct engine *read_exporting_config()
231
232 tmp_ci_list = (struct connector_instance_list *)callocz(1, sizeof(struct connector_instance_list));
233 memcpy(&tmp_ci_list->local_ci, &local_ci, sizeof(local_ci));
239 - tmp_ci_list->next = ci_list[backend_type];
240 - ci_list[backend_type] = tmp_ci_list;
234 + tmp_ci_list->backend_type = exporting_select_type(local_ci.connector_name);
235 + tmp_ci_list->next = tmp_ci_list_prev;
236 + tmp_ci_list_prev = tmp_ci_list;
237 instances_to_activate++;
238 } else
239 info("Instance (%s) on connector (%s) is not enabled", local_ci.instance_name, local_ci.connector_name);
@@ -245,11 +241,10 @@ struct engine *read_exporting_config()
241
242 if (unlikely(!instances_to_activate)) {
243 info("No connector instances to activate");
248 - freez(ci_list);
244 return NULL;
245 }
246
252 - engine = (struct engine *)calloc(1, sizeof(struct engine));
247 + engine = (struct engine *)callocz(1, sizeof(struct engine));
248 // TODO: Check and fill engine fields if actually needed
249
250 if (exporting_config_exists) {
@@ -260,107 +255,118 @@ struct engine *read_exporting_config()
255 exporter_get_number(CONFIG_SECTION_EXPORTING, EXPORTER_UPDATE_EVERY, EXPORTER_UPDATE_EVERY_DEFAULT);
256 }
257
263 - for (size_t i = 0; i < BACKEND_TYPE_NUM; i++) {
264 - // For each connector build list
265 - tmp_ci_list = ci_list[i];
258 + while (tmp_ci_list) {
259 + struct instance *tmp_instance;
260 + char *instance_name;
261
267 - // If we have a list of instances for this connector then build it
268 - if (tmp_ci_list) {
269 - struct connector *tmp_connector;
262 + info("Instance %s on %s", tmp_ci_list->local_ci.instance_name, tmp_ci_list->local_ci.connector_name);
263
271 - tmp_connector = (struct connector *)calloc(1, sizeof(struct connector));
272 - tmp_connector->next = engine->connector_root;
273 - engine->connector_root = tmp_connector;
264 + if (tmp_ci_list->backend_type == BACKEND_TYPE_UNKNOWN) {
265 + error("Unknown exporting connector type");
266 + goto next_connector_instance;
267 + }
268
275 - tmp_connector->config.type = i;
276 - tmp_connector->engine = engine;
269 +#ifndef HAVE_KINESIS
270 + if (tmp_ci_list->backend_type == BACKEND_TYPE_KINESIS) {
271 + error("AWS Kinesis support isn't compiled");
272 + goto next_connector_instance;
273 + }
274 +#endif
275 +
276 + tmp_instance = (struct instance *)callocz(1, sizeof(struct instance));
277 + tmp_instance->next = engine->instance_root;
278 + engine->instance_root = tmp_instance;
279 +
280 + tmp_instance->engine = engine;
281 + tmp_instance->config.type = tmp_ci_list->backend_type;
282 +
283 + instance_name = tmp_ci_list->local_ci.instance_name;
284
278 - while (tmp_ci_list) {
279 - struct instance *tmp_instance;
280 - char *instance_name;
285 + tmp_instance->config.name = strdupz(tmp_ci_list->local_ci.instance_name);
286
282 - info("Instance %s on %s", tmp_ci_list->local_ci.instance_name, tmp_ci_list->local_ci.connector_name);
287 + tmp_instance->config.destination =
288 + strdupz(exporter_get(instance_name, EXPORTER_DESTINATION, EXPORTER_DESTINATION_DEFAULT));
289
284 - tmp_instance = (struct instance *)calloc(1, sizeof(struct instance));
285 - tmp_instance->connector = engine->connector_root;
286 - tmp_instance->next = engine->connector_root->instance_root;
287 - engine->connector_root->instance_root = tmp_instance;
288 - tmp_instance->connector = engine->connector_root;
290 + tmp_instance->config.update_every =
291 + exporter_get_number(instance_name, EXPORTER_UPDATE_EVERY, EXPORTER_UPDATE_EVERY_DEFAULT);
292
290 - instance_name = tmp_ci_list->local_ci.instance_name;
293 + tmp_instance->config.buffer_on_failures =
294 + exporter_get_number(instance_name, EXPORTER_BUF_ONFAIL, EXPORTER_BUF_ONFAIL_DEFAULT);
295
292 - tmp_instance->config.name = strdupz(tmp_ci_list->local_ci.instance_name);
296 + tmp_instance->config.timeoutms =
297 + exporter_get_number(instance_name, EXPORTER_TIMEOUT_MS, EXPORTER_TIMEOUT_MS_DEFAULT);
298
294 - tmp_instance->config.destination =
295 - strdupz(exporter_get(instance_name, EXPORTER_DESTINATION, EXPORTER_DESTINATION_DEFAULT));
299 + tmp_instance->config.charts_pattern = simple_pattern_create(
300 + exporter_get(instance_name, EXPORTER_SEND_CHART_MATCH, EXPORTER_SEND_CHART_MATCH_DEFAULT),
301 + NULL,
302 + SIMPLE_PATTERN_EXACT);
303
297 - tmp_instance->config.update_every =
298 - exporter_get_number(instance_name, EXPORTER_UPDATE_EVERY, EXPORTER_UPDATE_EVERY_DEFAULT);
304 + tmp_instance->config.hosts_pattern = simple_pattern_create(
305 + exporter_get(instance_name, EXPORTER_SEND_HOST_MATCH, EXPORTER_SEND_HOST_MATCH_DEFAULT),
306 + NULL,
307 + SIMPLE_PATTERN_EXACT);
308
300 - tmp_instance->config.buffer_on_failures =
301 - exporter_get_number(instance_name, EXPORTER_BUF_ONFAIL, EXPORTER_BUF_ONFAIL_DEFAULT);
309 + char *data_source =
310 + exporter_get(instance_name, EXPORTER_DATA_SOURCE, EXPORTER_DATA_SOURCE_DEFAULT);
311
303 - tmp_instance->config.timeoutms =
304 - exporter_get_number(instance_name, EXPORTER_TIMEOUT_MS, EXPORTER_TIMEOUT_MS_DEFAULT);
312 + tmp_instance->config.options = exporting_parse_data_source(data_source, tmp_instance->config.options);
313
306 - tmp_instance->config.charts_pattern = simple_pattern_create(
307 - exporter_get(instance_name, EXPORTER_SEND_CHART_MATCH, EXPORTER_SEND_CHART_MATCH_DEFAULT),
308 - NULL,
309 - SIMPLE_PATTERN_EXACT);
314 + if (exporter_get_boolean(
315 + instance_name, EXPORTER_SEND_CONFIGURED_LABELS, EXPORTER_SEND_CONFIGURED_LABELS_DEFAULT))
316 + tmp_instance->config.options |= EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
317 + else
318 + tmp_instance->config.options &= ~EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
319
311 - tmp_instance->config.hosts_pattern = simple_pattern_create(
312 - exporter_get(instance_name, EXPORTER_SEND_HOST_MATCH, EXPORTER_SEND_HOST_MATCH_DEFAULT),
313 - NULL,
314 - SIMPLE_PATTERN_EXACT);
320 + if (exporter_get_boolean(
321 + instance_name, EXPORTER_SEND_AUTOMATIC_LABELS, EXPORTER_SEND_AUTOMATIC_LABELS_DEFAULT))
322 + tmp_instance->config.options |= EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
323 + else
324 + tmp_instance->config.options &= ~EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
325
316 - char *data_source =
317 - exporter_get(instance_name, EXPORTER_DATA_SOURCE, EXPORTER_DATA_SOURCE_DEFAULT);
326 + if (exporter_get_boolean(instance_name, EXPORTER_SEND_NAMES, EXPORTER_SEND_NAMES_DEFAULT))
327 + tmp_instance->config.options |= EXPORTING_OPTION_SEND_NAMES;
328 + else
329 + tmp_instance->config.options &= ~EXPORTING_OPTION_SEND_NAMES;
330
319 - tmp_instance->config.options = exporting_parse_data_source(data_source, tmp_instance->config.options);
331 + if (tmp_instance->config.type == BACKEND_TYPE_KINESIS) {
332 + struct aws_kinesis_specific_config *connector_specific_config =
333 + callocz(1, sizeof(struct aws_kinesis_specific_config));
334
321 - if (exporter_get_boolean(
322 - instance_name, EXPORTER_SEND_CONFIGURED_LABELS, EXPORTER_SEND_CONFIGURED_LABELS_DEFAULT))
323 - tmp_instance->config.options |= EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
324 - else
325 - tmp_instance->config.options &= ~EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
335 + tmp_instance->config.connector_specific_config = connector_specific_config;
336
327 - if (exporter_get_boolean(
328 - instance_name, EXPORTER_SEND_AUTOMATIC_LABELS, EXPORTER_SEND_AUTOMATIC_LABELS_DEFAULT))
329 - tmp_instance->config.options |= EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
330 - else
331 - tmp_instance->config.options &= ~EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
337 + connector_specific_config->stream_name = strdupz(exporter_get(
338 + instance_name, EXPORTER_KINESIS_STREAM_NAME, EXPORTER_KINESIS_STREAM_NAME_DEFAULT));
339
333 - if (exporter_get_boolean(instance_name, EXPORTER_SEND_NAMES, EXPORTER_SEND_NAMES_DEFAULT))
334 - tmp_instance->config.options |= EXPORTING_OPTION_SEND_NAMES;
335 - else
336 - tmp_instance->config.options &= ~EXPORTING_OPTION_SEND_NAMES;
340 + connector_specific_config->auth_key_id = strdupz(exporter_get(
341 + instance_name, EXPORTER_AWS_ACCESS_KEY_ID, ""));
342 +
343 + connector_specific_config->secure_key = strdupz(exporter_get(
344 + instance_name, EXPORTER_AWS_SECRET_ACCESS_KEY, ""));
345 + }
346
347 #ifdef NETDATA_INTERNAL_CHECKS
339 - info(
340 - " Dest=[%s], upd=[%d], buffer=[%d] timeout=[%ld] options=[%u]",
341 - tmp_instance->config.destination,
342 - tmp_instance->config.update_every,
343 - tmp_instance->config.buffer_on_failures,
344 - tmp_instance->config.timeoutms,
345 - tmp_instance->config.options);
348 + info(
349 + " Dest=[%s], upd=[%d], buffer=[%d] timeout=[%ld] options=[%u]",
350 + tmp_instance->config.destination,
351 + tmp_instance->config.update_every,
352 + tmp_instance->config.buffer_on_failures,
353 + tmp_instance->config.timeoutms,
354 + tmp_instance->config.options);
355 #endif
356
348 - if (unlikely(!exporting_config_exists) && !engine->config.hostname) {
349 - engine->config.hostname =
350 - strdupz(config_get(instance_name, "hostname", netdata_configured_hostname));
351 - engine->config.prefix = strdupz(config_get(instance_name, "prefix", "netdata"));
352 - engine->config.update_every =
353 - config_get_number(instance_name, EXPORTER_UPDATE_EVERY, EXPORTER_UPDATE_EVERY_DEFAULT);
354 - }
355 -
356 - tmp_ci_list1 = tmp_ci_list->next;
357 - freez(tmp_ci_list);
358 - tmp_ci_list = tmp_ci_list1;
359 - }
357 + if (unlikely(!exporting_config_exists) && !engine->config.hostname) {
358 + engine->config.hostname =
359 + strdupz(config_get(instance_name, "hostname", netdata_configured_hostname));
360 + engine->config.prefix = strdupz(config_get(instance_name, "prefix", "netdata"));
361 + engine->config.update_every =
362 + config_get_number(instance_name, EXPORTER_UPDATE_EVERY, EXPORTER_UPDATE_EVERY_DEFAULT);
363 }
361 - }
364
363 - freez(ci_list);
365 +next_connector_instance:
366 + tmp_ci_list1 = tmp_ci_list->next;
367 + freez(tmp_ci_list);
368 + tmp_ci_list = tmp_ci_list1;
369 + }
370
371 return engine;
372 }
exporting/send_data.c
+1 -1
@@ -145,7 +145,7 @@ void simple_connector_worker(void *instance_p)
145 {
146 struct instance *instance = (struct instance*)instance_p;
147
148 - struct simple_connector_config *connector_specific_config = instance->connector->config.connector_specific_config;
148 + struct simple_connector_config *connector_specific_config = instance->config.connector_specific_config;
149 struct stats *stats = &instance->stats;
150
151 int sock = -1;
exporting/tests/exporting_doubles.c
+46 -6
@@ -16,13 +16,11 @@ struct engine *__mock_read_exporting_config()
16 engine->config.hostname = strdupz("test-host");
17 engine->config.update_every = 3;
18
19 - engine->connector_root = calloc(1, sizeof(struct connector));
20 - engine->connector_root->config.type = BACKEND_TYPE_GRAPHITE;
21 - engine->connector_root->engine = engine;
19
23 - engine->connector_root->instance_root = calloc(1, sizeof(struct instance));
24 - struct instance *instance = engine->connector_root->instance_root;
25 - instance->connector = engine->connector_root;
20 + engine->instance_root = calloc(1, sizeof(struct instance));
21 + struct instance *instance = engine->instance_root;
22 + instance->engine = engine;
23 + instance->config.type = BACKEND_TYPE_GRAPHITE;
24 instance->config.name = strdupz("instance_name");
25 instance->config.destination = strdupz("localhost");
26 instance->config.update_every = 1;
@@ -160,3 +158,45 @@ int __mock_end_batch_formatting(struct instance *instance)
158 check_expected_ptr(instance);
159 return mock_type(int);
160 }
161 +
162 +#if HAVE_KINESIS
163 +void __wrap_aws_sdk_init()
164 +{
165 + function_called();
166 +}
167 +
168 +void __wrap_kinesis_init(
169 + void *kinesis_specific_data_p, const char *region, const char *access_key_id, const char *secret_key,
170 + const long timeout)
171 +{
172 + function_called();
173 + check_expected_ptr(kinesis_specific_data_p);
174 + check_expected_ptr(region);
175 + check_expected_ptr(access_key_id);
176 + check_expected_ptr(secret_key);
177 + check_expected(timeout);
178 +}
179 +
180 +void __wrap_kinesis_put_record(
181 + void *kinesis_specific_data_p, const char *stream_name, const char *partition_key, const char *data,
182 + size_t data_len)
183 +{
184 + function_called();
185 + check_expected_ptr(kinesis_specific_data_p);
186 + check_expected_ptr(stream_name);
187 + check_expected_ptr(partition_key);
188 + check_expected_ptr(data);
189 + check_expected_ptr(data);
190 + check_expected(data_len);
191 +}
192 +
193 +int __wrap_kinesis_get_result(void *request_outcomes_p, char *error_message, size_t *sent_bytes, size_t *lost_bytes)
194 +{
195 + function_called();
196 + check_expected_ptr(request_outcomes_p);
197 + check_expected_ptr(error_message);
198 + check_expected_ptr(sent_bytes);
199 + check_expected_ptr(lost_bytes);
200 + return mock_type(int);
201 +}
202 +#endif //HAVE_KINESIS
exporting/tests/exporting_fixtures.c
+3 -5
@@ -15,15 +15,13 @@ int teardown_configured_engine(void **state)
15 {
16 struct engine *engine = *state;
17
18 - struct instance *instance = engine->connector_root->instance_root;
18 + struct instance *instance = engine->instance_root;
19 free((void *)instance->config.destination);
20 free((void *)instance->config.name);
21 simple_pattern_free(instance->config.charts_pattern);
22 simple_pattern_free(instance->config.hosts_pattern);
23 free(instance);
24
25 - free(engine->connector_root);
26 -
25 free((void *)engine->config.prefix);
26 free((void *)engine->config.hostname);
27 free(engine);
@@ -122,8 +120,8 @@ int teardown_initialized_engine(void **state)
120 struct engine *engine = *state;
121
122 teardown_rrdhost();
125 - buffer_free(engine->connector_root->instance_root->labels);
126 - buffer_free(engine->connector_root->instance_root->buffer);
123 + buffer_free(engine->instance_root->labels);
124 + buffer_free(engine->instance_root->buffer);
125 teardown_configured_engine(state);
126
127 return 0;
exporting/tests/test_exporting_engine.c
+188 -82
@@ -21,16 +21,16 @@ void init_connectors_in_tests(struct engine *engine)
21
22 expect_function_call(__wrap_uv_thread_create);
23
24 - expect_value(__wrap_uv_thread_create, thread, &engine->connector_root->instance_root->thread);
24 + expect_value(__wrap_uv_thread_create, thread, &engine->instance_root->thread);
25 expect_value(__wrap_uv_thread_create, worker, simple_connector_worker);
26 - expect_value(__wrap_uv_thread_create, arg, engine->connector_root->instance_root);
26 + expect_value(__wrap_uv_thread_create, arg, engine->instance_root);
27
28 expect_function_call(__wrap_uv_thread_set_name_np);
29
30 assert_int_equal(__real_init_connectors(engine), 0);
31
32 assert_int_equal(engine->now, 2);
33 - assert_int_equal(engine->connector_root->instance_root->after, 2);
33 + assert_int_equal(engine->instance_root->after, 2);
34 }
35
36 static void test_exporting_engine(void **state)
@@ -82,16 +82,12 @@ static void test_read_exporting_config(void **state)
82 assert_int_equal(engine->config.update_every, 3);
83 assert_int_equal(engine->instance_num, 0);
84
85 - struct connector *connector = engine->connector_root;
86 - assert_ptr_not_equal(connector, NULL);
87 - assert_ptr_equal(connector->next, NULL);
88 - assert_ptr_equal(connector->engine, engine);
89 - assert_int_equal(connector->config.type, BACKEND_TYPE_GRAPHITE);
85
91 - struct instance *instance = connector->instance_root;
86 + struct instance *instance = engine->instance_root;
87 assert_ptr_not_equal(instance, NULL);
88 assert_ptr_equal(instance->next, NULL);
94 - assert_ptr_equal(instance->connector, connector);
89 + assert_ptr_equal(instance->engine, engine);
90 + assert_int_equal(instance->config.type, BACKEND_TYPE_GRAPHITE);
91 assert_string_equal(instance->config.destination, "localhost");
92 assert_int_equal(instance->config.update_every, 1);
93 assert_int_equal(instance->config.buffer_on_failures, 10);
@@ -111,16 +107,15 @@ static void test_init_connectors(void **state)
107
108 assert_int_equal(engine->instance_num, 1);
109
114 - struct connector *connector = engine->connector_root;
115 - assert_ptr_equal(connector->next, NULL);
116 - assert_ptr_equal(connector->worker, simple_connector_worker);
110 + struct instance *instance = engine->instance_root;
111
118 - struct simple_connector_config *connector_specific_config = connector->config.connector_specific_config;
119 - assert_int_equal(connector_specific_config->default_port, 2003);
120 -
121 - struct instance *instance = connector->instance_root;
112 assert_ptr_equal(instance->next, NULL);
113 assert_int_equal(instance->index, 0);
114 +
115 + struct simple_connector_config *connector_specific_config = instance->config.connector_specific_config;
116 + assert_int_equal(connector_specific_config->default_port, 2003);
117 +
118 + assert_ptr_equal(instance->worker, simple_connector_worker);
119 assert_ptr_equal(instance->start_batch_formatting, NULL);
120 assert_ptr_equal(instance->start_host_formatting, format_host_labels_graphite_plaintext);
121 assert_ptr_equal(instance->start_chart_formatting, NULL);
@@ -138,16 +133,13 @@ static void test_init_connectors(void **state)
133 static void test_init_graphite_instance(void **state)
134 {
135 struct engine *engine = *state;
141 - struct connector *connector = engine->connector_root;
142 - struct instance *instance = connector->instance_root;
143 -
144 - init_graphite_connector(connector);
145 - assert_int_equal(
146 - ((struct simple_connector_config *)(connector->config.connector_specific_config))->default_port, 2003);
147 - freez(connector->config.connector_specific_config);
136 + struct instance *instance = engine->instance_root;
137
138 instance->config.options = EXPORTING_SOURCE_DATA_AS_COLLECTED | EXPORTING_OPTION_SEND_NAMES;
139 assert_int_equal(init_graphite_instance(instance), 0);
140 + assert_int_equal(
141 + ((struct simple_connector_config *)(instance->config.connector_specific_config))->default_port, 2003);
142 + freez(instance->config.connector_specific_config);
143 assert_ptr_equal(instance->metric_formatting, format_dimension_collected_graphite_plaintext);
144 assert_ptr_not_equal(instance->buffer, NULL);
145 buffer_free(instance->buffer);
@@ -160,16 +152,13 @@ static void test_init_graphite_instance(void **state)
152 static void test_init_json_instance(void **state)
153 {
154 struct engine *engine = *state;
163 - struct connector *connector = engine->connector_root;
164 - struct instance *instance = connector->instance_root;
165 -
166 - init_json_connector(connector);
167 - assert_int_equal(
168 - ((struct simple_connector_config *)(connector->config.connector_specific_config))->default_port, 5448);
169 - freez(connector->config.connector_specific_config);
155 + struct instance *instance = engine->instance_root;
156
157 instance->config.options = EXPORTING_SOURCE_DATA_AS_COLLECTED | EXPORTING_OPTION_SEND_NAMES;
158 assert_int_equal(init_json_instance(instance), 0);
159 + assert_int_equal(
160 + ((struct simple_connector_config *)(instance->config.connector_specific_config))->default_port, 5448);
161 + freez(instance->config.connector_specific_config);
162 assert_ptr_equal(instance->metric_formatting, format_dimension_collected_json_plaintext);
163 assert_ptr_not_equal(instance->buffer, NULL);
164 buffer_free(instance->buffer);
@@ -179,25 +168,16 @@ static void test_init_json_instance(void **state)
168 assert_ptr_equal(instance->metric_formatting, format_dimension_stored_json_plaintext);
169 }
170
182 -static void test_init_opentsdb_connector(void **state)
183 -{
184 - struct engine *engine = *state;
185 - struct connector *connector = engine->connector_root;
186 -
187 - init_opentsdb_connector(connector);
188 - assert_int_equal(
189 - ((struct simple_connector_config *)(connector->config.connector_specific_config))->default_port, 4242);
190 - freez(connector->config.connector_specific_config);
191 -}
192 -
171 static void test_init_opentsdb_telnet_instance(void **state)
172 {
173 struct engine *engine = *state;
196 - struct connector *connector = engine->connector_root;
197 - struct instance *instance = connector->instance_root;
174 + struct instance *instance = engine->instance_root;
175
176 instance->config.options = EXPORTING_SOURCE_DATA_AS_COLLECTED | EXPORTING_OPTION_SEND_NAMES;
177 assert_int_equal(init_opentsdb_telnet_instance(instance), 0);
178 + assert_int_equal(
179 + ((struct simple_connector_config *)(instance->config.connector_specific_config))->default_port, 4242);
180 + freez(instance->config.connector_specific_config);
181 assert_ptr_equal(instance->metric_formatting, format_dimension_collected_opentsdb_telnet);
182 assert_ptr_not_equal(instance->buffer, NULL);
183 buffer_free(instance->buffer);
@@ -210,11 +190,13 @@ static void test_init_opentsdb_telnet_instance(void **state)
190 static void test_init_opentsdb_http_instance(void **state)
191 {
192 struct engine *engine = *state;
213 - struct connector *connector = engine->connector_root;
214 - struct instance *instance = connector->instance_root;
193 + struct instance *instance = engine->instance_root;
194
195 instance->config.options = EXPORTING_SOURCE_DATA_AS_COLLECTED | EXPORTING_OPTION_SEND_NAMES;
196 assert_int_equal(init_opentsdb_http_instance(instance), 0);
197 + assert_int_equal(
198 + ((struct simple_connector_config *)(instance->config.connector_specific_config))->default_port, 4242);
199 + freez(instance->config.connector_specific_config);
200 assert_ptr_equal(instance->metric_formatting, format_dimension_collected_opentsdb_http);
201 assert_ptr_not_equal(instance->buffer, NULL);
202 buffer_free(instance->buffer);
@@ -230,7 +212,7 @@ static void test_mark_scheduled_instances(void **state)
212
213 assert_int_equal(__real_mark_scheduled_instances(engine), 1);
214
233 - struct instance *instance = engine->connector_root->instance_root;
215 + struct instance *instance = engine->instance_root;
216 assert_int_equal(instance->scheduled, 1);
217 assert_int_equal(instance->before, 2);
218 }
@@ -238,7 +220,7 @@ static void test_mark_scheduled_instances(void **state)
220 static void test_rrdhost_is_exportable(void **state)
221 {
222 struct engine *engine = *state;
241 - struct instance *instance = engine->connector_root->instance_root;
223 + struct instance *instance = engine->instance_root;
224
225 expect_function_call(__wrap_info_int);
226
@@ -255,7 +237,7 @@ static void test_rrdhost_is_exportable(void **state)
237 static void test_false_rrdhost_is_exportable(void **state)
238 {
239 struct engine *engine = *state;
258 - struct instance *instance = engine->connector_root->instance_root;
240 + struct instance *instance = engine->instance_root;
241
242 simple_pattern_free(instance->config.hosts_pattern);
243 instance->config.hosts_pattern = simple_pattern_create("!*", NULL, SIMPLE_PATTERN_EXACT);
@@ -275,7 +257,7 @@ static void test_false_rrdhost_is_exportable(void **state)
257 static void test_rrdset_is_exportable(void **state)
258 {
259 struct engine *engine = *state;
278 - struct instance *instance = engine->connector_root->instance_root;
260 + struct instance *instance = engine->instance_root;
261 RRDSET *st = localhost->rrdset_root;
262
263 assert_ptr_equal(st->exporting_flags, NULL);
@@ -289,7 +271,7 @@ static void test_rrdset_is_exportable(void **state)
271 static void test_false_rrdset_is_exportable(void **state)
272 {
273 struct engine *engine = *state;
292 - struct instance *instance = engine->connector_root->instance_root;
274 + struct instance *instance = engine->instance_root;
275 RRDSET *st = localhost->rrdset_root;
276
277 simple_pattern_free(instance->config.charts_pattern);
@@ -306,7 +288,7 @@ static void test_false_rrdset_is_exportable(void **state)
288 static void test_exporting_calculate_value_from_stored_data(void **state)
289 {
290 struct engine *engine = *state;
309 - struct instance *instance = engine->connector_root->instance_root;
291 + struct instance *instance = engine->instance_root;
292 RRDDIM *rd = localhost->rrdset_root->dimensions;
293 time_t timestamp;
294
@@ -344,7 +326,7 @@ static void test_exporting_calculate_value_from_stored_data(void **state)
326 static void test_prepare_buffers(void **state)
327 {
328 struct engine *engine = *state;
347 - struct instance *instance = engine->connector_root->instance_root;
329 + struct instance *instance = engine->instance_root;
330
331 instance->start_batch_formatting = __mock_start_batch_formatting;
332 instance->start_host_formatting = __mock_start_host_formatting;
@@ -434,9 +416,9 @@ static void test_format_dimension_collected_graphite_plaintext(void **state)
416 struct engine *engine = *state;
417
418 RRDDIM *rd = localhost->rrdset_root->dimensions;
437 - assert_int_equal(format_dimension_collected_graphite_plaintext(engine->connector_root->instance_root, rd), 0);
419 + assert_int_equal(format_dimension_collected_graphite_plaintext(engine->instance_root, rd), 0);
420 assert_string_equal(
439 - buffer_tostring(engine->connector_root->instance_root->buffer),
421 + buffer_tostring(engine->instance_root->buffer),
422 "netdata.test-host.chart_name.dimension_name;TAG1=VALUE1 TAG2=VALUE2 123000321 15051\n");
423 }
424
@@ -448,9 +430,9 @@ static void test_format_dimension_stored_graphite_plaintext(void **state)
430 will_return(__wrap_exporting_calculate_value_from_stored_data, pack_storage_number(27, SN_EXISTS));
431
432 RRDDIM *rd = localhost->rrdset_root->dimensions;
451 - assert_int_equal(format_dimension_stored_graphite_plaintext(engine->connector_root->instance_root, rd), 0);
433 + assert_int_equal(format_dimension_stored_graphite_plaintext(engine->instance_root, rd), 0);
434 assert_string_equal(
453 - buffer_tostring(engine->connector_root->instance_root->buffer),
435 + buffer_tostring(engine->instance_root->buffer),
436 "netdata.test-host.chart_name.dimension_name;TAG1=VALUE1 TAG2=VALUE2 690565856.0000000 15052\n");
437 }
438
@@ -459,9 +441,9 @@ static void test_format_dimension_collected_json_plaintext(void **state)
441 struct engine *engine = *state;
442
443 RRDDIM *rd = localhost->rrdset_root->dimensions;
462 - assert_int_equal(format_dimension_collected_json_plaintext(engine->connector_root->instance_root, rd), 0);
444 + assert_int_equal(format_dimension_collected_json_plaintext(engine->instance_root, rd), 0);
445 assert_string_equal(
464 - buffer_tostring(engine->connector_root->instance_root->buffer),
446 + buffer_tostring(engine->instance_root->buffer),
447 "{\"prefix\":\"netdata\",\"hostname\":\"test-host\",\"host_tags\":\"TAG1=VALUE1 TAG2=VALUE2\","
448 "\"chart_id\":\"chart_id\",\"chart_name\":\"chart_name\",\"chart_family\":\"(null)\","
449 "\"chart_context\":\"(null)\",\"chart_type\":\"(null)\",\"units\":\"(null)\",\"id\":\"dimension_id\","
@@ -476,9 +458,9 @@ static void test_format_dimension_stored_json_plaintext(void **state)
458 will_return(__wrap_exporting_calculate_value_from_stored_data, pack_storage_number(27, SN_EXISTS));
459
460 RRDDIM *rd = localhost->rrdset_root->dimensions;
479 - assert_int_equal(format_dimension_stored_json_plaintext(engine->connector_root->instance_root, rd), 0);
461 + assert_int_equal(format_dimension_stored_json_plaintext(engine->instance_root, rd), 0);
462 assert_string_equal(
481 - buffer_tostring(engine->connector_root->instance_root->buffer),
463 + buffer_tostring(engine->instance_root->buffer),
464 "{\"prefix\":\"netdata\",\"hostname\":\"test-host\",\"host_tags\":\"TAG1=VALUE1 TAG2=VALUE2\","
465 "\"chart_id\":\"chart_id\",\"chart_name\":\"chart_name\",\"chart_family\":\"(null)\"," \
466 "\"chart_context\": \"(null)\",\"chart_type\":\"(null)\",\"units\": \"(null)\",\"id\":\"dimension_id\","
@@ -490,9 +472,9 @@ static void test_format_dimension_collected_opentsdb_telnet(void **state)
472 struct engine *engine = *state;
473
474 RRDDIM *rd = localhost->rrdset_root->dimensions;
493 - assert_int_equal(format_dimension_collected_opentsdb_telnet(engine->connector_root->instance_root, rd), 0);
475 + assert_int_equal(format_dimension_collected_opentsdb_telnet(engine->instance_root, rd), 0);
476 assert_string_equal(
495 - buffer_tostring(engine->connector_root->instance_root->buffer),
477 + buffer_tostring(engine->instance_root->buffer),
478 "put netdata.chart_name.dimension_name 15051 123000321 host=test-host TAG1=VALUE1 TAG2=VALUE2\n");
479 }
480
@@ -504,9 +486,9 @@ static void test_format_dimension_stored_opentsdb_telnet(void **state)
486 will_return(__wrap_exporting_calculate_value_from_stored_data, pack_storage_number(27, SN_EXISTS));
487
488 RRDDIM *rd = localhost->rrdset_root->dimensions;
507 - assert_int_equal(format_dimension_stored_opentsdb_telnet(engine->connector_root->instance_root, rd), 0);
489 + assert_int_equal(format_dimension_stored_opentsdb_telnet(engine->instance_root, rd), 0);
490 assert_string_equal(
509 - buffer_tostring(engine->connector_root->instance_root->buffer),
491 + buffer_tostring(engine->instance_root->buffer),
492 "put netdata.chart_name.dimension_name 15052 690565856.0000000 host=test-host TAG1=VALUE1 TAG2=VALUE2\n");
493 }
494
@@ -515,9 +497,9 @@ static void test_format_dimension_collected_opentsdb_http(void **state)
497 struct engine *engine = *state;
498
499 RRDDIM *rd = localhost->rrdset_root->dimensions;
518 - assert_int_equal(format_dimension_collected_opentsdb_http(engine->connector_root->instance_root, rd), 0);
500 + assert_int_equal(format_dimension_collected_opentsdb_http(engine->instance_root, rd), 0);
501 assert_string_equal(
520 - buffer_tostring(engine->connector_root->instance_root->buffer),
502 + buffer_tostring(engine->instance_root->buffer),
503 "POST /api/put HTTP/1.1\r\n"
504 "Host: test-host\r\n"
505 "Content-Type: application/json\r\n"
@@ -536,9 +518,9 @@ static void test_format_dimension_stored_opentsdb_http(void **state)
518 will_return(__wrap_exporting_calculate_value_from_stored_data, pack_storage_number(27, SN_EXISTS));
519
520 RRDDIM *rd = localhost->rrdset_root->dimensions;
539 - assert_int_equal(format_dimension_stored_opentsdb_http(engine->connector_root->instance_root, rd), 0);
521 + assert_int_equal(format_dimension_stored_opentsdb_http(engine->instance_root, rd), 0);
522 assert_string_equal(
541 - buffer_tostring(engine->connector_root->instance_root->buffer),
523 + buffer_tostring(engine->instance_root->buffer),
524 "POST /api/put HTTP/1.1\r\n"
525 "Host: test-host\r\n"
526 "Content-Type: application/json\r\n"
@@ -558,7 +540,7 @@ static void test_exporting_discard_response(void **state)
540
541 expect_function_call(__wrap_info_int);
542
561 - assert_int_equal(exporting_discard_response(response, engine->connector_root->instance_root), 0);
543 + assert_int_equal(exporting_discard_response(response, engine->instance_root), 0);
544
545 assert_string_equal(
546 log_line,
@@ -572,7 +554,7 @@ static void test_exporting_discard_response(void **state)
554 static void test_simple_connector_receive_response(void **state)
555 {
556 struct engine *engine = *state;
575 - struct instance *instance = engine->connector_root->instance_root;
557 + struct instance *instance = engine->instance_root;
558 struct stats *stats = &instance->stats;
559
560 int sock = 1;
@@ -599,7 +581,7 @@ static void test_simple_connector_receive_response(void **state)
581 static void test_simple_connector_send_buffer(void **state)
582 {
583 struct engine *engine = *state;
602 - struct instance *instance = engine->connector_root->instance_root;
584 + struct instance *instance = engine->instance_root;
585 struct stats *stats = &instance->stats;
586 BUFFER *buffer = instance->buffer;
587
@@ -645,7 +627,7 @@ static void test_simple_connector_send_buffer(void **state)
627 static void test_simple_connector_worker(void **state)
628 {
629 struct engine *engine = *state;
648 - struct instance *instance = engine->connector_root->instance_root;
630 + struct instance *instance = engine->instance_root;
631 BUFFER *buffer = instance->buffer;
632
633 __real_mark_scheduled_instances(engine);
@@ -721,7 +703,7 @@ static void test_sanitize_opentsdb_label_value(void **state)
703 static void test_format_host_labels_json_plaintext(void **state)
704 {
705 struct engine *engine = *state;
724 - struct instance *instance = engine->connector_root->instance_root;
706 + struct instance *instance = engine->instance_root;
707
708 instance->config.options |= EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
709 instance->config.options |= EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
@@ -733,7 +715,7 @@ static void test_format_host_labels_json_plaintext(void **state)
715 static void test_format_host_labels_graphite_plaintext(void **state)
716 {
717 struct engine *engine = *state;
736 - struct instance *instance = engine->connector_root->instance_root;
718 + struct instance *instance = engine->instance_root;
719
720 instance->config.options |= EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
721 instance->config.options |= EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
@@ -745,7 +727,7 @@ static void test_format_host_labels_graphite_plaintext(void **state)
727 static void test_format_host_labels_opentsdb_telnet(void **state)
728 {
729 struct engine *engine = *state;
748 - struct instance *instance = engine->connector_root->instance_root;
730 + struct instance *instance = engine->instance_root;
731
732 instance->config.options |= EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
733 instance->config.options |= EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
@@ -757,7 +739,7 @@ static void test_format_host_labels_opentsdb_telnet(void **state)
739 static void test_format_host_labels_opentsdb_http(void **state)
740 {
741 struct engine *engine = *state;
760 - struct instance *instance = engine->connector_root->instance_root;
742 + struct instance *instance = engine->instance_root;
743
744 instance->config.options |= EXPORTING_OPTION_SEND_CONFIGURED_LABELS;
745 instance->config.options |= EXPORTING_OPTION_SEND_AUTOMATIC_LABELS;
@@ -769,7 +751,7 @@ static void test_format_host_labels_opentsdb_http(void **state)
751 static void test_flush_host_labels(void **state)
752 {
753 struct engine *engine = *state;
772 - struct instance *instance = engine->connector_root->instance_root;
754 + struct instance *instance = engine->instance_root;
755
756 instance->labels = buffer_create(12);
757 buffer_strcat(instance->labels, "check string");
@@ -779,6 +761,117 @@ static void test_flush_host_labels(void **state)
761 assert_int_equal(buffer_strlen(instance->labels), 0);
762 }
763
764 +#if HAVE_KINESIS
765 +static void test_init_aws_kinesis_instance(void **state)
766 +{
767 + struct engine *engine = *state;
768 + struct instance *instance = engine->instance_root;
769 +
770 + instance->config.options = EXPORTING_SOURCE_DATA_AS_COLLECTED | EXPORTING_OPTION_SEND_NAMES;
771 +
772 + struct aws_kinesis_specific_config *connector_specific_config =
773 + callocz(1, sizeof(struct aws_kinesis_specific_config));
774 + instance->config.connector_specific_config = connector_specific_config;
775 + connector_specific_config->stream_name = strdupz("test_stream");
776 + connector_specific_config->auth_key_id = strdupz("test_auth_key_id");
777 + connector_specific_config->secure_key = strdupz("test_secure_key");
778 +
779 + expect_function_call(__wrap_aws_sdk_init);
780 + expect_function_call(__wrap_kinesis_init);
781 + expect_not_value(__wrap_kinesis_init, kinesis_specific_data_p, NULL);
782 + expect_string(__wrap_kinesis_init, region, "localhost");
783 + expect_string(__wrap_kinesis_init, access_key_id, "test_auth_key_id");
784 + expect_string(__wrap_kinesis_init, secret_key, "test_secure_key");
785 + expect_value(__wrap_kinesis_init, timeout, 10000);
786 + assert_int_equal(init_aws_kinesis_instance(instance), 0);
787 +
788 + assert_ptr_equal(instance->worker, aws_kinesis_connector_worker);
789 + assert_ptr_equal(instance->start_batch_formatting, NULL);
790 + assert_ptr_equal(instance->start_host_formatting, format_host_labels_json_plaintext);
791 + assert_ptr_equal(instance->start_chart_formatting, NULL);
792 + assert_ptr_equal(instance->metric_formatting, format_dimension_collected_json_plaintext);
793 + assert_ptr_equal(instance->end_chart_formatting, NULL);
794 + assert_ptr_equal(instance->end_host_formatting, flush_host_labels);
795 + assert_ptr_equal(instance->end_batch_formatting, NULL);
796 + assert_ptr_not_equal(instance->buffer, NULL);
797 + buffer_free(instance->buffer);
798 + assert_ptr_not_equal(instance->connector_specific_data, NULL);
799 + freez(instance->connector_specific_data);
800 +
801 + instance->config.options = EXPORTING_SOURCE_DATA_AVERAGE | EXPORTING_OPTION_SEND_NAMES;
802 +
803 + expect_function_call(__wrap_kinesis_init);
804 + expect_not_value(__wrap_kinesis_init, kinesis_specific_data_p, NULL);
805 + expect_string(__wrap_kinesis_init, region, "localhost");
806 + expect_string(__wrap_kinesis_init, access_key_id, "test_auth_key_id");
807 + expect_string(__wrap_kinesis_init, secret_key, "test_secure_key");
808 + expect_value(__wrap_kinesis_init, timeout, 10000);
809 +
810 + assert_int_equal(init_aws_kinesis_instance(instance), 0);
811 + assert_ptr_equal(instance->metric_formatting, format_dimension_stored_json_plaintext);
812 +
813 + free(connector_specific_config->stream_name);
814 + free(connector_specific_config->auth_key_id);
815 + free(connector_specific_config->secure_key);
816 +}
817 +
818 +static void test_aws_kinesis_connector_worker(void **state)
819 +{
820 + struct engine *engine = *state;
821 + struct instance *instance = engine->instance_root;
822 + BUFFER *buffer = instance->buffer;
823 +
824 + __real_mark_scheduled_instances(engine);
825 +
826 + expect_function_call(__wrap_rrdhost_is_exportable);
827 + expect_value(__wrap_rrdhost_is_exportable, instance, instance);
828 + expect_value(__wrap_rrdhost_is_exportable, host, localhost);
829 + will_return(__wrap_rrdhost_is_exportable, 1);
830 +
831 + RRDSET *st = localhost->rrdset_root;
832 + expect_function_call(__wrap_rrdset_is_exportable);
833 + expect_value(__wrap_rrdset_is_exportable, instance, instance);
834 + expect_value(__wrap_rrdset_is_exportable, st, st);
835 + will_return(__wrap_rrdset_is_exportable, 1);
836 +
837 + __real_prepare_buffers(engine);
838 +
839 + struct aws_kinesis_specific_config *connector_specific_config =
840 + callocz(1, sizeof(struct aws_kinesis_specific_config));
841 + instance->config.connector_specific_config = connector_specific_config;
842 + connector_specific_config->stream_name = strdupz("test_stream");
843 + connector_specific_config->auth_key_id = strdupz("test_auth_key_id");
844 + connector_specific_config->secure_key = strdupz("test_secure_key");
845 +
846 + struct aws_kinesis_specific_data *connector_specific_data = callocz(1, sizeof(struct aws_kinesis_specific_data));
847 + instance->connector_specific_data = (void *)connector_specific_data;
848 +
849 + expect_function_call(__wrap_kinesis_put_record);
850 + expect_not_value(__wrap_kinesis_put_record, kinesis_specific_data_p, NULL);
851 + expect_string(__wrap_kinesis_put_record, stream_name, "test_stream");
852 + expect_string(__wrap_kinesis_put_record, partition_key, "netdata_0");
853 + expect_value(__wrap_kinesis_put_record, data, buffer_tostring(buffer));
854 + // The buffer is prepated by Graphite exporting connector
855 + expect_string(
856 + __wrap_kinesis_put_record, data,
857 + "netdata.test-host.chart_name.dimension_name;TAG1=VALUE1 TAG2=VALUE2 123000321 15051\n");
858 + expect_value(__wrap_kinesis_put_record, data_len, 84);
859 +
860 + expect_function_call(__wrap_kinesis_get_result);
861 + expect_value(__wrap_kinesis_get_result, request_outcomes_p, NULL);
862 + expect_not_value(__wrap_kinesis_get_result, error_message, NULL);
863 + expect_not_value(__wrap_kinesis_get_result, sent_bytes, NULL);
864 + expect_not_value(__wrap_kinesis_get_result, lost_bytes, NULL);
865 + will_return(__wrap_kinesis_get_result, 0);
866 +
867 + aws_kinesis_connector_worker(instance);
868 +
869 + free(connector_specific_config->stream_name);
870 + free(connector_specific_config->auth_key_id);
871 + free(connector_specific_config->secure_key);
872 +}
873 +#endif // HAVE_KINESIS
874 +
875 int main(void)
876 {
877 const struct CMUnitTest tests[] = {
@@ -789,8 +882,6 @@ int main(void)
882 test_init_graphite_instance, setup_configured_engine, teardown_configured_engine),
883 cmocka_unit_test_setup_teardown(
884 test_init_json_instance, setup_configured_engine, teardown_configured_engine),
792 - cmocka_unit_test_setup_teardown(
793 - test_init_opentsdb_connector, setup_configured_engine, teardown_configured_engine),
885 cmocka_unit_test_setup_teardown(
886 test_init_opentsdb_telnet_instance, setup_configured_engine, teardown_configured_engine),
887 cmocka_unit_test_setup_teardown(
@@ -850,6 +941,21 @@ int main(void)
941 cmocka_unit_test_setup_teardown(test_flush_host_labels, setup_initialized_engine, teardown_initialized_engine),
942 };
943
853 - return cmocka_run_group_tests_name("exporting_engine", tests, NULL, NULL) +
854 - cmocka_run_group_tests_name("labels_in_exporting_engine", label_tests, NULL, NULL);
944 +#if HAVE_KINESIS
945 + const struct CMUnitTest kinesis_tests[] = {
946 + cmocka_unit_test_setup_teardown(
947 + test_init_aws_kinesis_instance, setup_configured_engine, teardown_configured_engine),
948 + cmocka_unit_test_setup_teardown(
949 + test_aws_kinesis_connector_worker, setup_initialized_engine, teardown_initialized_engine),
950 + };
951 +#endif
952 +
953 + int test_res = cmocka_run_group_tests_name("exporting_engine", tests, NULL, NULL) +
954 + cmocka_run_group_tests_name("labels_in_exporting_engine", label_tests, NULL, NULL);
955 +
956 +#if HAVE_KINESIS
957 + test_res += cmocka_run_group_tests_name("kinesis_exporting_connector", kinesis_tests, NULL, NULL);
958 +#endif
959 +
960 + return test_res;
961 }
exporting/tests/test_exporting_engine.h
+10
@@ -9,6 +9,7 @@
9 #include "exporting/graphite/graphite.h"
10 #include "exporting/json/json.h"
11 #include "exporting/opentsdb/opentsdb.h"
12 +#include "exporting/aws_kinesis/aws_kinesis.h"
13
14 #include <stdarg.h>
15 #include <stddef.h>
@@ -95,6 +96,15 @@ int __mock_end_chart_formatting(struct instance *instance, RRDSET *st);
96 int __mock_end_host_formatting(struct instance *instance, RRDHOST *host);
97 int __mock_end_batch_formatting(struct instance *instance);
98
99 +void __wrap_aws_sdk_init();
100 +void __wrap_kinesis_init(
101 + void *kinesis_specific_data_p, const char *region, const char *access_key_id, const char *secret_key,
102 + const long timeout);
103 +void __wrap_kinesis_put_record(
104 + void *kinesis_specific_data_p, const char *stream_name, const char *partition_key, const char *data,
105 + size_t data_len);
106 +int __wrap_kinesis_get_result(void *request_outcomes_p, char *error_message, size_t *sent_bytes, size_t *lost_bytes);
107 +
108 // -----------------------------------------------------------------------
109 // fixtures
110