/src/fluent-bit/plugins/out_kinesis_firehose/firehose_api.c
Line | Count | Source |
1 | | /* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */ |
2 | | |
3 | | /* Fluent Bit |
4 | | * ========== |
5 | | * Copyright (C) 2015-2026 The Fluent Bit Authors |
6 | | * |
7 | | * Licensed under the Apache License, Version 2.0 (the "License"); |
8 | | * you may not use this file except in compliance with the License. |
9 | | * You may obtain a copy of the License at |
10 | | * |
11 | | * http://www.apache.org/licenses/LICENSE-2.0 |
12 | | * |
13 | | * Unless required by applicable law or agreed to in writing, software |
14 | | * distributed under the License is distributed on an "AS IS" BASIS, |
15 | | * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
16 | | * See the License for the specific language governing permissions and |
17 | | * limitations under the License. |
18 | | */ |
19 | | |
20 | | #include <fluent-bit/flb_compat.h> |
21 | | #include <fluent-bit/flb_info.h> |
22 | | #include <fluent-bit/flb_output.h> |
23 | | #include <fluent-bit/flb_utils.h> |
24 | | #include <fluent-bit/flb_slist.h> |
25 | | #include <fluent-bit/flb_time.h> |
26 | | #include <fluent-bit/flb_pack.h> |
27 | | #include <fluent-bit/flb_macros.h> |
28 | | #include <fluent-bit/flb_config_map.h> |
29 | | #include <fluent-bit/flb_output_plugin.h> |
30 | | #include <fluent-bit/flb_log_event_decoder.h> |
31 | | |
32 | | #include <fluent-bit/flb_sds.h> |
33 | | #include <fluent-bit/flb_aws_credentials.h> |
34 | | #include <fluent-bit/flb_aws_util.h> |
35 | | #include <fluent-bit/flb_mem.h> |
36 | | #include <fluent-bit/flb_http_client.h> |
37 | | #include <fluent-bit/flb_utils.h> |
38 | | |
39 | | #include <fluent-bit/flb_base64.h> |
40 | | #include <fluent-bit/aws/flb_aws_compress.h> |
41 | | #include <fluent-bit/aws/flb_aws_aggregation.h> |
42 | | |
43 | | #include <monkey/mk_core.h> |
44 | | #include <msgpack.h> |
45 | | #include <string.h> |
46 | | #include <stdio.h> |
47 | | #include <stdlib.h> |
48 | | |
49 | | #ifndef FLB_SYSTEM_WINDOWS |
50 | | #include <unistd.h> |
51 | | #endif |
52 | | |
53 | | #include "firehose_api.h" |
54 | | |
55 | 0 | #define ERR_CODE_SERVICE_UNAVAILABLE "ServiceUnavailableException" |
56 | | |
57 | | /* Forward declarations */ |
58 | | static int send_log_events(struct flb_firehose *ctx, struct flush *buf); |
59 | | |
60 | | static struct flb_aws_header put_record_batch_header = { |
61 | | .key = "X-Amz-Target", |
62 | | .key_len = 12, |
63 | | .val = "Firehose_20150804.PutRecordBatch", |
64 | | .val_len = 32, |
65 | | }; |
66 | | |
67 | | static inline int try_to_write(char *buf, int *off, size_t left, |
68 | | const char *str, size_t str_len) |
69 | 0 | { |
70 | 0 | if (str_len <= 0){ |
71 | 0 | str_len = strlen(str); |
72 | 0 | } |
73 | 0 | if (left <= *off+str_len) { |
74 | 0 | return FLB_FALSE; |
75 | 0 | } |
76 | 0 | memcpy(buf+*off, str, str_len); |
77 | 0 | *off += str_len; |
78 | 0 | return FLB_TRUE; |
79 | 0 | } |
80 | | |
81 | | /* |
82 | | * Writes the "header" for a put_record_batch payload |
83 | | */ |
84 | | static int init_put_payload(struct flb_firehose *ctx, struct flush *buf, |
85 | | int *offset) |
86 | 0 | { |
87 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
88 | 0 | "{\"DeliveryStreamName\":\"", 23)) { |
89 | 0 | goto error; |
90 | 0 | } |
91 | | |
92 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
93 | 0 | ctx->delivery_stream, 0)) { |
94 | 0 | goto error; |
95 | 0 | } |
96 | | |
97 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
98 | 0 | "\",\"Records\":[", 13)) { |
99 | 0 | goto error; |
100 | 0 | } |
101 | 0 | return 0; |
102 | | |
103 | 0 | error: |
104 | 0 | return -1; |
105 | 0 | } |
106 | | |
107 | | /* |
108 | | * Writes a log event to the output buffer |
109 | | */ |
110 | | static int write_event(struct flb_firehose *ctx, struct flush *buf, |
111 | | struct firehose_event *event, int *offset) |
112 | 0 | { |
113 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
114 | 0 | "{\"Data\":\"", 9)) { |
115 | 0 | goto error; |
116 | 0 | } |
117 | | |
118 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
119 | 0 | event->json, event->len)) { |
120 | 0 | goto error; |
121 | 0 | } |
122 | | |
123 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
124 | 0 | "\"}", 2)) { |
125 | 0 | goto error; |
126 | 0 | } |
127 | | |
128 | 0 | return 0; |
129 | | |
130 | 0 | error: |
131 | 0 | return -1; |
132 | 0 | } |
133 | | |
134 | | /* Terminates a PutRecordBatch payload */ |
135 | | static int end_put_payload(struct flb_firehose *ctx, struct flush *buf, |
136 | | int *offset) |
137 | 0 | { |
138 | 0 | if (!try_to_write(buf->out_buf, offset, buf->out_buf_size, |
139 | 0 | "]}", 2)) { |
140 | 0 | return -1; |
141 | 0 | } |
142 | 0 | buf->out_buf[*offset] = '\0'; |
143 | |
|
144 | 0 | return 0; |
145 | 0 | } |
146 | | |
147 | | |
148 | | /* |
149 | | * Process event with simple aggregation (Firehose version) |
150 | | * Uses shared aggregation implementation |
151 | | */ |
152 | | static int process_event_simple_aggregation(struct flb_firehose *ctx, struct flush *buf, |
153 | | const msgpack_object *obj, struct flb_time *tms, |
154 | | struct flb_config *config) |
155 | 0 | { |
156 | 0 | return flb_aws_aggregation_process_event(&buf->agg_buf, |
157 | 0 | buf->tmp_buf, |
158 | 0 | buf->tmp_buf_size, |
159 | 0 | &buf->tmp_buf_offset, |
160 | 0 | obj, |
161 | 0 | tms, |
162 | 0 | config, |
163 | 0 | ctx->ins, |
164 | 0 | ctx->delivery_stream, |
165 | 0 | ctx->log_key, |
166 | 0 | ctx->time_key, |
167 | 0 | ctx->time_key_format, |
168 | 0 | MAX_EVENT_SIZE); |
169 | 0 | } |
170 | | |
171 | | /* |
172 | | * Processes the msgpack object |
173 | | * -1 = failure, record not added |
174 | | * 0 = success, record added |
175 | | * 1 = we ran out of space, send and retry |
176 | | * 2 = record could not be processed, discard it |
177 | | * Returns 0 on success, -1 on general errors, |
178 | | * and 1 if we ran out of space to write the event |
179 | | * which means a send must occur |
180 | | */ |
181 | | static int process_event(struct flb_firehose *ctx, struct flush *buf, |
182 | | const msgpack_object *obj, struct flb_time *tms, |
183 | | struct flb_config *config) |
184 | 0 | { |
185 | 0 | size_t written = 0; |
186 | 0 | int ret; |
187 | 0 | size_t size; |
188 | 0 | size_t b64_len; |
189 | 0 | struct firehose_event *event; |
190 | 0 | char *tmp_buf_ptr; |
191 | 0 | char *time_key_ptr; |
192 | 0 | struct tm time_stamp; |
193 | 0 | struct tm *tmp; |
194 | 0 | size_t len; |
195 | 0 | size_t tmp_size; |
196 | 0 | void *compressed_tmp_buf; |
197 | 0 | char *out_buf; |
198 | |
|
199 | 0 | tmp_buf_ptr = buf->tmp_buf + buf->tmp_buf_offset; |
200 | 0 | ret = flb_msgpack_to_json(tmp_buf_ptr, |
201 | 0 | buf->tmp_buf_size - buf->tmp_buf_offset, |
202 | 0 | obj, config->json_escape_unicode); |
203 | 0 | if (ret <= 0) { |
204 | | /* |
205 | | * negative value means failure to write to buffer, |
206 | | * which means we ran out of space, and must send the logs |
207 | | * |
208 | | * TODO: This could also incorrectly be triggered if the record |
209 | | * is larger than MAX_EVENT_SIZE |
210 | | */ |
211 | 0 | return 1; |
212 | 0 | } |
213 | 0 | written = (size_t) ret; |
214 | | |
215 | | /* Discard empty messages (written == 2 means '""') */ |
216 | 0 | if (written <= 2) { |
217 | 0 | flb_plg_debug(ctx->ins, "Found empty log message, %s", ctx->delivery_stream); |
218 | 0 | return 2; |
219 | 0 | } |
220 | | |
221 | 0 | if (ctx->log_key) { |
222 | | /* |
223 | | * flb_msgpack_to_json will encase the value in quotes |
224 | | * We don't want that for log_key, so we ignore the first |
225 | | * and last character |
226 | | */ |
227 | 0 | written -= 2; |
228 | 0 | tmp_buf_ptr++; /* pass over the opening quote */ |
229 | 0 | buf->tmp_buf_offset++; |
230 | 0 | } |
231 | | |
232 | | /* is (written + 1) because we still have to append newline */ |
233 | 0 | if ((written + 1) >= MAX_EVENT_SIZE) { |
234 | 0 | flb_plg_warn(ctx->ins, "[size=%zu] Discarding record which is larger than " |
235 | 0 | "max size allowed by Firehose, %s", written + 1, |
236 | 0 | ctx->delivery_stream); |
237 | 0 | return 2; |
238 | 0 | } |
239 | | |
240 | 0 | if (ctx->time_key) { |
241 | | /* append time_key to end of json string */ |
242 | 0 | tmp = gmtime_r(&tms->tm.tv_sec, &time_stamp); |
243 | 0 | if (!tmp) { |
244 | 0 | flb_plg_error(ctx->ins, "Could not create time stamp for %lu unix " |
245 | 0 | "seconds, discarding record, %s", tms->tm.tv_sec, |
246 | 0 | ctx->delivery_stream); |
247 | 0 | return 2; |
248 | 0 | } |
249 | | |
250 | | /* format time output and return the length */ |
251 | 0 | len = flb_aws_strftime_precision(&out_buf, ctx->time_key_format, tms); |
252 | | |
253 | | /* how much space do we have left */ |
254 | 0 | tmp_size = (buf->tmp_buf_size - buf->tmp_buf_offset) - written; |
255 | 0 | if (len > tmp_size) { |
256 | | /* not enough space - tell caller to retry */ |
257 | 0 | flb_free(out_buf); |
258 | 0 | return 1; |
259 | 0 | } |
260 | | |
261 | 0 | if (len == 0) { |
262 | | /* |
263 | | * when the length of out_buf is not enough for time_key_format, |
264 | | * time_key will not be added to record. |
265 | | */ |
266 | 0 | flb_plg_error(ctx->ins, "Failed to add time_key %s to record, %s", |
267 | 0 | ctx->time_key, ctx->delivery_stream); |
268 | 0 | flb_free(out_buf); |
269 | 0 | } |
270 | 0 | else { |
271 | 0 | time_key_ptr = tmp_buf_ptr + written - 1; |
272 | 0 | memcpy(time_key_ptr, ",", 1); |
273 | 0 | time_key_ptr++; |
274 | 0 | memcpy(time_key_ptr, "\"", 1); |
275 | 0 | time_key_ptr++; |
276 | 0 | memcpy(time_key_ptr, ctx->time_key, strlen(ctx->time_key)); |
277 | 0 | time_key_ptr += strlen(ctx->time_key); |
278 | 0 | memcpy(time_key_ptr, "\":\"", 3); |
279 | 0 | time_key_ptr += 3; |
280 | 0 | tmp_size = buf->tmp_buf_size - buf->tmp_buf_offset; |
281 | 0 | tmp_size -= (time_key_ptr - tmp_buf_ptr); |
282 | | |
283 | | /* merge out_buf to time_key_ptr */ |
284 | 0 | memcpy(time_key_ptr, out_buf, len); |
285 | 0 | flb_free(out_buf); |
286 | 0 | time_key_ptr += len; |
287 | 0 | memcpy(time_key_ptr, "\"}", 2); |
288 | 0 | time_key_ptr += 2; |
289 | 0 | written = (time_key_ptr - tmp_buf_ptr); |
290 | 0 | } |
291 | 0 | } |
292 | | |
293 | | /* is (written + 1) because we still have to append newline */ |
294 | 0 | if ((written + 1) >= MAX_EVENT_SIZE) { |
295 | 0 | flb_plg_warn(ctx->ins, "[size=%zu] Discarding record which is larger than " |
296 | 0 | "max size allowed by Firehose, %s", written + 1, |
297 | 0 | ctx->delivery_stream); |
298 | 0 | return 2; |
299 | 0 | } |
300 | | |
301 | | /* append newline to record */ |
302 | | |
303 | 0 | tmp_size = (buf->tmp_buf_size - buf->tmp_buf_offset) - written; |
304 | 0 | if (tmp_size <= 1) { |
305 | | /* no space left- tell caller to retry */ |
306 | 0 | return 1; |
307 | 0 | } |
308 | | |
309 | 0 | memcpy(tmp_buf_ptr + written, "\n", 1); |
310 | 0 | written++; |
311 | |
|
312 | 0 | if (ctx->compression == FLB_AWS_COMPRESS_NONE) { |
313 | | /* |
314 | | * check if event_buf is initialized and big enough |
315 | | * Base64 encoding will increase size by ~4/3 |
316 | | */ |
317 | 0 | size = (written * 1.5) + 4; |
318 | 0 | if (buf->event_buf == NULL || buf->event_buf_size < size) { |
319 | 0 | flb_free(buf->event_buf); |
320 | 0 | buf->event_buf = flb_malloc(size); |
321 | 0 | buf->event_buf_size = size; |
322 | 0 | if (buf->event_buf == NULL) { |
323 | 0 | flb_errno(); |
324 | 0 | return -1; |
325 | 0 | } |
326 | 0 | } |
327 | | |
328 | 0 | tmp_buf_ptr = buf->tmp_buf + buf->tmp_buf_offset; |
329 | | |
330 | 0 | ret = flb_base64_encode((unsigned char *) buf->event_buf, size, &b64_len, |
331 | 0 | (unsigned char *) tmp_buf_ptr, written); |
332 | 0 | if (ret != 0) { |
333 | 0 | flb_errno(); |
334 | 0 | return -1; |
335 | 0 | } |
336 | 0 | written = b64_len; |
337 | 0 | } |
338 | 0 | else { |
339 | | /* |
340 | | * compress event, truncating input if needed |
341 | | * replace event buffer with compressed buffer |
342 | | */ |
343 | 0 | ret = flb_aws_compression_b64_truncate_compress(ctx->compression, |
344 | 0 | MAX_B64_EVENT_SIZE, |
345 | 0 | tmp_buf_ptr, |
346 | 0 | written, &compressed_tmp_buf, |
347 | 0 | &size); /* evaluate size */ |
348 | 0 | if (ret == -1) { |
349 | 0 | flb_plg_error(ctx->ins, "Unable to compress record, discarding, " |
350 | 0 | "%s", ctx->delivery_stream); |
351 | 0 | return 2; |
352 | 0 | } |
353 | 0 | flb_free(buf->event_buf); |
354 | 0 | buf->event_buf = compressed_tmp_buf; |
355 | 0 | compressed_tmp_buf = NULL; |
356 | 0 | written = size; |
357 | 0 | } |
358 | | |
359 | 0 | tmp_buf_ptr = buf->tmp_buf + buf->tmp_buf_offset; |
360 | 0 | if ((buf->tmp_buf_size - buf->tmp_buf_offset) < written) { |
361 | | /* not enough space, send logs */ |
362 | 0 | return 1; |
363 | 0 | } |
364 | | |
365 | | /* copy serialized json to tmp_buf */ |
366 | 0 | memcpy(tmp_buf_ptr, buf->event_buf, written); |
367 | |
|
368 | 0 | buf->tmp_buf_offset += written; |
369 | 0 | event = &buf->events[buf->event_index]; |
370 | 0 | event->json = tmp_buf_ptr; |
371 | 0 | event->len = written; |
372 | 0 | event->timestamp.tv_sec = tms->tm.tv_sec; |
373 | 0 | event->timestamp.tv_nsec = tms->tm.tv_nsec; |
374 | |
|
375 | 0 | return 0; |
376 | 0 | } |
377 | | |
378 | | /* Resets or inits a flush struct */ |
379 | 0 | static void reset_flush_buf(struct flb_firehose *ctx, struct flush *buf) { |
380 | 0 | buf->event_index = 0; |
381 | 0 | buf->tmp_buf_offset = 0; |
382 | 0 | buf->data_size = PUT_RECORD_BATCH_HEADER_LEN + PUT_RECORD_BATCH_FOOTER_LEN; |
383 | 0 | buf->data_size += strlen(ctx->delivery_stream); |
384 | 0 | } |
385 | | |
386 | | /* Finalize and send aggregated record */ |
387 | 0 | static int send_aggregated_record(struct flb_firehose *ctx, struct flush *buf) { |
388 | 0 | int ret; |
389 | 0 | size_t agg_size; |
390 | 0 | size_t b64_len; |
391 | 0 | struct firehose_event *event; |
392 | 0 | void *compressed_tmp_buf; |
393 | 0 | size_t compressed_size; |
394 | | |
395 | | /* Finalize the aggregated record */ |
396 | 0 | ret = flb_aws_aggregation_finalize(&buf->agg_buf, 1, &agg_size); |
397 | 0 | if (ret < 0) { |
398 | | /* No data to finalize */ |
399 | 0 | return 0; |
400 | 0 | } |
401 | | |
402 | | /* Handle compression if enabled */ |
403 | 0 | if (ctx->compression != FLB_AWS_COMPRESS_NONE) { |
404 | 0 | ret = flb_aws_compression_b64_truncate_compress(ctx->compression, |
405 | 0 | MAX_B64_EVENT_SIZE, |
406 | 0 | buf->agg_buf.agg_buf, |
407 | 0 | agg_size, |
408 | 0 | &compressed_tmp_buf, |
409 | 0 | &compressed_size); |
410 | 0 | if (ret == -1) { |
411 | 0 | flb_plg_error(ctx->ins, "Unable to compress aggregated record, discarding, %s", |
412 | 0 | ctx->delivery_stream); |
413 | 0 | flb_aws_aggregation_reset(&buf->agg_buf); |
414 | 0 | return 0; |
415 | 0 | } |
416 | | |
417 | | /* Ensure event_buf is large enough */ |
418 | 0 | if (buf->event_buf == NULL || buf->event_buf_size < compressed_size) { |
419 | 0 | flb_free(buf->event_buf); |
420 | 0 | buf->event_buf = compressed_tmp_buf; |
421 | 0 | buf->event_buf_size = compressed_size; |
422 | 0 | compressed_tmp_buf = NULL; |
423 | 0 | } else { |
424 | 0 | memcpy(buf->event_buf, compressed_tmp_buf, compressed_size); |
425 | 0 | flb_free(compressed_tmp_buf); |
426 | 0 | } |
427 | 0 | agg_size = compressed_size; |
428 | 0 | } |
429 | 0 | else { |
430 | | /* Base64 encode the aggregated record */ |
431 | 0 | size_t size = (agg_size * 1.5) + 4; |
432 | 0 | if (buf->event_buf == NULL || buf->event_buf_size < size) { |
433 | 0 | flb_free(buf->event_buf); |
434 | 0 | buf->event_buf = flb_malloc(size); |
435 | 0 | buf->event_buf_size = size; |
436 | 0 | if (buf->event_buf == NULL) { |
437 | 0 | flb_errno(); |
438 | 0 | return -1; |
439 | 0 | } |
440 | 0 | } |
441 | | |
442 | 0 | ret = flb_base64_encode((unsigned char *) buf->event_buf, size, &b64_len, |
443 | 0 | (unsigned char *) buf->agg_buf.agg_buf, agg_size); |
444 | 0 | if (ret != 0) { |
445 | 0 | flb_errno(); |
446 | 0 | return -1; |
447 | 0 | } |
448 | 0 | agg_size = b64_len; |
449 | 0 | } |
450 | | |
451 | | /* Copy to tmp_buf for sending */ |
452 | 0 | if (buf->tmp_buf_size < agg_size) { |
453 | 0 | flb_plg_error(ctx->ins, "Aggregated record too large for buffer"); |
454 | 0 | flb_aws_aggregation_reset(&buf->agg_buf); |
455 | 0 | return 0; |
456 | 0 | } |
457 | | |
458 | 0 | memcpy(buf->tmp_buf, buf->event_buf, agg_size); |
459 | | |
460 | | /* Create event record */ |
461 | 0 | event = &buf->events[0]; |
462 | 0 | event->json = buf->tmp_buf; |
463 | 0 | event->len = agg_size; |
464 | 0 | event->timestamp.tv_sec = 0; |
465 | 0 | event->timestamp.tv_nsec = 0; |
466 | 0 | buf->event_index = 1; |
467 | | |
468 | | /* Calculate data_size for the payload */ |
469 | 0 | buf->data_size = PUT_RECORD_BATCH_HEADER_LEN + PUT_RECORD_BATCH_FOOTER_LEN; |
470 | 0 | buf->data_size += strlen(ctx->delivery_stream); |
471 | 0 | buf->data_size += agg_size + PUT_RECORD_BATCH_PER_RECORD_LEN; |
472 | | |
473 | | /* Send the aggregated record */ |
474 | 0 | ret = send_log_events(ctx, buf); |
475 | | |
476 | | /* Reset aggregation buffer */ |
477 | 0 | flb_aws_aggregation_reset(&buf->agg_buf); |
478 | |
|
479 | 0 | return ret; |
480 | 0 | } |
481 | | |
482 | | /* constructs a put payload, and then sends */ |
483 | 0 | static int send_log_events(struct flb_firehose *ctx, struct flush *buf) { |
484 | 0 | int ret; |
485 | 0 | int offset; |
486 | 0 | int i; |
487 | 0 | struct firehose_event *event; |
488 | |
|
489 | 0 | if (buf->event_index <= 0) { |
490 | | /* |
491 | | * event_index should always be 1 more than the actual last event index |
492 | | * when this function is called. |
493 | | * Except in the case where send_log_events() is called at the end of |
494 | | * process_and_send. If all records were already sent, event_index |
495 | | * will be 0. Hence this check. |
496 | | */ |
497 | 0 | return 0; |
498 | 0 | } |
499 | | |
500 | | /* alloc out_buf if needed */ |
501 | 0 | if (buf->out_buf == NULL || buf->out_buf_size < buf->data_size) { |
502 | 0 | if (buf->out_buf != NULL) { |
503 | 0 | flb_free(buf->out_buf); |
504 | 0 | } |
505 | 0 | buf->out_buf = flb_malloc(buf->data_size + 1); |
506 | 0 | if (!buf->out_buf) { |
507 | 0 | flb_errno(); |
508 | 0 | return -1; |
509 | 0 | } |
510 | 0 | buf->out_buf_size = buf->data_size; |
511 | 0 | } |
512 | | |
513 | 0 | offset = 0; |
514 | 0 | ret = init_put_payload(ctx, buf, &offset); |
515 | 0 | if (ret < 0) { |
516 | 0 | flb_plg_error(ctx->ins, "Failed to initialize PutRecordBatch payload, %s", |
517 | 0 | ctx->delivery_stream); |
518 | 0 | return -1; |
519 | 0 | } |
520 | | |
521 | 0 | for (i = 0; i < buf->event_index; i++) { |
522 | 0 | event = &buf->events[i]; |
523 | 0 | ret = write_event(ctx, buf, event, &offset); |
524 | 0 | if (ret < 0) { |
525 | 0 | flb_plg_error(ctx->ins, "Failed to write log record %d to " |
526 | 0 | "payload buffer, %s", i, ctx->delivery_stream); |
527 | 0 | return -1; |
528 | 0 | } |
529 | 0 | if (i != (buf->event_index -1)) { |
530 | 0 | if (!try_to_write(buf->out_buf, &offset, buf->out_buf_size, |
531 | 0 | ",", 1)) { |
532 | 0 | flb_plg_error(ctx->ins, "Could not terminate record with ','"); |
533 | 0 | return -1; |
534 | 0 | } |
535 | 0 | } |
536 | 0 | } |
537 | | |
538 | 0 | ret = end_put_payload(ctx, buf, &offset); |
539 | 0 | if (ret < 0) { |
540 | 0 | flb_plg_error(ctx->ins, "Could not complete PutRecordBatch payload"); |
541 | 0 | return -1; |
542 | 0 | } |
543 | 0 | flb_plg_debug(ctx->ins, "firehose:PutRecordBatch: events=%d, payload=%d bytes", i, offset); |
544 | 0 | ret = put_record_batch(ctx, buf, (size_t) offset, i); |
545 | 0 | if (ret < 0) { |
546 | 0 | flb_plg_error(ctx->ins, "Failed to send log records"); |
547 | 0 | return -1; |
548 | 0 | } |
549 | 0 | buf->records_sent += i; |
550 | |
|
551 | 0 | return 0; |
552 | 0 | } |
553 | | |
554 | | /* |
555 | | * Processes the msgpack object, sends the current batch if needed |
556 | | */ |
557 | | static int add_event(struct flb_firehose *ctx, struct flush *buf, |
558 | | const msgpack_object *obj, struct flb_time *tms, |
559 | | struct flb_config *config) |
560 | 0 | { |
561 | 0 | int ret; |
562 | 0 | struct firehose_event *event; |
563 | 0 | int retry_add = FLB_FALSE; |
564 | 0 | size_t event_bytes = 0; |
565 | |
|
566 | 0 | if (buf->event_index == 0) { |
567 | | /* init */ |
568 | 0 | reset_flush_buf(ctx, buf); |
569 | 0 | } |
570 | | |
571 | | /* Use simple aggregation if enabled */ |
572 | 0 | if (ctx->simple_aggregation) { |
573 | 0 | retry_add_event_agg: |
574 | 0 | retry_add = FLB_FALSE; |
575 | 0 | ret = process_event_simple_aggregation(ctx, buf, obj, tms, config); |
576 | 0 | if (ret < 0) { |
577 | 0 | return -1; |
578 | 0 | } |
579 | 0 | else if (ret == 1) { |
580 | | /* Buffer full - check if buffer was empty before sending (record too large) */ |
581 | 0 | if (buf->agg_buf.agg_buf_offset == 0) { |
582 | 0 | flb_plg_warn(ctx->ins, "Discarding unprocessable record (too large for aggregation buffer), %s", |
583 | 0 | ctx->delivery_stream); |
584 | 0 | reset_flush_buf(ctx, buf); |
585 | 0 | return 0; |
586 | 0 | } |
587 | | |
588 | | /* Send aggregated record and retry */ |
589 | 0 | ret = send_aggregated_record(ctx, buf); |
590 | 0 | reset_flush_buf(ctx, buf); |
591 | 0 | if (ret < 0) { |
592 | 0 | return -1; |
593 | 0 | } |
594 | 0 | retry_add = FLB_TRUE; |
595 | 0 | } |
596 | 0 | else if (ret == 2) { |
597 | | /* Discard this record */ |
598 | 0 | flb_plg_warn(ctx->ins, "Discarding large or unprocessable record, %s", |
599 | 0 | ctx->delivery_stream); |
600 | 0 | return 0; |
601 | 0 | } |
602 | | |
603 | 0 | if (retry_add == FLB_TRUE) { |
604 | 0 | goto retry_add_event_agg; |
605 | 0 | } |
606 | | |
607 | 0 | return 0; |
608 | 0 | } |
609 | | |
610 | | /* Normal processing without aggregation */ |
611 | 0 | retry_add_event: |
612 | 0 | retry_add = FLB_FALSE; |
613 | 0 | ret = process_event(ctx, buf, obj, tms, config); |
614 | 0 | if (ret < 0) { |
615 | 0 | return -1; |
616 | 0 | } |
617 | 0 | else if (ret == 1) { |
618 | 0 | if (buf->event_index <= 0) { |
619 | | /* somehow the record was larger than our entire request buffer */ |
620 | 0 | flb_plg_warn(ctx->ins, "Discarding massive log record, %s", |
621 | 0 | ctx->delivery_stream); |
622 | 0 | return 0; /* discard this record and return to caller */ |
623 | 0 | } |
624 | | /* send logs and then retry the add */ |
625 | 0 | retry_add = FLB_TRUE; |
626 | 0 | goto send; |
627 | 0 | } else if (ret == 2) { |
628 | | /* discard this record and return to caller */ |
629 | 0 | flb_plg_warn(ctx->ins, "Discarding large or unprocessable record, %s", |
630 | 0 | ctx->delivery_stream); |
631 | 0 | return 0; |
632 | 0 | } |
633 | | |
634 | 0 | event = &buf->events[buf->event_index]; |
635 | 0 | event_bytes = event->len + PUT_RECORD_BATCH_PER_RECORD_LEN; |
636 | |
|
637 | 0 | if ((buf->data_size + event_bytes) > PUT_RECORD_BATCH_PAYLOAD_SIZE) { |
638 | 0 | if (buf->event_index <= 0) { |
639 | | /* somehow the record was larger than our entire request buffer */ |
640 | 0 | flb_plg_warn(ctx->ins, "[size=%zu] Discarding massive log record, %s", |
641 | 0 | event_bytes, ctx->delivery_stream); |
642 | 0 | return 0; /* discard this record and return to caller */ |
643 | 0 | } |
644 | | /* do not send this event */ |
645 | 0 | retry_add = FLB_TRUE; |
646 | 0 | goto send; |
647 | 0 | } |
648 | | |
649 | | /* send is not needed yet, return to caller */ |
650 | 0 | buf->data_size += event_bytes; |
651 | 0 | buf->event_index++; |
652 | |
|
653 | 0 | if (buf->event_index == MAX_EVENTS_PER_PUT) { |
654 | 0 | goto send; |
655 | 0 | } |
656 | | |
657 | 0 | return 0; |
658 | | |
659 | 0 | send: |
660 | 0 | ret = send_log_events(ctx, buf); |
661 | 0 | reset_flush_buf(ctx, buf); |
662 | 0 | if (ret < 0) { |
663 | 0 | return -1; |
664 | 0 | } |
665 | | |
666 | 0 | if (retry_add == FLB_TRUE) { |
667 | 0 | goto retry_add_event; |
668 | 0 | } |
669 | | |
670 | 0 | return 0; |
671 | 0 | } |
672 | | |
673 | | /* |
674 | | * Main routine- processes msgpack and sends in batches |
675 | | * return value is the number of events processed (number sent is stored in buf) |
676 | | */ |
677 | | int process_and_send_records(struct flb_firehose *ctx, struct flush *buf, |
678 | | const char *data, size_t bytes, |
679 | | struct flb_config *config) |
680 | 0 | { |
681 | | // size_t off = 0; |
682 | 0 | int i = 0; |
683 | 0 | size_t map_size; |
684 | | // msgpack_unpacked result; |
685 | | // msgpack_object *obj; |
686 | 0 | msgpack_object map; |
687 | | // msgpack_object root; |
688 | 0 | msgpack_object_kv *kv; |
689 | 0 | msgpack_object key; |
690 | 0 | msgpack_object val; |
691 | 0 | char *key_str = NULL; |
692 | 0 | size_t key_str_size = 0; |
693 | 0 | int j; |
694 | 0 | int ret; |
695 | 0 | int check = FLB_FALSE; |
696 | 0 | int found = FLB_FALSE; |
697 | | // struct flb_time tms; |
698 | 0 | struct flb_log_event_decoder log_decoder; |
699 | 0 | struct flb_log_event log_event; |
700 | |
|
701 | 0 | ret = flb_log_event_decoder_init(&log_decoder, (char *) data, bytes); |
702 | |
|
703 | 0 | if (ret != FLB_EVENT_DECODER_SUCCESS) { |
704 | 0 | flb_plg_error(ctx->ins, |
705 | 0 | "Log event decoder initialization error : %d", ret); |
706 | |
|
707 | 0 | return -1; |
708 | 0 | } |
709 | | |
710 | | |
711 | 0 | while ((ret = flb_log_event_decoder_next( |
712 | 0 | &log_decoder, |
713 | 0 | &log_event)) == FLB_EVENT_DECODER_SUCCESS) { |
714 | 0 | map = *log_event.body; |
715 | 0 | map_size = map.via.map.size; |
716 | |
|
717 | 0 | if (ctx->log_key) { |
718 | 0 | key_str = NULL; |
719 | 0 | key_str_size = 0; |
720 | 0 | check = FLB_FALSE; |
721 | 0 | found = FLB_FALSE; |
722 | |
|
723 | 0 | kv = map.via.map.ptr; |
724 | |
|
725 | 0 | for(j=0; j < map_size; j++) { |
726 | 0 | key = (kv+j)->key; |
727 | 0 | if (key.type == MSGPACK_OBJECT_BIN) { |
728 | 0 | key_str = (char *) key.via.bin.ptr; |
729 | 0 | key_str_size = key.via.bin.size; |
730 | 0 | check = FLB_TRUE; |
731 | 0 | } |
732 | 0 | if (key.type == MSGPACK_OBJECT_STR) { |
733 | 0 | key_str = (char *) key.via.str.ptr; |
734 | 0 | key_str_size = key.via.str.size; |
735 | 0 | check = FLB_TRUE; |
736 | 0 | } |
737 | |
|
738 | 0 | if (check == FLB_TRUE) { |
739 | 0 | if (strncmp(ctx->log_key, key_str, key_str_size) == 0) { |
740 | 0 | found = FLB_TRUE; |
741 | 0 | val = (kv+j)->val; |
742 | 0 | ret = add_event(ctx, buf, &val, &log_event.timestamp, config); |
743 | 0 | if (ret < 0 ) { |
744 | 0 | goto error; |
745 | 0 | } |
746 | 0 | } |
747 | 0 | } |
748 | |
|
749 | 0 | } |
750 | 0 | if (found == FLB_FALSE) { |
751 | 0 | flb_plg_error(ctx->ins, "Could not find log_key '%s' in record, %s", |
752 | 0 | ctx->log_key, ctx->delivery_stream); |
753 | 0 | } |
754 | 0 | else { |
755 | 0 | i++; |
756 | 0 | } |
757 | 0 | continue; |
758 | 0 | } |
759 | | |
760 | 0 | ret = add_event(ctx, buf, &map, &log_event.timestamp, config); |
761 | 0 | if (ret < 0 ) { |
762 | 0 | goto error; |
763 | 0 | } |
764 | 0 | i++; |
765 | 0 | } |
766 | 0 | flb_log_event_decoder_destroy(&log_decoder); |
767 | | |
768 | | /* send any remaining events */ |
769 | 0 | if (ctx->simple_aggregation) { |
770 | | /* Send any remaining aggregated data */ |
771 | 0 | ret = send_aggregated_record(ctx, buf); |
772 | 0 | } |
773 | 0 | else { |
774 | 0 | ret = send_log_events(ctx, buf); |
775 | 0 | } |
776 | 0 | reset_flush_buf(ctx, buf); |
777 | |
|
778 | 0 | if (ret < 0) { |
779 | 0 | return -1; |
780 | 0 | } |
781 | | |
782 | | /* return number of events processed */ |
783 | 0 | buf->records_processed = i; |
784 | |
|
785 | 0 | return i; |
786 | | |
787 | 0 | error: |
788 | 0 | flb_log_event_decoder_destroy(&log_decoder); |
789 | |
|
790 | 0 | return -1; |
791 | 0 | } |
792 | | |
793 | | /* |
794 | | * Returns number of failed records on success, -1 on failure |
795 | | */ |
796 | | static int process_api_response(struct flb_firehose *ctx, |
797 | | struct flb_http_client *c) |
798 | 0 | { |
799 | 0 | int i; |
800 | 0 | int k; |
801 | 0 | int w; |
802 | 0 | int ret; |
803 | 0 | int failed_records = -1; |
804 | 0 | int root_type; |
805 | 0 | char *out_buf; |
806 | 0 | int throughput_exceeded = FLB_FALSE; |
807 | 0 | size_t off = 0; |
808 | 0 | size_t out_size; |
809 | 0 | msgpack_unpacked result; |
810 | 0 | msgpack_object root; |
811 | 0 | msgpack_object key; |
812 | 0 | msgpack_object val; |
813 | 0 | msgpack_object response; |
814 | 0 | msgpack_object response_key; |
815 | 0 | msgpack_object response_val; |
816 | |
|
817 | 0 | if (strstr(c->resp.payload, "\"FailedPutCount\":0")) { |
818 | 0 | return 0; |
819 | 0 | } |
820 | | |
821 | | /* Convert JSON payload to msgpack */ |
822 | 0 | ret = flb_pack_json(c->resp.payload, c->resp.payload_size, |
823 | 0 | &out_buf, &out_size, &root_type, NULL); |
824 | 0 | if (ret == -1) { |
825 | 0 | flb_plg_error(ctx->ins, "could not pack/validate JSON API response\n%s", |
826 | 0 | c->resp.payload); |
827 | 0 | return -1; |
828 | 0 | } |
829 | | |
830 | | /* Lookup error field */ |
831 | 0 | msgpack_unpacked_init(&result); |
832 | 0 | ret = msgpack_unpack_next(&result, out_buf, out_size, &off); |
833 | 0 | if (ret != MSGPACK_UNPACK_SUCCESS) { |
834 | 0 | flb_plg_error(ctx->ins, "Cannot unpack response to find error\n%s", |
835 | 0 | c->resp.payload); |
836 | 0 | failed_records = -1; |
837 | 0 | goto done; |
838 | 0 | } |
839 | | |
840 | 0 | root = result.data; |
841 | 0 | if (root.type != MSGPACK_OBJECT_MAP) { |
842 | 0 | flb_plg_error(ctx->ins, "unexpected payload type=%i", |
843 | 0 | root.type); |
844 | 0 | failed_records = -1; |
845 | 0 | goto done; |
846 | 0 | } |
847 | | |
848 | 0 | for (i = 0; i < root.via.map.size; i++) { |
849 | 0 | key = root.via.map.ptr[i].key; |
850 | 0 | if (key.type != MSGPACK_OBJECT_STR) { |
851 | 0 | flb_plg_error(ctx->ins, "unexpected key type=%i", |
852 | 0 | key.type); |
853 | 0 | failed_records = -1; |
854 | 0 | goto done; |
855 | 0 | } |
856 | | |
857 | 0 | if (key.via.str.size >= 14 && |
858 | 0 | strncmp(key.via.str.ptr, "FailedPutCount", 14) == 0) { |
859 | 0 | val = root.via.map.ptr[i].val; |
860 | 0 | if (val.type != MSGPACK_OBJECT_POSITIVE_INTEGER) { |
861 | 0 | flb_plg_error(ctx->ins, "unexpected 'FailedPutCount' value type=%i", |
862 | 0 | val.type); |
863 | 0 | failed_records = -1; |
864 | 0 | goto done; |
865 | 0 | } |
866 | | |
867 | 0 | failed_records = val.via.u64; |
868 | 0 | if (failed_records == 0) { |
869 | | /* no need to check RequestResponses field */ |
870 | 0 | goto done; |
871 | 0 | } |
872 | 0 | } |
873 | | |
874 | 0 | if (key.via.str.size >= 14 && |
875 | 0 | strncmp(key.via.str.ptr, "RequestResponses", 16) == 0) { |
876 | 0 | val = root.via.map.ptr[i].val; |
877 | 0 | if (val.type != MSGPACK_OBJECT_ARRAY) { |
878 | 0 | flb_plg_error(ctx->ins, "unexpected 'RequestResponses' value type=%i", |
879 | 0 | val.type); |
880 | 0 | failed_records = -1; |
881 | 0 | goto done; |
882 | 0 | } |
883 | | |
884 | 0 | if (val.via.array.size == 0) { |
885 | 0 | flb_plg_error(ctx->ins, "'RequestResponses' field in response is empty"); |
886 | 0 | failed_records = -1; |
887 | 0 | goto done; |
888 | 0 | } |
889 | | |
890 | 0 | for (k = 0; k < val.via.array.size; k++) { |
891 | | /* iterate through the responses */ |
892 | 0 | response = val.via.array.ptr[k]; |
893 | 0 | if (response.type != MSGPACK_OBJECT_MAP) { |
894 | 0 | flb_plg_error(ctx->ins, "unexpected 'RequestResponses[%d]' value type=%i", |
895 | 0 | k, response.type); |
896 | 0 | failed_records = -1; |
897 | 0 | goto done; |
898 | 0 | } |
899 | 0 | for (w = 0; w < response.via.map.size; w++) { |
900 | | /* iterate through the response's keys */ |
901 | 0 | response_key = response.via.map.ptr[w].key; |
902 | 0 | if (response_key.type != MSGPACK_OBJECT_STR) { |
903 | 0 | flb_plg_error(ctx->ins, "unexpected key type=%i", |
904 | 0 | response_key.type); |
905 | 0 | failed_records = -1; |
906 | 0 | goto done; |
907 | 0 | } |
908 | 0 | if (response_key.via.str.size >= 9 && |
909 | 0 | strncmp(response_key.via.str.ptr, "ErrorCode", 9) == 0) { |
910 | 0 | response_val = response.via.map.ptr[w].val; |
911 | 0 | if (!throughput_exceeded && |
912 | 0 | response_val.via.str.size >= 27 && |
913 | 0 | (strncmp(response_val.via.str.ptr, |
914 | 0 | ERR_CODE_SERVICE_UNAVAILABLE, 27) == 0)) { |
915 | 0 | throughput_exceeded = FLB_TRUE; |
916 | 0 | flb_plg_error(ctx->ins, "Throughput limits may have been exceeded, %s", |
917 | 0 | ctx->delivery_stream); |
918 | 0 | } |
919 | 0 | flb_plg_debug(ctx->ins, "Record %i failed with err_code=%.*s", |
920 | 0 | k, response_val.via.str.size, |
921 | 0 | response_val.via.str.ptr); |
922 | 0 | } |
923 | 0 | if (response_key.via.str.size >= 12 && |
924 | 0 | strncmp(response_key.via.str.ptr, "ErrorMessage", 12) == 0) { |
925 | 0 | response_val = response.via.map.ptr[w].val; |
926 | 0 | flb_plg_debug(ctx->ins, "Record %i failed with err_msg=%.*s", |
927 | 0 | k, response_val.via.str.size, |
928 | 0 | response_val.via.str.ptr); |
929 | 0 | } |
930 | 0 | } |
931 | 0 | } |
932 | 0 | } |
933 | 0 | } |
934 | | |
935 | 0 | done: |
936 | 0 | flb_free(out_buf); |
937 | 0 | msgpack_unpacked_destroy(&result); |
938 | 0 | return failed_records; |
939 | 0 | } |
940 | | |
941 | | static int plugin_under_test() |
942 | 0 | { |
943 | 0 | if (getenv("FLB_FIREHOSE_PLUGIN_UNDER_TEST") != NULL) { |
944 | 0 | return FLB_TRUE; |
945 | 0 | } |
946 | | |
947 | 0 | return FLB_FALSE; |
948 | 0 | } |
949 | | |
950 | | static char *mock_error_response(char *error_env_var) |
951 | 0 | { |
952 | 0 | char *err_val = NULL; |
953 | 0 | char *error = NULL; |
954 | 0 | int len = 0; |
955 | |
|
956 | 0 | err_val = getenv(error_env_var); |
957 | 0 | if (err_val != NULL && strlen(err_val) > 0) { |
958 | 0 | error = flb_malloc(strlen(err_val) + sizeof(char)); |
959 | 0 | if (error == NULL) { |
960 | 0 | flb_errno(); |
961 | 0 | return NULL; |
962 | 0 | } |
963 | | |
964 | 0 | len = strlen(err_val); |
965 | 0 | memcpy(error, err_val, len); |
966 | 0 | error[len] = '\0'; |
967 | 0 | return error; |
968 | 0 | } |
969 | | |
970 | 0 | return NULL; |
971 | 0 | } |
972 | | |
973 | | int partial_success() |
974 | 0 | { |
975 | 0 | char *err_val = NULL; |
976 | |
|
977 | 0 | err_val = getenv("PARTIAL_SUCCESS_CASE"); |
978 | 0 | if (err_val != NULL && strlen(err_val) > 0) { |
979 | 0 | return FLB_TRUE; |
980 | 0 | } |
981 | | |
982 | 0 | return FLB_FALSE; |
983 | 0 | } |
984 | | |
985 | | static struct flb_http_client *mock_http_call(char *error_env_var) |
986 | 0 | { |
987 | | /* create an http client so that we can set the response */ |
988 | 0 | struct flb_http_client *c = NULL; |
989 | 0 | char *error = mock_error_response(error_env_var); |
990 | |
|
991 | 0 | c = flb_calloc(1, sizeof(struct flb_http_client)); |
992 | 0 | if (!c) { |
993 | 0 | flb_errno(); |
994 | 0 | flb_free(error); |
995 | 0 | return NULL; |
996 | 0 | } |
997 | 0 | mk_list_init(&c->headers); |
998 | |
|
999 | 0 | if (error != NULL) { |
1000 | 0 | c->resp.status = 400; |
1001 | | /* resp.data is freed on destroy, payload is supposed to reference it */ |
1002 | 0 | c->resp.data = error; |
1003 | 0 | c->resp.payload = c->resp.data; |
1004 | 0 | c->resp.payload_size = strlen(error); |
1005 | 0 | } |
1006 | 0 | else { |
1007 | 0 | c->resp.status = 200; |
1008 | 0 | c->resp.payload = ""; |
1009 | 0 | c->resp.payload_size = 0; |
1010 | 0 | if (partial_success() == FLB_TRUE) { |
1011 | | /* mocked partial failure response */ |
1012 | 0 | c->resp.payload = "{\"Encrypted\": false,\"FailedPutCount\": 1,\"RequestResponses\":[{\"RecordId\": \"Me0CqhxK3BK3MiBWgy/AydQrVUg7vbc40Z4zNds3jiiJDscqGtWFz9bJugbrAoN70YCaxpXgmyR9R+LFxS2rleDepqFljYArBtXnRmVzSMOAzTJZlwsO84+757kBvA5RUycF3wC3XZjFtUFP0Q4QTdhuD8HMJBvKGiBY9Yy5jBUmZuKhXxCLQ/YTwKQaQKn4fnc5iISxaErPXsWMI7OApHZ1eFGvcHVZ\"},{\"RecordId\": \"NRAZVkblYgWWDSvTAF/9jBR4MlciEUFV+QIjb1D8uar7YbC3wqeLQuSZ0GEopGlE/8JAK9h9aAyTub5lH5V+bZuR3SeKKABWoJ788/tI455Kup9oRzmXTKWiXeklxmAe9MtsSz0y4t3oIrSLq8e3QVH9DJKWdhDkIXd8lXK1wuJi8tKmnNgxFob/Cz398kQFXPc4JwKj3Dv3Ou0qibZiusko6f7yBUve\",\"ErrorCode\":\"ServiceUnavailableException\",\"ErrorMessage\": \"Catsssss\"},{\"RecordId\": \"InFGTFvML/MGCLtnC3moI/zCISrKSScu/D8oCGmeIIeVaYUfywHpr2NmsQiZsxUL9+4ThOm2ypxqFGudZvgXQ45gUWMG+R4Y5xzS03N+vQ71+UaL392jY6HUs2SxYkZQe6vpdK+xHaJJ1b8uE++Laxg9rmsXtNt193WjmH3FhU1veu9pnSiGZgqC7czpyVgvZBNeWc+hTjEVicj3VAHBg/9yRN0sC30C\",\"ErrorCode\":\"ServiceUnavailableException\",\"ErrorMessage\": \"Catsssss 2\"},{\"RecordId\":\"KufmrRJ2z8zAgYAYGz6rm4BQC8SA7g87lQJQl2DQ+Be5EiEpr5bG33ilnQVvo1Q05BJuQBnjbw2cm919Ya72awapxfOBdZcPPKJN7KDZV/n1DFCDDrJ2vgyNK4qhKdo3Mr7nyrBpkLIs93PdxOdrTh11Y9HHEaFtim0cHJYpKCSZBjNObfWjfjHx5TuB7L3PHQqMKMu0MT5L9gPgVXHElGalqKZGTcfB\"}]}"; |
1013 | 0 | c->resp.payload_size = strlen(c->resp.payload); |
1014 | 0 | } |
1015 | 0 | else { |
1016 | | /* mocked success response */ |
1017 | 0 | c->resp.payload = "{\"Encrypted\": false,\"FailedPutCount\": 0,\"RequestResponses\":[{\"RecordId\": \"Me0CqhxK3BK3MiBWgy/AydQrVUg7vbc40Z4zNds3jiiJDscqGtWFz9bJugbrAoN70YCaxpXgmyR9R+LFxS2rleDepqFljYArBtXnRmVzSMOAzTJZlwsO84+757kBvA5RUycF3wC3XZjFtUFP0Q4QTdhuD8HMJBvKGiBY9Yy5jBUmZuKhXxCLQ/YTwKQaQKn4fnc5iISxaErPXsWMI7OApHZ1eFGvcHVZ\"},{\"RecordId\": \"NRAZVkblYgWWDSvTAF/9jBR4MlciEUFV+QIjb1D8uar7YbC3wqeLQuSZ0GEopGlE/8JAK9h9aAyTub5lH5V+bZuR3SeKKABWoJ788/tI455Kup9oRzmXTKWiXeklxmAe9MtsSz0y4t3oIrSLq8e3QVH9DJKWdhDkIXd8lXK1wuJi8tKmnNgxFob/Cz398kQFXPc4JwKj3Dv3Ou0qibZiusko6f7yBUve\"},{\"RecordId\": \"InFGTFvML/MGCLtnC3moI/zCISrKSScu/D8oCGmeIIeVaYUfywHpr2NmsQiZsxUL9+4ThOm2ypxqFGudZvgXQ45gUWMG+R4Y5xzS03N+vQ71+UaL392jY6HUs2SxYkZQe6vpdK+xHaJJ1b8uE++Laxg9rmsXtNt193WjmH3FhU1veu9pnSiGZgqC7czpyVgvZBNeWc+hTjEVicj3VAHBg/9yRN0sC30C\"},{\"RecordId\": \"KufmrRJ2z8zAgYAYGz6rm4BQC8SA7g87lQJQl2DQ+Be5EiEpr5bG33ilnQVvo1Q05BJuQBnjbw2cm919Ya72awapxfOBdZcPPKJN7KDZV/n1DFCDDrJ2vgyNK4qhKdo3Mr7nyrBpkLIs93PdxOdrTh11Y9HHEaFtim0cHJYpKCSZBjNObfWjfjHx5TuB7L3PHQqMKMu0MT5L9gPgVXHElGalqKZGTcfB\"}]}"; |
1018 | 0 | c->resp.payload_size = strlen(c->resp.payload); |
1019 | 0 | } |
1020 | 0 | } |
1021 | |
|
1022 | 0 | return c; |
1023 | 0 | } |
1024 | | |
1025 | | |
1026 | | /* |
1027 | | * Returns -1 on failure, 0 on success |
1028 | | */ |
1029 | | int put_record_batch(struct flb_firehose *ctx, struct flush *buf, |
1030 | | size_t payload_size, int num_records) |
1031 | 0 | { |
1032 | |
|
1033 | 0 | struct flb_http_client *c = NULL; |
1034 | 0 | struct flb_aws_client *firehose_client; |
1035 | 0 | flb_sds_t error; |
1036 | 0 | int failed_records = 0; |
1037 | |
|
1038 | 0 | flb_plg_debug(ctx->ins, "Sending log records to delivery stream %s", |
1039 | 0 | ctx->delivery_stream); |
1040 | |
|
1041 | 0 | if (plugin_under_test() == FLB_TRUE) { |
1042 | 0 | c = mock_http_call("TEST_PUT_RECORD_BATCH_ERROR"); |
1043 | 0 | } |
1044 | 0 | else { |
1045 | 0 | firehose_client = ctx->firehose_client; |
1046 | 0 | c = firehose_client->client_vtable->request(firehose_client, FLB_HTTP_POST, |
1047 | 0 | "/", buf->out_buf, payload_size, |
1048 | 0 | &put_record_batch_header, 1); |
1049 | 0 | } |
1050 | |
|
1051 | 0 | if (c) { |
1052 | 0 | flb_plg_debug(ctx->ins, "PutRecordBatch http status=%d", c->resp.status); |
1053 | |
|
1054 | 0 | if (c->resp.status == 200) { |
1055 | | /* Firehose API can return partial success- check response */ |
1056 | 0 | if (c->resp.payload_size > 0) { |
1057 | 0 | failed_records = process_api_response(ctx, c); |
1058 | 0 | if (failed_records < 0) { |
1059 | 0 | flb_plg_error(ctx->ins, "PutRecordBatch response " |
1060 | 0 | "could not be parsed, %s", |
1061 | 0 | c->resp.payload); |
1062 | 0 | flb_http_client_destroy(c); |
1063 | 0 | return -1; |
1064 | 0 | } |
1065 | 0 | if (failed_records == num_records) { |
1066 | 0 | flb_plg_error(ctx->ins, "PutRecordBatch request returned " |
1067 | 0 | "with no records successfully recieved, %s", |
1068 | 0 | ctx->delivery_stream); |
1069 | 0 | flb_http_client_destroy(c); |
1070 | 0 | return -1; |
1071 | 0 | } |
1072 | 0 | if (failed_records > 0) { |
1073 | 0 | flb_plg_error(ctx->ins, "%d out of %d records failed to be " |
1074 | 0 | "delivered, will retry this batch, %s", |
1075 | 0 | failed_records, num_records, |
1076 | 0 | ctx->delivery_stream); |
1077 | 0 | flb_http_client_destroy(c); |
1078 | 0 | return -1; |
1079 | 0 | } |
1080 | 0 | } |
1081 | 0 | flb_plg_debug(ctx->ins, "Sent events to %s", ctx->delivery_stream); |
1082 | 0 | flb_http_client_destroy(c); |
1083 | 0 | return 0; |
1084 | 0 | } |
1085 | | |
1086 | | /* Check error */ |
1087 | 0 | if (c->resp.payload_size > 0) { |
1088 | 0 | error = flb_aws_error(c->resp.payload, c->resp.payload_size); |
1089 | 0 | if (error != NULL) { |
1090 | 0 | if (strcmp(error, ERR_CODE_SERVICE_UNAVAILABLE) == 0) { |
1091 | 0 | flb_plg_error(ctx->ins, "Throughput limits for %s " |
1092 | 0 | "may have been exceeded.", |
1093 | 0 | ctx->delivery_stream); |
1094 | 0 | } |
1095 | 0 | if (strncmp(error, "SerializationException", 22) == 0) { |
1096 | | /* |
1097 | | * If this happens, we habe a bug in the code |
1098 | | * User should send us the output to debug |
1099 | | */ |
1100 | 0 | flb_plg_error(ctx->ins, "<<------Bug in Code------>>"); |
1101 | 0 | printf("Malformed request: %s", buf->out_buf); |
1102 | 0 | } |
1103 | 0 | flb_aws_print_error(c->resp.payload, c->resp.payload_size, |
1104 | 0 | "PutRecordBatch", ctx->ins); |
1105 | 0 | flb_sds_destroy(error); |
1106 | 0 | } |
1107 | 0 | else { |
1108 | | /* error could not be parsed, print raw response to debug */ |
1109 | 0 | flb_plg_debug(ctx->ins, "Raw response: %s", c->resp.payload); |
1110 | 0 | } |
1111 | 0 | } |
1112 | 0 | } |
1113 | | |
1114 | 0 | flb_plg_error(ctx->ins, "Failed to send log records to %s", ctx->delivery_stream); |
1115 | 0 | if (c) { |
1116 | 0 | flb_http_client_destroy(c); |
1117 | 0 | } |
1118 | 0 | return -1; |
1119 | 0 | } |
1120 | | |
1121 | | |
1122 | | void flush_destroy(struct flush *buf) |
1123 | 0 | { |
1124 | 0 | if (buf) { |
1125 | 0 | if (buf->agg_buf_initialized) { |
1126 | 0 | flb_aws_aggregation_destroy(&buf->agg_buf); |
1127 | 0 | } |
1128 | 0 | flb_free(buf->tmp_buf); |
1129 | 0 | flb_free(buf->out_buf); |
1130 | 0 | flb_free(buf->events); |
1131 | 0 | flb_free(buf->event_buf); |
1132 | 0 | flb_free(buf); |
1133 | 0 | } |
1134 | 0 | } |