Coverage Report

Created: 2026-09-28 07:36

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/fluent-bit/plugins/in_opentelemetry/opentelemetry_logs.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_input_plugin.h>
21
#include <fluent-bit/flb_mem.h>
22
#include <fluent-bit/flb_sds.h>
23
#include <fluent-bit/flb_pack.h>
24
#include <fluent-bit/flb_log_event_encoder.h>
25
#include <fluent-bit/flb_time.h>
26
#include <fluent-bit/flb_opentelemetry.h>
27
#include <fluent-otel-proto/fluent-otel.h>
28
29
#include <cfl/cfl_arena.h>
30
31
#include "opentelemetry_protobuf.h"
32
#include "opentelemetry.h"
33
#include "opentelemetry_utils.h"
34
35
/*
36
 * Small requests use protobuf-c's default allocator to avoid reserving an
37
 * arena chunk. Realistic batched requests use the arena to amortize the large
38
 * number of short-lived allocations performed by protobuf-c.
39
 */
40
0
#define OTEL_PROTOBUF_ARENA_MIN_PAYLOAD_SIZE 8192
41
0
#define OTEL_PROTOBUF_ARENA_INITIAL_CHUNK_SIZE 4096
42
0
#define OTEL_PROTOBUF_ARENA_MAX_CHUNK_SIZE 65536
43
44
static void *protobuf_arena_chunk_malloc(void *context, size_t size)
45
0
{
46
0
    (void) context;
47
48
0
    return flb_malloc(size);
49
0
}
50
51
static void protobuf_arena_chunk_free(void *context, void *pointer)
52
0
{
53
0
    (void) context;
54
55
0
    flb_free(pointer);
56
0
}
57
58
static void *protobuf_arena_alloc(void *allocator_data, size_t size)
59
0
{
60
0
    struct cfl_arena *arena;
61
62
0
    arena = allocator_data;
63
0
    if (size == 0) {
64
0
        size = 1;
65
0
    }
66
67
0
    return cfl_arena_malloc(arena, size);
68
0
}
69
70
static void protobuf_arena_free(void *allocator_data, void *pointer)
71
0
{
72
0
    (void) allocator_data;
73
0
    (void) pointer;
74
0
}
75
76
static struct cfl_arena *protobuf_arena_create(void)
77
0
{
78
0
    struct cfl_arena_options options;
79
80
0
    cfl_arena_options_init(&options);
81
0
    options.chunk_size = OTEL_PROTOBUF_ARENA_INITIAL_CHUNK_SIZE;
82
0
    options.maximum_chunk_size = OTEL_PROTOBUF_ARENA_MAX_CHUNK_SIZE;
83
0
    options.malloc_fn = protobuf_arena_chunk_malloc;
84
0
    options.free_fn = protobuf_arena_chunk_free;
85
86
0
    return cfl_arena_create_with_options(&options);
87
0
}
88
89
static Opentelemetry__Proto__Collector__Logs__V1__ExportLogsServiceRequest *
90
protobuf_logs_unpack(ProtobufCAllocator *allocator, size_t size, const uint8_t *data)
91
0
{
92
0
    if (opentelemetry_protobuf_validate(
93
0
            &opentelemetry__proto__collector__logs__v1__export_logs_service_request__descriptor,
94
0
            data, size) != 0) {
95
0
        return NULL;
96
0
    }
97
0
    return opentelemetry__proto__collector__logs__v1__export_logs_service_request__unpack(
98
0
        allocator, size, data);
99
0
}
100
101
static void protobuf_logs_free(
102
    Opentelemetry__Proto__Collector__Logs__V1__ExportLogsServiceRequest *logs,
103
    ProtobufCAllocator *allocator)
104
0
{
105
0
    opentelemetry__proto__collector__logs__v1__export_logs_service_request__free_unpacked(
106
0
        logs, allocator);
107
0
}
108
109
/*
110
 * OTLP encoding functions to pack the log records as msgpack
111
 * ----------------------------------------------------------
112
 */
113
static int otlp_pack_any_value(msgpack_packer *mp_pck, Opentelemetry__Proto__Common__V1__AnyValue *body);
114
115
static int otel_pack_string(msgpack_packer *mp_pck, char *str)
116
0
{
117
0
    return msgpack_pack_str_with_body(mp_pck, str, strlen(str));
118
0
}
119
120
static int otel_pack_bool(msgpack_packer *mp_pck, bool val)
121
0
{
122
0
    if (val) {
123
0
        return msgpack_pack_true(mp_pck);
124
0
    }
125
0
    else {
126
0
        return msgpack_pack_false(mp_pck);
127
0
    }
128
0
}
129
130
static int otel_pack_int(msgpack_packer *mp_pck, int64_t val)
131
0
{
132
0
    return msgpack_pack_int64(mp_pck, val);
133
0
}
134
135
static int otel_pack_double(msgpack_packer *mp_pck, double val)
136
0
{
137
0
    return msgpack_pack_double(mp_pck, val);
138
0
}
139
140
static int otel_pack_kvarray(msgpack_packer *mp_pck,
141
                             Opentelemetry__Proto__Common__V1__KeyValue **kv_array,
142
                             size_t kv_count)
