From 00c4d3ba804ff9f0e050acccc979bf4c1eba0f97 Mon Sep 17 00:00:00 2001 From: Bernd Ahlers Date: Thu, 16 Jul 2026 13:23:08 +0200 Subject: [PATCH 1/3] Extend KafkaTransportIT with tests for compression options (#26686) Adding tests for gzip, snappy, lz4, and zstd compression options. --- .../graylog/testing/kafka/KafkaContainer.java | 21 +++++- .../inputs/transports/KafkaTransportIT.java | 71 ++++++++++++++++--- 2 files changed, 80 insertions(+), 12 deletions(-) diff --git a/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java b/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java index 1f6dd3504c35..fdaf0fc5693c 100644 --- a/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java +++ b/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java @@ -170,17 +170,36 @@ protected void containerIsStarted(InspectContainerResponse containerInfo) { /** * Returns a new {@code KafkaProducer} instance that is connected to the Kafka container. + * The producer doesn't compress record batches. * * @return the new producer instance */ public KafkaProducer createByteArrayProducer() { + return createByteArrayProducer("none"); + } + + /** + * Returns a new {@code KafkaProducer} instance that is connected to the Kafka container and + * compresses record batches using the given compression type. + *

+ * To make the producer actually build a compressed record batch that contains more than a single record, + * {@code linger.ms} is raised so that records sent in quick succession are grouped into the same batch. + * + * @param compressionType the {@code compression.type} producer setting (e.g. {@code none}, {@code gzip}, + * {@code snappy}, {@code lz4} or {@code zstd}) + * @return the new producer instance + */ + public KafkaProducer createByteArrayProducer(String compressionType) { final var props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:" + getKafkaPort()); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); props.put(ProducerConfig.CLIENT_ID_CONFIG, "graylog-node-" + UUID.randomUUID()); props.put(ProducerConfig.ACKS_CONFIG, "1"); - props.put(ProducerConfig.LINGER_MS_CONFIG, 0); + props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, requireNonNull(compressionType, "compressionType cannot be null")); + // Give records a chance to accumulate in a single (compressed) record batch instead of being sent + // immediately in a batch of their own. + props.put(ProducerConfig.LINGER_MS_CONFIG, "none".equals(compressionType) ? 0 : 100); return new KafkaProducer<>(props); } diff --git a/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java b/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java index dc2435a2b90b..98fea01a6870 100644 --- a/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java +++ b/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java @@ -30,6 +30,8 @@ import org.graylog2.shared.SuppressForbidden; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import org.mockito.ArgumentCaptor; import org.mockito.Captor; import org.mockito.junit.jupiter.MockitoExtension; @@ -37,10 +39,12 @@ import org.testcontainers.junit.jupiter.Testcontainers; import java.nio.charset.StandardCharsets; +import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import java.util.stream.IntStream; import static org.assertj.core.api.Assertions.assertThat; import static org.graylog2.shared.utilities.StringUtils.f; @@ -59,17 +63,66 @@ class KafkaTransportIT { ArgumentCaptor messageCaptor; @Test - @SuppressForbidden("Executors.newSingleThreadScheduledExecutor is okay in tests") void basicConsumer() throws Exception { - KAFKA.createTopic("test"); + final var topic = "test"; + KAFKA.createTopic(topic); final var messageValue = UUID.randomUUID().toString().getBytes(StandardCharsets.UTF_8); - final ProducerRecord record = new ProducerRecord<>("test", messageValue); + final ProducerRecord record = new ProducerRecord<>(topic, messageValue); try (KafkaProducer producer = KAFKA.createByteArrayProducer()) { producer.send(record).get(30, TimeUnit.SECONDS); } + final var input = launchTransport(topic, "basic-consumer"); + + verify(input, timeout(5_000).times(1)).processRawMessage(messageCaptor.capture()); + + assertThat(messageCaptor.getValue()).isNotNull().satisfies(rawMessage -> { + assertThat(rawMessage.getId()).isNotNull(); + assertThat(rawMessage.getPayload()).isEqualTo(messageValue); + }); + } + + /** + * Verifies that the transport can consume record batches that were compressed by the producer. Consuming + * compressed batches requires the matching compression codec (and its native library) to be present on the + * classpath, so this test guards against a codec library going missing at runtime. + * + * @see PR #26674 + */ + @ParameterizedTest + @ValueSource(strings = {"gzip", "snappy", "lz4", "zstd"}) + void compressedConsumer(String compressionType) throws Exception { + final var topic = f("test-%s", compressionType); + KAFKA.createTopic(topic); + + // Produce a batch of records with compressible (repetitive) payloads so the codec actually kicks in. + final List messageValues = IntStream.range(0, 10) + .mapToObj(i -> f("%s-compressed-message-%d-%s", compressionType, i, "x".repeat(256))) + .toList(); + + try (KafkaProducer producer = KAFKA.createByteArrayProducer(compressionType)) { + for (final String messageValue : messageValues) { + producer.send(new ProducerRecord<>(topic, messageValue.getBytes(StandardCharsets.UTF_8))); + } + producer.flush(); + } + + final var input = launchTransport(topic, f("compressed-consumer-%s", compressionType)); + + verify(input, timeout(10_000).times(messageValues.size())).processRawMessage(messageCaptor.capture()); + + // Compare the payloads as strings because AssertJ compares byte[] elements by reference, not by content. + final List receivedPayloads = messageCaptor.getAllValues().stream() + .map(rawMessage -> new String(rawMessage.getPayload(), StandardCharsets.UTF_8)) + .toList(); + + assertThat(receivedPayloads).containsExactlyInAnyOrderElementsOf(messageValues); + } + + @SuppressForbidden("Executors.newSingleThreadScheduledExecutor is okay in tests") + private MessageInput launchTransport(String topicFilter, String groupId) throws Exception { final var serverStatus = mock(ServerStatus.class); final var config = new Configuration(Map.of( KafkaTransport.CK_LEGACY, false, @@ -77,8 +130,9 @@ void basicConsumer() throws Exception { KafkaTransport.CK_BOOTSTRAP, f("localhost:%d", KAFKA.getKafkaPort()), KafkaTransport.CK_FETCH_MIN_BYTES, 1, KafkaTransport.CK_FETCH_WAIT_MAX, 100, - KafkaTransport.CK_TOPIC_FILTER, "test", - KafkaTransport.CK_OFFSET_RESET, "smallest" + KafkaTransport.CK_TOPIC_FILTER, topicFilter, + KafkaTransport.CK_OFFSET_RESET, "smallest", + KafkaTransport.CK_GROUP_ID, groupId )); final var transport = new KafkaTransport( config, @@ -94,11 +148,6 @@ KafkaTransport.CK_BOOTSTRAP, f("localhost:%d", KAFKA.getKafkaPort()), transport.lifecycleStateChange(Lifecycle.RUNNING); // Required to set paused=false transport.launch(input); - verify(input, timeout(5_000).times(1)).processRawMessage(messageCaptor.capture()); - - assertThat(messageCaptor.getValue()).isNotNull().satisfies(rawMessage -> { - assertThat(rawMessage.getId()).isNotNull(); - assertThat(rawMessage.getPayload()).isEqualTo(messageValue); - }); + return input; } } From 3c88b2c0ae0ad0e4a516efb0e7630af983a88727 Mon Sep 17 00:00:00 2001 From: Bernd Ahlers Date: Thu, 16 Jul 2026 16:59:43 +0200 Subject: [PATCH 2/3] Add SlowTest and SlowParameterizedTest annotations (#26686) Should be use to annotate tests which run multiple seconds. --- .../testing/SlowParameterizedTest.java | 39 +++++++++++++++++++ .../java/org/graylog/testing/SlowTest.java | 38 ++++++++++++++++++ 2 files changed, 77 insertions(+) create mode 100644 graylog2-server/src/test/java/org/graylog/testing/SlowParameterizedTest.java create mode 100644 graylog2-server/src/test/java/org/graylog/testing/SlowTest.java diff --git a/graylog2-server/src/test/java/org/graylog/testing/SlowParameterizedTest.java b/graylog2-server/src/test/java/org/graylog/testing/SlowParameterizedTest.java new file mode 100644 index 000000000000..c5f0892bbf39 --- /dev/null +++ b/graylog2-server/src/test/java/org/graylog/testing/SlowParameterizedTest.java @@ -0,0 +1,39 @@ +/* + * Copyright (C) 2020 Graylog, Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the Server Side Public License, version 1, + * as published by MongoDB, Inc. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * Server Side Public License for more details. + * + * You should have received a copy of the Server Side Public License + * along with this program. If not, see + * . + */ +package org.graylog.testing; + +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.platform.commons.annotation.Testable; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +/** + * Use this annotation instead of the {@link ParameterizedTest} annotation to tag slow tests. Rule of thumb: everything that takes + * multiple seconds. + */ +@Retention(RetentionPolicy.RUNTIME) +@Target(ElementType.METHOD) +@Testable +@Tag("slow-test") +@ParameterizedTest +public @interface SlowParameterizedTest { +} diff --git a/graylog2-server/src/test/java/org/graylog/testing/SlowTest.java b/graylog2-server/src/test/java/org/graylog/testing/SlowTest.java new file mode 100644 index 000000000000..7e651310b686 --- /dev/null +++ b/graylog2-server/src/test/java/org/graylog/testing/SlowTest.java @@ -0,0 +1,38 @@ +/* + * Copyright (C) 2020 Graylog, Inc. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the Server Side Public License, version 1, + * as published by MongoDB, Inc. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * Server Side Public License for more details. + * + * You should have received a copy of the Server Side Public License + * along with this program. If not, see + * . + */ +package org.graylog.testing; + +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.platform.commons.annotation.Testable; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +/** + * Use this annotation instead of the {@link Test} annotation to tag slow tests. Rule of thumb: everything that takes + * multiple seconds. + */ +@Retention(RetentionPolicy.RUNTIME) +@Target(ElementType.METHOD) +@Testable +@Tag("slow-test") +@Test +public @interface SlowTest { +} From 80909e5d55fca32ec7253963fa8be3a253ccf16b Mon Sep 17 00:00:00 2001 From: Bernd Ahlers Date: Thu, 16 Jul 2026 17:00:23 +0200 Subject: [PATCH 3/3] Run Kafka transport tests for versions 3.7 - 4.3 (#26686) --- .../graylog/testing/kafka/KafkaContainer.java | 14 +- .../inputs/transports/KafkaTransportIT.java | 127 ++++++++++++++++-- 2 files changed, 123 insertions(+), 18 deletions(-) diff --git a/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java b/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java index fdaf0fc5693c..fb64be6e4032 100644 --- a/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java +++ b/graylog2-server/src/test/java/org/graylog/testing/kafka/KafkaContainer.java @@ -45,11 +45,17 @@ /** * An Apache Kafka container that is optimized for fast startup. The container is using the - * bitnami/kafka image. + * apache/kafka image. */ public class KafkaContainer extends GenericContainer { public enum Version { - V34("3.7.0"); + V37("3.7.2"), + V38("3.8.1"), + V39("3.9.2"), + V40("4.0.2"), + V41("4.1.2"), + V42("4.2.1"), + V43("4.3.1"); private final String version; @@ -63,7 +69,7 @@ public String getVersion() { } private static final Logger LOG = LoggerFactory.getLogger(KafkaContainer.class); - private static final Version DEFAULT_VERSION = Version.V34; + private static final Version DEFAULT_VERSION = Version.V37; private static final String KAFKA_ADVERTISED_LISTENERS_FILE = "/.env-kafka-advertised-listeners"; private Admin adminClient = null; @@ -186,7 +192,7 @@ public KafkaProducer createByteArrayProducer() { * {@code linger.ms} is raised so that records sent in quick succession are grouped into the same batch. * * @param compressionType the {@code compression.type} producer setting (e.g. {@code none}, {@code gzip}, - * {@code snappy}, {@code lz4} or {@code zstd}) + * {@code snappy}, {@code lz4} or {@code zstd}) * @return the new producer instance */ public KafkaProducer createByteArrayProducer(String compressionType) { diff --git a/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java b/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java index 98fea01a6870..78cd396a47ae 100644 --- a/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java +++ b/graylog2-server/src/test/java/org/graylog2/inputs/transports/KafkaTransportIT.java @@ -19,6 +19,8 @@ import com.google.common.eventbus.EventBus; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; +import org.graylog.testing.SlowParameterizedTest; +import org.graylog.testing.SlowTest; import org.graylog.testing.kafka.KafkaContainer; import org.graylog2.plugin.LocalMetricRegistry; import org.graylog2.plugin.ServerStatus; @@ -28,9 +30,8 @@ import org.graylog2.plugin.lifecycles.Lifecycle; import org.graylog2.plugin.system.SimpleNodeId; import org.graylog2.shared.SuppressForbidden; -import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.extension.ExtendWith; -import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; import org.mockito.ArgumentCaptor; import org.mockito.Captor; @@ -39,6 +40,7 @@ import org.testcontainers.junit.jupiter.Testcontainers; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.UUID; @@ -53,24 +55,43 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; -@Testcontainers +/** + * Runs the Kafka transport tests against a specific Kafka version. Concrete subclasses (see below) each provide + * their own {@link KafkaContainer} of a fixed version, so the same set of tests runs against every supported + * Kafka broker version. + */ @ExtendWith(MockitoExtension.class) -class KafkaTransportIT { - @Container - private static final KafkaContainer KAFKA = KafkaContainer.create(); - +abstract class KafkaTransportIT { @Captor ArgumentCaptor messageCaptor; - @Test + private final List launchedTransports = new ArrayList<>(); + + /** + * Returns the Kafka container to run the tests against. Each subclass provides a container for a specific + * Kafka version. + * + * @return the Kafka container + */ + protected abstract KafkaContainer kafka(); + + @AfterEach + void stopTransports() { + // Stop the transports so their consumer threads shut down. Otherwise leaked consumers keep retrying + // against the (torn down) broker and flood the log with connection warnings. + launchedTransports.forEach(KafkaTransport::stop); + launchedTransports.clear(); + } + + @SlowTest void basicConsumer() throws Exception { final var topic = "test"; - KAFKA.createTopic(topic); + kafka().createTopic(topic); final var messageValue = UUID.randomUUID().toString().getBytes(StandardCharsets.UTF_8); final ProducerRecord record = new ProducerRecord<>(topic, messageValue); - try (KafkaProducer producer = KAFKA.createByteArrayProducer()) { + try (KafkaProducer producer = kafka().createByteArrayProducer()) { producer.send(record).get(30, TimeUnit.SECONDS); } @@ -91,18 +112,18 @@ void basicConsumer() throws Exception { * * @see PR #26674 */ - @ParameterizedTest + @SlowParameterizedTest @ValueSource(strings = {"gzip", "snappy", "lz4", "zstd"}) void compressedConsumer(String compressionType) throws Exception { final var topic = f("test-%s", compressionType); - KAFKA.createTopic(topic); + kafka().createTopic(topic); // Produce a batch of records with compressible (repetitive) payloads so the codec actually kicks in. final List messageValues = IntStream.range(0, 10) .mapToObj(i -> f("%s-compressed-message-%d-%s", compressionType, i, "x".repeat(256))) .toList(); - try (KafkaProducer producer = KAFKA.createByteArrayProducer(compressionType)) { + try (KafkaProducer producer = kafka().createByteArrayProducer(compressionType)) { for (final String messageValue : messageValues) { producer.send(new ProducerRecord<>(topic, messageValue.getBytes(StandardCharsets.UTF_8))); } @@ -127,7 +148,7 @@ private MessageInput launchTransport(String topicFilter, String groupId) throws final var config = new Configuration(Map.of( KafkaTransport.CK_LEGACY, false, KafkaTransport.CK_THREADS, 1, - KafkaTransport.CK_BOOTSTRAP, f("localhost:%d", KAFKA.getKafkaPort()), + KafkaTransport.CK_BOOTSTRAP, f("localhost:%d", kafka().getKafkaPort()), KafkaTransport.CK_FETCH_MIN_BYTES, 1, KafkaTransport.CK_FETCH_WAIT_MAX, 100, KafkaTransport.CK_TOPIC_FILTER, topicFilter, @@ -147,7 +168,85 @@ KafkaTransport.CK_BOOTSTRAP, f("localhost:%d", KAFKA.getKafkaPort()), transport.lifecycleStateChange(Lifecycle.RUNNING); // Required to set paused=false transport.launch(input); + launchedTransports.add(transport); return input; } } + +@Testcontainers +class KafkaTransport37IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V37); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +} + +@Testcontainers +class KafkaTransport38IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V38); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +} + +@Testcontainers +class KafkaTransport39IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V39); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +} + +@Testcontainers +class KafkaTransport40IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V40); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +} + +@Testcontainers +class KafkaTransport41IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V41); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +} + +@Testcontainers +class KafkaTransport42IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V42); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +} + +@Testcontainers +class KafkaTransport43IT extends KafkaTransportIT { + @Container + private static final KafkaContainer KAFKA = KafkaContainer.create(KafkaContainer.Version.V43); + + @Override + protected KafkaContainer kafka() { + return KAFKA; + } +}