Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 15 additions & 6 deletions tests/kafkatest/tests/core/downgrade_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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):

Expand Down Expand Up @@ -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()

Expand All @@ -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"])
Expand All @@ -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
Expand All @@ -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")
Expand Down