143
0
{
144
0
    int result;
145
0
    int index;
146
147
0
    result = msgpack_pack_map(mp_pck, kv_count);
148
149
0
    if (result != 0) {
150
0
        return result;
151
0
    }
152
153
0
    for (index = 0; index < kv_count && result == 0; index++) {
154
0
        result = otel_pack_string(mp_pck, kv_array[index]->key);
155
156
0
        if(result == 0) {
157
0
           result = otlp_pack_any_value(mp_pck, kv_array[index]->value);
158
0
        }
159
0
    }
160
161
0
    return result;
162
0
}
163
164
static int otel_pack_kvlist(msgpack_packer *mp_pck,
165
                            Opentelemetry__Proto__Common__V1__KeyValueList *kv_list)
166
0
{
167
0
    int kv_index;
168
0
    int ret;
169
0
    char *key;
170
0
    Opentelemetry__Proto__Common__V1__AnyValue *value;
171
172
0
    ret = msgpack_pack_map(mp_pck, kv_list->n_values);
173
0
    if (ret != 0) {
174
0
        return ret;
175
0
    }
176
177
0
    for (kv_index = 0; kv_index < kv_list->n_values && ret == 0; kv_index++) {
178
0
        key = kv_list->values[kv_index]->key;
179
0
        value = kv_list->values[kv_index]->value;
180
181
0
        ret = otel_pack_string(mp_pck, key);
182
183
0
        if(ret == 0) {
184
0
           ret = otlp_pack_any_value(mp_pck, value);
185
0
        }
186
0
    }
187
188
0
    return ret;
189
0
}
190
191
static int otel_pack_array(msgpack_packer *mp_pck,
192
                           Opentelemetry__Proto__Common__V1__ArrayValue *array)
193
0
{
194
0
    int ret;
195
0
    int array_index;
196
197
0
    ret = msgpack_pack_array(mp_pck, array->n_values);
198
199
0
    if (ret != 0) {
200
0
        return ret;
201
0
    }
202
203
0
    for (array_index = 0; array_index < array->n_values && ret == 0; array_index++) {
204
0
        ret = otlp_pack_any_value(mp_pck, array->values[array_index]);
205
0
    }
206
207
0
    return ret;
208
0
}
209
210
static int otel_pack_bytes(msgpack_packer *mp_pck,
211
                           ProtobufCBinaryData bytes)
212
0
{
213
0
    return msgpack_pack_bin_with_body(mp_pck, bytes.data, bytes.len);
214
0
}
215
216
static int otlp_pack_any_value(msgpack_packer *mp_pck,
217
                               Opentelemetry__Proto__Common__V1__AnyValue *body)
218
0
{
219
0
    int result;
220
221
0
    result = -2;
222
223
0
    if (body == NULL) {
224
0
        msgpack_pack_nil(mp_pck);
225
0
        return 0;
226
0
    }
227
228
0
    switch(body->value_case){
229
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_STRING_VALUE:
230
0
            result = otel_pack_string(mp_pck, body->string_value);
231
0
            break;
232
233
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_STRING_VALUE_STRINDEX:
234
            /* Profiling-only string dictionary reference: ignore in logs. */
235
0
            result = msgpack_pack_nil(mp_pck);
236
0
            break;
237
238
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_BOOL_VALUE:
239
0
            result =  otel_pack_bool(mp_pck, body->bool_value);
240
0
            break;
241
242
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_INT_VALUE:
243
0
            result = otel_pack_int(mp_pck, body->int_value);
244
0
            break;
245
246
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_DOUBLE_VALUE:
247
0
            result = otel_pack_double(mp_pck, body->double_value);
248
0
            break;
249
250
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_ARRAY_VALUE:
251
0
            result = otel_pack_array(mp_pck, body->array_value);
252
0
            break;
253
254
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_KVLIST_VALUE:
255
0
            result = otel_pack_kvlist(mp_pck, body->kvlist_value);
256
0
            break;
257
258
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_BYTES_VALUE:
259
0
            result = otel_pack_bytes(mp_pck, body->bytes_value);
260
0
            break;
261
262
0
        case OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE__NOT_SET:
263
            /* treat an unset value as null */
264
0
            result = msgpack_pack_nil(mp_pck);
265
0
            break;
266
267
0
        default:
268
0
            break;
269
0
    }
270
271
0
    if (result == -2) {
272
0
        flb_error("[otel]: invalid value type in pack_any_value");
273
0
        result = -1;
274
0
    }
275
276
0
    return result;
277
0
}
278
279
/* https://opentelemetry.io/docs/specs/otel/logs/data-model/#log-and-event-record-definition */
280
static int otel_pack_v1_metadata(struct flb_opentelemetry *ctx,
281
                                 msgpack_packer *mp_pck,
282
                                 struct Opentelemetry__Proto__Logs__V1__LogRecord *log_record,
283
                                 Opentelemetry__Proto__Resource__V1__Resource *resource,
284
                                 Opentelemetry__Proto__Common__V1__InstrumentationScope *scope)
