From dc988a09a9a4168cf6b4c4d0034d03e5655d263c Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 14:46:37 +0900 Subject: [PATCH 1/8] http_client: Set up recycle connection properly Signed-off-by: Hiroshi Hatake --- src/flb_http_client.c | 1 + 1 file changed, 1 insertion(+) diff --git a/src/flb_http_client.c b/src/flb_http_client.c index 3b55b86f6ad..f42c625b90b 100644 --- a/src/flb_http_client.c +++ b/src/flb_http_client.c @@ -2560,6 +2560,7 @@ struct flb_http_client_session *flb_http_client_session_begin(struct flb_http_cl if (protocol_version == HTTP_PROTOCOL_VERSION_20) { flb_stream_disable_keepalive(&upstream->base); + flb_upstream_conn_recycle(connection, FLB_FALSE); } session = flb_http_client_session_create(client, protocol_version, connection); From 30b7f20cc777941fdd9b7e3f1e7fdf8b0d681d50 Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 14:47:37 +0900 Subject: [PATCH 2/8] opentelemetry: Implement otel batch operations for metrics Signed-off-by: Hiroshi Hatake --- include/fluent-bit/flb_opentelemetry.h | 35 + src/CMakeLists.txt | 1 + .../flb_opentelemetry_metrics_batch.c | 838 ++++++++++++++++++ 3 files changed, 874 insertions(+) create mode 100644 src/opentelemetry/flb_opentelemetry_metrics_batch.c diff --git a/include/fluent-bit/flb_opentelemetry.h b/include/fluent-bit/flb_opentelemetry.h index cdb529e0aef..3e65d9643fc 100644 --- a/include/fluent-bit/flb_opentelemetry.h +++ b/include/fluent-bit/flb_opentelemetry.h @@ -95,6 +95,16 @@ struct flb_otel_error_map { struct cmt; struct ctrace; + +struct flb_opentelemetry_metrics_proto_batch { + flb_sds_t payload; + size_t data_point_count; +}; + +struct flb_opentelemetry_metrics_proto_batches { + size_t count; + struct flb_opentelemetry_metrics_proto_batch *entries; +}; struct flb_log_event; enum flb_opentelemetry_otlp_json_result { @@ -263,6 +273,31 @@ flb_sds_t flb_opentelemetry_metrics_msgpack_to_otlp_proto(const void *data, size_t size, int *result); +/* + * Split an encoded OTLP ExportMetricsServiceRequest without changing metric, + * resource, or scope metadata. A zero limit returns the original request as a + * single batch. The returned payloads are owned by the batch collection. + */ +struct flb_opentelemetry_metrics_proto_batches * +flb_opentelemetry_metrics_proto_batches_create(const void *payload, + size_t payload_size, + size_t max_data_points, + int *result); + +struct flb_opentelemetry_metrics_proto_batches * +flb_opentelemetry_metrics_to_otlp_proto_batches(struct cmt *context, + size_t max_data_points, + int *result); + +struct flb_opentelemetry_metrics_proto_batches * +flb_opentelemetry_metrics_msgpack_to_otlp_proto_batches(const void *data, + size_t size, + size_t max_data_points, + int *result); + +void flb_opentelemetry_metrics_proto_batches_destroy( + struct flb_opentelemetry_metrics_proto_batches *batches); + flb_sds_t flb_opentelemetry_logs_to_otlp_proto(const void *event_chunk_data, size_t event_chunk_size, struct flb_opentelemetry_otlp_logs_options *options, diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 178fe16a534..7a10aae0bb0 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -122,6 +122,7 @@ set(src ${src} opentelemetry/flb_opentelemetry_logs.c opentelemetry/flb_opentelemetry_metrics.c + opentelemetry/flb_opentelemetry_metrics_batch.c opentelemetry/flb_opentelemetry_otlp_json.c opentelemetry/flb_opentelemetry_otlp_proto.c opentelemetry/flb_opentelemetry_traces.c diff --git a/src/opentelemetry/flb_opentelemetry_metrics_batch.c b/src/opentelemetry/flb_opentelemetry_metrics_batch.c new file mode 100644 index 00000000000..ea33ec4c536 --- /dev/null +++ b/src/opentelemetry/flb_opentelemetry_metrics_batch.c @@ -0,0 +1,838 @@ +/* -*- Mode: C; tab-width: 4; indent-tabs-mode: nil; c-basic-offset: 4 -*- */ + +/* Fluent Bit + * ========== + * Copyright (C) 2015-2026 The Fluent Bit Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include +#include +#include + +#include +#include +#include + +#include + +struct metrics_batch_view { + size_t data_point_count; + Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest request; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *active_resource; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *active_scope; + const Opentelemetry__Proto__Metrics__V1__ResourceMetrics *source_resource; + const Opentelemetry__Proto__Metrics__V1__ScopeMetrics *source_scope; +}; + +static void set_result(int *result, int value) +{ + if (result != NULL) { + *result = value; + } +} + +static void metrics_batch_view_init(struct metrics_batch_view *batch) +{ + memset(batch, 0, sizeof(struct metrics_batch_view)); + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__init( + &batch->request); +} + +static void metric_view_destroy(Opentelemetry__Proto__Metrics__V1__Metric *metric) +{ + if (metric == NULL) { + return; + } + + if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE) { + flb_free(metric->gauge); + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM) { + flb_free(metric->sum); + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM) { + flb_free(metric->histogram); + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM) { + flb_free(metric->exponential_histogram); + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY) { + flb_free(metric->summary); + } + + flb_free(metric); +} + +static void metrics_batch_view_destroy(struct metrics_batch_view *batch) +{ + size_t resource_index; + size_t scope_index; + size_t metric_index; + Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scope; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resource; + + for (resource_index = 0; + resource_index < batch->request.n_resource_metrics; + resource_index++) { + resource = batch->request.resource_metrics[resource_index]; + + for (scope_index = 0; + scope_index < resource->n_scope_metrics; + scope_index++) { + scope = resource->scope_metrics[scope_index]; + + for (metric_index = 0; + metric_index < scope->n_metrics; + metric_index++) { + metric = scope->metrics[metric_index]; + metric_view_destroy(metric); + } + + flb_free(scope->metrics); + flb_free(scope); + } + + flb_free(resource->scope_metrics); + flb_free(resource); + } + + flb_free(batch->request.resource_metrics); + metrics_batch_view_init(batch); +} + +void flb_opentelemetry_metrics_proto_batches_destroy( + struct flb_opentelemetry_metrics_proto_batches *batches) +{ + size_t index; + + if (batches == NULL) { + return; + } + + for (index = 0; index < batches->count; index++) { + flb_sds_destroy(batches->entries[index].payload); + } + + flb_free(batches->entries); + flb_free(batches); +} + +static int batches_append(struct flb_opentelemetry_metrics_proto_batches *batches, + flb_sds_t payload, + size_t data_point_count) +{ + size_t count; + struct flb_opentelemetry_metrics_proto_batch *entries; + + count = batches->count; + if (count >= SIZE_MAX / sizeof(struct flb_opentelemetry_metrics_proto_batch)) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + entries = flb_realloc( + batches->entries, + (count + 1) * sizeof(struct flb_opentelemetry_metrics_proto_batch)); + if (entries == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + entries[count].payload = payload; + entries[count].data_point_count = data_point_count; + batches->entries = entries; + batches->count++; + + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; +} + +static int batch_add_resource( + struct metrics_batch_view *batch, + const Opentelemetry__Proto__Metrics__V1__ResourceMetrics *source) +{ + size_t count; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resource; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics **resources; + + resource = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__ResourceMetrics)); + if (resource == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + *resource = *source; + resource->n_scope_metrics = 0; + resource->scope_metrics = NULL; + + count = batch->request.n_resource_metrics; + if (count >= SIZE_MAX / sizeof(Opentelemetry__Proto__Metrics__V1__ResourceMetrics *)) { + flb_free(resource); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + resources = flb_realloc( + batch->request.resource_metrics, + (count + 1) * + sizeof(Opentelemetry__Proto__Metrics__V1__ResourceMetrics *)); + if (resources == NULL) { + flb_free(resource); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + resources[count] = resource; + batch->request.resource_metrics = resources; + batch->request.n_resource_metrics++; + batch->active_resource = resource; + batch->active_scope = NULL; + batch->source_resource = source; + batch->source_scope = NULL; + + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; +} + +static int batch_ensure_resource( + struct metrics_batch_view *batch, + const Opentelemetry__Proto__Metrics__V1__ResourceMetrics *source) +{ + if (batch->source_resource == source && batch->active_resource != NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + } + + return batch_add_resource(batch, source); +} + +static int batch_add_scope( + struct metrics_batch_view *batch, + const Opentelemetry__Proto__Metrics__V1__ScopeMetrics *source) +{ + size_t count; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scope; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics **scopes; + + scope = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__ScopeMetrics)); + if (scope == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + *scope = *source; + scope->n_metrics = 0; + scope->metrics = NULL; + + count = batch->active_resource->n_scope_metrics; + if (count >= SIZE_MAX / sizeof(Opentelemetry__Proto__Metrics__V1__ScopeMetrics *)) { + flb_free(scope); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + scopes = flb_realloc( + batch->active_resource->scope_metrics, + (count + 1) * sizeof(Opentelemetry__Proto__Metrics__V1__ScopeMetrics *)); + if (scopes == NULL) { + flb_free(scope); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + scopes[count] = scope; + batch->active_resource->scope_metrics = scopes; + batch->active_resource->n_scope_metrics++; + batch->active_scope = scope; + batch->source_scope = source; + + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; +} + +static int batch_ensure_scope( + struct metrics_batch_view *batch, + const Opentelemetry__Proto__Metrics__V1__ResourceMetrics *source_resource, + const Opentelemetry__Proto__Metrics__V1__ScopeMetrics *source_scope) +{ + int result; + + result = batch_ensure_resource(batch, source_resource); + if (result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + return result; + } + + if (batch->source_scope == source_scope && batch->active_scope != NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + } + + return batch_add_scope(batch, source_scope); +} + +static int metric_data_point_count( + const Opentelemetry__Proto__Metrics__V1__Metric *metric, + size_t *count) +{ + *count = 0; + + if (metric == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + + if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA__NOT_SET) { + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE) { + if (metric->gauge == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + *count = metric->gauge->n_data_points; + if (*count > 0 && metric->gauge->data_points == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM) { + if (metric->sum == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + *count = metric->sum->n_data_points; + if (*count > 0 && metric->sum->data_points == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM) { + if (metric->histogram == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + *count = metric->histogram->n_data_points; + if (*count > 0 && metric->histogram->data_points == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM) { + if (metric->exponential_histogram == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + *count = metric->exponential_histogram->n_data_points; + if (*count > 0 && metric->exponential_histogram->data_points == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY) { + if (metric->summary == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + *count = metric->summary->n_data_points; + if (*count > 0 && metric->summary->data_points == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + } + else { + /* Preserve metric types added by newer OTLP schemas. */ + *count = 0; + } + + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; +} + +static Opentelemetry__Proto__Metrics__V1__Metric *metric_view_create( + const Opentelemetry__Proto__Metrics__V1__Metric *source, + size_t offset, + size_t count) +{ + Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__Gauge *gauge; + Opentelemetry__Proto__Metrics__V1__Sum *sum; + Opentelemetry__Proto__Metrics__V1__Histogram *histogram; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogram *exp_histogram; + Opentelemetry__Proto__Metrics__V1__Summary *summary; + + metric = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__Metric)); + if (metric == NULL) { + return NULL; + } + + *metric = *source; + + if (source->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE) { + gauge = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__Gauge)); + if (gauge == NULL) { + flb_free(metric); + return NULL; + } + *gauge = *source->gauge; + gauge->n_data_points = count; + gauge->data_points = count > 0 ? source->gauge->data_points + offset : NULL; + metric->gauge = gauge; + } + else if (source->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM) { + sum = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__Sum)); + if (sum == NULL) { + flb_free(metric); + return NULL; + } + *sum = *source->sum; + sum->n_data_points = count; + sum->data_points = count > 0 ? source->sum->data_points + offset : NULL; + metric->sum = sum; + } + else if (source->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM) { + histogram = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__Histogram)); + if (histogram == NULL) { + flb_free(metric); + return NULL; + } + *histogram = *source->histogram; + histogram->n_data_points = count; + histogram->data_points = count > 0 ? source->histogram->data_points + offset : NULL; + metric->histogram = histogram; + } + else if (source->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM) { + exp_histogram = flb_calloc( + 1, + sizeof(Opentelemetry__Proto__Metrics__V1__ExponentialHistogram)); + if (exp_histogram == NULL) { + flb_free(metric); + return NULL; + } + *exp_histogram = *source->exponential_histogram; + exp_histogram->n_data_points = count; + exp_histogram->data_points = count > 0 ? + source->exponential_histogram->data_points + offset : + NULL; + metric->exponential_histogram = exp_histogram; + } + else if (source->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY) { + summary = flb_calloc(1, sizeof(Opentelemetry__Proto__Metrics__V1__Summary)); + if (summary == NULL) { + flb_free(metric); + return NULL; + } + *summary = *source->summary; + summary->n_data_points = count; + summary->data_points = count > 0 ? source->summary->data_points + offset : NULL; + metric->summary = summary; + } + + return metric; +} + +static int batch_add_metric( + struct metrics_batch_view *batch, + const Opentelemetry__Proto__Metrics__V1__ResourceMetrics *source_resource, + const Opentelemetry__Proto__Metrics__V1__ScopeMetrics *source_scope, + const Opentelemetry__Proto__Metrics__V1__Metric *source_metric, + size_t offset, + size_t count) +{ + int result; + size_t metric_count; + Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__Metric **metrics; + + result = batch_ensure_scope(batch, source_resource, source_scope); + if (result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + return result; + } + + metric = metric_view_create(source_metric, offset, count); + if (metric == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + metric_count = batch->active_scope->n_metrics; + if (metric_count >= SIZE_MAX / sizeof(Opentelemetry__Proto__Metrics__V1__Metric *)) { + metric_view_destroy(metric); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + metrics = flb_realloc( + batch->active_scope->metrics, + (metric_count + 1) * sizeof(Opentelemetry__Proto__Metrics__V1__Metric *)); + if (metrics == NULL) { + metric_view_destroy(metric); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + metrics[metric_count] = metric; + batch->active_scope->metrics = metrics; + batch->active_scope->n_metrics++; + batch->data_point_count += count; + + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; +} + +static int batch_pack(struct metrics_batch_view *batch, + struct flb_opentelemetry_metrics_proto_batches *batches) +{ + int result; + size_t payload_size; + size_t packed_size; + flb_sds_t payload; + + if (batch->request.n_resource_metrics == 0) { + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + } + + payload_size = + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__get_packed_size( + &batch->request); + payload = flb_sds_create_size(payload_size); + if (payload == NULL) { + metrics_batch_view_destroy(batch); + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + packed_size = + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__pack( + &batch->request, + (uint8_t *) payload); + if (packed_size != payload_size) { + flb_sds_destroy(payload); + metrics_batch_view_destroy(batch); + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + + flb_sds_len_set(payload, packed_size); + result = batches_append(batches, payload, batch->data_point_count); + if (result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + flb_sds_destroy(payload); + } + + metrics_batch_view_destroy(batch); + + return result; +} + +static int request_data_point_count( + Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest *request, + size_t *total) +{ + int result; + size_t resource_index; + size_t scope_index; + size_t metric_index; + size_t count; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scope; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resource; + + *total = 0; + + for (resource_index = 0; resource_index < request->n_resource_metrics; resource_index++) { + resource = request->resource_metrics[resource_index]; + if (resource == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + + for (scope_index = 0; scope_index < resource->n_scope_metrics; scope_index++) { + scope = resource->scope_metrics[scope_index]; + if (scope == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT; + } + + for (metric_index = 0; metric_index < scope->n_metrics; metric_index++) { + result = metric_data_point_count(scope->metrics[metric_index], &count); + if (result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + return result; + } + if (count > SIZE_MAX - *total) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + *total += count; + } + } + } + + return FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; +} + +static int append_original_payload( + struct flb_opentelemetry_metrics_proto_batches *batches, + const void *payload, + size_t payload_size, + size_t data_point_count) +{ + int result; + flb_sds_t copy; + + copy = flb_sds_create_size(payload_size); + if (copy == NULL) { + return FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED; + } + + memcpy(copy, payload, payload_size); + flb_sds_len_set(copy, payload_size); + result = batches_append(batches, copy, data_point_count); + if (result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + flb_sds_destroy(copy); + } + + return result; +} + +struct flb_opentelemetry_metrics_proto_batches * +flb_opentelemetry_metrics_proto_batches_create(const void *payload, + size_t payload_size, + size_t max_data_points, + int *result) +{ + int local_result; + size_t resource_index; + size_t scope_index; + size_t metric_index; + size_t data_point_count; + size_t offset; + size_t remaining; + size_t batch_count; + size_t total_data_points; + Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scope; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resource; + Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest *request; + struct flb_opentelemetry_metrics_proto_batches *batches; + struct metrics_batch_view batch; + + if (payload == NULL || payload_size == 0) { + errno = EINVAL; + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT); + return NULL; + } + + batches = flb_calloc(1, sizeof(struct flb_opentelemetry_metrics_proto_batches)); + if (batches == NULL) { + errno = ENOMEM; + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED); + return NULL; + } + + request = + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__unpack( + NULL, + payload_size, + (const uint8_t *) payload); + if (request == NULL) { + errno = EINVAL; + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT); + flb_opentelemetry_metrics_proto_batches_destroy(batches); + return NULL; + } + + local_result = request_data_point_count(request, &total_data_points); + if (local_result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + goto error; + } + + if (max_data_points == 0 || total_data_points <= max_data_points) { + local_result = append_original_payload(batches, + payload, + payload_size, + total_data_points); + if (local_result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + goto error; + } + + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__free_unpacked( + request, + NULL); + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + return batches; + } + + local_result = FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + metrics_batch_view_init(&batch); + + for (resource_index = 0; + resource_index < request->n_resource_metrics && + local_result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + resource_index++) { + resource = request->resource_metrics[resource_index]; + + if (resource->n_scope_metrics == 0) { + local_result = batch_ensure_resource(&batch, resource); + continue; + } + + for (scope_index = 0; + scope_index < resource->n_scope_metrics && + local_result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + scope_index++) { + scope = resource->scope_metrics[scope_index]; + + if (scope->n_metrics == 0) { + local_result = batch_ensure_scope(&batch, resource, scope); + continue; + } + + for (metric_index = 0; + metric_index < scope->n_metrics && + local_result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS; + metric_index++) { + metric = scope->metrics[metric_index]; + local_result = metric_data_point_count(metric, &data_point_count); + if (local_result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + break; + } + + if (data_point_count == 0) { + local_result = batch_add_metric(&batch, + resource, + scope, + metric, + 0, + 0); + continue; + } + + offset = 0; + while (offset < data_point_count && + local_result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + if (batch.data_point_count == max_data_points) { + local_result = batch_pack(&batch, batches); + if (local_result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + break; + } + } + + remaining = data_point_count - offset; + batch_count = max_data_points - batch.data_point_count; + if (batch_count > remaining) { + batch_count = remaining; + } + + local_result = batch_add_metric(&batch, + resource, + scope, + metric, + offset, + batch_count); + offset += batch_count; + } + } + } + } + + if (local_result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS && + batch.request.n_resource_metrics > 0) { + local_result = batch_pack(&batch, batches); + } + else { + metrics_batch_view_destroy(&batch); + } + + if (local_result != FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS) { + goto error; + } + + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__free_unpacked( + request, + NULL); + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + return batches; + +error: + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__free_unpacked( + request, + NULL); + flb_opentelemetry_metrics_proto_batches_destroy(batches); + if (local_result == FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED) { + errno = ENOMEM; + } + else { + errno = EINVAL; + } + set_result(result, local_result); + return NULL; +} + +static struct flb_opentelemetry_metrics_proto_batches *empty_batches_create(int *result) +{ + struct flb_opentelemetry_metrics_proto_batches *batches; + + batches = flb_calloc(1, sizeof(struct flb_opentelemetry_metrics_proto_batches)); + if (batches == NULL) { + errno = ENOMEM; + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED); + return NULL; + } + + set_result(result, FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + return batches; +} + +struct flb_opentelemetry_metrics_proto_batches * +flb_opentelemetry_metrics_to_otlp_proto_batches(struct cmt *context, + size_t max_data_points, + int *result) +{ + int encode_result; + flb_sds_t payload; + struct flb_opentelemetry_metrics_proto_batches *batches; + + payload = flb_opentelemetry_metrics_to_otlp_proto(context, &encode_result); + if (payload == NULL) { + set_result(result, encode_result); + return NULL; + } + + if (flb_sds_len(payload) == 0) { + flb_opentelemetry_metrics_proto_destroy(payload); + return empty_batches_create(result); + } + + batches = flb_opentelemetry_metrics_proto_batches_create(payload, + flb_sds_len(payload), + max_data_points, + result); + flb_opentelemetry_metrics_proto_destroy(payload); + + return batches; +} + +struct flb_opentelemetry_metrics_proto_batches * +flb_opentelemetry_metrics_msgpack_to_otlp_proto_batches(const void *data, + size_t size, + size_t max_data_points, + int *result) +{ + int encode_result; + flb_sds_t payload; + struct flb_opentelemetry_metrics_proto_batches *batches; + + payload = flb_opentelemetry_metrics_msgpack_to_otlp_proto(data, + size, + &encode_result); + if (payload == NULL) { + set_result(result, encode_result); + return NULL; + } + + if (flb_sds_len(payload) == 0) { + flb_opentelemetry_metrics_proto_destroy(payload); + return empty_batches_create(result); + } + + batches = flb_opentelemetry_metrics_proto_batches_create(payload, + flb_sds_len(payload), + max_data_points, + result); + flb_opentelemetry_metrics_proto_destroy(payload); + + return batches; +} From 623eb4a9f11190e3063f5f593e5e66fa9980c808 Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 14:48:50 +0900 Subject: [PATCH 3/8] out_opentelemetry: Support batch operations for metrics with max data points Signed-off-by: Hiroshi Hatake --- plugins/out_opentelemetry/opentelemetry.c | 80 +++++++++++++++++++++-- plugins/out_opentelemetry/opentelemetry.h | 4 ++ 2 files changed, 79 insertions(+), 5 deletions(-) diff --git a/plugins/out_opentelemetry/opentelemetry.c b/plugins/out_opentelemetry/opentelemetry.c index 50bd7e05edc..bc94f6db6cb 100644 --- a/plugins/out_opentelemetry/opentelemetry.c +++ b/plugins/out_opentelemetry/opentelemetry.c @@ -25,6 +25,7 @@ #include #include #include +#include #include #include @@ -730,6 +731,67 @@ static int opentelemetry_format_test(struct flb_config *config, return 0; } +static int post_metrics_payload(struct opentelemetry_context *ctx, + struct flb_event_chunk *event_chunk, + flb_sds_t payload) +{ + int result; + int encode_result; + size_t index; + struct flb_opentelemetry_metrics_proto_batches *batches; + + if (ctx->metrics_max_datapoints == 0) { + return opentelemetry_post(ctx, + payload, + flb_sds_len(payload), + event_chunk->tag, + flb_sds_len(event_chunk->tag), + ctx->metrics_uri_sanitized, + ctx->grpc_metrics_uri); + } + + batches = flb_opentelemetry_metrics_proto_batches_create( + payload, + flb_sds_len(payload), + (size_t) ctx->metrics_max_datapoints, + &encode_result); + if (batches == NULL) { + flb_plg_error(ctx->ins, + "could not split metric payload into batches: %i", + encode_result); + if (encode_result == FLB_OPENTELEMETRY_OTLP_PROTO_NOT_SUPPORTED) { + return FLB_RETRY; + } + return FLB_ERROR; + } + + result = FLB_OK; + for (index = 0; index < batches->count; index++) { + result = opentelemetry_post(ctx, + batches->entries[index].payload, + flb_sds_len(batches->entries[index].payload), + event_chunk->tag, + flb_sds_len(event_chunk->tag), + ctx->metrics_uri_sanitized, + ctx->grpc_metrics_uri); + if (result != FLB_OK) { + if (result == FLB_RETRY && index > 0) { + flb_plg_warn(ctx->ins, + "metric payload partially succeeded (%zu/%zu batches); " + "skipping retry to avoid resending accepted data", + index, + batches->count); + result = FLB_OK; + } + break; + } + } + + flb_opentelemetry_metrics_proto_batches_destroy(batches); + + return result; +} + static int process_metrics(struct flb_event_chunk *event_chunk, struct flb_output_flush *out1_flush, struct flb_input_instance *ins, void *out_context, @@ -796,11 +858,7 @@ static int process_metrics(struct flb_event_chunk *event_chunk, flb_plg_debug(ctx->ins, "final payload size: %lu", flb_sds_len(buf)); if (buf && flb_sds_len(buf) > 0) { /* Send HTTP request */ - result = opentelemetry_post(ctx, buf, flb_sds_len(buf), - event_chunk->tag, - flb_sds_len(event_chunk->tag), - ctx->metrics_uri_sanitized, - ctx->grpc_metrics_uri); + result = post_metrics_payload(ctx, event_chunk, buf); /* Debug http_post() result statuses */ if (result == FLB_OK) { @@ -1020,6 +1078,12 @@ static int cb_opentelemetry_init(struct flb_output_instance *ins, ctx->batch_size = atoi(DEFAULT_LOG_RECORD_BATCH_SIZE); } + if (ctx->metrics_max_datapoints < 0) { + flb_plg_error(ins, "metrics_max_datapoints must be zero or greater"); + flb_opentelemetry_context_destroy(ctx); + return -1; + } + flb_output_set_context(ins, ctx); /* @@ -1117,6 +1181,12 @@ static struct flb_config_map config_map[] = { 0, FLB_TRUE, offsetof(struct opentelemetry_context, grpc_metrics_uri), "Specify an optional gRPC URI for the target OTel endpoint." }, + { + FLB_CONFIG_MAP_INT, "metrics_max_datapoints", DEFAULT_METRICS_MAX_DATAPOINTS, + 0, FLB_TRUE, offsetof(struct opentelemetry_context, metrics_max_datapoints), + "Set the maximum number of metric data points per OTLP export request " + "(0 disables the limit; default: 0)" + }, { FLB_CONFIG_MAP_INT, "batch_size", DEFAULT_LOG_RECORD_BATCH_SIZE, diff --git a/plugins/out_opentelemetry/opentelemetry.h b/plugins/out_opentelemetry/opentelemetry.h index b45b97c92e6..bfec5824d54 100644 --- a/plugins/out_opentelemetry/opentelemetry.h +++ b/plugins/out_opentelemetry/opentelemetry.h @@ -44,6 +44,7 @@ #define DEFAULT_LOG_RECORD_BATCH_SIZE "1000" #define DEFAULT_MAX_RESOURCE_EXPORT "0" /* no resource limits */ #define DEFAULT_MAX_SCOPE_EXPORT "0" /* no scope limits */ +#define DEFAULT_METRICS_MAX_DATAPOINTS "0" /* no data point limit */ struct opentelemetry_body_key { flb_sds_t key; @@ -145,6 +146,9 @@ struct opentelemetry_context { /* Number of logs to flush at a time */ int batch_size; + /* Maximum number of metric data points per OTLP export request */ + int metrics_max_datapoints; + /* Maximum number of resources per OTLP export */ int max_resources; From c945439431f8a468688ad62299d82dc863486cd2 Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 14:49:40 +0900 Subject: [PATCH 4/8] tests: internal: Add a test case for batch request with data points' limitation Signed-off-by: Hiroshi Hatake --- tests/internal/opentelemetry.c | 391 +++++++++++++++++++++++++++++++++ 1 file changed, 391 insertions(+) diff --git a/tests/internal/opentelemetry.c b/tests/internal/opentelemetry.c index 13e9add6538..7434df1fe89 100644 --- a/tests/internal/opentelemetry.c +++ b/tests/internal/opentelemetry.c @@ -2706,6 +2706,391 @@ void test_opentelemetry_metrics_msgpack_otlp_proto_merges_contexts() destroy_metrics_context_list(&contexts); } +void test_opentelemetry_metrics_otlp_proto_data_point_batches() +{ + int index; + int result; + int ret; + int seen[11]; + size_t batch_index; + size_t resource_index; + size_t scope_index; + size_t metric_index; + size_t point_index; + size_t total_data_points; + uint64_t timestamp; + char *label_keys[] = {"series"}; + char *label_values[1]; + char *series[] = { + "series-0", "series-1", "series-2", "series-3", "series-4", "series-5", + "series-6", "series-7", "series-8", "series-9", "series-10" + }; + struct cmt *context; + struct cmt_gauge *gauge; + struct flb_opentelemetry_metrics_proto_batches *batches; + Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scope; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resource; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint *point; + Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest *decoded; + + memset(seen, 0, sizeof(seen)); + context = cmt_create(); + TEST_CHECK(context != NULL); + if (context == NULL) { + return; + } + + gauge = cmt_gauge_create(context, + "test", + "batch", + "value", + "batching test", + 1, + label_keys); + TEST_CHECK(gauge != NULL); + if (gauge == NULL) { + cmt_destroy(context); + return; + } + + for (index = 0; index < 11; index++) { + label_values[0] = series[index]; + ret = cmt_gauge_set(gauge, + (uint64_t) index + 1, + (double) index, + 1, + label_values); + TEST_CHECK(ret == 0); + } + + batches = flb_opentelemetry_metrics_to_otlp_proto_batches(context, 4, &result); + TEST_CHECK(result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + TEST_CHECK(batches != NULL); + if (batches == NULL) { + cmt_destroy(context); + return; + } + + TEST_CHECK(batches->count == 3); + TEST_CHECK(batches->entries[0].data_point_count == 4); + TEST_CHECK(batches->entries[1].data_point_count == 4); + TEST_CHECK(batches->entries[2].data_point_count == 3); + + total_data_points = 0; + for (batch_index = 0; batch_index < batches->count; batch_index++) { + decoded = + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__unpack( + NULL, + flb_sds_len(batches->entries[batch_index].payload), + (uint8_t *) batches->entries[batch_index].payload); + TEST_CHECK(decoded != NULL); + if (decoded == NULL) { + continue; + } + + for (resource_index = 0; + resource_index < decoded->n_resource_metrics; + resource_index++) { + resource = decoded->resource_metrics[resource_index]; + for (scope_index = 0; scope_index < resource->n_scope_metrics; scope_index++) { + scope = resource->scope_metrics[scope_index]; + for (metric_index = 0; metric_index < scope->n_metrics; metric_index++) { + metric = scope->metrics[metric_index]; + TEST_CHECK(metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE); + if (metric->data_case != + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE) { + continue; + } + + for (point_index = 0; + point_index < metric->gauge->n_data_points; + point_index++) { + point = metric->gauge->data_points[point_index]; + timestamp = point->time_unix_nano; + TEST_CHECK(timestamp >= 1 && timestamp <= 11); + if (timestamp >= 1 && timestamp <= 11) { + seen[timestamp - 1]++; + } + total_data_points++; + } + } + } + } + + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__free_unpacked( + decoded, + NULL); + } + + TEST_CHECK(total_data_points == 11); + for (index = 0; index < 11; index++) { + TEST_CHECK(seen[index] == 1); + } + + flb_opentelemetry_metrics_proto_batches_destroy(batches); + + batches = flb_opentelemetry_metrics_to_otlp_proto_batches(context, 11, &result); + TEST_CHECK(result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + TEST_CHECK(batches != NULL); + if (batches != NULL) { + TEST_CHECK(batches->count == 1); + TEST_CHECK(batches->entries[0].data_point_count == 11); + flb_opentelemetry_metrics_proto_batches_destroy(batches); + } + + batches = flb_opentelemetry_metrics_to_otlp_proto_batches(context, 0, &result); + TEST_CHECK(result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + TEST_CHECK(batches != NULL); + if (batches != NULL) { + TEST_CHECK(batches->count == 1); + TEST_CHECK(batches->entries[0].data_point_count == 11); + flb_opentelemetry_metrics_proto_batches_destroy(batches); + } + + cmt_destroy(context); + + batches = flb_opentelemetry_metrics_proto_batches_create("invalid", 7, 4, &result); + TEST_CHECK(batches == NULL); + TEST_CHECK(result == FLB_OPENTELEMETRY_OTLP_PROTO_INVALID_ARGUMENT); +} + +void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() +{ + int result; + int gauge_seen; + int sum_seen; + int histogram_seen; + int exp_histogram_seen; + int summary_seen; + size_t payload_size; + size_t batch_index; + size_t resource_index; + size_t scope_index; + size_t metric_index; + size_t total_data_points; + flb_sds_t payload; + Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__Metric metrics[5]; + Opentelemetry__Proto__Metrics__V1__Metric *metric_entries[5]; + Opentelemetry__Proto__Metrics__V1__Gauge gauge; + Opentelemetry__Proto__Metrics__V1__Sum sum; + Opentelemetry__Proto__Metrics__V1__Histogram histogram; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogram exp_histogram; + Opentelemetry__Proto__Metrics__V1__Summary summary; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint gauge_point; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint sum_point; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint *gauge_points[1]; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint *sum_points[1]; + Opentelemetry__Proto__Metrics__V1__HistogramDataPoint histogram_point; + Opentelemetry__Proto__Metrics__V1__HistogramDataPoint *histogram_points[1]; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint exp_histogram_point; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint *exp_histogram_points[1]; + Opentelemetry__Proto__Metrics__V1__SummaryDataPoint summary_point; + Opentelemetry__Proto__Metrics__V1__SummaryDataPoint *summary_points[1]; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics scope; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scopes[1]; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics resource; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resources[1]; + Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest request; + Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest *decoded; + struct flb_opentelemetry_metrics_proto_batches *batches; + + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__init( + &request); + opentelemetry__proto__metrics__v1__resource_metrics__init(&resource); + opentelemetry__proto__metrics__v1__scope_metrics__init(&scope); + opentelemetry__proto__metrics__v1__gauge__init(&gauge); + opentelemetry__proto__metrics__v1__sum__init(&sum); + opentelemetry__proto__metrics__v1__histogram__init(&histogram); + opentelemetry__proto__metrics__v1__exponential_histogram__init(&exp_histogram); + opentelemetry__proto__metrics__v1__summary__init(&summary); + opentelemetry__proto__metrics__v1__number_data_point__init(&gauge_point); + opentelemetry__proto__metrics__v1__number_data_point__init(&sum_point); + opentelemetry__proto__metrics__v1__histogram_data_point__init(&histogram_point); + opentelemetry__proto__metrics__v1__exponential_histogram_data_point__init( + &exp_histogram_point); + opentelemetry__proto__metrics__v1__summary_data_point__init(&summary_point); + + for (metric_index = 0; metric_index < 5; metric_index++) { + opentelemetry__proto__metrics__v1__metric__init(&metrics[metric_index]); + metric_entries[metric_index] = &metrics[metric_index]; + } + + gauge_points[0] = &gauge_point; + gauge.n_data_points = 1; + gauge.data_points = gauge_points; + metrics[0].name = "gauge"; + metrics[0].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE; + metrics[0].gauge = &gauge; + + sum_points[0] = &sum_point; + sum.n_data_points = 1; + sum.data_points = sum_points; + metrics[1].name = "sum"; + metrics[1].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM; + metrics[1].sum = ∑ + + histogram_points[0] = &histogram_point; + histogram.n_data_points = 1; + histogram.data_points = histogram_points; + metrics[2].name = "histogram"; + metrics[2].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM; + metrics[2].histogram = &histogram; + + exp_histogram_points[0] = &exp_histogram_point; + exp_histogram.n_data_points = 1; + exp_histogram.data_points = exp_histogram_points; + metrics[3].name = "exponential_histogram"; + metrics[3].data_case = + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM; + metrics[3].exponential_histogram = &exp_histogram; + + summary_points[0] = &summary_point; + summary.n_data_points = 1; + summary.data_points = summary_points; + metrics[4].name = "summary"; + metrics[4].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY; + metrics[4].summary = &summary; + + scope.n_metrics = 5; + scope.metrics = metric_entries; + scopes[0] = &scope; + resource.n_scope_metrics = 1; + resource.scope_metrics = scopes; + resources[0] = &resource; + request.n_resource_metrics = 1; + request.resource_metrics = resources; + + payload_size = + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__get_packed_size( + &request); + payload = flb_sds_create_size(payload_size); + TEST_CHECK(payload != NULL); + if (payload == NULL) { + return; + } + + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__pack( + &request, + (uint8_t *) payload); + flb_sds_len_set(payload, payload_size); + + batches = flb_opentelemetry_metrics_proto_batches_create(payload, + payload_size, + 2, + &result); + flb_sds_destroy(payload); + TEST_CHECK(result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + TEST_CHECK(batches != NULL); + if (batches == NULL) { + return; + } + + TEST_CHECK(batches->count == 3); + TEST_CHECK(batches->entries[0].data_point_count == 2); + TEST_CHECK(batches->entries[1].data_point_count == 2); + TEST_CHECK(batches->entries[2].data_point_count == 1); + + gauge_seen = 0; + sum_seen = 0; + histogram_seen = 0; + exp_histogram_seen = 0; + summary_seen = 0; + total_data_points = 0; + + for (batch_index = 0; batch_index < batches->count; batch_index++) { + decoded = + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__unpack( + NULL, + flb_sds_len(batches->entries[batch_index].payload), + (uint8_t *) batches->entries[batch_index].payload); + TEST_CHECK(decoded != NULL); + if (decoded == NULL) { + continue; + } + + for (resource_index = 0; + resource_index < decoded->n_resource_metrics; + resource_index++) { + resource = *decoded->resource_metrics[resource_index]; + for (scope_index = 0; + scope_index < resource.n_scope_metrics; + scope_index++) { + scope = *resource.scope_metrics[scope_index]; + for (metric_index = 0; metric_index < scope.n_metrics; metric_index++) { + metric = scope.metrics[metric_index]; + if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE) { + gauge_seen++; + total_data_points += metric->gauge->n_data_points; + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM) { + sum_seen++; + total_data_points += metric->sum->n_data_points; + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM) { + histogram_seen++; + total_data_points += metric->histogram->n_data_points; + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM) { + exp_histogram_seen++; + total_data_points += metric->exponential_histogram->n_data_points; + } + else if (metric->data_case == + OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY) { + summary_seen++; + total_data_points += metric->summary->n_data_points; + } + } + } + } + + opentelemetry__proto__collector__metrics__v1__export_metrics_service_request__free_unpacked( + decoded, + NULL); + } + + TEST_CHECK(total_data_points == 5); + TEST_CHECK(gauge_seen == 1); + TEST_CHECK(sum_seen == 1); + TEST_CHECK(histogram_seen == 1); + TEST_CHECK(exp_histogram_seen == 1); + TEST_CHECK(summary_seen == 1); + + flb_opentelemetry_metrics_proto_batches_destroy(batches); +} + +void test_opentelemetry_metrics_otlp_proto_batches_empty_context() +{ + int result; + struct cmt *context; + struct flb_opentelemetry_metrics_proto_batches *batches; + + context = cmt_create(); + TEST_CHECK(context != NULL); + if (context == NULL) { + return; + } + + batches = flb_opentelemetry_metrics_to_otlp_proto_batches(context, 4, &result); + TEST_CHECK(result == FLB_OPENTELEMETRY_OTLP_PROTO_SUCCESS); + TEST_CHECK(batches != NULL); + if (batches != NULL) { + TEST_CHECK(batches->count <= 1); + if (batches->count == 1) { + TEST_CHECK(batches->entries[0].data_point_count == 0); + } + flb_opentelemetry_metrics_proto_batches_destroy(batches); + } + + cmt_destroy(context); +} + void test_opentelemetry_traces_otlp_proto_roundtrip() { int result; @@ -2780,5 +3165,11 @@ TEST_LIST = { test_opentelemetry_metrics_otlp_proto_roundtrip }, { "opentelemetry_metrics_msgpack_otlp_proto_merges_contexts", test_opentelemetry_metrics_msgpack_otlp_proto_merges_contexts }, + { "opentelemetry_metrics_otlp_proto_data_point_batches", + test_opentelemetry_metrics_otlp_proto_data_point_batches }, + { "opentelemetry_metrics_otlp_proto_batches_all_metric_types", + test_opentelemetry_metrics_otlp_proto_batches_all_metric_types }, + { "opentelemetry_metrics_otlp_proto_batches_empty_context", + test_opentelemetry_metrics_otlp_proto_batches_empty_context }, { 0 } }; From 6d2839873593f2b4c5ecf818f03fae322cafd0f7 Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 14:50:17 +0900 Subject: [PATCH 5/8] tests: integration: Add e2e cases for the limitation of max data points Signed-off-by: Hiroshi Hatake --- tests/integration/README.md | 1 + .../out_otel_grpc_metrics_max_datapoints.conf | 19 ++ .../out_otel_http_metrics_max_datapoints.conf | 17 ++ .../tests/test_out_opentelemetry_001.py | 186 ++++++++++++++++++ tests/integration/src/server/otlp_server.py | 22 ++- 5 files changed, 242 insertions(+), 3 deletions(-) create mode 100644 tests/integration/scenarios/out_opentelemetry/config/out_otel_grpc_metrics_max_datapoints.conf create mode 100644 tests/integration/scenarios/out_opentelemetry/config/out_otel_http_metrics_max_datapoints.conf diff --git a/tests/integration/README.md b/tests/integration/README.md index d7b9020cf64..e198987bae8 100644 --- a/tests/integration/README.md +++ b/tests/integration/README.md @@ -490,6 +490,7 @@ Covers: - metadata and message key mapping - `add_label` - `batch_size` +- `metrics_max_datapoints` - `logs_max_resources` - `logs_max_scopes` diff --git a/tests/integration/scenarios/out_opentelemetry/config/out_otel_grpc_metrics_max_datapoints.conf b/tests/integration/scenarios/out_opentelemetry/config/out_otel_grpc_metrics_max_datapoints.conf new file mode 100644 index 00000000000..9e476823877 --- /dev/null +++ b/tests/integration/scenarios/out_opentelemetry/config/out_otel_grpc_metrics_max_datapoints.conf @@ -0,0 +1,19 @@ +[SERVICE] + flush 1 + log_level info + http_server on + http_port ${FLUENT_BIT_HTTP_MONITORING_PORT} + +[INPUT] + name opentelemetry + port ${FLUENT_BIT_TEST_LISTENER_PORT} + +[OUTPUT] + name opentelemetry + match * + host 127.0.0.1 + port ${TEST_SUITE_HTTP_PORT} + http2 on + grpc on + grpc_metrics_uri /batched.metrics.v1.Metrics/Export + metrics_max_datapoints 4 diff --git a/tests/integration/scenarios/out_opentelemetry/config/out_otel_http_metrics_max_datapoints.conf b/tests/integration/scenarios/out_opentelemetry/config/out_otel_http_metrics_max_datapoints.conf new file mode 100644 index 00000000000..4094c2b39b7 --- /dev/null +++ b/tests/integration/scenarios/out_opentelemetry/config/out_otel_http_metrics_max_datapoints.conf @@ -0,0 +1,17 @@ +[SERVICE] + flush 1 + log_level info + http_server on + http_port ${FLUENT_BIT_HTTP_MONITORING_PORT} + +[INPUT] + name opentelemetry + port ${FLUENT_BIT_TEST_LISTENER_PORT} + +[OUTPUT] + name opentelemetry + match * + host 127.0.0.1 + port ${TEST_SUITE_HTTP_PORT} + metrics_uri /batched/metrics + metrics_max_datapoints 4 diff --git a/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py b/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py index 27ec8d14d6e..6d4d08ec8ef 100644 --- a/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py +++ b/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py @@ -22,6 +22,7 @@ ) from server.otlp_server import ( configure_otlp_grpc_methods, + configure_otlp_response, data_storage, otlp_server_run, stop_otlp_server, @@ -523,6 +524,100 @@ def scope(scope_name, scope_version, body, flags): return {"resource_logs": resource_logs} +def _build_batched_metrics_payload(): + def attributes(series_id): + return [ + { + "key": "series.id", + "value": { + "string_value": str(series_id), + }, + } + ] + + def resource(service_name, metric): + return { + "resource": { + "attributes": [ + { + "key": "service.name", + "value": { + "string_value": service_name, + }, + } + ], + }, + "scope_metrics": [ + { + "scope": { + "name": "batch-test", + "version": "1.0.0", + }, + "metrics": [metric], + } + ], + } + + gauge_points = [ + { + "attributes": attributes(series_id), + "time_unix_nano": str(1704067200000000000 + series_id), + "as_int": series_id, + } + for series_id in range(7) + ] + sum_points = [ + { + "attributes": attributes(series_id), + "start_time_unix_nano": "1704067200000000000", + "time_unix_nano": str(1704067200000000000 + series_id), + "as_int": series_id, + } + for series_id in range(7, 11) + ] + + return { + "resource_metrics": [ + resource( + "service-a", + { + "name": "batch.gauge", + "gauge": { + "data_points": gauge_points, + }, + }, + ), + resource( + "service-b", + { + "name": "batch.sum", + "sum": { + "data_points": sum_points, + "aggregation_temporality": 2, + "is_monotonic": True, + }, + }, + ), + ] + } + + +def iter_metric_points_with_resource(output): + data_keys = ("gauge", "sum", "histogram", "exponentialHistogram", "summary") + + for resource_metric in output.get("resourceMetrics", []): + resource_attributes = _attributes_to_dict( + resource_metric.get("resource", {}).get("attributes", []) + ) + for scope_metric in resource_metric.get("scopeMetrics", []): + for metric in scope_metric.get("metrics", []): + for data_key in data_keys: + if data_key in metric: + for point in metric[data_key].get("dataPoints", []): + yield point, resource_attributes + break + + def _assert_log_resource_attribution(logs_seen): output = json.loads(json_format.MessageToJson(logs_seen[0])) records = list(iter_log_records(output)) @@ -851,6 +946,97 @@ def test_out_opentelemetry_grpc_metrics_uri(): assert output["resourceMetrics"] +@pytest.mark.parametrize( + "config_file,receiver_mode,request_path", + [ + ("out_otel_http_metrics_max_datapoints.conf", "http", "/batched/metrics"), + ( + "out_otel_grpc_metrics_max_datapoints.conf", + "grpc", + "/batched.metrics.v1.Metrics/Export", + ), + ], + ids=["http", "grpc"], +) +def test_out_opentelemetry_metrics_max_datapoints( + config_file, + receiver_mode, + request_path, +): + grpc_methods = {"metrics": request_path} if receiver_mode == "grpc" else None + service = Service( + config_file, + receiver_mode=receiver_mode, + grpc_methods=grpc_methods, + ) + service.start() + service.send_payload_dict(_build_batched_metrics_payload(), "metrics") + metrics_seen = service.wait_for_signal("metrics", minimum_count=3, timeout=15) + requests_seen = service.wait_for_requests(3, timeout=15) + service.stop() + + assert len(metrics_seen) == 3 + assert len(requests_seen) == 3 + assert {request["path"] for request in requests_seen} == {request_path} + + batch_sizes = [] + observed_series = {} + for export_request in metrics_seen: + output = json.loads(json_format.MessageToJson(export_request)) + points = list(iter_metric_points_with_resource(output)) + batch_sizes.append(len(points)) + assert len(points) <= 4 + + for point, resource_attributes in points: + point_attributes = _attributes_to_dict(point.get("attributes", [])) + series_id = int(point_attributes["series.id"]) + assert series_id not in observed_series + observed_series[series_id] = resource_attributes["service.name"] + + assert sorted(batch_sizes) == [3, 4, 4] + assert observed_series == { + **{series_id: "service-a" for series_id in range(7)}, + **{series_id: "service-b" for series_id in range(7, 11)}, + } + + +def test_out_opentelemetry_metrics_partial_success_is_not_retried(): + payload = _build_batched_metrics_payload() + resource_metrics = payload["resource_metrics"] + resource_metrics[0]["scope_metrics"][0]["metrics"].extend( + resource_metrics[1]["scope_metrics"][0]["metrics"] + ) + payload["resource_metrics"] = [resource_metrics[0]] + + service = Service("out_otel_http_metrics_max_datapoints.conf") + service.start() + try: + configure_otlp_response(status_codes=[200, 503, 200]) + service.send_payload_dict(payload, "metrics") + metrics_seen = service.wait_for_signal("metrics", minimum_count=2, timeout=15) + _wait_for_log_message(service, "metric payload partially succeeded") + requests_seen = list(data_storage["requests"]) + finally: + service.stop() + + assert len(metrics_seen) == 2 + assert len(requests_seen) == 2 + + batch_series = [] + for export_request in metrics_seen: + output = json.loads(json_format.MessageToJson(export_request)) + points = list(iter_metric_points_with_resource(output)) + assert len(points) <= 4 + batch_series.append( + { + int(_attributes_to_dict(point.get("attributes", []))["series.id"]) + for point, _ in points + } + ) + + assert batch_series[0].isdisjoint(batch_series[1]) + + def test_out_opentelemetry_traces_uri(): service = Service("out_otel_http_traces.yaml") service.start() diff --git a/tests/integration/src/server/otlp_server.py b/tests/integration/src/server/otlp_server.py index 0724951f071..e4a755ac832 100644 --- a/tests/integration/src/server/otlp_server.py +++ b/tests/integration/src/server/otlp_server.py @@ -42,6 +42,7 @@ data_storage = {"traces": [], "metrics": [], "logs": [], "requests": []} response_config = { "status_code": 200, + "status_codes": [], "body": {"status": "received"}, "content_type": "application/json", "delay_seconds": 0, @@ -65,6 +66,7 @@ def reset_otlp_server_state(): response_config.update( { "status_code": 200, + "status_codes": [], "body": {"status": "received"}, "content_type": "application/json", "delay_seconds": 0, @@ -80,9 +82,18 @@ def reset_otlp_server_state(): shutdown_flag.clear() -def configure_otlp_response(*, status_code=None, body=None, content_type=None, delay_seconds=None): +def configure_otlp_response( + *, + status_code=None, + status_codes=None, + body=None, + content_type=None, + delay_seconds=None, +): if status_code is not None: response_config["status_code"] = status_code + if status_codes is not None: + response_config["status_codes"] = list(status_codes) if body is not None: response_config["body"] = body if content_type is not None: @@ -101,16 +112,21 @@ def configure_otlp_grpc_methods(*, logs=None, metrics=None, traces=None): def _build_response(): + status_code = response_config["status_code"] + if response_config["delay_seconds"]: time.sleep(response_config["delay_seconds"]) + if response_config["status_codes"]: + status_code = response_config["status_codes"].pop(0) + body = response_config["body"] if isinstance(body, (dict, list)): - return jsonify(body), response_config["status_code"] + return jsonify(body), status_code return Response( body, - status=response_config["status_code"], + status=status_code, content_type=response_config["content_type"], ) From aa2be0e00baa97477c97ba34821a68efb72c423e Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 15:08:05 +0900 Subject: [PATCH 6/8] tests: integration: Mark as skip if zstd command is absent Signed-off-by: Hiroshi Hatake --- .../out_opentelemetry/tests/test_out_opentelemetry_001.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py b/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py index 6d4d08ec8ef..bf6f7cebbb3 100644 --- a/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py +++ b/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py @@ -2,6 +2,7 @@ import json import logging import os +import shutil import socket import threading @@ -812,6 +813,10 @@ def test_out_opentelemetry_gzip_and_logs_body_key_attributes(): assert "message" not in attributes +@pytest.mark.skipif( + shutil.which("zstd") is None, + reason="zstd executable is required to decode the test payload", +) def test_out_opentelemetry_zstd_and_logs_body_key_attributes(): service = Service("out_otel_http_logs_zstd.yaml") service.start() From 445b42e0a4f8e3fdc9e4cce7ecd26573ffbd4baf Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 16:31:18 +0900 Subject: [PATCH 7/8] tests: integration: Validate resource and scope metadata preservation Signed-off-by: Hiroshi Hatake --- .../tests/test_out_opentelemetry_001.py | 50 +++++++++++++++---- 1 file changed, 41 insertions(+), 9 deletions(-) diff --git a/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py b/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py index bf6f7cebbb3..6da2a045f8b 100644 --- a/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py +++ b/tests/integration/scenarios/out_opentelemetry/tests/test_out_opentelemetry_001.py @@ -538,6 +538,7 @@ def attributes(series_id): def resource(service_name, metric): return { + "schema_url": f"https://example.com/resource/{service_name}/1.0.0", "resource": { "attributes": [ { @@ -551,9 +552,10 @@ def resource(service_name, metric): "scope_metrics": [ { "scope": { - "name": "batch-test", + "name": f"batch-test-{service_name}", "version": "1.0.0", }, + "schema_url": f"https://example.com/scope/{service_name}/1.0.0", "metrics": [metric], } ], @@ -610,12 +612,21 @@ def iter_metric_points_with_resource(output): resource_attributes = _attributes_to_dict( resource_metric.get("resource", {}).get("attributes", []) ) + resource_schema_url = resource_metric.get("schemaUrl") for scope_metric in resource_metric.get("scopeMetrics", []): + scope = scope_metric.get("scope", {}) + scope_schema_url = scope_metric.get("schemaUrl") for metric in scope_metric.get("metrics", []): for data_key in data_keys: if data_key in metric: for point in metric[data_key].get("dataPoints", []): - yield point, resource_attributes + yield ( + point, + resource_attributes, + resource_schema_url, + scope, + scope_schema_url, + ) break @@ -992,11 +1003,29 @@ def test_out_opentelemetry_metrics_max_datapoints( batch_sizes.append(len(points)) assert len(points) <= 4 - for point, resource_attributes in points: + for ( + point, + resource_attributes, + resource_schema_url, + scope, + scope_schema_url, + ) in points: point_attributes = _attributes_to_dict(point.get("attributes", [])) series_id = int(point_attributes["series.id"]) assert series_id not in observed_series - observed_series[series_id] = resource_attributes["service.name"] + service_name = "service-a" if series_id < 7 else "service-b" + assert resource_attributes["service.name"] == service_name + assert resource_schema_url == ( + f"https://example.com/resource/{service_name}/1.0.0" + ) + assert scope == { + "name": f"batch-test-{service_name}", + "version": "1.0.0", + } + assert scope_schema_url == ( + f"https://example.com/scope/{service_name}/1.0.0" + ) + observed_series[series_id] = service_name assert sorted(batch_sizes) == [3, 4, 4] assert observed_series == { @@ -1008,8 +1037,9 @@ def test_out_opentelemetry_metrics_max_datapoints( def test_out_opentelemetry_metrics_partial_success_is_not_retried(): payload = _build_batched_metrics_payload() resource_metrics = payload["resource_metrics"] - resource_metrics[0]["scope_metrics"][0]["metrics"].extend( - resource_metrics[1]["scope_metrics"][0]["metrics"] + gauge_points = resource_metrics[0]["scope_metrics"][0]["metrics"][0]["gauge"] + gauge_points["data_points"].extend( + resource_metrics[1]["scope_metrics"][0]["metrics"][0]["sum"]["data_points"] ) payload["resource_metrics"] = [resource_metrics[0]] @@ -1031,15 +1061,17 @@ def test_out_opentelemetry_metrics_partial_success_is_not_retried(): for export_request in metrics_seen: output = json.loads(json_format.MessageToJson(export_request)) points = list(iter_metric_points_with_resource(output)) - assert len(points) <= 4 + assert len(points) == 4 batch_series.append( { int(_attributes_to_dict(point.get("attributes", []))["series.id"]) - for point, _ in points + for point, _, _, _, _ in points } ) - assert batch_series[0].isdisjoint(batch_series[1]) + assert batch_series[0] == {0, 1, 2, 3} + assert batch_series[1] == {4, 5, 6, 7} + assert {8, 9, 10}.isdisjoint(set().union(*batch_series)) def test_out_opentelemetry_traces_uri(): From e17fe88379e7e708e052d1369df9371f097a2479 Mon Sep 17 00:00:00 2001 From: Hiroshi Hatake Date: Mon, 10 Aug 2026 16:32:47 +0900 Subject: [PATCH 8/8] tests: internal: Assert each decoded batch for all supported types Signed-off-by: Hiroshi Hatake --- tests/internal/opentelemetry.c | 358 +++++++++++++++++++++++++++++---- 1 file changed, 317 insertions(+), 41 deletions(-) diff --git a/tests/internal/opentelemetry.c b/tests/internal/opentelemetry.c index 7434df1fe89..21e749fb549 100644 --- a/tests/internal/opentelemetry.c +++ b/tests/internal/opentelemetry.c @@ -2864,14 +2864,26 @@ void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() int histogram_seen; int exp_histogram_seen; int summary_seen; + int gauge_points_seen[3] = {0}; + int sum_points_seen[3] = {0}; + int histogram_points_seen[3] = {0}; + int exp_histogram_points_seen[3] = {0}; + int summary_points_seen[3] = {0}; size_t payload_size; size_t batch_index; size_t resource_index; size_t scope_index; size_t metric_index; + size_t point_index; + size_t data_point_index; size_t total_data_points; flb_sds_t payload; Opentelemetry__Proto__Metrics__V1__Metric *metric; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint *number_point; + Opentelemetry__Proto__Metrics__V1__HistogramDataPoint *histogram_point; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint + *exp_histogram_point; + Opentelemetry__Proto__Metrics__V1__SummaryDataPoint *summary_point; Opentelemetry__Proto__Metrics__V1__Metric metrics[5]; Opentelemetry__Proto__Metrics__V1__Metric *metric_entries[5]; Opentelemetry__Proto__Metrics__V1__Gauge gauge; @@ -2879,20 +2891,26 @@ void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() Opentelemetry__Proto__Metrics__V1__Histogram histogram; Opentelemetry__Proto__Metrics__V1__ExponentialHistogram exp_histogram; Opentelemetry__Proto__Metrics__V1__Summary summary; - Opentelemetry__Proto__Metrics__V1__NumberDataPoint gauge_point; - Opentelemetry__Proto__Metrics__V1__NumberDataPoint sum_point; - Opentelemetry__Proto__Metrics__V1__NumberDataPoint *gauge_points[1]; - Opentelemetry__Proto__Metrics__V1__NumberDataPoint *sum_points[1]; - Opentelemetry__Proto__Metrics__V1__HistogramDataPoint histogram_point; - Opentelemetry__Proto__Metrics__V1__HistogramDataPoint *histogram_points[1]; - Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint exp_histogram_point; - Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint *exp_histogram_points[1]; - Opentelemetry__Proto__Metrics__V1__SummaryDataPoint summary_point; - Opentelemetry__Proto__Metrics__V1__SummaryDataPoint *summary_points[1]; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint gauge_point_values[3]; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint sum_point_values[3]; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint *gauge_points[3]; + Opentelemetry__Proto__Metrics__V1__NumberDataPoint *sum_points[3]; + Opentelemetry__Proto__Metrics__V1__HistogramDataPoint histogram_point_values[3]; + Opentelemetry__Proto__Metrics__V1__HistogramDataPoint *histogram_points[3]; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint + exp_histogram_point_values[3]; + Opentelemetry__Proto__Metrics__V1__ExponentialHistogramDataPoint + *exp_histogram_points[3]; + Opentelemetry__Proto__Metrics__V1__SummaryDataPoint summary_point_values[3]; + Opentelemetry__Proto__Metrics__V1__SummaryDataPoint *summary_points[3]; + Opentelemetry__Proto__Resource__V1__Resource resource_metadata; + Opentelemetry__Proto__Common__V1__InstrumentationScope scope_metadata; Opentelemetry__Proto__Metrics__V1__ScopeMetrics scope; Opentelemetry__Proto__Metrics__V1__ScopeMetrics *scopes[1]; Opentelemetry__Proto__Metrics__V1__ResourceMetrics resource; Opentelemetry__Proto__Metrics__V1__ResourceMetrics *resources[1]; + Opentelemetry__Proto__Metrics__V1__ScopeMetrics *decoded_scope; + Opentelemetry__Proto__Metrics__V1__ResourceMetrics *decoded_resource; Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest request; Opentelemetry__Proto__Collector__Metrics__V1__ExportMetricsServiceRequest *decoded; struct flb_opentelemetry_metrics_proto_batches *batches; @@ -2906,57 +2924,136 @@ void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() opentelemetry__proto__metrics__v1__histogram__init(&histogram); opentelemetry__proto__metrics__v1__exponential_histogram__init(&exp_histogram); opentelemetry__proto__metrics__v1__summary__init(&summary); - opentelemetry__proto__metrics__v1__number_data_point__init(&gauge_point); - opentelemetry__proto__metrics__v1__number_data_point__init(&sum_point); - opentelemetry__proto__metrics__v1__histogram_data_point__init(&histogram_point); - opentelemetry__proto__metrics__v1__exponential_histogram_data_point__init( - &exp_histogram_point); - opentelemetry__proto__metrics__v1__summary_data_point__init(&summary_point); + opentelemetry__proto__resource__v1__resource__init(&resource_metadata); + opentelemetry__proto__common__v1__instrumentation_scope__init(&scope_metadata); for (metric_index = 0; metric_index < 5; metric_index++) { opentelemetry__proto__metrics__v1__metric__init(&metrics[metric_index]); metric_entries[metric_index] = &metrics[metric_index]; } - gauge_points[0] = &gauge_point; - gauge.n_data_points = 1; + for (point_index = 0; point_index < 3; point_index++) { + opentelemetry__proto__metrics__v1__number_data_point__init( + &gauge_point_values[point_index]); + opentelemetry__proto__metrics__v1__number_data_point__init( + &sum_point_values[point_index]); + opentelemetry__proto__metrics__v1__histogram_data_point__init( + &histogram_point_values[point_index]); + opentelemetry__proto__metrics__v1__exponential_histogram_data_point__init( + &exp_histogram_point_values[point_index]); + opentelemetry__proto__metrics__v1__summary_data_point__init( + &summary_point_values[point_index]); + + gauge_points[point_index] = &gauge_point_values[point_index]; + gauge_point_values[point_index].start_time_unix_nano = 100 + point_index; + gauge_point_values[point_index].time_unix_nano = 1000 + point_index; + gauge_point_values[point_index].flags = 1 + point_index; + gauge_point_values[point_index].value_case = + OPENTELEMETRY__PROTO__METRICS__V1__NUMBER_DATA_POINT__VALUE_AS_INT; + gauge_point_values[point_index].as_int = 10 + point_index; + + sum_points[point_index] = &sum_point_values[point_index]; + sum_point_values[point_index].start_time_unix_nano = 200 + point_index; + sum_point_values[point_index].time_unix_nano = 2000 + point_index; + sum_point_values[point_index].flags = 11 + point_index; + sum_point_values[point_index].value_case = + OPENTELEMETRY__PROTO__METRICS__V1__NUMBER_DATA_POINT__VALUE_AS_DOUBLE; + sum_point_values[point_index].as_double = 20.5 + point_index; + + histogram_points[point_index] = &histogram_point_values[point_index]; + histogram_point_values[point_index].start_time_unix_nano = 300 + point_index; + histogram_point_values[point_index].time_unix_nano = 3000 + point_index; + histogram_point_values[point_index].count = 30 + point_index; + histogram_point_values[point_index].has_sum = FLB_TRUE; + histogram_point_values[point_index].sum = 30.5 + point_index; + histogram_point_values[point_index].flags = 21 + point_index; + histogram_point_values[point_index].has_min = FLB_TRUE; + histogram_point_values[point_index].min = 3.5 + point_index; + histogram_point_values[point_index].has_max = FLB_TRUE; + histogram_point_values[point_index].max = 35.5 + point_index; + + exp_histogram_points[point_index] = &exp_histogram_point_values[point_index]; + exp_histogram_point_values[point_index].start_time_unix_nano = 400 + point_index; + exp_histogram_point_values[point_index].time_unix_nano = 4000 + point_index; + exp_histogram_point_values[point_index].count = 40 + point_index; + exp_histogram_point_values[point_index].has_sum = FLB_TRUE; + exp_histogram_point_values[point_index].sum = 40.5 + point_index; + exp_histogram_point_values[point_index].scale = 4 + point_index; + exp_histogram_point_values[point_index].zero_count = 40 + point_index; + exp_histogram_point_values[point_index].flags = 31 + point_index; + exp_histogram_point_values[point_index].has_min = FLB_TRUE; + exp_histogram_point_values[point_index].min = 4.5 + point_index; + exp_histogram_point_values[point_index].has_max = FLB_TRUE; + exp_histogram_point_values[point_index].max = 45.5 + point_index; + exp_histogram_point_values[point_index].zero_threshold = 0.5 + point_index; + + summary_points[point_index] = &summary_point_values[point_index]; + summary_point_values[point_index].start_time_unix_nano = 500 + point_index; + summary_point_values[point_index].time_unix_nano = 5000 + point_index; + summary_point_values[point_index].count = 50 + point_index; + summary_point_values[point_index].sum = 50.5 + point_index; + summary_point_values[point_index].flags = 41 + point_index; + } + + gauge.n_data_points = 3; gauge.data_points = gauge_points; metrics[0].name = "gauge"; + metrics[0].description = "gauge description"; + metrics[0].unit = "gauge unit"; metrics[0].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE; metrics[0].gauge = &gauge; - sum_points[0] = &sum_point; - sum.n_data_points = 1; + sum.n_data_points = 3; sum.data_points = sum_points; + sum.aggregation_temporality = + OPENTELEMETRY__PROTO__METRICS__V1__AGGREGATION_TEMPORALITY__AGGREGATION_TEMPORALITY_DELTA; + sum.is_monotonic = FLB_TRUE; metrics[1].name = "sum"; + metrics[1].description = "sum description"; + metrics[1].unit = "sum unit"; metrics[1].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM; metrics[1].sum = ∑ - histogram_points[0] = &histogram_point; - histogram.n_data_points = 1; + histogram.n_data_points = 3; histogram.data_points = histogram_points; + histogram.aggregation_temporality = + OPENTELEMETRY__PROTO__METRICS__V1__AGGREGATION_TEMPORALITY__AGGREGATION_TEMPORALITY_CUMULATIVE; metrics[2].name = "histogram"; + metrics[2].description = "histogram description"; + metrics[2].unit = "histogram unit"; metrics[2].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM; metrics[2].histogram = &histogram; - exp_histogram_points[0] = &exp_histogram_point; - exp_histogram.n_data_points = 1; + exp_histogram.n_data_points = 3; exp_histogram.data_points = exp_histogram_points; + exp_histogram.aggregation_temporality = + OPENTELEMETRY__PROTO__METRICS__V1__AGGREGATION_TEMPORALITY__AGGREGATION_TEMPORALITY_DELTA; metrics[3].name = "exponential_histogram"; + metrics[3].description = "exponential histogram description"; + metrics[3].unit = "exponential histogram unit"; metrics[3].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM; metrics[3].exponential_histogram = &exp_histogram; - summary_points[0] = &summary_point; - summary.n_data_points = 1; + summary.n_data_points = 3; summary.data_points = summary_points; metrics[4].name = "summary"; + metrics[4].description = "summary description"; + metrics[4].unit = "summary unit"; metrics[4].data_case = OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY; metrics[4].summary = &summary; + scope_metadata.name = "splitter scope"; + scope_metadata.version = "2.0.0"; + scope_metadata.dropped_attributes_count = 9; + scope.scope = &scope_metadata; + scope.schema_url = "https://example.com/scope/2.0.0"; scope.n_metrics = 5; scope.metrics = metric_entries; scopes[0] = &scope; + resource_metadata.dropped_attributes_count = 7; + resource.resource = &resource_metadata; + resource.schema_url = "https://example.com/resource/1.0.0"; resource.n_scope_metrics = 1; resource.scope_metrics = scopes; resources[0] = &resource; @@ -2988,10 +3085,15 @@ void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() return; } - TEST_CHECK(batches->count == 3); - TEST_CHECK(batches->entries[0].data_point_count == 2); - TEST_CHECK(batches->entries[1].data_point_count == 2); - TEST_CHECK(batches->entries[2].data_point_count == 1); + TEST_CHECK(batches->count == 8); + for (batch_index = 0; batch_index < batches->count; batch_index++) { + if (batch_index < 7) { + TEST_CHECK(batches->entries[batch_index].data_point_count == 2); + } + else { + TEST_CHECK(batches->entries[batch_index].data_point_count == 1); + } + } gauge_seen = 0; sum_seen = 0; @@ -3011,40 +3113,207 @@ void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() continue; } + TEST_CHECK(decoded->n_resource_metrics == 1); for (resource_index = 0; resource_index < decoded->n_resource_metrics; resource_index++) { - resource = *decoded->resource_metrics[resource_index]; + decoded_resource = decoded->resource_metrics[resource_index]; + TEST_CHECK(decoded_resource->resource != NULL); + TEST_CHECK(strcmp(decoded_resource->schema_url, + "https://example.com/resource/1.0.0") == 0); + if (decoded_resource->resource != NULL) { + TEST_CHECK(decoded_resource->resource->dropped_attributes_count == 7); + } + TEST_CHECK(decoded_resource->n_scope_metrics == 1); + for (scope_index = 0; - scope_index < resource.n_scope_metrics; + scope_index < decoded_resource->n_scope_metrics; scope_index++) { - scope = *resource.scope_metrics[scope_index]; - for (metric_index = 0; metric_index < scope.n_metrics; metric_index++) { - metric = scope.metrics[metric_index]; + decoded_scope = decoded_resource->scope_metrics[scope_index]; + TEST_CHECK(decoded_scope->scope != NULL); + TEST_CHECK(strcmp(decoded_scope->schema_url, + "https://example.com/scope/2.0.0") == 0); + if (decoded_scope->scope != NULL) { + TEST_CHECK(strcmp(decoded_scope->scope->name, "splitter scope") == 0); + TEST_CHECK(strcmp(decoded_scope->scope->version, "2.0.0") == 0); + TEST_CHECK(decoded_scope->scope->dropped_attributes_count == 9); + } + + for (metric_index = 0; + metric_index < decoded_scope->n_metrics; + metric_index++) { + metric = decoded_scope->metrics[metric_index]; if (metric->data_case == OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_GAUGE) { gauge_seen++; + TEST_CHECK(strcmp(metric->name, "gauge") == 0); + TEST_CHECK(strcmp(metric->description, "gauge description") == 0); + TEST_CHECK(strcmp(metric->unit, "gauge unit") == 0); + TEST_CHECK(metric->gauge->n_data_points > 0); + TEST_CHECK(metric->gauge->n_data_points <= 2); total_data_points += metric->gauge->n_data_points; + + for (point_index = 0; + point_index < metric->gauge->n_data_points; + point_index++) { + number_point = metric->gauge->data_points[point_index]; + TEST_CHECK(number_point->time_unix_nano >= 1000); + TEST_CHECK(number_point->time_unix_nano < 1003); + if (number_point->time_unix_nano < 1000 || + number_point->time_unix_nano >= 1003) { + continue; + } + data_point_index = number_point->time_unix_nano - 1000; + gauge_points_seen[data_point_index]++; + TEST_CHECK(number_point->start_time_unix_nano == + 100 + data_point_index); + TEST_CHECK(number_point->flags == 1 + data_point_index); + TEST_CHECK(number_point->value_case == + OPENTELEMETRY__PROTO__METRICS__V1__NUMBER_DATA_POINT__VALUE_AS_INT); + TEST_CHECK(number_point->as_int == 10 + data_point_index); + } } else if (metric->data_case == OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUM) { sum_seen++; + TEST_CHECK(strcmp(metric->name, "sum") == 0); + TEST_CHECK(strcmp(metric->description, "sum description") == 0); + TEST_CHECK(strcmp(metric->unit, "sum unit") == 0); + TEST_CHECK(metric->sum->aggregation_temporality == + sum.aggregation_temporality); + TEST_CHECK(metric->sum->is_monotonic == sum.is_monotonic); + TEST_CHECK(metric->sum->n_data_points > 0); + TEST_CHECK(metric->sum->n_data_points <= 2); total_data_points += metric->sum->n_data_points; + + for (point_index = 0; + point_index < metric->sum->n_data_points; + point_index++) { + number_point = metric->sum->data_points[point_index]; + TEST_CHECK(number_point->time_unix_nano >= 2000); + TEST_CHECK(number_point->time_unix_nano < 2003); + if (number_point->time_unix_nano < 2000 || + number_point->time_unix_nano >= 2003) { + continue; + } + data_point_index = number_point->time_unix_nano - 2000; + sum_points_seen[data_point_index]++; + TEST_CHECK(number_point->start_time_unix_nano == + 200 + data_point_index); + TEST_CHECK(number_point->flags == 11 + data_point_index); + TEST_CHECK(number_point->value_case == + OPENTELEMETRY__PROTO__METRICS__V1__NUMBER_DATA_POINT__VALUE_AS_DOUBLE); + TEST_CHECK(number_point->as_double == 20.5 + data_point_index); + } } else if (metric->data_case == OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_HISTOGRAM) { histogram_seen++; + TEST_CHECK(strcmp(metric->name, "histogram") == 0); + TEST_CHECK(strcmp(metric->description, "histogram description") == 0); + TEST_CHECK(strcmp(metric->unit, "histogram unit") == 0); + TEST_CHECK(metric->histogram->aggregation_temporality == + histogram.aggregation_temporality); + TEST_CHECK(metric->histogram->n_data_points > 0); + TEST_CHECK(metric->histogram->n_data_points <= 2); total_data_points += metric->histogram->n_data_points; + + for (point_index = 0; + point_index < metric->histogram->n_data_points; + point_index++) { + histogram_point = metric->histogram->data_points[point_index]; + TEST_CHECK(histogram_point->time_unix_nano >= 3000); + TEST_CHECK(histogram_point->time_unix_nano < 3003); + if (histogram_point->time_unix_nano < 3000 || + histogram_point->time_unix_nano >= 3003) { + continue; + } + data_point_index = histogram_point->time_unix_nano - 3000; + histogram_points_seen[data_point_index]++; + TEST_CHECK(histogram_point->start_time_unix_nano == + 300 + data_point_index); + TEST_CHECK(histogram_point->count == 30 + data_point_index); + TEST_CHECK(histogram_point->has_sum == FLB_TRUE); + TEST_CHECK(histogram_point->sum == 30.5 + data_point_index); + TEST_CHECK(histogram_point->flags == 21 + data_point_index); + TEST_CHECK(histogram_point->has_min == FLB_TRUE); + TEST_CHECK(histogram_point->min == 3.5 + data_point_index); + TEST_CHECK(histogram_point->has_max == FLB_TRUE); + TEST_CHECK(histogram_point->max == 35.5 + data_point_index); + } } else if (metric->data_case == OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_EXPONENTIAL_HISTOGRAM) { exp_histogram_seen++; + TEST_CHECK(strcmp(metric->name, "exponential_histogram") == 0); + TEST_CHECK(strcmp(metric->description, + "exponential histogram description") == 0); + TEST_CHECK(strcmp(metric->unit, + "exponential histogram unit") == 0); + TEST_CHECK(metric->exponential_histogram->aggregation_temporality == + exp_histogram.aggregation_temporality); + TEST_CHECK(metric->exponential_histogram->n_data_points > 0); + TEST_CHECK(metric->exponential_histogram->n_data_points <= 2); total_data_points += metric->exponential_histogram->n_data_points; + + for (point_index = 0; + point_index < metric->exponential_histogram->n_data_points; + point_index++) { + exp_histogram_point = + metric->exponential_histogram->data_points[point_index]; + TEST_CHECK(exp_histogram_point->time_unix_nano >= 4000); + TEST_CHECK(exp_histogram_point->time_unix_nano < 4003); + if (exp_histogram_point->time_unix_nano < 4000 || + exp_histogram_point->time_unix_nano >= 4003) { + continue; + } + data_point_index = exp_histogram_point->time_unix_nano - 4000; + exp_histogram_points_seen[data_point_index]++; + TEST_CHECK(exp_histogram_point->start_time_unix_nano == + 400 + data_point_index); + TEST_CHECK(exp_histogram_point->count == 40 + data_point_index); + TEST_CHECK(exp_histogram_point->has_sum == FLB_TRUE); + TEST_CHECK(exp_histogram_point->sum == 40.5 + data_point_index); + TEST_CHECK(exp_histogram_point->scale == 4 + data_point_index); + TEST_CHECK(exp_histogram_point->zero_count == + 40 + data_point_index); + TEST_CHECK(exp_histogram_point->flags == 31 + data_point_index); + TEST_CHECK(exp_histogram_point->has_min == FLB_TRUE); + TEST_CHECK(exp_histogram_point->min == 4.5 + data_point_index); + TEST_CHECK(exp_histogram_point->has_max == FLB_TRUE); + TEST_CHECK(exp_histogram_point->max == 45.5 + data_point_index); + TEST_CHECK(exp_histogram_point->zero_threshold == + 0.5 + data_point_index); + } } else if (metric->data_case == OPENTELEMETRY__PROTO__METRICS__V1__METRIC__DATA_SUMMARY) { summary_seen++; + TEST_CHECK(strcmp(metric->name, "summary") == 0); + TEST_CHECK(strcmp(metric->description, "summary description") == 0); + TEST_CHECK(strcmp(metric->unit, "summary unit") == 0); + TEST_CHECK(metric->summary->n_data_points > 0); + TEST_CHECK(metric->summary->n_data_points <= 2); total_data_points += metric->summary->n_data_points; + + for (point_index = 0; + point_index < metric->summary->n_data_points; + point_index++) { + summary_point = metric->summary->data_points[point_index]; + TEST_CHECK(summary_point->time_unix_nano >= 5000); + TEST_CHECK(summary_point->time_unix_nano < 5003); + if (summary_point->time_unix_nano < 5000 || + summary_point->time_unix_nano >= 5003) { + continue; + } + data_point_index = summary_point->time_unix_nano - 5000; + summary_points_seen[data_point_index]++; + TEST_CHECK(summary_point->start_time_unix_nano == + 500 + data_point_index); + TEST_CHECK(summary_point->count == 50 + data_point_index); + TEST_CHECK(summary_point->sum == 50.5 + data_point_index); + TEST_CHECK(summary_point->flags == 41 + data_point_index); + } } } } @@ -3055,12 +3324,19 @@ void test_opentelemetry_metrics_otlp_proto_batches_all_metric_types() NULL); } - TEST_CHECK(total_data_points == 5); - TEST_CHECK(gauge_seen == 1); - TEST_CHECK(sum_seen == 1); - TEST_CHECK(histogram_seen == 1); - TEST_CHECK(exp_histogram_seen == 1); - TEST_CHECK(summary_seen == 1); + TEST_CHECK(total_data_points == 15); + TEST_CHECK(gauge_seen == 2); + TEST_CHECK(sum_seen == 2); + TEST_CHECK(histogram_seen == 2); + TEST_CHECK(exp_histogram_seen == 2); + TEST_CHECK(summary_seen == 2); + for (point_index = 0; point_index < 3; point_index++) { + TEST_CHECK(gauge_points_seen[point_index] == 1); + TEST_CHECK(sum_points_seen[point_index] == 1); + TEST_CHECK(histogram_points_seen[point_index] == 1); + TEST_CHECK(exp_histogram_points_seen[point_index] == 1); + TEST_CHECK(summary_points_seen[point_index] == 1); + } flb_opentelemetry_metrics_proto_batches_destroy(batches); }