diff --git a/benchmarks/pipeline_compare.py b/benchmarks/pipeline_compare.py new file mode 100644 index 00000000000..38115897f56 --- /dev/null +++ b/benchmarks/pipeline_compare.py @@ -0,0 +1,283 @@ +#!/usr/bin/env python3 +"""Measure completed HTTP/Forward/OTLP -> body content_modifier -> null pipelines. + +OTLP protobuf requires the Python dependencies in tests/integration/requirements.txt. +OTLP logs include resource/scope groups, which retain the CFL processing path. + +Linux only. Uses equal CPU affinity, fresh processes, alternating run order, +10,000 warm-up records, four clients, and exact output-counter checks. CPU +seconds count only Fluent Bit; wall time ends after output completion. The +Forward workload uses PackedForward with V2 events and no compression/acks. +Field counts represent body width. Compare each input and width separately. +""" + +import argparse +import concurrent.futures +import csv +import datetime +import hashlib +import http.client +import json +import os +from pathlib import Path +import platform +import socket +import struct +import subprocess +import tempfile +import time + + +WARMUP_RECORDS = 10000 +BATCH_RECORDS = 1000 +CLIENTS = 4 + + +def free_port(): + with socket.socket() as listener: + listener.bind(("127.0.0.1", 0)) + return listener.getsockname()[1] + + +class MetricsNotReady(RuntimeError): + pass + + +def get_metrics(port): + connection = http.client.HTTPConnection("127.0.0.1", port, timeout=2) + try: + connection.request("GET", "/api/v1/metrics") + response = connection.getresponse() + if response.status == 404: + raise MetricsNotReady("Metrics snapshot is not available yet") + if response.status != 200: + raise RuntimeError(f"Metrics HTTP status: {response.status}") + return json.loads(response.read()) + finally: + connection.close() + + +def wait_for_count(port, expected, process): + deadline = time.monotonic() + 120 + while time.monotonic() < deadline: + if process.poll() is not None: + raise RuntimeError("Fluent Bit exited before output completion") + try: + metrics = get_metrics(port) + except MetricsNotReady: + time.sleep(0.01) + continue + output = metrics.get("output", {}).get("null.0", {}) + count = output.get("proc_records", 0) + if count >= expected: + if count != expected or output.get("dropped_records", 0) != 0: + raise RuntimeError(f"Unexpected output counters: {output}") + return + time.sleep(0.01) + raise TimeoutError(f"Output did not reach {expected} records") + + +def cpu_seconds(pid): + fields = Path(f"/proc/{pid}/stat").read_text().rsplit(")", 1)[1].split() + return (int(fields[11]) + int(fields[12])) / os.sysconf("SC_CLK_TCK") + + +def send_http(port, payload, batches, path="/bench", content_type="application/json"): + connection = http.client.HTTPConnection("127.0.0.1", port, timeout=30) + try: + for _ in range(batches): + connection.request("POST", path, payload, {"Content-Type": content_type}) + response = connection.getresponse() + body = response.read() + if response.status != 201: + raise RuntimeError(f"Input response: {response.status} {body!r}") + finally: + connection.close() + + +def pack_string(value): + value = value.encode() + return b"\xdb" + struct.pack(">I", len(value)) + value + + +def send_forward(port, payload, batches): + with socket.create_connection(("127.0.0.1", port), timeout=30) as connection: + for _ in range(batches): + connection.sendall(payload) + + +def make_payload(input_name, fields): + record = {f"field{index}": "representative log attribute" for index in range(fields)} + if input_name == "http": + return json.dumps([record] * BATCH_RECORDS, separators=(",", ":")).encode() + if input_name.startswith("otlp-"): + log = {"timeUnixNano": "1700000000000000000", + "body": {"kvlistValue": {"values": [ + {"key": key, "value": {"stringValue": value}} for key, value in record.items() + ]}}} + request = {"resourceLogs": [{"resource": {"attributes": [ + {"key": "service.name", "value": {"stringValue": "perf"}}]}, + "scopeLogs": [{"scope": {"name": "benchmark"}, "logRecords": [log] * BATCH_RECORDS}]}]} + if input_name == "otlp-json": + return json.dumps(request, separators=(",", ":")).encode() + from google.protobuf.json_format import ParseDict + from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ExportLogsServiceRequest + return ParseDict(request, ExportLogsServiceRequest()).SerializeToString() + body = b"\xde" + struct.pack(">H", len(record)) + body += b"".join(pack_string(key) + pack_string(value) for key, value in record.items()) + entry = b"\x92\x92\xce" + struct.pack(">I", 1700000000) + b"\x80" + body + entries = entry * BATCH_RECORDS + return b"\x92" + pack_string("bench") + b"\xc6" + struct.pack(">I", len(entries)) + entries + + +def run(binary, records, fields, cpu, input_name, modifiers, mixed, location): + input_port = free_port() + metrics_port = free_port() + payload = make_payload(input_name, fields) + if input_name.startswith("otlp-"): + content_type = "application/x-protobuf" if input_name == "otlp-protobuf" else "application/json" + def sender(port, data, batches): + return send_http(port, data, batches, "/v1/logs", content_type) + else: + sender = send_http if input_name == "http" else send_forward + plugin_name = "opentelemetry" if input_name.startswith("otlp-") else input_name + actions = "".join( + " - name: content_modifier\n" + " context: body\n" + " action: upsert\n" + f" key: environment{index}\n" + " value: production\n" + for index in range(modifiers) + ) + if mixed: + actions += (" - name: content_modifier\n" + " context: body\n" + " action: hash\n" + " key: environment0\n") + processor_config = "" + if actions: + processor_config = " processors:\n logs:\n" + actions + input_processors = processor_config if location == "input" else "" + output_processors = processor_config if location == "output" else "" + config = f"""service: + flush: 0.1 + grace: 1 + log_level: error + http_server: on + http_listen: 127.0.0.1 + http_port: {metrics_port} +pipeline: + inputs: + - name: {plugin_name} + listen: 127.0.0.1 + port: {input_port} + tag: bench + buffer_max_size: 16M +{input_processors} outputs: + - name: null + match: '*' +{output_processors}""" + with tempfile.TemporaryDirectory(prefix="flb-pipeline-perf-") as directory: + path = Path(directory) + config_path = path / "fluent-bit.yaml" + config_path.write_text(config) + with (path / "stderr.log").open("w+") as log: + process = subprocess.Popen( + ["taskset", "-c", str(cpu), str(binary), "-c", str(config_path)], + stdout=log, stderr=log, + ) + try: + deadline = time.monotonic() + 15 + while True: + try: + get_metrics(metrics_port) + break + except (OSError, ValueError, RuntimeError): + if time.monotonic() > deadline or process.poll() is not None: + log.seek(0) + raise RuntimeError(log.read()) + time.sleep(0.05) + sender(input_port, payload, WARMUP_RECORDS // BATCH_RECORDS) + wait_for_count(metrics_port, WARMUP_RECORDS, process) + start_cpu = cpu_seconds(process.pid) + start = time.monotonic() + with concurrent.futures.ThreadPoolExecutor(max_workers=CLIENTS) as pool: + tasks = [pool.submit(sender, input_port, payload, + records // (CLIENTS * BATCH_RECORDS)) + for _ in range(CLIENTS)] + for task in tasks: + task.result() + wait_for_count(metrics_port, records + WARMUP_RECORDS, process) + elapsed = time.monotonic() - start + cpu_used = cpu_seconds(process.pid) - start_cpu + return {"input": input_name, "records": records, "fields": fields, + "modifiers": modifiers, "mixed": mixed, "location": location, + "elapsed_s": elapsed, "cpu_s": cpu_used, + "records_per_second": records / elapsed, + "cpu_s_per_million": cpu_used * 1e6 / records, + "payload_sha256": hashlib.sha256(payload).hexdigest(), + "output_records": records} + finally: + process.terminate() + try: + process.wait(timeout=10) + except subprocess.TimeoutExpired: + process.kill() + process.wait() + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--baseline", type=Path, required=True) + parser.add_argument("--candidate", type=Path, required=True) + parser.add_argument("--output", type=Path, required=True) + parser.add_argument("--records", type=int, default=1000000) + parser.add_argument("--runs", type=int, default=5) + parser.add_argument("--input", choices=("http", "forward", "otlp-json", "otlp-protobuf"), default="http") + parser.add_argument("--fields", type=int, nargs="+", default=[12, 128]) + parser.add_argument("--modifiers", type=int, default=1) + parser.add_argument("--location", choices=("input", "output"), default="input") + parser.add_argument("--mixed", action="store_true", help="Append a CFL-only hash operation") + args = parser.parse_args() + if args.records <= 0 or args.records % (CLIENTS * BATCH_RECORDS) or args.runs <= 0: + parser.error("records must be a positive multiple of 4000; runs must be positive") + if args.modifiers < 0 or (args.mixed and args.modifiers == 0): + parser.error("modifiers must be nonnegative; mixed requires at least one modifier") + if any(fields < 1 or fields > 128 for fields in args.fields): + parser.error("fields must be between 1 and 128 (16M maximum batch size)") + cpu = min(os.sched_getaffinity(0)) + binaries = {"baseline": args.baseline.resolve(), "candidate": args.candidate.resolve()} + hashes = {name: hashlib.sha256(binary.read_bytes()).hexdigest() + for name, binary in binaries.items()} + manifest = {"created_utc": datetime.datetime.now(datetime.timezone.utc).isoformat(), + "platform": platform.platform(), "cpu": cpu, + "cpu_info": Path("/proc/cpuinfo").read_text(), + "arguments": {key: str(value) if isinstance(value, Path) else value + for key, value in vars(args).items()}, + "binaries": {name: str(binary) for name, binary in binaries.items()}, + "binary_sha256": hashes, "warmup_records": WARMUP_RECORDS, + "batch_records": BATCH_RECORDS, "clients": CLIENTS, + "source_sha256": hashlib.sha256(Path(__file__).read_bytes()).hexdigest()} + args.output.with_suffix(".manifest.json").write_text(json.dumps(manifest, indent=2) + "\n") + with args.output.open("w") as output: + writer = None + for fields in args.fields: + for repetition in range(1, args.runs + 1): + order = ("baseline", "candidate") if repetition % 2 else ("candidate", "baseline") + for system in order: + if hashlib.sha256(binaries[system].read_bytes()).hexdigest() != hashes[system]: + raise RuntimeError(f"The {system} binary changed during the campaign") + row = run(binaries[system], args.records, fields, cpu, args.input, + args.modifiers, args.mixed, args.location) + row.update(system=system, repetition=repetition, cpu=cpu, + binary_sha256=hashes[system]) + if writer is None: + writer = csv.DictWriter(output, fieldnames=list(row)) + writer.writeheader() + writer.writerow(row) + output.flush() + print(json.dumps(row), flush=True) + + +if __name__ == "__main__": + main() diff --git a/include/fluent-bit/flb_processor.h b/include/fluent-bit/flb_processor.h index 2a56531cbf5..824cf0e7106 100644 --- a/include/fluent-bit/flb_processor.h +++ b/include/fluent-bit/flb_processor.h @@ -34,6 +34,21 @@ #define FLB_PROCESSOR_SUCCESS 0 #define FLB_PROCESSOR_FAILURE -1 +/* + * Optional raw callbacks execute a complete native segment in configuration + * order and preserve record cardinality. All units must use the same callback + * and enable logs_raw_enabled during init, with no conditions. The runner + * locks all participating units and accounts + * for each separately. Existing CFL callbacks remain the fallback for the whole + * segment. Raw callbacks never modify or release the input buffer. + * MODIFIED transfers a separately allocated output buffer to the caller. + * NOTOUCH, UNSUPPORTED and FAILURE leave output NULL/zero. UNSUPPORTED must + * discard any partial output and invoke no externally visible side effects. + */ +#define FLB_PROCESSOR_RAW_NOTOUCH 0 +#define FLB_PROCESSOR_RAW_MODIFIED 1 +#define FLB_PROCESSOR_RAW_UNSUPPORTED 2 + /* Processor event types */ #define FLB_PROCESSOR_LOGS 1 #define FLB_PROCESSOR_METRICS 2 @@ -166,6 +181,11 @@ struct flb_processor_plugin { const char *, int); + int (*cb_process_logs_raw) (struct flb_processor_instance **, size_t, + const void *, size_t, + void **, size_t *, + const char *, int); + int (*cb_process_metrics) (struct flb_processor_instance *, struct cmt *, /* in */ struct cmt **, /* out */ @@ -195,6 +215,7 @@ struct flb_processor_instance { int id; /* instance id */ int log_level; /* instance log level */ int event_type; /* event type */ + int logs_raw_enabled; /* set during init for compatible configurations */ char name[32]; /* numbered name */ char *alias; /* alias name */ void *context; /* Instance local context */ diff --git a/plugins/processor_content_modifier/CMakeLists.txt b/plugins/processor_content_modifier/CMakeLists.txt index 681d2cd58df..6ccbad6e202 100644 --- a/plugins/processor_content_modifier/CMakeLists.txt +++ b/plugins/processor_content_modifier/CMakeLists.txt @@ -1,6 +1,7 @@ set(src cm_config.c cm_logs.c + cm_msgpack.c cm_metrics.c cm_traces.c cm_opentelemetry.c diff --git a/plugins/processor_content_modifier/cm.c b/plugins/processor_content_modifier/cm.c index 22180dbca17..a8903f1e0f6 100644 --- a/plugins/processor_content_modifier/cm.c +++ b/plugins/processor_content_modifier/cm.c @@ -41,6 +41,7 @@ static int cb_init(struct flb_processor_instance *ins, void *source_plugin_insta } flb_processor_instance_set_context(ins, ctx); + ins->logs_raw_enabled = cm_logs_raw_supported(ins); return FLB_PROCESSOR_SUCCESS; } @@ -163,6 +164,7 @@ struct flb_processor_plugin processor_content_modifier_plugin = { .description = "Modify the content of Logs, Metrics and Traces", .cb_init = cb_init, .cb_process_logs = cb_process_logs, + .cb_process_logs_raw = cm_logs_process_raw, .cb_process_metrics = cb_process_metrics, .cb_process_traces = cb_process_traces, .cb_exit = cb_exit, diff --git a/plugins/processor_content_modifier/cm.h b/plugins/processor_content_modifier/cm.h index 65e7664d755..90cd3446f56 100644 --- a/plugins/processor_content_modifier/cm.h +++ b/plugins/processor_content_modifier/cm.h @@ -107,6 +107,13 @@ struct content_modifier_ctx { }; /* Export telemetry functions */ +int cm_logs_raw_supported(struct flb_processor_instance *ins); + +int cm_logs_process_raw(struct flb_processor_instance **instances, size_t instance_count, + const void *data, size_t bytes, + void **out_buf, size_t *out_size, + const char *tag, int tag_len); + int cm_logs_process(struct flb_processor_instance *ins, struct content_modifier_ctx *ctx, struct flb_mp_chunk_cobj *chunk_cobj, diff --git a/plugins/processor_content_modifier/cm_msgpack.c b/plugins/processor_content_modifier/cm_msgpack.c new file mode 100644 index 00000000000..e09a1f0211d --- /dev/null +++ b/plugins/processor_content_modifier/cm_msgpack.c @@ -0,0 +1,232 @@ +/* -*- 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 "cm.h" +#include +#include +#include +#include + +/* Fall back for shapes that the existing CFL conversion cannot represent. */ +static int object_supported(msgpack_object *object, size_t depth) +{ + size_t index; + + if (depth > 32 || object->type == MSGPACK_OBJECT_EXT) { + return FLB_FALSE; + } + if (object->type == MSGPACK_OBJECT_MAP) { + for (index = 0; index < object->via.map.size; index++) { + if (object->via.map.ptr[index].key.type != MSGPACK_OBJECT_STR || + !object_supported(&object->via.map.ptr[index].val, depth + 1)) { + return FLB_FALSE; + } + } + } + else if (object->type == MSGPACK_OBJECT_ARRAY) { + for (index = 0; index < object->via.array.size; index++) { + if (!object_supported(&object->via.array.ptr[index], depth + 1)) { + return FLB_FALSE; + } + } + } + return FLB_TRUE; +} + +static void string_object(msgpack_object *object, cfl_sds_t value) +{ + object->type = MSGPACK_OBJECT_STR; + object->via.str.ptr = value; + object->via.str.size = cfl_sds_len(value); +} + +int cm_logs_raw_supported(struct flb_processor_instance *ins) +{ + struct content_modifier_ctx *ctx; + + ctx = ins->context; + return ctx != NULL && ctx->context_type == CM_CONTEXT_LOG_BODY && + (ctx->action_type == CM_ACTION_INSERT || ctx->action_type == CM_ACTION_UPSERT || + ctx->action_type == CM_ACTION_DELETE || ctx->action_type == CM_ACTION_RENAME); +} + +int cm_logs_process_raw(struct flb_processor_instance **instances, size_t instance_count, + const void *data, size_t bytes, + void **out_buf, size_t *out_size, + const char *tag, int tag_len) +{ + int ret; + int result; + int changed; + int record_changed; + int record_type; + size_t index; + size_t consumed; + size_t action_index; + size_t match; + size_t capacity; + size_t required; + size_t count; + size_t key_length; + struct content_modifier_ctx *ctx; + struct flb_log_event_decoder decoder; + struct flb_log_event event; + msgpack_sbuffer output; + msgpack_packer packer; + msgpack_object root; + msgpack_object elements[2]; + msgpack_object_kv *pairs; + msgpack_object_kv *resized; + + (void) tag; + (void) tag_len; + if (instance_count == 0 || instance_count > UINT32_MAX) { + return FLB_PROCESSOR_RAW_UNSUPPORTED; + } + for (action_index = 0; action_index < instance_count; action_index++) { + if (!cm_logs_raw_supported(instances[action_index])) { + return FLB_PROCESSOR_RAW_UNSUPPORTED; + } + } + + ret = flb_log_event_decoder_init(&decoder, (char *) data, bytes); + if (ret != FLB_EVENT_DECODER_SUCCESS) { + return FLB_PROCESSOR_RAW_UNSUPPORTED; + } + flb_log_event_decoder_read_groups(&decoder, FLB_TRUE); + msgpack_sbuffer_init(&output); + msgpack_packer_init(&packer, &output, msgpack_sbuffer_write); + pairs = NULL; + capacity = 0; + changed = FLB_FALSE; + consumed = 0; + result = FLB_PROCESSOR_RAW_UNSUPPORTED; + + while ((ret = flb_log_event_decoder_next(&decoder, &event)) == FLB_EVENT_DECODER_SUCCESS) { + /* The decoder can skip invalid markers. Let the existing path handle them. */ + if (decoder.record_base != (const char *) data + consumed) { + goto cleanup; + } + consumed += decoder.record_length; + ret = flb_log_event_decoder_get_record_type(&event, &record_type); + if (ret != FLB_EVENT_DECODER_SUCCESS || record_type != FLB_LOG_EVENT_NORMAL || + event.format != FLB_LOG_EVENT_FORMAT_FLUENT_BIT_V2 || + event.body->type != MSGPACK_OBJECT_MAP || + !object_supported(event.body, 0) || !object_supported(event.metadata, 0)) { + goto cleanup; + } + + /* An action can add at most one field. Reuse the largest scratch map. */ + if (instance_count > UINT32_MAX - event.body->via.map.size) { + result = FLB_PROCESSOR_FAILURE; + goto cleanup; + } + required = (size_t) event.body->via.map.size + instance_count; + if (required > SIZE_MAX / sizeof(msgpack_object_kv)) { + result = FLB_PROCESSOR_FAILURE; + goto cleanup; + } + if (required > capacity) { + resized = flb_realloc(pairs, required * sizeof(msgpack_object_kv)); + if (resized == NULL) { + flb_errno(); + result = FLB_PROCESSOR_FAILURE; + goto cleanup; + } + pairs = resized; + capacity = required; + } + count = event.body->via.map.size; + if (count > 0) { + memcpy(pairs, event.body->via.map.ptr, count * sizeof(msgpack_object_kv)); + } + record_changed = FLB_FALSE; + for (action_index = 0; action_index < instance_count; action_index++) { + ctx = instances[action_index]->context; + key_length = cfl_sds_len(ctx->key); + match = SIZE_MAX; + for (index = 0; index < count; index++) { + if (pairs[index].key.via.str.size == key_length && + strncmp(pairs[index].key.via.str.ptr, ctx->key, key_length) == 0) { + match = index; + break; + } + } + if ((ctx->action_type == CM_ACTION_INSERT && match != SIZE_MAX) || + ((ctx->action_type == CM_ACTION_DELETE || ctx->action_type == CM_ACTION_RENAME) && + match == SIZE_MAX)) { + continue; + } + if (ctx->action_type == CM_ACTION_RENAME) { + string_object(&pairs[match].key, ctx->value); + } + else { + if (match != SIZE_MAX) { + memmove(&pairs[match], &pairs[match + 1], + (count - match - 1) * sizeof(msgpack_object_kv)); + count--; + } + if (ctx->action_type == CM_ACTION_INSERT || ctx->action_type == CM_ACTION_UPSERT) { + string_object(&pairs[count].key, ctx->key); + string_object(&pairs[count].val, ctx->value); + count++; + } + } + record_changed = FLB_TRUE; + } + if (!record_changed) { + if (msgpack_sbuffer_write(&output, decoder.record_base, decoder.record_length) != 0) { + result = FLB_PROCESSOR_FAILURE; + goto cleanup; + } + continue; + } + + root = *event.root; + memcpy(elements, root.via.array.ptr, sizeof(elements)); + elements[1] = *event.body; + elements[1].via.map.ptr = pairs; + elements[1].via.map.size = count; + root.via.array.ptr = elements; + if (msgpack_pack_object(&packer, root) != 0) { + result = FLB_PROCESSOR_FAILURE; + goto cleanup; + } + changed = FLB_TRUE; + } + + if (consumed != bytes || + flb_log_event_decoder_get_last_result(&decoder) != FLB_EVENT_DECODER_SUCCESS) { + goto cleanup; + } + result = FLB_PROCESSOR_RAW_NOTOUCH; + if (changed) { + *out_buf = output.data; + *out_size = output.size; + output.data = NULL; + result = FLB_PROCESSOR_RAW_MODIFIED; + } + +cleanup: + msgpack_sbuffer_destroy(&output); + flb_free(pairs); + flb_log_event_decoder_destroy(&decoder); + return result; +} diff --git a/src/flb_processor.c b/src/flb_processor.c index 7f04146ff6f..8b81e47df6b 100644 --- a/src/flb_processor.c +++ b/src/flb_processor.c @@ -1308,6 +1308,118 @@ int flb_processor_is_active(struct flb_processor *proc) #include +/* + * A native segment shares one representation. Only optimize it when every + * unit accepts the same raw callback; otherwise keep the existing CFL chain. + */ +static int run_logs_raw_segment(struct flb_processor *proc, + struct flb_processor_unit *first, + const void *data, size_t bytes, int records, + const char *tag, size_t tag_len, + void **out_buf, size_t *out_size, + struct mk_list **last) +{ + int result; + int out_records; + int error; + size_t count; + size_t index; + size_t locked; + struct mk_list *head; + struct flb_processor_unit *unit; + struct flb_processor_instance *initial; + struct flb_processor_instance *single; + struct flb_processor_instance **instances; + const char *scope; + const char *owner; + + initial = first->ctx; + if (first->condition != NULL || !initial->logs_raw_enabled || + initial->p->cb_process_logs_raw == NULL) { + return FLB_PROCESSOR_RAW_UNSUPPORTED; + } + count = 0; + for (head = &first->_head; head != &proc->logs; head = head->next) { + unit = mk_list_entry(head, struct flb_processor_unit, _head); + if (unit->unit_type == FLB_PROCESSOR_UNIT_FILTER) { + break; + } + if (unit->condition != NULL || + !((struct flb_processor_instance *) unit->ctx)->logs_raw_enabled || + ((struct flb_processor_instance *) unit->ctx)->p->cb_process_logs_raw != + initial->p->cb_process_logs_raw) { + return FLB_PROCESSOR_RAW_UNSUPPORTED; + } + count++; + } + single = initial; + instances = &single; + if (count > 1) { + instances = flb_calloc(count, sizeof(struct flb_processor_instance *)); + if (instances == NULL) { + flb_errno(); + return FLB_PROCESSOR_RAW_UNSUPPORTED; + } + } + index = 0; + locked = 0; + result = FLB_PROCESSOR_FAILURE; + for (head = &first->_head; index < count; head = head->next) { + unit = mk_list_entry(head, struct flb_processor_unit, _head); + instances[index] = unit->ctx; + if (acquire_lock(&unit->lock, FLB_PROCESSOR_LOCK_RETRY_LIMIT, + FLB_PROCESSOR_LOCK_RETRY_DELAY) != FLB_TRUE) { + goto cleanup; + } + locked++; + index++; + *last = head; + } + result = initial->p->cb_process_logs_raw(instances, count, data, bytes, + out_buf, out_size, tag, tag_len); + if (result == FLB_PROCESSOR_RAW_MODIFIED) { + /* A fused segment cannot attribute additions/drops to individual units. */ + if (*out_buf == NULL || *out_buf == data || + flb_mp_count_log_records(*out_buf, *out_size) != records) { + if (*out_buf != data) { + flb_free(*out_buf); + } + *out_buf = NULL; + *out_size = 0; + result = FLB_PROCESSOR_FAILURE; + } + } + else if (*out_buf != NULL || *out_size != 0 || + (result != FLB_PROCESSOR_RAW_NOTOUCH && + result != FLB_PROCESSOR_RAW_UNSUPPORTED)) { + if (*out_buf != data) { + flb_free(*out_buf); + } + *out_buf = NULL; + *out_size = 0; + result = FLB_PROCESSOR_FAILURE; + } + +cleanup: + scope = processor_metrics_scope(proc); + owner = processor_metrics_owner(proc); + error = (result == FLB_PROCESSOR_FAILURE); + out_records = records; + for (index = 0; index < locked; index++) { + unit = instances[index]->pu; + if (result != FLB_PROCESSOR_RAW_UNSUPPORTED) { + processor_metrics_update(proc, unit, scope, owner, "logs", + records, out_records, FLB_TRUE, error); + } + release_lock(&unit->lock, FLB_PROCESSOR_LOCK_RETRY_LIMIT, + FLB_PROCESSOR_LOCK_RETRY_DELAY); + } + if (count > 1) { + flb_free(instances); + } + return result; +} + /* * This function will run all the processor units for the given tag and data, note * that depending of the 'type', 'data' can reference a msgpack for logs, a CMetrics @@ -1327,12 +1439,14 @@ int flb_processor_run(struct flb_processor *proc, int unit_out_items; int item_metrics_available; int finalize; + int raw_result; void *cur_buf = NULL; size_t cur_size; void *tmp_buf = NULL; size_t tmp_size; struct mk_list *head; struct mk_list *list = NULL; + struct mk_list *raw_last; struct flb_processor_unit *pu; struct flb_processor_unit *pu_next; struct flb_filter_instance *f_ins; @@ -1395,11 +1509,41 @@ int flb_processor_run(struct flb_processor *proc, item_metrics_available = FLB_FALSE; if (type == FLB_PROCESSOR_LOGS) { - unit_in_items = flb_mp_count_log_records(cur_buf, cur_size); + if (chunk_cobj != NULL && + chunk_cobj->log_decoder->offset == chunk_cobj->log_decoder->length) { + unit_in_items = flb_mp_chunk_cobj_count_log_records(chunk_cobj); + } + else { + unit_in_items = flb_mp_count_log_records(cur_buf, cur_size); + } unit_out_items = unit_in_items; item_metrics_available = FLB_TRUE; } + if (type == FLB_PROCESSOR_LOGS && chunk_cobj == NULL && + pu->unit_type == FLB_PROCESSOR_UNIT_NATIVE) { + raw_last = head; + raw_result = run_logs_raw_segment(proc, pu, cur_buf, cur_size, unit_in_items, + tag, tag_len, &tmp_buf, &tmp_size, &raw_last); + if (raw_result == FLB_PROCESSOR_FAILURE) { + if (cur_buf != data) { + flb_free(cur_buf); + } + return -1; + } + if (raw_result != FLB_PROCESSOR_RAW_UNSUPPORTED) { + if (raw_result == FLB_PROCESSOR_RAW_MODIFIED) { + if (cur_buf != data) { + flb_free(cur_buf); + } + cur_buf = tmp_buf; + cur_size = tmp_size; + } + head = raw_last; + continue; + } + } + ret = acquire_lock(&pu->lock, FLB_PROCESSOR_LOCK_RETRY_LIMIT, FLB_PROCESSOR_LOCK_RETRY_DELAY); diff --git a/tests/integration/scenarios/processor_content_modifier/tests/test_content_modifier.py b/tests/integration/scenarios/processor_content_modifier/tests/test_content_modifier.py new file mode 100644 index 00000000000..118a76ba906 --- /dev/null +++ b/tests/integration/scenarios/processor_content_modifier/tests/test_content_modifier.py @@ -0,0 +1,368 @@ +"""Content modifier correctness across input/output processing and fan-out.""" + +from copy import deepcopy +from concurrent.futures import ThreadPoolExecutor +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +import json +import hashlib +import re +import threading + +import pytest +import requests +import yaml + +from utils.test_service import FluentBitTestService + + +RECORDS = [ + {"id": 1, "target": "original", "nested": {"array": [None, True, 1, 1.5]}}, + {"id": 2, "other": "untouched"}, +] + + +def processor_metrics(text, scope, owner="http.0"): + values = {} + for line in text.splitlines(): + if not line.startswith("fluentbit_processor_") or "{" not in line: + continue + name, body = line.split("{", 1) + labels, value = body.split("}", 1) + labels = dict(re.findall(r'(\w+)="([^"]*)"', labels)) + if labels.get("scope") == scope and labels.get("owner") == owner: + values[(name, int(labels["stage"]))] = float(value.split()[0]) + return values + + +def run_pipeline(tmp_path, processors, expected, *, location="input", storage="memory", + records=None, copies=1, stage_counts=None, input_workers=1, + output_workers=2, input_threaded=False, clients=4, unique_copies=False): + records = RECORDS if records is None else records + total_records = len(records) * copies + stage_counts = ([(len(records), len(records))] * len(processors) + if stage_counts is None else stage_counts) + assert len(stage_counts) == len(processors) + route_counts = [len(expected) * copies, + len(expected if location == "input" else records) * copies] + received = [[], []] + lock = threading.Lock() + + class Sink(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def do_POST(self): + records = json.loads(self.rfile.read(int(self.headers["Content-Length"]))) + with lock: + received[int(self.path[1:])].extend(records) + self.send_response(200) + self.send_header("Content-Length", "0") + self.end_headers() + + def log_message(self, *_args): + pass + + class SinkServer(ThreadingHTTPServer): + request_queue_size = 128 + + sink = SinkServer(("127.0.0.1", 0), Sink) + thread = threading.Thread(target=sink.serve_forever, daemon=True) + thread.start() + input_config = { + "name": "http", "listen": "127.0.0.1", + "port": "${FLUENT_BIT_TEST_LISTENER_PORT}", "tag": "test", + "storage.type": storage, "workers": input_workers, "threaded": input_threaded, + } + outputs = [{ + "name": "http", "match": "*", "host": "127.0.0.1", + "port": sink.server_port, "uri": f"/{route}", "format": "json", + "json_date_key": False, "workers": output_workers, + } for route in range(2)] + if location == "input": + input_config["processors"] = {"logs": processors} + else: + outputs[0]["processors"] = {"logs": processors} + config = { + "service": {"flush": 0.1, "grace": 1, "log_level": "error", + "storage.path": str(tmp_path / "storage"), + "http_server": True, "http_listen": "127.0.0.1", + "http_port": "${FLUENT_BIT_HTTP_MONITORING_PORT}"}, + "pipeline": {"inputs": [input_config], "outputs": outputs}, + } + config_path = tmp_path / "fluent-bit.yaml" + config_path.write_text(yaml.safe_dump(config)) + service = FluentBitTestService(str(config_path)) + try: + service.start() + def send(index): + batch = deepcopy(records) + if unique_copies: + for record in batch: + record["id"] += index * len(records) + record["request"] = index + response = requests.post(f"http://127.0.0.1:{service.flb_listener_port}/test.{index}", + json=batch, timeout=30) + assert response.status_code == 201, response.text + + with ThreadPoolExecutor(max_workers=min(copies, clients)) as pool: + list(pool.map(send, range(copies))) + service.wait_for_condition( + lambda: all(len(route_records) >= count + for route_records, count in zip(received, route_counts)), + timeout=30, interval=0.1, description="both output routes to finish", + ) + scope = "input" if location == "input" else "output" + + def metrics_ready(): + response = requests.get( + f"http://127.0.0.1:{service.flb.http_monitoring_port}/api/v2/metrics/prometheus", + timeout=10, + ) + if response.status_code == 404: + return None + response.raise_for_status() + values = processor_metrics(response.text, scope) + for stage, counts in enumerate(stage_counts): + for metric, count in zip(["items_in_total", "items_out_total"], counts): + if values.get(("fluentbit_processor_" + metric, stage), 0) != count * copies: + return None + return values + + values = service.wait_for_condition(metrics_ready, timeout=30, interval=0.2, + description="per-processor accounting") + invocations = values[("fluentbit_processor_invocations_total", 0)] + if location == "input" and input_workers == 1: + assert invocations == total_records + elif copies == 1: + assert invocations == 1 + else: + assert 1 <= invocations <= total_records + for stage in range(len(processors)): + assert values[("fluentbit_processor_invocations_total", stage)] == invocations + assert values.get(("fluentbit_processor_errors_total", stage), 0) == 0 + before, after = stage_counts[stage] + assert values.get(("fluentbit_processor_items_drop_total", stage), 0) == ( + max(0, before - after) * copies + ) + assert values.get(("fluentbit_processor_items_add_total", stage), 0) == ( + max(0, after - before) * copies + ) + + finally: + try: + service.stop() + finally: + sink.shutdown() + sink.server_close() + thread.join(timeout=5) + def copied_records(source): + result = [] + for index in range(copies): + batch = deepcopy(source) + if unique_copies: + for record in batch: + record["id"] += index * len(records) + record["request"] = index + result.extend(batch) + return sorted(result, key=lambda record: record["id"]) + + assert sorted(received[0], key=lambda record: record["id"]) == copied_records(expected) + control = expected if location == "input" else records + assert sorted(received[1], key=lambda record: record["id"]) == copied_records(control) + + +@pytest.mark.parametrize("action", ["insert", "upsert", "delete", "rename"]) +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("storage", ["memory", "filesystem"]) +def test_body_actions(tmp_path, action, location, storage): + processor = {"name": "content_modifier", "context": "body", "action": action, + "key": "target"} + if action != "delete": + processor["value"] = "replacement" + expected = deepcopy(RECORDS) + for record in expected: + if action == "insert": + record.setdefault("target", "replacement") + elif action == "upsert": + record["target"] = "replacement" + elif action == "delete": + record.pop("target", None) + elif "target" in record: + record["replacement"] = record.pop("target") + run_pipeline(tmp_path, [processor], expected, location=location, storage=storage) + + +@pytest.mark.parametrize("location", ["input", "output"]) +def test_mixed_processor_chain_and_condition(tmp_path, location): + processors = [ + {"name": "content_modifier", "action": "insert", "key": "inserted", "value": "yes"}, + # A filter boundary consumes raw buffers and materializes them independently. + {"name": "modify", "add": "filtered yes"}, + {"name": "content_modifier", "action": "upsert", "key": "conditional", "value": "yes", + "condition": {"op": "and", "rules": [{"field": "$target", "op": "eq", "value": "original"}]}}, + # Later native units must see the materialized condition result. + {"name": "content_modifier", "action": "rename", "key": "target", "value": "renamed"}, + {"name": "content_modifier", "action": "delete", "key": "inserted"}, + ] + expected = deepcopy(RECORDS) + for record in expected: + record["filtered"] = "yes" + expected[0]["conditional"] = "yes" + expected[0]["renamed"] = expected[0].pop("target") + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem") + + +@pytest.mark.parametrize("action,key", [("insert", "id"), ("delete", "missing"), + ("rename", "missing")]) +@pytest.mark.parametrize("location", ["input", "output"]) +def test_body_noop_preserves_shared_buffer(tmp_path, action, key, location): + processor = {"name": "content_modifier", "action": action, "key": key} + if action != "delete": + processor["value"] = "replacement" + run_pipeline(tmp_path, [processor], RECORDS, location=location, storage="filesystem") + + +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("mixed", [False, True], ids=["fused", "cfl-fallback"]) +def test_ordered_native_segment(tmp_path, location, mixed): + processors = [ + {"name": "content_modifier", "action": "insert", "key": "stage0", "value": "initial"}, + {"name": "content_modifier", "action": "rename", "key": "stage0", "value": "stage1"}, + {"name": "content_modifier", "action": "upsert", "key": "stage1", "value": "updated"}, + {"name": "content_modifier", "action": "delete", "key": "target"}, + {"name": "content_modifier", "action": "insert", "key": "target", "value": "done"}, + {"name": "content_modifier", "action": "delete", "key": "stage1"}, + ] + expected = deepcopy(RECORDS) + for record in expected: + record["target"] = "done" + if mixed: + processors.append({"name": "content_modifier", "action": "hash", "key": "target"}) + for record in expected: + record["target"] = hashlib.sha256(b"done").hexdigest() + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem") + + +@pytest.mark.parametrize("location", ["input", "output"]) +def test_concurrent_fused_segments(tmp_path, location): + processors = [{"name": "content_modifier", "action": "upsert", "key": f"field{index}", + "value": f"value{index}"} for index in range(6)] + records = [{"id": index, "nested": {"values": [True, None, index]}} + for index in range(256)] + expected = deepcopy(records) + for record in expected: + record.update({f"field{index}": f"value{index}" for index in range(6)}) + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem", + records=records, copies=8) + + +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("keep_id", [1, 99], ids=["partial-drop", "empty-batch"]) +def test_native_drop_accounting(tmp_path, location, keep_id): + # SQL and both modifiers share CFL, including when SQL removes every record. + processors = [ + {"name": "content_modifier", "action": "insert", "key": "before", "value": "yes"}, + {"name": "sql", "query": f"SELECT * FROM STREAM WHERE id = {keep_id};"}, + {"name": "content_modifier", "action": "insert", "key": "after", "value": "yes"}, + ] + expected = [dict(record, before="yes", after="yes") + for record in RECORDS if record["id"] == keep_id] + remaining = len(expected) + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem", + stage_counts=[(2, 2), (2, remaining), (remaining, remaining)]) + + +@pytest.mark.parametrize("location", ["input", "output"]) +def test_multiple_fused_segments_and_noop_tail(tmp_path, location): + processors = [ + {"name": "content_modifier", "action": "insert", "key": "stage0", "value": "yes"}, + {"name": "content_modifier", "action": "rename", "key": "stage0", "value": "stage1"}, + {"name": "modify", "rename": "stage1 filtered"}, + {"name": "content_modifier", "action": "rename", "key": "filtered", "value": "stage2"}, + {"name": "content_modifier", "action": "upsert", "key": "stage2", "value": "done"}, + {"name": "modify", "add": "boundary yes"}, + # NOTOUCH must retain the buffer allocated by the preceding segment/filter. + {"name": "content_modifier", "action": "insert", "key": "id", "value": "unused"}, + {"name": "content_modifier", "action": "delete", "key": "missing"}, + ] + expected = [dict(record, stage2="done", boundary="yes") for record in RECORDS] + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem") + + +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("matches", [True, False], ids=["condition-matches", "condition-false"]) +def test_condition_observes_previous_native_edit(tmp_path, location, matches): + processors = [ + {"name": "content_modifier", "action": "insert", "key": "gate", "value": "ready"}, + {"name": "content_modifier", "action": "upsert", "key": "conditional", "value": "yes", + "condition": {"op": "and", "rules": [ + {"field": "$gate", "op": "eq", "value": "ready" if matches else "absent"}]}}, + {"name": "content_modifier", "action": "rename", "key": "gate", "value": "renamed"}, + ] + expected = [dict(record, renamed="ready") for record in RECORDS] + if matches: + for record in expected: + record["conditional"] = "yes" + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem") + + +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("depth", [4, 40], ids=["raw-compatible", "deep-cfl-fallback"]) +def test_nested_and_boundary_values(tmp_path, location, depth): + nested = {"values": [None, True, False, -9223372036854775808, 18446744073709551615, + 1.25, "embedded\0nul", "Unicode: café 日本語", [], {}]} + for _ in range(depth): + nested = {"nested": nested} + # On output, a deep second record forces fallback after the first was rewritten. + records = [{"id": 1, "target": "original", "payload": {}, + "large": "x" * 65536}, {"id": 2, "payload": nested}] + processors = [ + {"name": "content_modifier", "action": "upsert", "key": "target", "value": "updated"}, + {"name": "content_modifier", "action": "rename", "key": "target", "value": "renamed"}, + ] + expected = [dict(record, renamed="updated") for record in records] + for record in expected: + record.pop("target", None) + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem", + records=records) + + +@pytest.mark.parametrize("location", ["input", "output"]) +def test_fused_mixed_noop_and_growing_maps(tmp_path, location): + records = [ + {"id": 1, "target": {"nested": [True, None, "preserved"]}}, + dict({"id": 2}, **{f"field{index}": f"value{index}" for index in range(128)}), + {"id": 3, "target": "last-record"}, + {"id": 4}, + ] + processors = [ + {"name": "content_modifier", "action": "insert", "key": "id", "value": "unused"}, + {"name": "content_modifier", "action": "rename", "key": "target", "value": "renamed"}, + {"name": "content_modifier", "action": "delete", "key": "missing"}, + ] + expected = deepcopy(records) + for record in expected: + if "target" in record: + record["renamed"] = record.pop("target") + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem", + records=records) + + +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("threaded", [False, True], ids=["engine-input", "threaded-input"]) +@pytest.mark.parametrize("chain", ["fused", "filter-boundary"]) +def test_worker_threads_preserve_independent_batches(tmp_path, location, threaded, chain): + records = [dict({"id": index, "target": f"record-{index}"}, + **{f"field{field}": "payload" * 16 for field in range(32)}) + for index in range(128)] + processors = [{"name": "content_modifier", "action": "upsert", "key": f"stage{index}", + "value": f"value{index}"} for index in range(6)] + expected = deepcopy(records) + for record in expected: + record.update({f"stage{index}": f"value{index}" for index in range(6)}) + if chain == "filter-boundary": + processors.insert(3, {"name": "modify", "add": "boundary yes"}) + for record in expected: + record["boundary"] = "yes" + run_pipeline(tmp_path, processors, expected, location=location, storage="filesystem", + records=records, copies=32, clients=8, unique_copies=True, + input_workers=4, output_workers=4, input_threaded=threaded) diff --git a/tests/integration/scenarios/processor_content_modifier/tests/test_grouped_logs.py b/tests/integration/scenarios/processor_content_modifier/tests/test_grouped_logs.py new file mode 100644 index 00000000000..e99111f4861 --- /dev/null +++ b/tests/integration/scenarios/processor_content_modifier/tests/test_grouped_logs.py @@ -0,0 +1,175 @@ +"""Grouped OTLP fallback must preserve envelope data and logical-record accounting.""" + +from copy import deepcopy +import hashlib +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +import threading + +import pytest +import requests +import yaml +from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ExportLogsServiceRequest + +from utils.test_service import FluentBitTestService +from test_content_modifier import processor_metrics + + +def grouped_payload(): + payload = ExportLogsServiceRequest() + for group in range(2): + resource = payload.resource_logs.add() + resource.schema_url = f"https://schemas.example/resource/{group}" + attribute = resource.resource.attributes.add(key="service.name") + attribute.value.string_value = f"service-{group}" + scope = resource.scope_logs.add() + scope.schema_url = f"https://schemas.example/scope/{group}" + scope.scope.name = f"scope-{group}" + scope.scope.version = "1.2.3" + scope.scope.attributes.add(key="scope-key").value.string_value = "scope-value" + for index in range(2): + record_id = group * 2 + index + 1 + record = scope.log_records.add() + record.time_unix_nano = 1700000000000000000 + record_id + record.observed_time_unix_nano = 1700000001000000000 + record_id + record.severity_number = 9 + record.severity_text = "INFO" + record.trace_id = bytes([record_id]) * 16 + record.span_id = bytes([record_id]) * 8 + record.flags = record_id + record.attributes.add(key="original-attribute").value.string_value = "preserved" + record.body.kvlist_value.values.add(key="id").value.int_value = record_id + record.body.kvlist_value.values.add(key="target").value.string_value = "original" + return payload + + +def flatten(payloads): + """Associate every log with its complete resource and scope envelopes.""" + result = [] + for payload in payloads: + for resource in payload.resource_logs: + for scope in resource.scope_logs: + for record in scope.log_records: + result.append((resource.schema_url, resource.resource, + scope.schema_url, scope.scope, record)) + return sorted(result, key=lambda entry: next( + item.value.int_value for item in entry[4].body.kvlist_value.values if item.key == "id" + )) + + +@pytest.mark.parametrize("location", ["input", "output"]) +@pytest.mark.parametrize("mixed", [False, True], ids=["group-fallback", "mixed-native-fallback"]) +def test_grouped_fallback_envelopes_and_accounting(tmp_path, location, mixed): + source = grouped_payload() + expected = deepcopy(source) + for resource in expected.resource_logs: + for scope in resource.scope_logs: + for record in scope.log_records: + for item in record.body.kvlist_value.values: + if item.key == "target": + item.key = "renamed" + item.value.string_value = (hashlib.sha256(b"updated").hexdigest() + if mixed else "updated") + record.body.kvlist_value.values.add(key="after").value.string_value = "yes" + remaining = 4 + processors = [ + {"name": "content_modifier", "action": "upsert", "key": "target", "value": "updated"}, + {"name": "content_modifier", "action": "rename", "key": "target", "value": "renamed"}, + ] + counts = [(4, 4), (4, 4)] + if mixed: + processors.append({"name": "content_modifier", "action": "hash", "key": "renamed"}) + counts.append((4, 4)) + processors.append({"name": "content_modifier", "action": "insert", "key": "after", + "value": "yes"}) + counts.append((remaining, remaining)) + received = [[], []] + lock = threading.Lock() + + class Sink(BaseHTTPRequestHandler): + def do_POST(self): + payload = ExportLogsServiceRequest() + payload.ParseFromString(self.rfile.read(int(self.headers["Content-Length"]))) + with lock: + received[int(self.path[-1])].append(payload) + self.send_response(200) + self.send_header("Content-Type", "application/x-protobuf") + self.send_header("Content-Length", "0") + self.end_headers() + + def log_message(self, *_args): + pass + + sink = ThreadingHTTPServer(("127.0.0.1", 0), Sink) + thread = threading.Thread(target=sink.serve_forever, daemon=True) + thread.start() + input_config = {"name": "opentelemetry", "listen": "127.0.0.1", + "port": "${FLUENT_BIT_TEST_LISTENER_PORT}", "tag": "test", + "storage.type": "filesystem"} + outputs = [{"name": "opentelemetry", "match": "*", "host": "127.0.0.1", + "port": sink.server_port, "logs_uri": f"/logs/{route}"} + for route in range(2)] + if location == "input": + input_config["processors"] = {"logs": processors} + else: + outputs[0]["processors"] = {"logs": processors} + config = { + "service": {"flush": 0.1, "grace": 1, "log_level": "error", + "storage.path": str(tmp_path / "storage"), "http_server": True, + "http_listen": "127.0.0.1", + "http_port": "${FLUENT_BIT_HTTP_MONITORING_PORT}"}, + "pipeline": {"inputs": [input_config], "outputs": outputs}, + } + config_path = tmp_path / "fluent-bit.yaml" + config_path.write_text(yaml.safe_dump(config)) + service = FluentBitTestService(str(config_path)) + try: + service.start() + response = requests.post(f"http://127.0.0.1:{service.flb_listener_port}/v1/logs", + data=source.SerializeToString(), + headers={"Content-Type": "application/x-protobuf"}, timeout=30) + assert response.status_code == 201, response.text + control_count = remaining if location == "input" else 4 + service.wait_for_condition( + lambda: len(flatten(received[0])) >= remaining and + len(flatten(received[1])) >= control_count, + timeout=30, interval=0.1, description="both grouped output routes to finish", + ) + for route, payloads in enumerate(received): + (tmp_path / f"route-{route}.txt").write_text("\n".join(str(p) for p in payloads)) + + def metrics_ready(): + response = requests.get( + f"http://127.0.0.1:{service.flb.http_monitoring_port}/api/v2/metrics/prometheus", + timeout=10, + ) + response.raise_for_status() + values = processor_metrics(response.text, location, owner="opentelemetry.0") + for stage, (before, after) in enumerate(counts): + if values.get(("fluentbit_processor_items_in_total", stage), 0) != before: + return None + if values.get(("fluentbit_processor_items_out_total", stage), 0) != after: + return None + return values + + values = service.wait_for_condition(metrics_ready, timeout=30, interval=0.2, + description="group markers excluded from accounting") + for stage, (before, after) in enumerate(counts): + assert values[("fluentbit_processor_invocations_total", stage)] == 1 + assert values.get(("fluentbit_processor_errors_total", stage), 0) == 0 + assert values.get(("fluentbit_processor_items_drop_total", stage), 0) == before - after + assert values.get(("fluentbit_processor_items_add_total", stage), 0) == 0 + finally: + try: + service.stop() + finally: + sink.shutdown() + sink.server_close() + thread.join(timeout=5) + assert flatten(received[0]) == flatten([expected]) + assert flatten(received[1]) == flatten([expected if location == "input" else source]) + # Group envelopes must remain paired with their logical records. + for route in received: + for payload in route: + for resource in payload.resource_logs: + assert resource.scope_logs + assert all(scope.log_records for scope in resource.scope_logs) diff --git a/tests/internal/processor.c b/tests/internal/processor.c index e2ebcd7bf86..17ec0536d15 100644 --- a/tests/internal/processor.c +++ b/tests/internal/processor.c @@ -22,6 +22,8 @@ #include #include #include +#include +#include #include #include #include @@ -625,7 +627,120 @@ static void processor_metrics_counters() flb_sds_destroy(regex_prop_key); } +static int native_drop_second_record(struct flb_processor_instance *ins, + void *data, const char *tag, int tag_len) +{ + int ret; + int type; + int count; + struct flb_mp_chunk_cobj *chunk; + struct flb_mp_chunk_record *record; + + chunk = data; + count = 0; + while ((ret = flb_mp_chunk_cobj_record_next(chunk, &record)) == FLB_MP_CHUNK_RECORD_OK) { + TEST_CHECK(flb_log_event_decoder_get_record_type(&record->event, &type) == 0); + if (type == FLB_LOG_EVENT_NORMAL && ++count == 2) { + TEST_CHECK(flb_mp_chunk_cobj_record_destroy(chunk, record) == 0); + } + } + return ret == FLB_MP_CHUNK_RECORD_EOF ? FLB_PROCESSOR_SUCCESS : FLB_PROCESSOR_FAILURE; +} + +static void processor_native_chain_counters(void) +{ + int grouped; + int root_type; + char *data; + char *resized; + size_t bytes; + size_t out_size; + void *out; + double value; + struct flb_config *config; + struct flb_processor *proc; + struct flb_processor_unit *first; + struct flb_processor_unit *second; + struct flb_processor_instance *ins; + struct flb_processor_plugin *original; + struct flb_processor_plugin drop; + struct flb_input_instance owner; + const char *json = "[[1700000000,{}],{\"message\":\"first\"}]" + "[[1700000001,{}],{\"message\":\"second\"}]"; + + flb_init_env(); + for (grouped = 0; grouped <= 1; grouped++) { + config = flb_config_init(); + TEST_CHECK(config != NULL); + memset(&owner, 0, sizeof(owner)); + snprintf(owner.name, sizeof(owner.name), "unit_input.0"); + owner.cmt = cmt_create(); + TEST_CHECK(owner.cmt != NULL); + proc = flb_processor_create(config, "native_counters", &owner, FLB_PLUGIN_INPUT); + TEST_CHECK(proc != NULL); + first = flb_processor_unit_create(proc, FLB_PROCESSOR_LOGS, "content_modifier"); + second = flb_processor_unit_create(proc, FLB_PROCESSOR_LOGS, "content_modifier"); + TEST_CHECK(first != NULL && second != NULL); + TEST_CHECK(flb_processor_unit_set_property_str(first, "action", "insert") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(first, "key", "unused") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(first, "value", "unused") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(second, "action", "insert") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(second, "key", "after_drop") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(second, "value", "yes") == 0); + TEST_CHECK(flb_processor_init(proc) == 0); + ins = first->ctx; + original = ins->p; + drop = *original; + drop.cb_process_logs_raw = NULL; + drop.cb_process_logs = native_drop_second_record; + ins->p = &drop; + if (grouped) { + TEST_CHECK(create_grouped_msgpack_records(&data, &bytes) == 0); + resized = flb_realloc(data, bytes * 2); + TEST_CHECK(resized != NULL); + data = resized; + memcpy(data + bytes, data, bytes); + bytes *= 2; + } + else { + TEST_CHECK(flb_pack_json(json, strlen(json), &data, &bytes, &root_type, NULL) == 0); + } + TEST_CHECK(flb_mp_count_log_records(data, bytes) == 2); + out = NULL; + TEST_CHECK(flb_processor_run(proc, 0, FLB_PROCESSOR_LOGS, "test", 4, + data, bytes, &out, &out_size) == 0); + TEST_CHECK(flb_mp_count_log_records(out, out_size) == 1); + TEST_CHECK(get_counter_value_5(proc->cmt_items_in, "input", "unit_input.0", + "content_modifier", "0", "logs", &value) == 0); + TEST_CHECK(value == 2.0); + TEST_CHECK(get_counter_value_5(proc->cmt_items_out, "input", "unit_input.0", + "content_modifier", "0", "logs", &value) == 0); + TEST_CHECK(value == 1.0); + TEST_CHECK(get_counter_value_5(proc->cmt_items_in, "input", "unit_input.0", + "content_modifier", "1", "logs", &value) == 0); + TEST_CHECK(value == 1.0); + TEST_CHECK(get_counter_value_5(proc->cmt_items_out, "input", "unit_input.0", + "content_modifier", "1", "logs", &value) == 0); + TEST_CHECK(value == 1.0); + TEST_CHECK(get_counter_value_5(proc->cmt_items_drop, "input", "unit_input.0", + "content_modifier", "0", "logs", &value) == 0); + TEST_CHECK(value == 1.0); + TEST_CHECK(get_counter_value_5_or_zero(proc->cmt_items_drop, "input", "unit_input.0", + "content_modifier", "1", "logs", &value) == 0); + TEST_CHECK(value == 0.0); + ins->p = original; + if (out != data) { + flb_free(out); + } + flb_free(data); + flb_processor_destroy(proc); + cmt_destroy(owner.cmt); + flb_config_exit(config); + } +} + TEST_LIST = { + { "processor_native_chain_counters", processor_native_chain_counters }, { "processor_private_inputs_use_main_loop", processor_private_inputs_use_main_loop }, { "processor", processor }, { "processor_grouped_filter_counters", processor_grouped_filter_counters }, diff --git a/tests/runtime/processor_content_modifier.c b/tests/runtime/processor_content_modifier.c index fd494c5e6fd..5ecc1c6e8c7 100644 --- a/tests/runtime/processor_content_modifier.c +++ b/tests/runtime/processor_content_modifier.c @@ -1876,7 +1876,312 @@ static void flb_logs_otel_log_attributes_invalid_otlp_metadata() processor_test_destroy(ctx); } +/* Compare the optimized path with the existing CFL path on identical batches. */ +static void check_raw_chain_equivalence(const char *action, const char *json, + int expected_result, size_t modifiers) +{ + int ret; + int root_type; + int raw_result; + int expected_ret; + char *data; + char *saved; + char *actual_json; + char *expected_json; + size_t bytes; + size_t actual_size; + size_t expected_size; + size_t probe_size; + size_t index; + struct flb_processor_instance *instances[6]; + void *actual; + void *expected; + void *probe; + flb_ctx_t *flb; + struct flb_processor *proc; + struct flb_processor_unit *unit; + struct flb_processor_instance *ins; + struct flb_processor_plugin *original; + struct flb_processor_plugin fallback; + struct flb_log_event_decoder actual_decoder; + struct flb_log_event_decoder expected_decoder; + struct flb_log_event actual_event; + struct flb_log_event expected_event; + + flb = flb_create(); + TEST_CHECK(flb != NULL); + proc = flb_processor_create(flb->config, "raw_equivalence", NULL, 0); + TEST_CHECK(proc != NULL); + TEST_CHECK(modifiers > 0 && modifiers <= 6); + for (index = 0; index < modifiers; index++) { + unit = flb_processor_unit_create(proc, FLB_PROCESSOR_LOGS, "content_modifier"); + TEST_CHECK(unit != NULL); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "action", action) == 0); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "context", "body") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "key", "target") == 0); + if (strcmp(action, "delete") != 0) { + TEST_CHECK(flb_processor_unit_set_property_str(unit, "value", "replacement") == 0); + } + instances[index] = unit->ctx; + } + TEST_CHECK(flb_processor_init(proc) == 0); + ins = unit->ctx; + original = ins->p; + fallback = *original; + fallback.cb_process_logs_raw = NULL; + + ret = flb_pack_json(json, strlen(json), &data, &bytes, &root_type, NULL); + TEST_CHECK(ret == 0); + saved = flb_malloc(bytes); + TEST_CHECK(saved != NULL); + memcpy(saved, data, bytes); + probe = NULL; + probe_size = 0; + raw_result = original->cb_process_logs_raw(instances, modifiers, data, bytes, + &probe, &probe_size, "test", 4); + TEST_CHECK(raw_result == expected_result); + TEST_CHECK(memcmp(data, saved, bytes) == 0); + flb_free(probe); + + actual = NULL; + expected = NULL; + TEST_CHECK(flb_processor_run(proc, 0, FLB_PROCESSOR_LOGS, "test", 4, + data, bytes, &actual, &actual_size) == 0); + ins->p = &fallback; + TEST_CHECK(flb_processor_run(proc, 0, FLB_PROCESSOR_LOGS, "test", 4, + data, bytes, &expected, &expected_size) == 0); + ins->p = original; + TEST_CHECK(memcmp(data, saved, bytes) == 0); + + TEST_CHECK(flb_log_event_decoder_init(&actual_decoder, actual, actual_size) == 0); + TEST_CHECK(flb_log_event_decoder_init(&expected_decoder, expected, expected_size) == 0); + flb_log_event_decoder_read_groups(&actual_decoder, FLB_TRUE); + flb_log_event_decoder_read_groups(&expected_decoder, FLB_TRUE); + while ((ret = flb_log_event_decoder_next(&actual_decoder, &actual_event)) == 0) { + expected_ret = flb_log_event_decoder_next(&expected_decoder, &expected_event); + TEST_CHECK(expected_ret == 0); + if (expected_ret != 0) { + break; + } + TEST_CHECK(actual_event.timestamp.tm.tv_sec == expected_event.timestamp.tm.tv_sec); + TEST_CHECK(actual_event.timestamp.tm.tv_nsec == expected_event.timestamp.tm.tv_nsec); + actual_json = flb_msgpack_to_json_str(4096, actual_event.body, FLB_TRUE); + expected_json = flb_msgpack_to_json_str(4096, expected_event.body, FLB_TRUE); + TEST_CHECK(actual_json != NULL && expected_json != NULL); + TEST_CHECK(strcmp(actual_json, expected_json) == 0); + flb_free(actual_json); + flb_free(expected_json); + actual_json = flb_msgpack_to_json_str(4096, actual_event.metadata, FLB_TRUE); + expected_json = flb_msgpack_to_json_str(4096, expected_event.metadata, FLB_TRUE); + TEST_CHECK(actual_json != NULL && expected_json != NULL); + TEST_CHECK(strcmp(actual_json, expected_json) == 0); + flb_free(actual_json); + flb_free(expected_json); + } + TEST_CHECK(flb_log_event_decoder_get_last_result(&actual_decoder) == 0); + TEST_CHECK(flb_log_event_decoder_next(&expected_decoder, &expected_event) != 0); + TEST_CHECK(flb_log_event_decoder_get_last_result(&expected_decoder) == 0); + flb_log_event_decoder_destroy(&actual_decoder); + flb_log_event_decoder_destroy(&expected_decoder); + if (actual != data) { + flb_free(actual); + } + if (expected != data) { + flb_free(expected); + } + flb_free(saved); + flb_free(data); + flb_processor_destroy(proc); + flb_destroy(flb); +} + +static void check_raw_equivalence(const char *action, const char *json, int expected_result) +{ + check_raw_chain_equivalence(action, json, expected_result, 1); +} + +static void flb_logs_raw_chain_equivalence(void) +{ + const char *actions[] = {"insert", "upsert", "delete", "rename"}; + const char *json = + "[[1700000000,{\"origin\":\"test\"}]," + "{\"target\":\"first\",\"target\":\"second\",\"replacement\":\"collision\"}]" + "[[1700000001,{}],{\"nested\":{\"values\":[true,null,1,1.5]}}]" + "[[1700000002,{}],{}]"; + size_t index; + size_t count; + + for (index = 0; index < sizeof(actions) / sizeof(actions[0]); index++) { + for (count = 3; count <= 6; count += 3) { + check_raw_chain_equivalence(actions[index], json, FLB_PROCESSOR_RAW_MODIFIED, count); + check_raw_chain_equivalence(actions[index], + "[[1700000000,{}],{\"target\":\"plain\"}]" + "[[-1,{}],{}]" + "[[1700000001,{}],{\"target\":\"grouped\"}]" + "[[-2,{}],{}]", FLB_PROCESSOR_RAW_UNSUPPORTED, count); + } + } +} + +static void flb_logs_raw_equivalence(void) +{ + const char *actions[] = {"insert", "upsert", "delete", "rename"}; + const char *batch = + "[[1700000000.125,{\"source\":\"test\"}]," + "{\"target\":\"first\",\"target\":\"second\",\"replacement\":\"existing\",\"nested\":{\"x\":[null,true,1,1.5]}}]" + "[[1700000001,{}],{\"other\":\"unchanged\"}]" + "[[1700000002,{}],{}]"; + size_t index; + + for (index = 0; index < sizeof(actions) / sizeof(actions[0]); index++) { + check_raw_equivalence(actions[index], batch, FLB_PROCESSOR_RAW_MODIFIED); + check_raw_equivalence(actions[index], "[1700000000,{\"target\":\"legacy\"}]", + FLB_PROCESSOR_RAW_UNSUPPORTED); + check_raw_equivalence(actions[index], + "[[-1,{\"resource\":\"test\"}],{\"scope\":\"group\"}]" + "[[1700000000,{}],{\"target\":\"grouped\"}]" + "[[-2,{}],{}]", FLB_PROCESSOR_RAW_UNSUPPORTED); + /* Unsupported records after a modified prefix must discard the private output. */ + check_raw_equivalence(actions[index], + "[[1700000000,{}],{\"target\":\"plain\"}]" + "[[-1,{}],{}]" + "[[1700000001,{}],{\"target\":\"grouped\"}]" + "[[-2,{}],{}]", FLB_PROCESSOR_RAW_UNSUPPORTED); + check_raw_equivalence(actions[index], + "[[1700000000,{}],{\"target\":\"valid\"}]" + "[[-3,{}],{\"target\":\"invalid marker\"}]", + FLB_PROCESSOR_RAW_UNSUPPORTED); + } + check_raw_equivalence("insert", "[[1700000000,{}],{\"target\":\"keep\"}]", + FLB_PROCESSOR_RAW_NOTOUCH); + check_raw_equivalence("delete", "[[1700000000,{}],{\"other\":\"keep\"}]", + FLB_PROCESSOR_RAW_NOTOUCH); + check_raw_equivalence("rename", "[[1700000000,{}],{}]", FLB_PROCESSOR_RAW_NOTOUCH); +} + +static void flb_logs_raw_invalid_batch(void) +{ + int ret; + size_t bytes; + size_t out_size; + int root_type; + char *data; + void *out; + flb_ctx_t *flb; + struct flb_processor *proc; + struct flb_processor_unit *unit; + struct flb_processor_instance *ins; + const char *json = "[[1700000000,{}],{\"target\":\"valid\"}]" + "[[1700000001,{}],{\"target\":\"truncated\"}]"; + + flb = flb_create(); + proc = flb_processor_create(flb->config, "raw_invalid", NULL, 0); + unit = flb_processor_unit_create(proc, FLB_PROCESSOR_LOGS, "content_modifier"); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "action", "upsert") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "key", "target") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "value", "replacement") == 0); + TEST_CHECK(flb_processor_init(proc) == 0); + ins = unit->ctx; + TEST_CHECK(flb_pack_json(json, strlen(json), &data, &bytes, &root_type, NULL) == 0); + out = NULL; + out_size = 0; + ret = ins->p->cb_process_logs_raw(&ins, 1, data, bytes - 1, &out, &out_size, "test", 4); + TEST_CHECK(ret == FLB_PROCESSOR_RAW_UNSUPPORTED); + TEST_CHECK(out == NULL && out_size == 0); + flb_free(data); + flb_processor_destroy(proc); + flb_destroy(flb); +} + +static int raw_test_result; + +static int raw_test_callback(struct flb_processor_instance **instances, size_t count, + const void *data, size_t bytes, + void **out_buf, size_t *out_size, + const char *tag, int tag_len) +{ + *out_buf = NULL; + *out_size = 0; + return raw_test_result; +} + +static void flb_logs_raw_results(void) +{ + int ret; + int root_type; + size_t bytes; + size_t out_size; + size_t index; + char *data; + char *json; + void *out; + flb_ctx_t *flb; + struct flb_processor *proc; + struct flb_processor_unit *unit; + struct flb_processor_instance *ins; + struct flb_processor_plugin *original; + struct flb_processor_plugin mock; + int results[] = {FLB_PROCESSOR_RAW_NOTOUCH, FLB_PROCESSOR_RAW_UNSUPPORTED, + FLB_PROCESSOR_RAW_MODIFIED, FLB_PROCESSOR_FAILURE}; + const char *input = "[[1700000000,{}],{\"message\":\"test\"}]"; + + for (index = 0; index < sizeof(results) / sizeof(results[0]); index++) { + flb = flb_create(); + proc = flb_processor_create(flb->config, "raw_results", NULL, 0); + unit = flb_processor_unit_create(proc, FLB_PROCESSOR_LOGS, "modify"); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "add", "first yes") == 0); + unit = flb_processor_unit_create(proc, FLB_PROCESSOR_LOGS, "content_modifier"); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "action", "insert") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "key", "second") == 0); + TEST_CHECK(flb_processor_unit_set_property_str(unit, "value", "yes") == 0); + TEST_CHECK(flb_processor_init(proc) == 0); + ins = unit->ctx; + original = ins->p; + mock = *original; + mock.cb_process_logs_raw = raw_test_callback; + ins->p = &mock; + raw_test_result = results[index]; + TEST_CHECK(flb_pack_json(input, strlen(input), &data, &bytes, &root_type, NULL) == 0); + out = NULL; + out_size = 0; + ret = flb_processor_run(proc, 0, FLB_PROCESSOR_LOGS, "test", 4, + data, bytes, &out, &out_size); + if (raw_test_result == FLB_PROCESSOR_FAILURE) { + TEST_CHECK(ret == -1); + TEST_CHECK(out == NULL); + } + else if (raw_test_result == FLB_PROCESSOR_RAW_MODIFIED) { + TEST_CHECK(ret == -1); + TEST_CHECK(out == NULL && out_size == 0); + } + else { + TEST_CHECK(ret == 0); + json = flb_msgpack_raw_to_json_sds(out, out_size, FLB_TRUE); + TEST_CHECK(json != NULL); + TEST_CHECK(strstr(json, "first") != NULL); + if (raw_test_result == FLB_PROCESSOR_RAW_UNSUPPORTED) { + TEST_CHECK(strstr(json, "second") != NULL); + } + else { + TEST_CHECK(strstr(json, "second") == NULL); + } + flb_sds_destroy(json); + } + ins->p = original; + if (out != data) { + flb_free(out); + } + flb_free(data); + flb_processor_destroy(proc); + flb_destroy(flb); + } +} + TEST_LIST = { + {"logs.raw.chain_equivalence", flb_logs_raw_chain_equivalence}, + {"logs.raw.results", flb_logs_raw_results}, + {"logs.raw.equivalence", flb_logs_raw_equivalence}, + {"logs.raw.invalid_batch", flb_logs_raw_invalid_batch}, {"logs.action.insert" , flb_logs_action_insert }, {"logs.action.delete" , flb_logs_action_delete }, {"logs.action.rename" , flb_logs_action_rename },