master
c 220 lines 7.74 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #include "aws_kinesis.h"
4
5 /**
6 * Clean AWS Kinesis *
7 */
8 void aws_kinesis_cleanup(struct instance *instance)
9 {
10 netdata_log_info("EXPORTING: cleaning up instance %s ...", instance->config.name);
11 kinesis_shutdown(instance->connector_specific_data);
12
13 freez(instance->connector_specific_data);
14
15 struct aws_kinesis_specific_config *connector_specific_config = instance->config.connector_specific_config;
16 if (connector_specific_config) {
17 freez(connector_specific_config->auth_key_id);
18 freez(connector_specific_config->secure_key);
19 freez(connector_specific_config->stream_name);
20
21 freez(connector_specific_config);
22 }
23
24 netdata_log_info("EXPORTING: instance %s exited", instance->config.name);
25 instance->exited = 1;
26 }
27
28 /**
29 * Initialize AWS Kinesis connector instance
30 *
31 * @param instance an instance data structure.
32 * @return Returns 0 on success, 1 on failure.
33 */
34 int init_aws_kinesis_instance(struct instance *instance)
35 {
36 instance->worker = aws_kinesis_connector_worker;
37
38 instance->start_batch_formatting = NULL;
39 instance->start_host_formatting = format_host_labels_json_plaintext;
40 instance->start_chart_formatting = NULL;
41
42 if (EXPORTING_OPTIONS_DATA_SOURCE(instance->config.options) == EXPORTING_SOURCE_DATA_AS_COLLECTED)
43 instance->metric_formatting = format_dimension_collected_json_plaintext;
44 else
45 instance->metric_formatting = format_dimension_stored_json_plaintext;
46
47 instance->end_chart_formatting = NULL;
48 instance->variables_formatting = NULL;
49 instance->end_host_formatting = flush_host_labels;
50 instance->end_batch_formatting = NULL;
51
52 instance->prepare_header = NULL;
53 instance->check_response = NULL;
54
55 instance->buffer = (void *)buffer_create(0, &netdata_buffers_statistics.buffers_exporters);
56 if (!instance->buffer) {
57 netdata_log_error("EXPORTING: cannot create buffer for AWS Kinesis exporting connector instance %s",
58 instance->config.name);
59 return 1;
60 }
61 if (netdata_mutex_init(&instance->mutex))
62 return 1;
63 if (netdata_cond_init(&instance->cond_var))
64 return 1;
65
66 if (!instance->engine->aws_sdk_initialized) {
67 aws_sdk_init();
68 instance->engine->aws_sdk_initialized = 1;
69 }
70
71 struct aws_kinesis_specific_config *connector_specific_config = instance->config.connector_specific_config;
72 struct aws_kinesis_specific_data *connector_specific_data = callocz(1, sizeof(struct aws_kinesis_specific_data));
73 instance->connector_specific_data = (void *)connector_specific_data;
74
75 if (!strcmp(connector_specific_config->stream_name, "")) {
76 netdata_log_error("stream name is a mandatory Kinesis parameter but it is not configured");
77 return 1;
78 }
79
80 kinesis_init(
81 (void *)connector_specific_data,
82 instance->config.destination,
83 connector_specific_config->auth_key_id,
84 connector_specific_config->secure_key,
85 instance->config.timeoutms);
86
87 return 0;
88 }
89
90 /**
91 * AWS Kinesis connector worker
92 *
93 * Runs in a separate thread for every instance.
94 *
95 * @param instance_p an instance data structure.
96 */
97 void *aws_kinesis_connector_worker(void *instance_p)
98 {
99 struct instance *instance = (struct instance *)instance_p;
100 struct aws_kinesis_specific_config *connector_specific_config = instance->config.connector_specific_config;
101 struct aws_kinesis_specific_data *connector_specific_data = instance->connector_specific_data;
102
103 while (!instance->engine->exit) {
104 unsigned long long partition_key_seq = 0;
105 struct stats *stats = &instance->stats;
106
107 netdata_mutex_lock(&instance->mutex);
108 while (!instance->data_is_ready)
109 netdata_cond_wait(&instance->cond_var, &instance->mutex);
110 instance->data_is_ready = 0;
111
112 if (unlikely(instance->engine->exit)) {
113 netdata_mutex_unlock(&instance->mutex);
114 break;
115 }
116
117 // reset the monitoring chart counters
118 stats->received_bytes =
119 stats->sent_bytes =
120 stats->sent_metrics =
121 stats->lost_metrics =
122 stats->receptions =
123 stats->transmission_successes =
124 stats->transmission_failures =
125 stats->data_lost_events =
126 stats->lost_bytes =
127 stats->reconnects = 0;
128
129 BUFFER *buffer = (BUFFER *)instance->buffer;
130 size_t buffer_len = buffer_strlen(buffer);
131
132 stats->buffered_bytes = buffer_len;
133
134 size_t sent = 0;
135
136 while (sent < buffer_len) {
137 char partition_key[KINESIS_PARTITION_KEY_MAX + 1];
138 snprintf(partition_key, KINESIS_PARTITION_KEY_MAX, "netdata_%llu", partition_key_seq++);
139 size_t partition_key_len = strnlen(partition_key, KINESIS_PARTITION_KEY_MAX);
140
141 const char *first_char = buffer_tostring(buffer) + sent;
142
143 size_t record_len = 0;
144
145 // split buffer into chunks of maximum allowed size
146 if (buffer_len - sent < KINESIS_RECORD_MAX - partition_key_len) {
147 record_len = buffer_len - sent;
148 } else {
149 record_len = KINESIS_RECORD_MAX - partition_key_len;
150 while (record_len && *(first_char + record_len - 1) != '\n')
151 record_len--;
152 }
153 char error_message[ERROR_LINE_MAX + 1] = "";
154
155 netdata_log_debug(D_EXPORTING,
156 "EXPORTING: kinesis_put_record(): dest = %s, id = %s, key = %s, stream = %s, partition_key = %s, "
157 "buffer = %zu, record = %zu",
158 instance->config.destination,
159 connector_specific_config->auth_key_id,
160 connector_specific_config->secure_key,
161 connector_specific_config->stream_name,
162 partition_key,
163 buffer_len,
164 record_len);
165
166 kinesis_put_record(
167 connector_specific_data, connector_specific_config->stream_name, partition_key, first_char, record_len);
168
169 sent += record_len;
170 stats->transmission_successes++;
171
172 size_t sent_bytes = 0, lost_bytes = 0;
173
174 if (unlikely(kinesis_get_result(
175 connector_specific_data->request_outcomes, error_message, &sent_bytes, &lost_bytes))) {
176 // oops! we couldn't send (all or some of the) data
177 netdata_log_error("EXPORTING: %s", error_message);
178 netdata_log_error("EXPORTING: failed to write data to external database '%s'. Willing to write %zu bytes, wrote %zu bytes.",
179 instance->config.destination,
180 sent_bytes,
181 sent_bytes - lost_bytes);
182
183 stats->transmission_failures++;
184 stats->data_lost_events++;
185 stats->lost_bytes += lost_bytes;
186
187 // estimate the number of lost metrics
188 stats->lost_metrics += (collected_number)(
189 stats->buffered_metrics *
190 (buffer_len && (lost_bytes > buffer_len) ? (double)lost_bytes / buffer_len : 1));
191
192 break;
193 } else {
194 stats->receptions++;
195 }
196
197 if (unlikely(instance->engine->exit))
198 break;
199 }
200
201 stats->sent_bytes += sent;
202 if (likely(sent == buffer_len))
203 stats->sent_metrics = stats->buffered_metrics;
204
205 buffer_flush(buffer);
206
207 send_internal_metrics(instance);
208
209 stats->buffered_metrics = 0;
210
211 netdata_mutex_unlock(&instance->mutex);
212
213 #ifdef UNIT_TESTING
214 break;
215 #endif
216 }
217
218 aws_kinesis_cleanup(instance);
219 return NULL;
220 }