diff --git a/plugins/out_prometheus_remote_write/remote_write.c b/plugins/out_prometheus_remote_write/remote_write.c index d9e264be671..6bbdef1d761 100644 --- a/plugins/out_prometheus_remote_write/remote_write.c +++ b/plugins/out_prometheus_remote_write/remote_write.c @@ -20,6 +20,7 @@ #include #include #include +#include #include #include @@ -71,6 +72,10 @@ static int http_post(struct prometheus_remote_write_context *ctx, ret = flb_gzip_compress((void *) body, body_len, &payload_buf, &payload_size); } + else if (strcasecmp(ctx->compression, "zstd") == 0) { + ret = flb_zstd_compress((void *) body, body_len, + &payload_buf, &payload_size); + } else { payload_buf = (void *) body; payload_size = body_len; @@ -134,6 +139,13 @@ static int http_post(struct prometheus_remote_write_context *ctx, "gzip", strlen("gzip")); } + else if (strcasecmp(ctx->compression, "zstd") == 0) { + flb_http_add_header(c, + "Content-Encoding", + strlen("Content-Encoding"), + "zstd", + strlen("zstd")); + } /* Basic Auth headers */ if (ctx->http_user && ctx->http_passwd) { @@ -415,7 +427,8 @@ static struct flb_config_map config_map[] = { { FLB_CONFIG_MAP_STR, "compression", "snappy", 0, FLB_TRUE, offsetof(struct prometheus_remote_write_context, compression), - "Compress the payload with either snappy, gzip if set" + "Set the payload compression mechanism. Options available are 'snappy', " + "'gzip' and 'zstd'." }, #ifdef FLB_HAVE_SIGNV4 diff --git a/plugins/out_prometheus_remote_write/remote_write_conf.c b/plugins/out_prometheus_remote_write/remote_write_conf.c index 8f02398e715..f6ee9d61e37 100644 --- a/plugins/out_prometheus_remote_write/remote_write_conf.c +++ b/plugins/out_prometheus_remote_write/remote_write_conf.c @@ -64,6 +64,26 @@ static int config_add_labels(struct flb_output_instance *ins, return 0; } +/* Validate the 'compression' property */ +static int config_validate_compression(struct flb_output_instance *ins, + struct prometheus_remote_write_context *ctx) +{ + if (!ctx->compression) { + return 0; + } + + if (strcasecmp(ctx->compression, "snappy") == 0 || + strcasecmp(ctx->compression, "gzip") == 0 || + strcasecmp(ctx->compression, "zstd") == 0) { + return 0; + } + + flb_plg_error(ins, "invalid 'compression' value '%s', it must be one of " + "'snappy', 'gzip' or 'zstd'", ctx->compression); + + return -1; +} + struct prometheus_remote_write_context *flb_prometheus_remote_write_context_create( struct flb_output_instance *ins, struct flb_config *config) { @@ -94,6 +114,13 @@ struct prometheus_remote_write_context *flb_prometheus_remote_write_context_crea return NULL; } + /* Validate 'compression' */ + ret = config_validate_compression(ins, ctx); + if (ret == -1) { + flb_free(ctx); + return NULL; + } + /* Parse 'add_label' */ ret = config_add_labels(ins, ctx); if (ret == -1) { diff --git a/tests/integration/scenarios/in_prometheus_remote_write/config/sender_compression_wire.yaml b/tests/integration/scenarios/in_prometheus_remote_write/config/sender_compression_wire.yaml new file mode 100644 index 00000000000..23b1db4955c --- /dev/null +++ b/tests/integration/scenarios/in_prometheus_remote_write/config/sender_compression_wire.yaml @@ -0,0 +1,21 @@ +service: + flush: 1 + grace: 1 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + +pipeline: + inputs: + - name: fluentbit_metrics + scrape_interval: 1 + scrape_on_start: true + + outputs: + - name: prometheus_remote_write + match: "*" + host: 127.0.0.1 + port: ${TEST_SUITE_HTTP_PORT} + tls: off + uri: /data + compression: ${PROM_RW_COMPRESSION} diff --git a/tests/integration/scenarios/in_prometheus_remote_write/config/sender_invalid_compression.yaml b/tests/integration/scenarios/in_prometheus_remote_write/config/sender_invalid_compression.yaml new file mode 100644 index 00000000000..4895f51a229 --- /dev/null +++ b/tests/integration/scenarios/in_prometheus_remote_write/config/sender_invalid_compression.yaml @@ -0,0 +1,20 @@ +service: + flush: 1 + grace: 1 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + +pipeline: + inputs: + - name: fluentbit_metrics + scrape_interval: 1 + scrape_on_start: true + + outputs: + - name: prometheus_remote_write + match: "*" + host: 127.0.0.1 + port: ${PROM_RW_RECEIVER_PORT} + uri: /write + compression: not_a_real_algorithm diff --git a/tests/integration/scenarios/in_prometheus_remote_write/config/sender_zstd.yaml b/tests/integration/scenarios/in_prometheus_remote_write/config/sender_zstd.yaml new file mode 100644 index 00000000000..6db7edd5aa1 --- /dev/null +++ b/tests/integration/scenarios/in_prometheus_remote_write/config/sender_zstd.yaml @@ -0,0 +1,20 @@ +service: + flush: 1 + grace: 1 + log_level: info + http_server: on + http_port: ${FLUENT_BIT_HTTP_MONITORING_PORT} + +pipeline: + inputs: + - name: fluentbit_metrics + scrape_interval: 1 + scrape_on_start: true + + outputs: + - name: prometheus_remote_write + match: "*" + host: 127.0.0.1 + port: ${PROM_RW_RECEIVER_PORT} + uri: /write + compression: zstd diff --git a/tests/integration/scenarios/in_prometheus_remote_write/tests/test_in_prometheus_remote_write_001.py b/tests/integration/scenarios/in_prometheus_remote_write/tests/test_in_prometheus_remote_write_001.py index ba18d7fe609..bf8bcff425d 100644 --- a/tests/integration/scenarios/in_prometheus_remote_write/tests/test_in_prometheus_remote_write_001.py +++ b/tests/integration/scenarios/in_prometheus_remote_write/tests/test_in_prometheus_remote_write_001.py @@ -2,8 +2,11 @@ import time import pytest +import requests -from utils.fluent_bit_manager import FluentBitManager +from server.http_server import configure_http_response, data_storage, http_server_run +from utils.fluent_bit_manager import FluentBitManager, FluentBitStartupError +from utils.test_service import FluentBitTestService from utils.input_pause_resume import ( assert_connection_closed, is_valgrind, @@ -152,6 +155,147 @@ def test_in_prometheus_remote_write_matrix(case, workers_enabled): service.stop() +class WireCaptureService: + """ + Sends to the test suite HTTP server instead of a Fluent Bit receiver, so + the outbound request can be inspected before anything decodes it. + """ + + def __init__(self, config_file, extra_env=None): + self.config_file = os.path.abspath( + os.path.join(os.path.dirname(__file__), "../config", config_file) + ) + self.service = FluentBitTestService( + self.config_file, + data_storage=data_storage, + data_keys=["payloads", "requests"], + extra_env=extra_env, + pre_start=self._start_receiver, + post_stop=self._stop_receiver, + ) + + def _start_receiver(self, service): + http_server_run(service.test_suite_http_port) + configure_http_response(status_code=200, body={}) + self.service.wait_for_http_endpoint( + f"http://127.0.0.1:{service.test_suite_http_port}/ping", + timeout=10, + interval=0.5, + ) + + def _stop_receiver(self, service): + try: + requests.post( + f"http://127.0.0.1:{service.test_suite_http_port}/shutdown", + timeout=2, + ) + except requests.RequestException: + pass + + def start(self): + self.service.start() + + def stop(self): + self.service.stop() + + def wait_for_requests(self, minimum_count, timeout=30): + if os.environ.get("VALGRIND"): + timeout = max(timeout * 3, 60) + + return self.service.wait_for_condition( + lambda: data_storage["requests"] + if len(data_storage["requests"]) >= minimum_count + else None, + timeout=timeout, + interval=0.5, + description=f"{minimum_count} outbound remote write requests", + ) + + +@pytest.mark.parametrize("compression", ["snappy", "gzip", "zstd"]) +def test_in_prometheus_remote_write_compression_content_encoding(compression): + """ + The sender must advertise the compression it actually applied. + + This is checked on the wire rather than at a Fluent Bit receiver, because + a receiver cannot distinguish the failure case. When a compression branch + is not applied the body is sent uncompressed and without a + Content-Encoding header, and an uncompressed remote write body still + decodes as protobuf, so the metrics would arrive and look correct. + + All three supported algorithms are covered because they share the same + fall through, so a regression in any of them fails the same silent way. + """ + service = WireCaptureService( + "sender_compression_wire.yaml", + extra_env={"PROM_RW_COMPRESSION": compression}, + ) + service.start() + + try: + requests_seen = service.wait_for_requests(1) + finally: + service.stop() + + headers = requests_seen[0]["headers"] + assert headers.get("Content-Encoding") == compression + assert headers.get("Content-Type") == "application/x-protobuf" + + +def test_in_prometheus_remote_write_zstd_compression(): + """ + The sender compresses the payload with zstd and the receiver decodes it. + + Metrics arriving is not sufficient on its own. When decompression fails, + flb_http_request_uncompress_body() leaves the body untouched and still + reports success, so an uncompressed payload sent with a zstd + Content-Encoding header would also decode as protobuf and produce metrics. + That case is only distinguishable by the decompression failure the + receiver logs, so assert it never appears. + """ + service = Service("receiver_http1_cleartext.yaml", "sender_zstd.yaml") + service.start() + + try: + receiver_log = service.wait_for_log( + service.receiver.log_file, + "fluentbit_input_metrics_scrapes_total", + timeout=40, + interval=1, + ) + assert f"listening on 127.0.0.1:{service.receiver_port}" in receiver_log + assert "fluentbit_input_metrics_scrapes_total" in receiver_log + + receiver_log = _read_file(service.receiver.log_file) + assert "[http zstd] decompression failed" not in receiver_log + finally: + service.stop() + + +def test_in_prometheus_remote_write_rejects_invalid_compression(): + """ + An unrecognized 'compression' value must fail during initialization. It + used to be treated as 'no compression' without logging anything, which + made a typo indistinguishable from a working configuration. + """ + service = Service( + "receiver_http1_cleartext.yaml", + "sender_invalid_compression.yaml", + ) + service.start(start_sender=False) + + try: + with pytest.raises(FluentBitStartupError): + service.start_sender() + + sender_log = _read_file(service.sender.log_file) + assert "invalid 'compression' value 'not_a_real_algorithm'" in sender_log + assert "it must be one of 'snappy', 'gzip' or 'zstd'" in sender_log + assert "failed to initialize 'prometheus_remote_write' plugin" in sender_log + finally: + service.stop() + + @pytest.mark.parametrize( "receiver_config", [