285
0
{
286
0
    int ret;
287
0
    int len;
288
0
    struct flb_mp_map_header mh;
289
0
    struct flb_mp_map_header otlp_mh;
290
291
0
    flb_mp_map_header_init(&otlp_mh, mp_pck);
292
293
0
    len = flb_sds_len(ctx->logs_metadata_key);
294
295
    /* otlp key start */
296
0
    flb_mp_map_header_append(&otlp_mh);
297
298
0
    msgpack_pack_str(mp_pck, len);
299
0
    msgpack_pack_str_body(mp_pck, ctx->logs_metadata_key, len);
300
301
0
    flb_mp_map_header_init(&mh, mp_pck);
302
303
0
    if (log_record->observed_time_unix_nano != 0) {
304
0
        flb_mp_map_header_append(&mh);
305
0
        msgpack_pack_str(mp_pck, 18);
306
0
        msgpack_pack_str_body(mp_pck, "observed_timestamp", 18);
307
0
        msgpack_pack_uint64(mp_pck, log_record->observed_time_unix_nano);
308
0
    }
309
310
    /* Value of 0 indicates unknown or missing timestamp. */
311
0
    if (log_record->time_unix_nano != 0) {
312
0
        flb_mp_map_header_append(&mh);
313
0
        msgpack_pack_str(mp_pck, 9);
314
0
        msgpack_pack_str_body(mp_pck, "timestamp", 9);
315
0
        msgpack_pack_uint64(mp_pck, log_record->time_unix_nano);
316
0
    }
317
318
    /* https://opentelemetry.io/docs/specs/otel/logs/data-model/#field-severitynumber */
319
0
    if (log_record->severity_number >= 1 && log_record->severity_number <= 24) {
320
0
        flb_mp_map_header_append(&mh);
321
0
        msgpack_pack_str(mp_pck, 15);
322
0
        msgpack_pack_str_body(mp_pck, "severity_number", 15);
323
0
        msgpack_pack_uint64(mp_pck, log_record->severity_number);
324
0
    }
325
326
0
    if (log_record->severity_text != NULL && strlen(log_record->severity_text) > 0) {
327
0
        flb_mp_map_header_append(&mh);
328
0
        msgpack_pack_str(mp_pck, 13);
329
0
        msgpack_pack_str_body(mp_pck, "severity_text", 13);
330
0
        msgpack_pack_str(mp_pck, strlen(log_record->severity_text));
331
0
        msgpack_pack_str_body(mp_pck, log_record->severity_text, strlen(log_record->severity_text));
332
0
    }
333
334
0
    if (log_record->n_attributes > 0) {
335
0
        flb_mp_map_header_append(&mh);
336
0
        msgpack_pack_str(mp_pck, 10);
337
0
        msgpack_pack_str_body(mp_pck, "attributes", 10);
338
0
        ret = otel_pack_kvarray(mp_pck,
339
0
                                log_record->attributes,
340
0
                                log_record->n_attributes);
341
0
        if (ret != 0) {
342
0
            return ret;
343
0
        }
344
0
    }
345
346
0
    if (log_record->dropped_attributes_count > 0) {
347
0
        flb_mp_map_header_append(&mh);
348
0
        msgpack_pack_str(mp_pck, 24);
349
0
        msgpack_pack_str_body(mp_pck, "dropped_attributes_count", 24);
350
0
        msgpack_pack_uint64(mp_pck, log_record->dropped_attributes_count);
351
0
    }
352
353
0
    if (log_record->trace_id.len > 0) {
354
0
        flb_mp_map_header_append(&mh);
355
0
        msgpack_pack_str(mp_pck, 8);
356
0
        msgpack_pack_str_body(mp_pck, "trace_id", 8);
357
0
        ret = otel_pack_bytes(mp_pck, log_record->trace_id);
358
0
        if (ret != 0) {
359
0
            return ret;
360
0
        }
361
0
    }
362
363
0
    if (log_record->span_id.len > 0) {
364
0
        flb_mp_map_header_append(&mh);
365
0
        msgpack_pack_str(mp_pck, 7);
366
0
        msgpack_pack_str_body(mp_pck, "span_id", 7);
367
0
        ret = otel_pack_bytes(mp_pck, log_record->span_id);
368
0
        if (ret != 0) {
369
0
            return ret;
370
0
        }
371
0
    }
372
373
0
    flb_mp_map_header_append(&mh);
374
0
    msgpack_pack_str(mp_pck, 11);
375
0
    msgpack_pack_str_body(mp_pck, "trace_flags", 11);
376
0
    msgpack_pack_uint8(mp_pck, (uint8_t) log_record->flags & 0xff);
377
378
0
    if (log_record->event_name != NULL && strlen(log_record->event_name) > 0) {
379
0
        flb_mp_map_header_append(&mh);
380
0
        msgpack_pack_str(mp_pck, 10);
381
0
        msgpack_pack_str_body(mp_pck, "event_name", 10);
382
0
        msgpack_pack_str(mp_pck, strlen(log_record->event_name));
383
0
        msgpack_pack_str_body(mp_pck, log_record->event_name, strlen(log_record->event_name));
384
0
    }
385
386
0
    flb_mp_map_header_end(&mh);
387
388
    /* otlp key end */
389
0
    flb_mp_map_header_end(&otlp_mh);
390
391
0
    return 0;
392
0
}
393
394
static int binary_payload_to_msgpack(struct flb_opentelemetry *ctx,
395
                                     struct flb_log_event_encoder *encoder,
396
                                     char *tag, size_t tag_len,
397
                                     uint8_t *in_buf,
398
                                     size_t in_size,
399
                                     size_t *record_count)
