Coverage Report

Created: 2026-09-28 07:36

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}