diff --git a/tests/kafkatest/tests/core/downgrade_test.py b/tests/kafkatest/tests/core/downgrade_test.py index 0a7322d2f763a..eccd6c1eb913f 100644 --- a/tests/kafkatest/tests/core/downgrade_test.py +++ b/tests/kafkatest/tests/core/downgrade_test.py @@ -13,7 +13,7 @@ # See the License for the specific language governing permissions and # limitations under the License. -from ducktape.mark import parametrize +from ducktape.mark import parametrize, matrix from ducktape.mark.resource import cluster from kafkatest.services.console_consumer import ConsoleConsumer @@ -23,7 +23,7 @@ from kafkatest.services.zookeeper import ZookeeperService from kafkatest.tests.end_to_end import EndToEndTest from kafkatest.utils import is_int -from kafkatest.version import LATEST_0_9, LATEST_0_10, LATEST_0_10_0, LATEST_0_10_1, LATEST_0_10_2, LATEST_0_11_0, LATEST_1_0, LATEST_1_1, LATEST_2_0, LATEST_2_1, LATEST_2_2, LATEST_2_3, V_0_9_0_0, V_0_11_0_0, DEV_BRANCH, KafkaVersion +from kafkatest.version import LATEST_0_9, LATEST_0_10, LATEST_0_10_0, LATEST_0_10_1, LATEST_0_10_2, LATEST_0_11_0, LATEST_1_0, LATEST_1_1, LATEST_2_0, LATEST_2_1, LATEST_2_2, LATEST_2_3, LATEST_2_4, LATEST_2_5, V_0_9_0_0, V_0_11_0_0, DEV_BRANCH, KafkaVersion class TestDowngrade(EndToEndTest): @@ -52,7 +52,7 @@ def downgrade_to(self, kafka_version): del node.config[config_property.MESSAGE_FORMAT_VERSION] self.kafka.start_node(node) - def setup_services(self, kafka_version, compression_types, security_protocol): + def setup_services(self, kafka_version, compression_types, security_protocol, static_membership): self.create_zookeeper() self.zk.start() @@ -68,10 +68,18 @@ def setup_services(self, kafka_version, compression_types, security_protocol): self.producer.start() self.create_consumer(log_level="DEBUG", - version=kafka_version) + version=kafka_version, + static_membership=static_membership) + self.consumer.start() @cluster(num_nodes=7) + @matrix(version=[str(LATEST_2_5)], compression_types=[["none"]], static_membership=[False, True]) + @parametrize(version=str(LATEST_2_5), compression_types=["zstd"], security_protocol="SASL_SSL") + # static membership was introduced with a buggy verifiable console consumer which + # required static membership to be enabled + @parametrize(version=str(LATEST_2_4), compression_types=["none"], static_membership=True) + @parametrize(version=str(LATEST_2_4), compression_types=["zstd"], security_protocol="SASL_SSL", static_membership=True) @parametrize(version=str(LATEST_2_3), compression_types=["none"]) @parametrize(version=str(LATEST_2_3), compression_types=["zstd"], security_protocol="SASL_SSL") @parametrize(version=str(LATEST_2_2), compression_types=["none"]) @@ -82,7 +90,8 @@ def setup_services(self, kafka_version, compression_types, security_protocol): @parametrize(version=str(LATEST_2_0), compression_types=["snappy"], security_protocol="SASL_SSL") @parametrize(version=str(LATEST_1_1), compression_types=["none"]) @parametrize(version=str(LATEST_1_1), compression_types=["lz4"], security_protocol="SASL_SSL") - def test_upgrade_and_downgrade(self, version, compression_types, security_protocol="PLAINTEXT"): + def test_upgrade_and_downgrade(self, version, compression_types, security_protocol="PLAINTEXT", + static_membership=False): """Test upgrade and downgrade of Kafka cluster from old versions to the current version `version` is the Kafka version to upgrade from and downgrade back to @@ -103,7 +112,7 @@ def test_upgrade_and_downgrade(self, version, compression_types, security_protoc """ kafka_version = KafkaVersion(version) - self.setup_services(kafka_version, compression_types, security_protocol) + self.setup_services(kafka_version, compression_types, security_protocol, static_membership) self.await_startup() self.logger.info("First pass bounce - rolling upgrade")