400
0
{
401
0
    int ret = 0;
402
0
    int len;
403
0
    int resource_logs_index;
404
0
    int scope_log_index;
405
0
    int log_record_index;
406
0
    char *logs_body_key;
407
0
    int scope_has_schema_url;
408
0
    struct cfl_arena *protobuf_arena;
409
0
    struct flb_mp_map_header mh;
410
0
    struct flb_mp_map_header mh_tmp;
411
0
    struct flb_time tm;
412
0
    ProtobufCAllocator arena_allocator;
413
0
    ProtobufCAllocator *protobuf_allocator;
414
415
0
    msgpack_packer *mp_pck;
416
0
    msgpack_packer *mp_pck_meta;
417
418
    /* OTel proto suff */
419
0
    Opentelemetry__Proto__Collector__Logs__V1__ExportLogsServiceRequest *input_logs;
420
0
    Opentelemetry__Proto__Logs__V1__ScopeLogs **scope_logs;
421
0
    Opentelemetry__Proto__Logs__V1__ScopeLogs *scope_log;
422
0
    Opentelemetry__Proto__Common__V1__InstrumentationScope *scope;
423
424
0
    Opentelemetry__Proto__Logs__V1__ResourceLogs **resource_logs;
425
0
    Opentelemetry__Proto__Logs__V1__ResourceLogs *resource_log;
426
0
    Opentelemetry__Proto__Logs__V1__LogRecord **log_records;
427
0
    Opentelemetry__Proto__Resource__V1__Resource *resource;
428
429
0
    mp_pck = &encoder->body.packer;
430
0
    mp_pck_meta = &encoder->metadata.packer;
431
0
    input_logs = NULL;
432
0
    protobuf_arena = NULL;
433
0
    protobuf_allocator = NULL;
434
0
    *record_count = 0;
435
436
0
    if (in_size >= OTEL_PROTOBUF_ARENA_MIN_PAYLOAD_SIZE) {
437
0
        protobuf_arena = protobuf_arena_create();
438
0
        if (protobuf_arena != NULL) {
439
0
            arena_allocator.alloc = protobuf_arena_alloc;
440
0
            arena_allocator.free = protobuf_arena_free;
441
0
            arena_allocator.allocator_data = protobuf_arena;
442
0
            protobuf_allocator = &arena_allocator;
443
0
        }
444
0
    }
445
446
    /* unpack logs from protobuf payload */
447
0
    input_logs = protobuf_logs_unpack(protobuf_allocator, in_size, in_buf);
448
0
    if (input_logs == NULL) {
449
0
        flb_plg_warn(ctx->ins, "failed to unpack input logs from OpenTelemetry payload");
450
0
        ret = -1;
451
0
        goto binary_payload_to_msgpack_end;
452
0
    }
453
454
0
    resource_logs = input_logs->resource_logs;
455
0
    if (input_logs->n_resource_logs == 0) {
456
0
        ret = 0;
457
0
        goto binary_payload_to_msgpack_end;
458
0
    }
459
460
0
    if (resource_logs == NULL) {
461
0
        flb_plg_warn(ctx->ins, "no resource logs found");
462
0
        ret = -1;
463
0
        goto binary_payload_to_msgpack_end;
464
0
    }
465
466
0
    for (resource_logs_index = 0; resource_logs_index < input_logs->n_resource_logs; resource_logs_index++) {
467
0
        resource_log = resource_logs[resource_logs_index];
468
0
        if (resource_log == NULL) {
469
0
            flb_plg_warn(ctx->ins, "null resource logs entry found");
470
0
            ret = -1;
471
0
            goto binary_payload_to_msgpack_end;
472
0
        }
473
474
0
        resource = resource_log->resource;
475
0
        scope_logs = resource_log->scope_logs;
476
477
0
        if (resource_log->n_scope_logs > 0 && scope_logs == NULL) {
478
0
            flb_plg_warn(ctx->ins, "no scope logs found");
479
0
            ret = -1;
480
0
            goto binary_payload_to_msgpack_end;
481
0
        }
482
483
0
        for (scope_log_index = 0; scope_log_index < resource_log->n_scope_logs; scope_log_index++) {
484
0
            scope_log = scope_logs[scope_log_index];
485
0
            if (scope_log == NULL) {
486
0
                flb_plg_warn(ctx->ins, "null scope logs entry found");
487
0
                ret = -1;
488
0
                goto binary_payload_to_msgpack_end;
489
0
            }
490
491
0
            log_records = scope_log->log_records;
492
493
0
            if (scope_log->n_log_records == 0) {
494
0
                continue;
495
0
            }
496
497
0
            if (log_records == NULL) {
498
0
                flb_plg_warn(ctx->ins, "no log records found");
499
0
                ret = -1;
500
0
                goto binary_payload_to_msgpack_end;
501
0
            }
502
503
0
            flb_log_event_encoder_group_init(encoder);
504
505
            /* pack schema (internal) */
506
0
            ret = flb_log_event_encoder_append_metadata_values(encoder,
507
0
                                                               FLB_LOG_EVENT_STRING_VALUE("schema", 6),
508
0
                                                               FLB_LOG_EVENT_STRING_VALUE("otlp", 4),
509
0
                                                               FLB_LOG_EVENT_STRING_VALUE("resource_id", 11),
510
0
                                                               FLB_LOG_EVENT_INT64_VALUE(resource_logs_index),
511
0
                                                               FLB_LOG_EVENT_STRING_VALUE("scope_id", 8),
512
0
                                                               FLB_LOG_EVENT_INT64_VALUE(scope_log_index));
513
514
515
0
            ret = flb_log_event_encoder_dynamic_field_reset(&encoder->body);
516
0
            if (ret != FLB_EVENT_ENCODER_SUCCESS) {
517
0
                flb_plg_error(ctx->ins, "failed to reset log event body: %s",
518
0
                              flb_log_event_encoder_get_error_description(ret));
519
0
                goto binary_payload_to_msgpack_end;
520
0
            }
521
522
0
            flb_mp_map_header_init(&mh, mp_pck);
523
524
            /* Resource */
525
0
            flb_mp_map_header_append(&mh);
526
0
            msgpack_pack_str(mp_pck, 8);
527
0
            msgpack_pack_str_body(mp_pck, "resource", 8);
528
529
0
            flb_mp_map_header_init(&mh_tmp, mp_pck);
530
0
            if (resource) {
531
                /* look for OTel resource attributes */
532
0
                if (resource->n_attributes > 0 && resource->attributes) {
533
0
                    flb_mp_map_header_append(&mh_tmp);
534
0
                    msgpack_pack_str(mp_pck, 10);
535
0
                    msgpack_pack_str_body(mp_pck, "attributes", 10);
536
537
0
                    ret = otel_pack_kvarray(mp_pck,
538
0
                                            resource->attributes,
539
0
                                            resource->n_attributes);
540
0
                    if (ret != 0) {
541
0
                        goto binary_payload_to_msgpack_end;
542
0
                    }
543
0
                }
544
545
0
                if (resource->dropped_attributes_count > 0) {
546
0
                    flb_mp_map_header_append(&mh_tmp);
547
0
                    msgpack_pack_str(mp_pck, 24);
548
0
                    msgpack_pack_str_body(mp_pck, "dropped_attributes_count", 24);
549
0
                    msgpack_pack_uint64(mp_pck, resource->dropped_attributes_count);
550
0
                }
551
552
0
                if (resource_log->schema_url) {
553
0
                    flb_mp_map_header_append(&mh_tmp);
554
0
                    msgpack_pack_str(mp_pck, 10);
555
0
                    msgpack_pack_str_body(mp_pck, "schema_url", 10);
556
557
0
                    len = strlen(resource_log->schema_url);
558
0
                    msgpack_pack_str(mp_pck, len);
559
0
                    msgpack_pack_str_body(mp_pck, resource_log->schema_url, len);
560
0
                }
561
0
            }
562
0
            flb_mp_map_header_end(&mh_tmp);
563
564
            /* scope */
565
0
            flb_mp_map_header_append(&mh);
566
0
            msgpack_pack_str(mp_pck, 5);
567
0
            msgpack_pack_str_body(mp_pck, "scope", 5);
568
569
            /* Scope */
570
0
            scope = scope_log->scope;
571
0
            scope_has_schema_url = FLB_FALSE;
572
573
0
            if (scope_log->schema_url && strlen(scope_log->schema_url) > 0) {
574
0
                scope_has_schema_url = FLB_TRUE;
575
0
            }
576
577
0
            if (scope && (scope->name || scope->version ||
578
0
                          scope->n_attributes > 0 || scope->dropped_attributes_count > 0 ||
579
0
                          scope_has_schema_url == FLB_TRUE)) {
580
0
                flb_mp_map_header_init(&mh_tmp, mp_pck);
581
582
0
                if (scope_has_schema_url == FLB_TRUE) {
583
0
                    flb_mp_map_header_append(&mh_tmp);
584
0
                    msgpack_pack_str(mp_pck, 10);
585
0
                    msgpack_pack_str_body(mp_pck, "schema_url", 10);
586
587
0
                    len = strlen(scope_log->schema_url);
588
0
                    msgpack_pack_str(mp_pck, len);
589
0
                    msgpack_pack_str_body(mp_pck, scope_log->schema_url, len);
590
0
                }
591
592
0
                if (scope->name && strlen(scope->name) > 0) {
593
0
                    flb_mp_map_header_append(&mh_tmp);
594
0
                    msgpack_pack_str(mp_pck, 4);
595
0
                    msgpack_pack_str_body(mp_pck, "name", 4);
596
597
0
                    len = strlen(scope->name);
598
0
                    msgpack_pack_str(mp_pck, len);
599
0
                    msgpack_pack_str_body(mp_pck, scope->name, len);
600
0
                }
601
0
                if (scope->version && strlen(scope->version) > 0) {
602
0
                    flb_mp_map_header_append(&mh_tmp);
603
604
0
                    msgpack_pack_str(mp_pck, 7);
605
0
                    msgpack_pack_str_body(mp_pck, "version", 7);
606
607
0
                    len = strlen(scope->version);
608
0
                    msgpack_pack_str(mp_pck, len);
609
0
                    msgpack_pack_str_body(mp_pck, scope->version, len);
610
0
                }
611
612
0
                if (scope->n_attributes > 0 && scope->attributes) {
613
0
                    flb_mp_map_header_append(&mh_tmp);
614
0
                    msgpack_pack_str(mp_pck, 10);
615
0
                    msgpack_pack_str_body(mp_pck, "attributes", 10);
616
0
                    ret = otel_pack_kvarray(mp_pck,
617
0
                                            scope->attributes,
618
0
                                            scope->n_attributes);
619
0
                    if (ret != 0) {
620
0
                        goto binary_payload_to_msgpack_end;
621
0
                    }
622
0
                }
623
624
0
                if (scope->dropped_attributes_count > 0) {
625
0
                    flb_mp_map_header_append(&mh_tmp);
626
0
                    msgpack_pack_str(mp_pck, 24);
627
0
                    msgpack_pack_str_body(mp_pck, "dropped_attributes_count", 24);
628
0
                    msgpack_pack_uint64(mp_pck, scope->dropped_attributes_count);
629
0
                }
630
631
0
                flb_mp_map_header_end(&mh_tmp);
632
0
            }
633
0
            else {
634
0
                flb_mp_map_header_init(&mh_tmp, mp_pck);
635
636
0
                if (scope_has_schema_url == FLB_TRUE) {
637
0
                    flb_mp_map_header_append(&mh_tmp);
638
0
                    msgpack_pack_str(mp_pck, 10);
639
0
                    msgpack_pack_str_body(mp_pck, "schema_url", 10);
640
641
0
                    len = strlen(scope_log->schema_url);
642
0
                    msgpack_pack_str(mp_pck, len);
643
0
                    msgpack_pack_str_body(mp_pck, scope_log->schema_url, len);
644
0
                }
645
646
0
                flb_mp_map_header_end(&mh_tmp);
647
0
            }
648
649
0
            flb_mp_map_header_end(&mh);
650
651
0
            ret = flb_log_event_encoder_dynamic_field_flush(&encoder->body);
652
0
            if (ret != FLB_EVENT_ENCODER_SUCCESS) {
653
0
                flb_plg_error(ctx->ins, "could not set group content metadata: %s",
654
0
                              flb_log_event_encoder_get_error_description(ret));
655
0
                goto binary_payload_to_msgpack_end;
656
0
            }
657
658
0
            flb_log_event_encoder_group_header_end(encoder);
659
660
0
            for (log_record_index=0; log_record_index < scope_log->n_log_records; log_record_index++) {
661
0
                ret = flb_log_event_encoder_begin_record(encoder);
662
663
0
                if (ret == FLB_EVENT_ENCODER_SUCCESS) {
664
0
                    if (log_records[log_record_index]->time_unix_nano > 0) {
665
0
                        ret = flb_time_from_uint64(
666
0
                                &tm,
667
0
                                log_records[log_record_index]->time_unix_nano);
668
0
                        if (ret == 0) {
669
0
                            ret = flb_log_event_encoder_set_timestamp(encoder, &tm);
670
0
                        }
671
0
                    }
672
0
                    else if (log_records[log_record_index]->observed_time_unix_nano > 0) {
673
0
                        ret = flb_time_from_uint64(
674
0
                                &tm,
675
0
                                log_records[log_record_index]->observed_time_unix_nano);
676
0
                        if (ret == 0) {
677
0
                            ret = flb_log_event_encoder_set_timestamp(encoder, &tm);
678
0
                        }
679
0
                    }
680
0
                    else {
681
0
                        flb_time_get(&tm);
682
0
                        ret = flb_log_event_encoder_set_timestamp(encoder, &tm);
683
0
                    }
684
0
                }
685
686
0
                if (ret == FLB_EVENT_ENCODER_SUCCESS) {
687
0
                    ret = flb_log_event_encoder_dynamic_field_reset(&encoder->metadata);
688
0
                    if (ret != FLB_EVENT_ENCODER_SUCCESS) {
689
0
                        flb_plg_error(ctx->ins, "failed to reset log event metadata: %s",
690
0
                                      flb_log_event_encoder_get_error_description(ret));
691
0
                        ret = FLB_EVENT_ENCODER_ERROR_SERIALIZATION_FAILURE;
692
0
                    }
693
0
                    else {
694
0
                        ret = otel_pack_v1_metadata(ctx,
695
0
                                                    mp_pck_meta,
696
0
                                                    log_records[log_record_index],
697
0
                                                    resource,
698
0
                                                    scope_log->scope);
699
0
                    }
700
701
0
                    if (ret != 0) {
702
0
                        flb_plg_error(ctx->ins, "failed to convert log record");
703
0
                        ret = FLB_EVENT_ENCODER_ERROR_SERIALIZATION_FAILURE;
704
0
                    }
705
0
                    else {
706
0
                        ret = flb_log_event_encoder_dynamic_field_flush(&encoder->metadata);
707
0
                    }
708
0
                }
709
710
0
                if (ret == FLB_EVENT_ENCODER_SUCCESS) {
711
0
                    ret = flb_log_event_encoder_dynamic_field_reset(&encoder->body);
712
0
                    if (ret != FLB_EVENT_ENCODER_SUCCESS) {
713
0
                        flb_plg_error(ctx->ins, "failed to reset log event body: %s",
714
0
                                      flb_log_event_encoder_get_error_description(ret));
715
0
                        ret = FLB_EVENT_ENCODER_ERROR_SERIALIZATION_FAILURE;
716
0
                    }
717
0
                    else if (ctx->logs_body_key == NULL &&
718
0
                             log_records[log_record_index]->body != NULL &&
719
0
                             log_records[log_record_index]->body->value_case ==
720
0
                             OPENTELEMETRY__PROTO__COMMON__V1__ANY_VALUE__VALUE_KVLIST_VALUE) {
721
0
                        ret = otlp_pack_any_value(
722
0
                                mp_pck,
723
0
                                log_records[log_record_index]->body);
724
0
                    }
725
0
                    else {
726
0
                        logs_body_key = ctx->logs_body_key;
727
0
                        if (logs_body_key == NULL) {
728
0
                            logs_body_key = "log";
729
0
                        }
730
0
                        ret = msgpack_pack_map(mp_pck, 1);
731
0
                        if (ret == 0) {
732
0
                            ret = msgpack_pack_str(mp_pck, strlen(logs_body_key));
733
0
                        }
734
0
                        if (ret == 0) {
735
0
                            ret = msgpack_pack_str_body(mp_pck,
736
0
                                                        logs_body_key,
737
0
                                                        strlen(logs_body_key));
738
0
                        }
739
0
                        if (ret == 0) {
740
0
                            ret = otlp_pack_any_value(
741
0
                                    mp_pck,
742
0
                                    log_records[log_record_index]->body);
743
0
                        }
744
0
                    }
745
746
0
                    if (ret != 0) {
747
0
                        flb_plg_error(ctx->ins, "failed to convert log record body");
748
0
                        ret = FLB_EVENT_ENCODER_ERROR_SERIALIZATION_FAILURE;
749
0
                    }
750
0
                    else {
751
0
                        ret = flb_log_event_encoder_dynamic_field_flush(&encoder->body);
752
0
                    }
753
0
                }
754
755
0
                if (ret == FLB_EVENT_ENCODER_SUCCESS) {
756
0
                    ret = flb_log_event_encoder_commit_record(encoder);
757
0
                    if (ret == FLB_EVENT_ENCODER_SUCCESS) {
758
0
                        (*record_count)++;
759
0
                    }
760
0
                }
761
0
                else {
762
0
                    flb_plg_error(ctx->ins, "marshalling error");
763
0
                    goto binary_payload_to_msgpack_end;
764
0
                }
765
0
            }
766
767
0
            flb_log_event_encoder_group_end(encoder);
768
769
0
        }
770
0
    }
771
772
0
 binary_payload_to_msgpack_end:
773
0
    if (input_logs) {
774
0
        protobuf_logs_free(input_logs, protobuf_allocator);
775
0
    }
776
0
    cfl_arena_destroy(protobuf_arena);
777
778
0
    if (ret != 0) {
779
0
        return -1;
780
0
    }
781
782
0
    return 0;
783
0
}
784
785
/*
786
 * Main function used from opentelemetry_prot.c to process logs either in JSON or Protobuf format.
787
 * -----------------------------------------------------------------------------------------------
788
 */
