/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 | } |