From 1b4aaae0c75aaf7de5c7f66c3f8656d1413b3c9b Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Sat, 19 Sep 2026 09:23:33 -0600 Subject: [PATCH 1/7] in_elasticsearch: respond to malformed bulk payloads Signed-off-by: Eduardo Silva --- plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c b/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c index 72f9d641b03..22cd260154d 100644 --- a/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c +++ b/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c @@ -609,6 +609,8 @@ static int process_payload_ng(struct flb_http_request *request, flb_sds_t tag, flb_sds_t *bulk_statuses) { + int ret; + if (request->content_type == NULL) { send_response_ng(response, 400, NULL, "error: header 'Content-Type' is not set\n"); @@ -627,8 +629,13 @@ static int process_payload_ng(struct flb_http_request *request, return -1; } - return parse_payload_ndjson(context, tag, request->body, - cfl_sds_len(request->body), bulk_statuses); + ret = parse_payload_ndjson(context, tag, request->body, + cfl_sds_len(request->body), bulk_statuses); + if (ret != 0 && ret != FLB_INPUT_INGRESS_BUSY) { + send_response_ng(response, 400, NULL, "error: invalid bulk payload\n"); + } + + return ret; } int in_elasticsearch_bulk_prot_handle_ng(struct flb_http_request *request, From b7a76748e234251dd870a9dc5c84f258c964aa69 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Sat, 19 Sep 2026 09:23:33 -0600 Subject: [PATCH 2/7] tests: integration: cover elasticsearch invalid bulk Signed-off-by: Eduardo Silva --- .../tests/test_elasticsearch_invalid_bulk.py | 135 ++++++++++++++++++ 1 file changed, 135 insertions(+) create mode 100644 tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py diff --git a/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py b/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py new file mode 100644 index 00000000000..a4a64857582 --- /dev/null +++ b/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py @@ -0,0 +1,135 @@ +"""Bounded nesting and recovery checks for network ingestion.""" +import contextlib +import http.client +import os +from pathlib import Path +import signal +import socket +import struct +import subprocess +import time + +import pytest + + + +def wait_for(predicate, timeout=30): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + if predicate(): + return + time.sleep(0.05) + assert predicate(), "Timed out waiting for Fluent Bit" + + +@contextlib.contextmanager +def daemon(tmp_path, mode): + with socket.socket() as listener: + listener.bind(("127.0.0.1", 0)) + port = listener.getsockname()[1] + plugin = mode.split("-")[0] + address = str(tmp_path / "input.sock") if plugin == "unix_socket" else ("127.0.0.1", port) + parser = tmp_path / "parsers.conf" + parser.write_text("[PARSER]\n Name json\n Format json\n") + command = [os.environ["FLUENT_BIT_BINARY"], "-f", "0.1", "-R", str(parser), "-i", plugin] + if plugin == "unix_socket": + command += ["-p", f"socket_path={address}"] + else: + command += ["-p", "listen=127.0.0.1", "-p", f"port={port}"] + if mode.endswith("-parser"): + command += ["-p", "format=none", "-p", "parser=json"] + if plugin == "syslog": + command += ["-p", "mode=tcp", "-p", "parser=json"] + command += ["-o", "stdout", "-m", "*", "-p", "format=json_lines"] + log = tmp_path / "fluent-bit.log" + memlog = tmp_path / "valgrind.log" + memory = os.environ.get("VALGRIND") == "1" + if memory: + command = ["valgrind", "--leak-check=full", "--show-leak-kinds=all", + "--errors-for-leak-kinds=definite,indirect", "--error-exitcode=99", + f"--log-file={memlog}"] + command + with log.open("w") as output: + process = subprocess.Popen(command, stdout=output, stderr=subprocess.STDOUT) + def ready(): + assert process.poll() is None, log.read_text() + return "[output:stdout:" in log.read_text() + try: + wait_for(ready) + yield address, process, log + finally: + if process.poll() is None: + process.send_signal(signal.SIGTERM) + try: + process.wait(timeout=30) + except subprocess.TimeoutExpired: + process.kill() + process.wait() + pytest.fail("Fluent Bit did not shut down cleanly") + assert process.returncode == 0, log.read_text() + (memlog.read_text() if memory else "") + if memory: + assert "ERROR SUMMARY: 0 errors" in memlog.read_text(), memlog.read_text() + + +def send(mode, address, payload): + plugin = mode.split("-")[0] + if plugin in ("http", "splunk", "elasticsearch"): + conn = http.client.HTTPConnection(*address, timeout=10) + try: + path = "/test" + if plugin == "splunk": + path = "/services/collector/event" + payload = b'{"event":' + payload + b"}" + elif plugin == "elasticsearch": + path = "/_bulk" + payload = b'{"index":{}}\n' + payload + b"\n" + conn.request("POST", path, payload, {"Content-Type": "application/json"}) + response = conn.getresponse() + response.read() + return response.status + finally: + conn.close() + return + if plugin == "udp": + with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as sock: + sock.sendto(payload + b"\n", address) + return + family = socket.AF_UNIX if plugin == "unix_socket" else socket.AF_INET + with socket.socket(family, socket.SOCK_STREAM) as sock: + sock.settimeout(5) + sock.connect(address) + if plugin == "mqtt": + # A regular MQTT CONNECT followed by a QoS 0 JSON publication. + sock.sendall(b"\x10\x10\x00\x04MQTT\x04\x02\x00\x0a\x00\x04test") + assert sock.recv(4)[0] == 0x20 + body = b"\x00\x01a" + payload + length = len(body) + encoded = bytearray() + while True: + digit = length % 128 + length //= 128 + encoded.append(digit | (0x80 if length else 0)) + if not length: + break + sock.sendall(b"\x30" + bytes(encoded) + body) + else: + sock.sendall(payload if plugin == "forward" else payload + b"\n") + + +@pytest.mark.parametrize("mode", ["elasticsearch"]) +@pytest.mark.parametrize("kind", ["array", "map", "mixed"]) +def test_nested_json_recovery(tmp_path, mode, kind): + # Bounded regression data, with no process-crash or exploit-chain behavior. + value = b"0" + for index in range(65): + value = (b'{"k":' + value + b"}") if kind == "map" or ( + kind == "mixed" and index % 2) else b"[" + value + b"]" + nested = b'{"nested":' + value + b"}" + with daemon(tmp_path, mode) as (address, process, log): + send(mode, address, b'{"marker":"before"}') + wait_for(lambda: '"marker":"before"' in log.read_text()) + assert send(mode, address, nested) == 400 + send(mode, address, b'{"marker":"after"}') + wait_for(lambda: '"marker":"after"' in log.read_text()) + assert process.poll() is None + # Parser inputs may preserve rejected JSON as a raw log string. + assert '"nested":' not in log.read_text() From f85eb2c039ff4b571e8fa0e7086f90d4741d58ba Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 23 Sep 2026 11:44:51 -0600 Subject: [PATCH 3/7] in_elasticsearch: keep bulk ingestion failures retryable Signed-off-by: Eduardo Silva --- .../in_elasticsearch_bulk_prot.c | 53 +++++++++++-------- 1 file changed, 31 insertions(+), 22 deletions(-) diff --git a/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c b/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c index 22cd260154d..94cc7d6ff21 100644 --- a/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c +++ b/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c @@ -256,10 +256,10 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char ret = flb_log_event_encoder_begin_record(encoder); if (ret != FLB_EVENT_ENCODER_SUCCESS) { - flb_sds_destroy(write_op); flb_plg_error(ctx->ins, "event encoder error : %d", ret); error_op = FLB_TRUE; - + ingest_result = -1; + flb_sds_destroy(write_op); break; } @@ -268,10 +268,10 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char &tm); if (ret != FLB_EVENT_ENCODER_SUCCESS) { - flb_sds_destroy(write_op); flb_plg_error(ctx->ins, "event encoder error : %d", ret); error_op = FLB_TRUE; - + ingest_result = -1; + flb_sds_destroy(write_op); break; } @@ -283,10 +283,10 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char } if (ret != FLB_EVENT_ENCODER_SUCCESS) { - flb_sds_destroy(write_op); flb_plg_error(ctx->ins, "event encoder error : %d", ret); error_op = FLB_TRUE; - + ingest_result = -1; + flb_sds_destroy(write_op); break; } } @@ -310,7 +310,8 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char if (ret != FLB_EVENT_ENCODER_SUCCESS) { flb_plg_error(ctx->ins, "event encoder error : %d", ret); error_op = FLB_TRUE; - + ingest_result = -1; + flb_sds_destroy(write_op); break; } @@ -319,7 +320,8 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char if (ret != FLB_EVENT_ENCODER_SUCCESS) { flb_plg_error(ctx->ins, "event encoder error : %d", ret); error_op = FLB_TRUE; - + ingest_result = -1; + flb_sds_destroy(write_op); break; } @@ -418,14 +420,23 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char else { flb_plg_error(ctx->ins, "skip record from invalid type: %i", result.data.type); + flb_sds_destroy(write_op); msgpack_unpacked_destroy(&result); if (destroy_local_encoder == FLB_TRUE) { flb_log_event_encoder_destroy(encoder); } - return -1; + return FLB_ERR_JSON_INVAL; } } + if (ingest_result != 0) { + msgpack_unpacked_destroy(&result); + if (destroy_local_encoder == FLB_TRUE) { + flb_log_event_encoder_destroy(encoder); + } + return ingest_result; + } + if (idx % 2 != 0) { flb_plg_warn(ctx->ins, "decode payload of Bulk API is failed"); msgpack_unpacked_destroy(&result); @@ -440,19 +451,11 @@ static int process_ndpack(struct flb_in_elasticsearch *ctx, flb_sds_t tag, char flb_log_event_encoder_destroy(encoder); } - return -1; + return FLB_ERR_JSON_INVAL; } msgpack_unpacked_destroy(&result); - if (ingest_result != 0) { - if (destroy_local_encoder == FLB_TRUE) { - flb_log_event_encoder_destroy(encoder); - } - - return ingest_result; - } - if (destroy_local_encoder == FLB_TRUE) { flb_log_event_encoder_destroy(encoder); } @@ -469,7 +472,10 @@ static ssize_t parse_payload_ndjson(struct flb_in_elasticsearch *ctx, flb_sds_t struct flb_pack_state pack_state; /* Initialize packer */ - flb_pack_state_init(&pack_state); + ret = flb_pack_state_init(&pack_state); + if (ret != 0) { + return -1; + } /* Pack JSON as msgpack */ ret = flb_pack_json_state(payload, size, @@ -479,11 +485,11 @@ static ssize_t parse_payload_ndjson(struct flb_in_elasticsearch *ctx, flb_sds_t /* Handle exceptions */ if (ret == FLB_ERR_JSON_PART) { flb_plg_warn(ctx->ins, "JSON data is incomplete, skipping"); - return -1; + return FLB_ERR_JSON_INVAL; } else if (ret == FLB_ERR_JSON_INVAL) { flb_plg_warn(ctx->ins, "invalid JSON message, skipping"); - return -1; + return FLB_ERR_JSON_INVAL; } else if (ret == -1) { return -1; @@ -631,9 +637,12 @@ static int process_payload_ng(struct flb_http_request *request, ret = parse_payload_ndjson(context, tag, request->body, cfl_sds_len(request->body), bulk_statuses); - if (ret != 0 && ret != FLB_INPUT_INGRESS_BUSY) { + if (ret == FLB_ERR_JSON_INVAL) { send_response_ng(response, 400, NULL, "error: invalid bulk payload\n"); } + else if (ret != 0 && ret != FLB_INPUT_INGRESS_BUSY) { + send_response_ng(response, 500, NULL, "error: could not ingest bulk payload\n"); + } return ret; } From 6b9ae6b7575a697a3b9b6e59c10d88198dd8a8f2 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 23 Sep 2026 11:44:51 -0600 Subject: [PATCH 4/7] tests: cover elasticsearch bulk append failures Signed-off-by: Eduardo Silva --- tests/runtime/in_elasticsearch.c | 45 ++++++++++++++++++++++++++++++++ 1 file changed, 45 insertions(+) diff --git a/tests/runtime/in_elasticsearch.c b/tests/runtime/in_elasticsearch.c index b373d2859cf..b60e9f8eaa6 100644 --- a/tests/runtime/in_elasticsearch.c +++ b/tests/runtime/in_elasticsearch.c @@ -25,6 +25,7 @@ #include #include #include +#include #include #include "flb_tests_runtime.h" @@ -882,7 +883,51 @@ void flb_test_in_elasticsearch_index_op_with_plugin_tag() flb_test_in_elasticsearch("index", 9210, "es.index"); } +void flb_test_in_elasticsearch_ingestion_failure() +{ + struct flb_lib_out_cb cb_data = {0}; + struct test_ctx *ctx; + struct flb_input_instance *input; + struct flb_http_client *client; + int ret; + size_t bytes_sent; + char *payload = "{\"index\":{}}\n{\"message\":\"valid\"}\n"; + + ctx = test_ctx_create(&cb_data); + if (!TEST_CHECK(ctx != NULL)) { + return; + } + + ret = flb_input_set(ctx->flb, ctx->i_ffd, "port", "9211", NULL); + TEST_CHECK(ret == 0); + ret = flb_output_set(ctx->flb, ctx->o_ffd, "match", "*", NULL); + TEST_CHECK(ret == 0); + ret = flb_start(ctx->flb); + TEST_CHECK(ret == 0); + + input = flb_input_get_instance(ctx->flb->config, ctx->i_ffd); + /* Force append failure without pausing the HTTP listener itself. */ + input->mem_buf_status = FLB_INPUT_PAUSED; + ctx->httpc = in_elasticsearch_client_ctx_create(9211); + client = flb_http_client(ctx->httpc->u_conn, FLB_HTTP_POST, "/_bulk", + payload, strlen(payload), "127.0.0.1", 9211, NULL, 0); + if (TEST_CHECK(client != NULL)) { + flb_http_add_header(client, FLB_HTTP_HEADER_CONTENT_TYPE, + strlen(FLB_HTTP_HEADER_CONTENT_TYPE), + NDJSON_CONTENT_TYPE, strlen(NDJSON_CONTENT_TYPE)); + ret = flb_http_do(client, &bytes_sent); + TEST_CHECK(ret == 0); + TEST_CHECK(client->resp.status == 500); + TEST_MSG("expected HTTP 500, got %d", client->resp.status); + flb_http_client_destroy(client); + } + input->mem_buf_status = FLB_INPUT_RUNNING; + flb_upstream_conn_release(ctx->httpc->u_conn); + test_ctx_destroy(ctx); +} + TEST_LIST = { + {"ingestion_failure", flb_test_in_elasticsearch_ingestion_failure}, {"version", flb_test_in_elasticsearch_version}, {"configured_version", flb_test_in_elasticsearch_version_configured}, {"index_op", flb_test_in_elasticsearch_index_op}, From 112cffb1f34573149f38d559953fdf0ae6c14303 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 23 Sep 2026 11:44:51 -0600 Subject: [PATCH 5/7] tests: integration: validate elasticsearch errors with shared memory checks Signed-off-by: Eduardo Silva --- .../tests/test_elasticsearch_invalid_bulk.py | 148 +++++++----------- 1 file changed, 57 insertions(+), 91 deletions(-) diff --git a/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py b/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py index a4a64857582..9dd265201b8 100644 --- a/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py +++ b/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py @@ -1,16 +1,14 @@ -"""Bounded nesting and recovery checks for network ingestion.""" +"""Reject malformed bulk payloads while keeping ingestion failures retryable.""" import contextlib import http.client -import os +import json from pathlib import Path -import signal import socket -import struct -import subprocess import time import pytest +from utils.fluent_bit_manager import FluentBitManager def wait_for(predicate, timeout=30): @@ -19,100 +17,51 @@ def wait_for(predicate, timeout=30): if predicate(): return time.sleep(0.05) - assert predicate(), "Timed out waiting for Fluent Bit" + assert predicate(), "Timed out waiting for Fluent Bit output" @contextlib.contextmanager -def daemon(tmp_path, mode): +def daemon(tmp_path, mode, input_options=None): with socket.socket() as listener: listener.bind(("127.0.0.1", 0)) port = listener.getsockname()[1] - plugin = mode.split("-")[0] - address = str(tmp_path / "input.sock") if plugin == "unix_socket" else ("127.0.0.1", port) - parser = tmp_path / "parsers.conf" - parser.write_text("[PARSER]\n Name json\n Format json\n") - command = [os.environ["FLUENT_BIT_BINARY"], "-f", "0.1", "-R", str(parser), "-i", plugin] - if plugin == "unix_socket": - command += ["-p", f"socket_path={address}"] - else: - command += ["-p", "listen=127.0.0.1", "-p", f"port={port}"] - if mode.endswith("-parser"): - command += ["-p", "format=none", "-p", "parser=json"] - if plugin == "syslog": - command += ["-p", "mode=tcp", "-p", "parser=json"] - command += ["-o", "stdout", "-m", "*", "-p", "format=json_lines"] - log = tmp_path / "fluent-bit.log" - memlog = tmp_path / "valgrind.log" - memory = os.environ.get("VALGRIND") == "1" - if memory: - command = ["valgrind", "--leak-check=full", "--show-leak-kinds=all", - "--errors-for-leak-kinds=definite,indirect", "--error-exitcode=99", - f"--log-file={memlog}"] + command - with log.open("w") as output: - process = subprocess.Popen(command, stdout=output, stderr=subprocess.STDOUT) - def ready(): - assert process.poll() is None, log.read_text() - return "[output:stdout:" in log.read_text() - try: - wait_for(ready) - yield address, process, log - finally: - if process.poll() is None: - process.send_signal(signal.SIGTERM) - try: - process.wait(timeout=30) - except subprocess.TimeoutExpired: - process.kill() - process.wait() - pytest.fail("Fluent Bit did not shut down cleanly") - assert process.returncode == 0, log.read_text() + (memlog.read_text() if memory else "") - if memory: - assert "ERROR SUMMARY: 0 errors" in memlog.read_text(), memlog.read_text() + input_config = {"name": mode, "listen": "127.0.0.1", "port": port} + if input_options: + input_config.update(input_options) + config = { + "service": { + "flush": 0.1, + "grace": 1, + "http_server": "on", + "http_listen": "127.0.0.1", + "http_port": "${FLUENT_BIT_HTTP_MONITORING_PORT}", + }, + "pipeline": { + "inputs": [input_config], + "outputs": [{"name": "stdout", "match": "*", "format": "json_lines"}], + }, + } + config_path = tmp_path / "fluent-bit.yaml" + config_path.write_text(json.dumps(config)) + manager = FluentBitManager(str(config_path)) + try: + manager.start() + log = Path(manager.log_file) + yield ("127.0.0.1", port), manager.process, log + finally: + manager.stop() def send(mode, address, payload): - plugin = mode.split("-")[0] - if plugin in ("http", "splunk", "elasticsearch"): - conn = http.client.HTTPConnection(*address, timeout=10) - try: - path = "/test" - if plugin == "splunk": - path = "/services/collector/event" - payload = b'{"event":' + payload + b"}" - elif plugin == "elasticsearch": - path = "/_bulk" - payload = b'{"index":{}}\n' + payload + b"\n" - conn.request("POST", path, payload, {"Content-Type": "application/json"}) - response = conn.getresponse() - response.read() - return response.status - finally: - conn.close() - return - if plugin == "udp": - with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as sock: - sock.sendto(payload + b"\n", address) - return - family = socket.AF_UNIX if plugin == "unix_socket" else socket.AF_INET - with socket.socket(family, socket.SOCK_STREAM) as sock: - sock.settimeout(5) - sock.connect(address) - if plugin == "mqtt": - # A regular MQTT CONNECT followed by a QoS 0 JSON publication. - sock.sendall(b"\x10\x10\x00\x04MQTT\x04\x02\x00\x0a\x00\x04test") - assert sock.recv(4)[0] == 0x20 - body = b"\x00\x01a" + payload - length = len(body) - encoded = bytearray() - while True: - digit = length % 128 - length //= 128 - encoded.append(digit | (0x80 if length else 0)) - if not length: - break - sock.sendall(b"\x30" + bytes(encoded) + body) - else: - sock.sendall(payload if plugin == "forward" else payload + b"\n") + conn = http.client.HTTPConnection(*address, timeout=10) + try: + payload = b'{"index":{}}\n' + payload + b"\n" + conn.request("POST", "/_bulk", payload, {"Content-Type": "application/json"}) + response = conn.getresponse() + response.read() + return response.status + finally: + conn.close() @pytest.mark.parametrize("mode", ["elasticsearch"]) @@ -131,5 +80,22 @@ def test_nested_json_recovery(tmp_path, mode, kind): send(mode, address, b'{"marker":"after"}') wait_for(lambda: '"marker":"after"' in log.read_text()) assert process.poll() is None - # Parser inputs may preserve rejected JSON as a raw log string. assert '"nested":' not in log.read_text() + + +@pytest.mark.parametrize("payload", [b'{"broken":', b'not-json', b'[]', b'42']) +def test_malformed_bulk_recovery(tmp_path, payload): + with daemon(tmp_path, "elasticsearch") as (address, process, log): + assert send("elasticsearch", address, payload) == 400 + assert send("elasticsearch", address, b'{"marker":"after"}') == 200 + wait_for(lambda: '"marker":"after"' in log.read_text()) + assert process.poll() is None + + +def test_busy_ingress_is_retryable(tmp_path): + options = {"http_server.workers": 2, "http_server.ingress_queue_byte_limit": "128"} + with daemon(tmp_path, "elasticsearch", options) as (address, process, log): + assert send("elasticsearch", address, json.dumps({"large": "x" * 1024}).encode()) == 503 + assert send("elasticsearch", address, b'{"marker":"after"}') == 200 + wait_for(lambda: '"marker":"after"' in log.read_text()) + assert process.poll() is None From c204e0b9061a99e69779dc865123dc371aace11d Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 23 Sep 2026 12:31:28 -0600 Subject: [PATCH 6/7] in_elasticsearch: reject unconsumed bulk payload data Signed-off-by: Eduardo Silva --- .../in_elasticsearch/in_elasticsearch_bulk_prot.c | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c b/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c index 94cc7d6ff21..b1f59e32946 100644 --- a/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c +++ b/plugins/in_elasticsearch/in_elasticsearch_bulk_prot.c @@ -468,6 +468,7 @@ static ssize_t parse_payload_ndjson(struct flb_in_elasticsearch *ctx, flb_sds_t { int ret; int out_size; + size_t offset; char *pack; struct flb_pack_state pack_state; @@ -480,6 +481,17 @@ static ssize_t parse_payload_ndjson(struct flb_in_elasticsearch *ctx, flb_sds_t /* Pack JSON as msgpack */ ret = flb_pack_json_state(payload, size, &pack, &out_size, &pack_state); + if (ret == 0) { + /* The packer can succeed with a complete prefix of an incomplete request. */ + for (offset = pack_state.last_byte; offset < size; offset++) { + if (payload[offset] != ' ' && payload[offset] != '\t' && + payload[offset] != '\r' && payload[offset] != '\n') { + flb_free(pack); + ret = FLB_ERR_JSON_INVAL; + break; + } + } + } flb_pack_state_reset(&pack_state); /* Handle exceptions */ From 24282eae4ce41d26e3467947fadfd558ce3d96dc Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 23 Sep 2026 12:31:28 -0600 Subject: [PATCH 7/7] tests: integration: cover trailing elasticsearch bulk data Signed-off-by: Eduardo Silva --- .../tests/test_elasticsearch_invalid_bulk.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py b/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py index 9dd265201b8..328eb50b4fc 100644 --- a/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py +++ b/tests/integration/scenarios/elasticsearch_invalid_bulk/tests/test_elasticsearch_invalid_bulk.py @@ -99,3 +99,22 @@ def test_busy_ingress_is_retryable(tmp_path): assert send("elasticsearch", address, b'{"marker":"after"}') == 200 wait_for(lambda: '"marker":"after"' in log.read_text()) assert process.poll() is None + + +@pytest.mark.parametrize("trailer", [b'{"index":', b'[', b'"unfinished', b'garbage']) +def test_trailing_data_rejected_before_ingestion(tmp_path, trailer): + with daemon(tmp_path, "elasticsearch") as (address, process, log): + payload = b'{"marker":"rejected"}\n' + trailer + assert send("elasticsearch", address, payload) == 400 + assert send("elasticsearch", address, b'{"marker":"after"}') == 200 + wait_for(lambda: '"marker":"after"' in log.read_text()) + assert '"marker":"rejected"' not in log.read_text() + assert process.poll() is None + + +@pytest.mark.parametrize("trailer", [b'', b' \t\r\n \t']) +def test_trailing_whitespace_accepted(tmp_path, trailer): + with daemon(tmp_path, "elasticsearch") as (address, process, log): + assert send("elasticsearch", address, b'{"marker":"accepted"}' + trailer) == 200 + wait_for(lambda: '"marker":"accepted"' in log.read_text()) + assert process.poll() is None