789
int opentelemetry_process_logs(struct flb_opentelemetry *ctx,
790
                               flb_sds_t content_type,
791
                               flb_sds_t tag,
792
                               size_t tag_len,
793
                               void *data, size_t size)
794
0
{
795
0
    int ret = -1;
796
0
    int is_proto = FLB_FALSE; /* default to JSON */
797
0
    int error_status = 0;
798
0
    char *buf;
799
0
    uint8_t *payload;
800
0
    uint64_t payload_size;
801
0
    size_t record_count;
802
0
    struct flb_log_event_encoder *encoder;
803
804
0
    buf = (char *) data;
805
0
    payload = data;
806
0
    payload_size = size;
807
0
    record_count = 0;
808
809
    /* Detect the type of payload */
810
0
    if (content_type) {
811
0
        if (opentelemetry_is_json_content_type(content_type) == FLB_TRUE) {
812
0
            if (opentelemetry_payload_starts_with_json_object(buf, size) != FLB_TRUE) {
813
0
                flb_plg_error(ctx->ins, "Invalid JSON payload");
814
0
                return -1;
815
0
            }
816
0
            is_proto = FLB_FALSE;
817
0
        }
818
0
        else if (opentelemetry_is_protobuf_content_type(content_type) == FLB_TRUE) {
819
0
            is_proto = FLB_TRUE;
820
0
        }
821
0
        else {
822
0
            flb_plg_error(ctx->ins, "Unsupported content type %s", content_type);
823
0
            return -1;
824
0
        }
825
0
    }
826
827
0
    encoder = flb_log_event_encoder_create(FLB_LOG_EVENT_FORMAT_FLUENT_BIT_V2);
828
0
    if (encoder == NULL) {
829
0
        return -1;
830
0
    }
831
832
0
    if (is_proto == FLB_TRUE) {
833
0
        ret = binary_payload_to_msgpack(ctx, encoder,
834
0
                                        tag, tag_len,
835
0
                                        (uint8_t *) payload, payload_size,
836
0
                                        &record_count);
837
0
        if (ret < 0) {
838
0
            flb_plg_error(ctx->ins, "failed to process logs from protobuf payload");
839
0
        }
840
0
    }
841
0
    else {
842
0
        ret = flb_opentelemetry_logs_json_to_msgpack(encoder,
843
0
                                                     (const char *) payload, payload_size,
844
0
                                                     ctx->logs_body_key,
845
0
                                                     &error_status);
846
0
        if (ret != 0) {
847
            /* we are printing the error for now, let's see what is the user's preference later */
848
0
            flb_plg_error(ctx->ins, "failed to process logs from JSON payload (%i) %s",
849
0
                          error_status,
850
0
                          flb_opentelemetry_error_to_string(error_status));
851
0
        }
852
853
0
    }
854
855
0
    if (ret >= 0) {
856
0
        if (opentelemetry_uses_worker_ingress_queue(ctx)) {
857
0
            size_t allocation_size;
858
0
            void *resized_buffer;
859
860
0
            allocation_size = encoder->buffer.alloc;
861
862
0
            if (allocation_size > encoder->output_length) {
863
0
                resized_buffer = flb_realloc(encoder->output_buffer,
864
0
                                             encoder->output_length);
865
0
                if (resized_buffer != NULL) {
866
0
                    encoder->buffer.data = resized_buffer;
867
0
                    encoder->output_buffer = resized_buffer;
868
0
                    encoder->buffer.alloc = encoder->output_length;
869
0
                    allocation_size = encoder->output_length;
870
0
                }
871
0
            }
872
873
0
            ret = opentelemetry_ingest_logs_take(ctx,
874
0
                                                 record_count,
875
0
                                                 tag,
876
0
                                                 flb_sds_len(tag),
877
0
                                                 encoder->output_buffer,
878
0
                                                 encoder->output_length,
879
0
                                                 allocation_size);
880
0
            flb_log_event_encoder_claim_internal_buffer_ownership(encoder);
881
0
        }
882
0
        else {
883
0
            ret = opentelemetry_ingest_logs(ctx,
884
0
                                            tag,
885
0
                                            flb_sds_len(tag),
886
0
                                            encoder->output_buffer,
887
0
                                            encoder->output_length);
888
0
        }
889
890
0
        if (ret != 0) {
891
0
            flb_plg_error(ctx->ins, "failed to append logs to the input buffer");
892
0
        }
893
0
    }
894
895
0
    flb_log_event_encoder_destroy(encoder);
896
0
    return ret;
897
0
}