mjsax commented on code in PR #23146: URL: https://github.com/apache/kafka/pull/23146#discussion_r3826187427
########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) Review Comment: Why is this necessary? Should be set to `1` by default. Guess this might be legacy code from 4.1 EA release, which just got c&p around. We might want to do some cleanup (maybe new PR -- if small enough, also ok to just add to this PR) to remove from other tests, too ########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) Review Comment: Claude also pointed out, that there is some other potential legacy code. `use_streams_groups also sets UNSTABLE_API_VERSIONS_ENABLE and UNSTABLE_FEATURE_VERSIONS_ENABLE in kafka.py` -- this also sounds like 4.1 EA stuff, which required to enable unable API/Features to allows us to enable "streams" protocol. ########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) + + def start_processor(self, client_version, extra_configs=None): + """ + Start a single Streams instance on the streams rebalance protocol and wait until it is + RUNNING. Reaching RUNNING is itself the core protocol assertion: it requires a successful + join, topology initialization, an assignment from the group coordinator, and task creation. + """ + processor = StreamsUpgradeTestJobRunnerService(self.test_context, self.kafka) + # An empty version string makes kafka-run-class.sh use the trunk build rather than an + # installed release, which is how the harness selects the DEV client. + processor.set_version("" if client_version == str(DEV_BRANCH) else client_version) + processor.set_config("group.protocol", "streams") + for key, value in (extra_configs or {}).items(): + processor.set_config(key, value) + + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + processor.start() + monitor.wait_until(self.RUNNING_LOG, + timeout_sec=120, + err_msg="Streams client (version '%s') never reached RUNNING on %s" + % (client_version, str(processor.node.account))) + return processor + + def count_in_file(self, node, literal, path): + """Number of lines in `path` containing `literal`, matched as a fixed string.""" + cmd = "grep -c -F -- '%s' %s 2>/dev/null || true" % (literal, path) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def count_unsupported_version_errors(self, node, path): + """ + Number of UnsupportedVersionException lines in `path`, ignoring the benign KIP-714 + telemetry probe (see BENIGN_TELEMETRY_UVE_LOG). + """ + cmd = "grep -F -- '%s' %s 2>/dev/null | grep -c -v -F -- '%s' || true" \ + % (self.UNSUPPORTED_VERSION_LOG, path, self.BENIGN_TELEMETRY_UVE_LOG) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def run_streams_groups_command(self, args, tool_version=None): + """ + Run bin/kafka-streams-groups.sh against the cluster, from the install of `tool_version` + (the trunk build when omitted). Returns the combined stdout/stderr; the command is allowed + to fail so callers can assert on how a failure is reported. + """ + node = self.kafka.nodes[0] + version = DEV_BRANCH if tool_version is None else KafkaVersion(tool_version) + script = self.kafka.path.script("kafka-streams-groups.sh", version) + cmd = "%s --bootstrap-server %s %s 2>&1" % (script, self.kafka.bootstrap_servers(), args) + self.logger.info("Running streams-groups command: %s" % cmd) + return node.account.ssh_output(cmd, allow_fail=True).decode("utf-8", errors="replace") + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_new_client_old_broker(self, broker_version, metadata_quorum): Review Comment: If we say "new client old broker" it seems `broker_version` above should only be 4.2/4.2 but not DEV_BRANCH? But the test below does assert that acceptable-recovery-lag is provided when broker is new -- so maybe the test name is just wrong? -- On the other hand, for 4.4/4.4 broker/client we can cover this in unit/integration tests, so it seems not necessary to add as a system test? ########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) + + def start_processor(self, client_version, extra_configs=None): + """ + Start a single Streams instance on the streams rebalance protocol and wait until it is + RUNNING. Reaching RUNNING is itself the core protocol assertion: it requires a successful + join, topology initialization, an assignment from the group coordinator, and task creation. + """ + processor = StreamsUpgradeTestJobRunnerService(self.test_context, self.kafka) + # An empty version string makes kafka-run-class.sh use the trunk build rather than an + # installed release, which is how the harness selects the DEV client. + processor.set_version("" if client_version == str(DEV_BRANCH) else client_version) + processor.set_config("group.protocol", "streams") + for key, value in (extra_configs or {}).items(): + processor.set_config(key, value) + + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + processor.start() + monitor.wait_until(self.RUNNING_LOG, + timeout_sec=120, + err_msg="Streams client (version '%s') never reached RUNNING on %s" + % (client_version, str(processor.node.account))) + return processor + + def count_in_file(self, node, literal, path): + """Number of lines in `path` containing `literal`, matched as a fixed string.""" + cmd = "grep -c -F -- '%s' %s 2>/dev/null || true" % (literal, path) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def count_unsupported_version_errors(self, node, path): + """ + Number of UnsupportedVersionException lines in `path`, ignoring the benign KIP-714 + telemetry probe (see BENIGN_TELEMETRY_UVE_LOG). + """ + cmd = "grep -F -- '%s' %s 2>/dev/null | grep -c -v -F -- '%s' || true" \ + % (self.UNSUPPORTED_VERSION_LOG, path, self.BENIGN_TELEMETRY_UVE_LOG) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def run_streams_groups_command(self, args, tool_version=None): + """ + Run bin/kafka-streams-groups.sh against the cluster, from the install of `tool_version` + (the trunk build when omitted). Returns the combined stdout/stderr; the command is allowed + to fail so callers can assert on how a failure is reported. + """ + node = self.kafka.nodes[0] + version = DEV_BRANCH if tool_version is None else KafkaVersion(tool_version) + script = self.kafka.path.script("kafka-streams-groups.sh", version) + cmd = "%s --bootstrap-server %s %s 2>&1" % (script, self.kafka.bootstrap_servers(), args) + self.logger.info("Running streams-groups command: %s" % cmd) + return node.account.ssh_output(cmd, allow_fail=True).decode("utf-8", errors="replace") + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_new_client_old_broker(self, broker_version, metadata_quorum): + """ + A 4.4 client must work against a 4.2/4.3 broker: it must not send a request the older + broker rejects, and it must tolerate a v0 response in which the new int64 + AcceptableRecoveryLag is absent and therefore reads back as its default of -1. + """ + self.setup_kafka(broker_version) + processor = self.start_processor(str(DEV_BRANCH)) + + lag_not_provided = self.count_in_file(processor.node, self.OLD_BROKER_LAG_LOG, processor.LOG_FILE) + if broker_version == str(DEV_BRANCH): + assert lag_not_provided == 0, \ Review Comment: Is this check sufficient? We only have a negative test, but not a positive test -- for a new broker, should we assert that we see that `acceptable.recover.lag=10000` (I believe 10K is the default) gets logged? ########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) + + def start_processor(self, client_version, extra_configs=None): + """ + Start a single Streams instance on the streams rebalance protocol and wait until it is + RUNNING. Reaching RUNNING is itself the core protocol assertion: it requires a successful + join, topology initialization, an assignment from the group coordinator, and task creation. + """ + processor = StreamsUpgradeTestJobRunnerService(self.test_context, self.kafka) + # An empty version string makes kafka-run-class.sh use the trunk build rather than an + # installed release, which is how the harness selects the DEV client. + processor.set_version("" if client_version == str(DEV_BRANCH) else client_version) + processor.set_config("group.protocol", "streams") + for key, value in (extra_configs or {}).items(): + processor.set_config(key, value) + + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + processor.start() + monitor.wait_until(self.RUNNING_LOG, + timeout_sec=120, + err_msg="Streams client (version '%s') never reached RUNNING on %s" + % (client_version, str(processor.node.account))) + return processor + + def count_in_file(self, node, literal, path): + """Number of lines in `path` containing `literal`, matched as a fixed string.""" + cmd = "grep -c -F -- '%s' %s 2>/dev/null || true" % (literal, path) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def count_unsupported_version_errors(self, node, path): + """ + Number of UnsupportedVersionException lines in `path`, ignoring the benign KIP-714 + telemetry probe (see BENIGN_TELEMETRY_UVE_LOG). + """ + cmd = "grep -F -- '%s' %s 2>/dev/null | grep -c -v -F -- '%s' || true" \ + % (self.UNSUPPORTED_VERSION_LOG, path, self.BENIGN_TELEMETRY_UVE_LOG) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def run_streams_groups_command(self, args, tool_version=None): + """ + Run bin/kafka-streams-groups.sh against the cluster, from the install of `tool_version` + (the trunk build when omitted). Returns the combined stdout/stderr; the command is allowed + to fail so callers can assert on how a failure is reported. + """ + node = self.kafka.nodes[0] + version = DEV_BRANCH if tool_version is None else KafkaVersion(tool_version) + script = self.kafka.path.script("kafka-streams-groups.sh", version) + cmd = "%s --bootstrap-server %s %s 2>&1" % (script, self.kafka.bootstrap_servers(), args) + self.logger.info("Running streams-groups command: %s" % cmd) + return node.account.ssh_output(cmd, allow_fail=True).decode("utf-8", errors="replace") + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_new_client_old_broker(self, broker_version, metadata_quorum): + """ + A 4.4 client must work against a 4.2/4.3 broker: it must not send a request the older + broker rejects, and it must tolerate a v0 response in which the new int64 + AcceptableRecoveryLag is absent and therefore reads back as its default of -1. + """ + self.setup_kafka(broker_version) + processor = self.start_processor(str(DEV_BRANCH)) + + lag_not_provided = self.count_in_file(processor.node, self.OLD_BROKER_LAG_LOG, processor.LOG_FILE) + if broker_version == str(DEV_BRANCH): + assert lag_not_provided == 0, \ + "A 4.4 broker negotiates heartbeat v1 and must supply acceptableRecoveryLag, " \ + "but the client logged it as not provided" + else: + assert lag_not_provided > 0, \ + "Against a %s broker the heartbeat response is v0 and carries no int64 " \ + "acceptableRecoveryLag, so the client should have logged it as not provided" \ + % broker_version + + assert self.count_unsupported_version_errors(processor.node, processor.LOG_FILE) == 0, \ + "The 4.4 client hit an UnsupportedVersionException against a %s broker" % broker_version + + processor.stop() Review Comment: Should we also assert that there was no error on the broker? Or would this be redundant? Not sure from top of my head, but if the broker received something unexpected, it would also log an error. I am particularly thinking about task-offset-sum / task-offset-end-sum in the HB request-- these fields are already present in v0, but should _never_ be set (and I believe 4.2/4.3 brokers check that both fields are always zero, and we did lift this check with 4.4 brokers). Worth checking if we made similar changes for other field. Also wondering about now populated fields in the response like `taskOffsetIntervalMs` ? ########## tests/kafkatest/tests/streams/streams_topology_description_plugin_test.py: ########## @@ -161,3 +163,42 @@ def test_topology_description_not_stored_without_plugin(self, metadata_quorum): assert int(next(pushed).strip()) == 0, \ "Client logged a successful push despite no plugin being configured on the broker" processor.stop() + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3)], + metadata_quorum=[quorum.combined_kraft]) + def test_no_push_solicited_by_old_broker(self, broker_version, metadata_quorum): + """ + Test a 4.4 client against a 4.2/4.3 broker. The broker solicits a topology description push + by setting topologyDescriptionRequired on the heartbeat response, but that field only exists + in response version 1 (KIP-1331), so a 4.2/4.3 broker can never ask. Such a broker also does + not know the plugin config and only warns about it, which this configures deliberately: even Review Comment: > Such a broker also does not know the plugin config and only warns about it Not sure what this means? Are you saying, if the config is provided broker side the broker just logs "unknown" -- sure, but this seems to be irrelevant for this test? > even with the config present, the older broker must not cause the newer client to push. Similar. An older broker really cannot "cause" a new client to push -- and the test is not about this. It's really about a buggy client which would push even if there was no signal from the broker, that the broker understand the RPC. ########## tests/kafkatest/tests/streams/streams_topology_description_plugin_test.py: ########## @@ -161,3 +163,42 @@ def test_topology_description_not_stored_without_plugin(self, metadata_quorum): assert int(next(pushed).strip()) == 0, \ "Client logged a successful push despite no plugin being configured on the broker" processor.stop() + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3)], + metadata_quorum=[quorum.combined_kraft]) + def test_no_push_solicited_by_old_broker(self, broker_version, metadata_quorum): + """ + Test a 4.4 client against a 4.2/4.3 broker. The broker solicits a topology description push Review Comment: `Test a 4.4 client` -- this comment is only "correct" on 4.4 branch, but on trunk it would already be 4.5. Can we phrase this differently, so the comment does not outdate, each time `trunk` version is bumped. Might apply to other comments, too -- just noticed here for the first time. ########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) + + def start_processor(self, client_version, extra_configs=None): + """ + Start a single Streams instance on the streams rebalance protocol and wait until it is + RUNNING. Reaching RUNNING is itself the core protocol assertion: it requires a successful + join, topology initialization, an assignment from the group coordinator, and task creation. + """ + processor = StreamsUpgradeTestJobRunnerService(self.test_context, self.kafka) + # An empty version string makes kafka-run-class.sh use the trunk build rather than an + # installed release, which is how the harness selects the DEV client. + processor.set_version("" if client_version == str(DEV_BRANCH) else client_version) + processor.set_config("group.protocol", "streams") + for key, value in (extra_configs or {}).items(): + processor.set_config(key, value) + + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + processor.start() + monitor.wait_until(self.RUNNING_LOG, + timeout_sec=120, + err_msg="Streams client (version '%s') never reached RUNNING on %s" + % (client_version, str(processor.node.account))) + return processor + + def count_in_file(self, node, literal, path): + """Number of lines in `path` containing `literal`, matched as a fixed string.""" + cmd = "grep -c -F -- '%s' %s 2>/dev/null || true" % (literal, path) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def count_unsupported_version_errors(self, node, path): + """ + Number of UnsupportedVersionException lines in `path`, ignoring the benign KIP-714 + telemetry probe (see BENIGN_TELEMETRY_UVE_LOG). + """ + cmd = "grep -F -- '%s' %s 2>/dev/null | grep -c -v -F -- '%s' || true" \ + % (self.UNSUPPORTED_VERSION_LOG, path, self.BENIGN_TELEMETRY_UVE_LOG) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def run_streams_groups_command(self, args, tool_version=None): + """ + Run bin/kafka-streams-groups.sh against the cluster, from the install of `tool_version` + (the trunk build when omitted). Returns the combined stdout/stderr; the command is allowed + to fail so callers can assert on how a failure is reported. + """ + node = self.kafka.nodes[0] + version = DEV_BRANCH if tool_version is None else KafkaVersion(tool_version) + script = self.kafka.path.script("kafka-streams-groups.sh", version) + cmd = "%s --bootstrap-server %s %s 2>&1" % (script, self.kafka.bootstrap_servers(), args) + self.logger.info("Running streams-groups command: %s" % cmd) + return node.account.ssh_output(cmd, allow_fail=True).decode("utf-8", errors="replace") + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_new_client_old_broker(self, broker_version, metadata_quorum): + """ + A 4.4 client must work against a 4.2/4.3 broker: it must not send a request the older + broker rejects, and it must tolerate a v0 response in which the new int64 + AcceptableRecoveryLag is absent and therefore reads back as its default of -1. + """ + self.setup_kafka(broker_version) + processor = self.start_processor(str(DEV_BRANCH)) + + lag_not_provided = self.count_in_file(processor.node, self.OLD_BROKER_LAG_LOG, processor.LOG_FILE) + if broker_version == str(DEV_BRANCH): + assert lag_not_provided == 0, \ + "A 4.4 broker negotiates heartbeat v1 and must supply acceptableRecoveryLag, " \ + "but the client logged it as not provided" + else: + assert lag_not_provided > 0, \ + "Against a %s broker the heartbeat response is v0 and carries no int64 " \ + "acceptableRecoveryLag, so the client should have logged it as not provided" \ + % broker_version + + assert self.count_unsupported_version_errors(processor.node, processor.LOG_FILE) == 0, \ + "The 4.4 client hit an UnsupportedVersionException against a %s broker" % broker_version + + processor.stop() + + @cluster(num_nodes=2) + @matrix(client_version=[str(LATEST_4_2), str(LATEST_4_3)], + metadata_quorum=[quorum.combined_kraft]) + def test_old_client_new_broker(self, client_version, metadata_quorum): + """ + A 4.2/4.3 client must keep working against a 4.4 broker. The broker sets the v1-only + AcceptableRecoveryLag unconditionally and relies on it being `ignorable` so that it is + dropped when the response is serialized at v0; an older client must not see a malformed + response. + """ + self.setup_kafka(str(DEV_BRANCH)) + processor = self.start_processor(client_version) + + assert self.count_unsupported_version_errors(processor.node, processor.LOG_FILE) == 0, \ + "The %s client hit an UnsupportedVersionException against a 4.4 broker" % client_version + + processor.stop() + + @cluster(num_nodes=2) + @matrix(client_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_missing_client_tags_status_gated_by_rpc_version(self, client_version, metadata_quorum): + """ + With rack-aware assignment tags required on the broker but absent on the client, the group + coordinator returns the MISSING_CLIENT_TAGS status (code 6). That status is gated on + heartbeat version >= 1, so only a 4.4 client may receive it: the Status enum shipped in + 4.2/4.3 defines codes 0-5 only. + """ + self.setup_kafka(str(DEV_BRANCH), + extra_server_prop_overrides=[ + ["group.streams.rack.aware.assignment.tags", self.RACK_AWARE_TAG] + ]) + # No client.tag.<tag> is configured, so the required tag is missing. + processor = self.start_processor(client_version) + + status_logged = self.count_in_file(processor.node, self.MISSING_CLIENT_TAGS_LOG, processor.LOG_FILE) + if client_version == str(DEV_BRANCH): + assert status_logged > 0, \ + "A 4.4 client negotiates heartbeat v1 and should have been told its required " \ + "client tags are missing" + else: + assert status_logged == 0, \ + "The broker returned the MISSING_CLIENT_TAGS status to a %s client, which " \ + "negotiates heartbeat v0 and does not understand status code 6" % client_version + + processor.stop() + + @cluster(num_nodes=2) + @matrix(metadata_quorum=[quorum.combined_kraft]) + def test_missing_client_tags_status_absent_when_tag_configured(self, metadata_quorum): Review Comment: Do we need this test as system test? Seems this can be covered (and hopefully is already covered) as integration test? ########## tests/kafkatest/tests/streams/streams_protocol_cross_version_test.py: ########## @@ -0,0 +1,336 @@ +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from ducktape.mark import matrix +from ducktape.mark.resource import cluster +from ducktape.tests.test import Test +from kafkatest.services.kafka import KafkaService, quorum +from kafkatest.services.streams import ( + INMEMORY_TOPOLOGY_DESCRIPTION_PLUGIN_CLASS, + StreamsUpgradeTestJobRunnerService, +) +from kafkatest.version import LATEST_4_2, LATEST_4_3, DEV_BRANCH, KafkaVersion + + +class StreamsProtocolCrossVersionTest(Test): + """ + Cross-version tests for the streams rebalance protocol (KIP-1071). + + AK 4.4 bumped StreamsGroupHeartbeat (apiKey 88) and StreamsGroupDescribe (apiKey 89) from + version 0 to version 1. The new fields are only exchanged when both the client and the broker + speak v1, so mixed-version deployments must keep working: + + - StreamsGroupHeartbeatRequest v1 is byte-identical to v0. It was bumped only so that + response v1 is negotiated, and so the broker may return the MISSING_CLIENT_TAGS status + (code 6) that v0 clients do not know about (KAFKA-20744). + - StreamsGroupHeartbeatResponse v1 replaces the v0 `AcceptableRecoveryLagLegacy` (int32) + with `AcceptableRecoveryLag` (int64, ignorable) and adds `TopologyDescriptionRequired` + (KIP-1331). + - StreamsGroupDescribeRequest v1 adds `IncludeTopologyDescription`, and the response adds + `TopologyDescription`, `TopologyDescriptionStatus` and `AssignorName` (KIP-1331, KIP-1357). + + 4.2 is the earliest release that speaks the protocol at all -- `streams.version` level 1 + requires metadata version 4.2-IV1 -- so 4.2 and 4.3 are the only older versions worth pairing + against a 4.4 broker or client. + + These tests drive the `StreamsUpgradeTest` harness, which is published in the 4.2/4.3 streams + test jars and mirrored for trunk, so the same topology can be run on either side of the bump. + """ + + # Topics the StreamsUpgradeTest harness reads from and writes to. + input_topic = "data" + output_topic = "echo" + + # application.id hard-coded by the StreamsUpgradeTest harness, in every version. + group_id = "StreamsUpgradeTest" + + RUNNING_LOG = "State transition from REBALANCING to RUNNING" + + # Logged by StreamsGroupHeartbeatRequestManager#describeConfig when the broker leaves + # acceptableRecoveryLag at its protocol default, i.e. when the response came back as v0. + # Only a 4.4+ client logs this line at all. + OLD_BROKER_LAG_LOG = "acceptableRecoveryLag=not provided (older broker)" + + # StatusDetail of the MISSING_CLIENT_TAGS status (code 6), which the group coordinator only + # returns on a v1 heartbeat. + MISSING_CLIENT_TAGS_LOG = "Missing required client tags for rack-aware standby assignment" + + UNSUPPORTED_VERSION_LOG = "UnsupportedVersionException" + + # On startup the client probes the broker for KIP-714 telemetry support; a broker without a + # client-metrics receiver rejects GET_TELEMETRY_SUBSCRIPTIONS and the client logs the + # rejection (with a stack trace) at DEBUG before disabling telemetry. That happens against + # brokers of every version, including dev, and is unrelated to the streams protocol, so it + # must not trip the UnsupportedVersionException guard. + BENIGN_TELEMETRY_UVE_LOG = "does not support GET_TELEMETRY_SUBSCRIPTIONS" + + RACK_AWARE_TAG = "zone" + + def __init__(self, test_context): + super(StreamsProtocolCrossVersionTest, self).__init__(test_context=test_context) + self.topics = { + self.input_topic: {"partitions": 1, "replication-factor": 1}, + self.output_topic: {"partitions": 1, "replication-factor": 1}, + } + + def setup_kafka(self, broker_version, extra_server_prop_overrides=None): + server_prop_overrides = [ + # The harness configures session.timeout.ms=10000, which is below the default lower bound. + ["group.streams.min.session.timeout.ms", "10000"], + ["group.streams.session.timeout.ms", "10000"], + ] + server_prop_overrides.extend(extra_server_prop_overrides or []) + + self.kafka = KafkaService(self.test_context, + num_nodes=1, + zk=None, + topics=self.topics, + use_streams_groups=True, + server_prop_overrides=server_prop_overrides) + self.kafka.set_version(KafkaVersion(broker_version)) + self.kafka.start() + self.kafka.run_features_command("upgrade", "streams.version", 1) + + def start_processor(self, client_version, extra_configs=None): + """ + Start a single Streams instance on the streams rebalance protocol and wait until it is + RUNNING. Reaching RUNNING is itself the core protocol assertion: it requires a successful + join, topology initialization, an assignment from the group coordinator, and task creation. + """ + processor = StreamsUpgradeTestJobRunnerService(self.test_context, self.kafka) + # An empty version string makes kafka-run-class.sh use the trunk build rather than an + # installed release, which is how the harness selects the DEV client. + processor.set_version("" if client_version == str(DEV_BRANCH) else client_version) + processor.set_config("group.protocol", "streams") + for key, value in (extra_configs or {}).items(): + processor.set_config(key, value) + + with processor.node.account.monitor_log(processor.LOG_FILE) as monitor: + processor.start() + monitor.wait_until(self.RUNNING_LOG, + timeout_sec=120, + err_msg="Streams client (version '%s') never reached RUNNING on %s" + % (client_version, str(processor.node.account))) + return processor + + def count_in_file(self, node, literal, path): + """Number of lines in `path` containing `literal`, matched as a fixed string.""" + cmd = "grep -c -F -- '%s' %s 2>/dev/null || true" % (literal, path) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def count_unsupported_version_errors(self, node, path): + """ + Number of UnsupportedVersionException lines in `path`, ignoring the benign KIP-714 + telemetry probe (see BENIGN_TELEMETRY_UVE_LOG). + """ + cmd = "grep -F -- '%s' %s 2>/dev/null | grep -c -v -F -- '%s' || true" \ + % (self.UNSUPPORTED_VERSION_LOG, path, self.BENIGN_TELEMETRY_UVE_LOG) + output = node.account.ssh_output(cmd, allow_fail=False).decode("utf-8", errors="replace").strip() + return int(output) if output.isdigit() else 0 + + def run_streams_groups_command(self, args, tool_version=None): + """ + Run bin/kafka-streams-groups.sh against the cluster, from the install of `tool_version` + (the trunk build when omitted). Returns the combined stdout/stderr; the command is allowed + to fail so callers can assert on how a failure is reported. + """ + node = self.kafka.nodes[0] + version = DEV_BRANCH if tool_version is None else KafkaVersion(tool_version) + script = self.kafka.path.script("kafka-streams-groups.sh", version) + cmd = "%s --bootstrap-server %s %s 2>&1" % (script, self.kafka.bootstrap_servers(), args) + self.logger.info("Running streams-groups command: %s" % cmd) + return node.account.ssh_output(cmd, allow_fail=True).decode("utf-8", errors="replace") + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_new_client_old_broker(self, broker_version, metadata_quorum): + """ + A 4.4 client must work against a 4.2/4.3 broker: it must not send a request the older + broker rejects, and it must tolerate a v0 response in which the new int64 + AcceptableRecoveryLag is absent and therefore reads back as its default of -1. + """ + self.setup_kafka(broker_version) + processor = self.start_processor(str(DEV_BRANCH)) + + lag_not_provided = self.count_in_file(processor.node, self.OLD_BROKER_LAG_LOG, processor.LOG_FILE) + if broker_version == str(DEV_BRANCH): + assert lag_not_provided == 0, \ + "A 4.4 broker negotiates heartbeat v1 and must supply acceptableRecoveryLag, " \ + "but the client logged it as not provided" + else: + assert lag_not_provided > 0, \ + "Against a %s broker the heartbeat response is v0 and carries no int64 " \ + "acceptableRecoveryLag, so the client should have logged it as not provided" \ + % broker_version + + assert self.count_unsupported_version_errors(processor.node, processor.LOG_FILE) == 0, \ + "The 4.4 client hit an UnsupportedVersionException against a %s broker" % broker_version + + processor.stop() + + @cluster(num_nodes=2) + @matrix(client_version=[str(LATEST_4_2), str(LATEST_4_3)], + metadata_quorum=[quorum.combined_kraft]) + def test_old_client_new_broker(self, client_version, metadata_quorum): + """ + A 4.2/4.3 client must keep working against a 4.4 broker. The broker sets the v1-only + AcceptableRecoveryLag unconditionally and relies on it being `ignorable` so that it is + dropped when the response is serialized at v0; an older client must not see a malformed + response. + """ + self.setup_kafka(str(DEV_BRANCH)) + processor = self.start_processor(client_version) + + assert self.count_unsupported_version_errors(processor.node, processor.LOG_FILE) == 0, \ + "The %s client hit an UnsupportedVersionException against a 4.4 broker" % client_version + + processor.stop() + + @cluster(num_nodes=2) + @matrix(client_version=[str(LATEST_4_2), str(LATEST_4_3), str(DEV_BRANCH)], + metadata_quorum=[quorum.combined_kraft]) + def test_missing_client_tags_status_gated_by_rpc_version(self, client_version, metadata_quorum): + """ + With rack-aware assignment tags required on the broker but absent on the client, the group + coordinator returns the MISSING_CLIENT_TAGS status (code 6). That status is gated on + heartbeat version >= 1, so only a 4.4 client may receive it: the Status enum shipped in + 4.2/4.3 defines codes 0-5 only. + """ + self.setup_kafka(str(DEV_BRANCH), + extra_server_prop_overrides=[ + ["group.streams.rack.aware.assignment.tags", self.RACK_AWARE_TAG] + ]) + # No client.tag.<tag> is configured, so the required tag is missing. + processor = self.start_processor(client_version) + + status_logged = self.count_in_file(processor.node, self.MISSING_CLIENT_TAGS_LOG, processor.LOG_FILE) + if client_version == str(DEV_BRANCH): + assert status_logged > 0, \ + "A 4.4 client negotiates heartbeat v1 and should have been told its required " \ + "client tags are missing" + else: + assert status_logged == 0, \ + "The broker returned the MISSING_CLIENT_TAGS status to a %s client, which " \ + "negotiates heartbeat v0 and does not understand status code 6" % client_version + + processor.stop() + + @cluster(num_nodes=2) + @matrix(metadata_quorum=[quorum.combined_kraft]) + def test_missing_client_tags_status_absent_when_tag_configured(self, metadata_quorum): + """ + Counterpart to the gating test: a 4.4 client that does configure the required tag must not + be told anything is missing. Without this, a version gate that never fires would look the + same as one that is correctly withholding the status. + """ + self.setup_kafka(str(DEV_BRANCH), + extra_server_prop_overrides=[ + ["group.streams.rack.aware.assignment.tags", self.RACK_AWARE_TAG] + ]) + processor = self.start_processor( + str(DEV_BRANCH), + extra_configs={"client.tag.%s" % self.RACK_AWARE_TAG: "eu-central-1a"}) + + assert self.count_in_file(processor.node, self.MISSING_CLIENT_TAGS_LOG, processor.LOG_FILE) == 0, \ + "The broker reported missing client tags even though the required tag was configured" + + processor.stop() + + @cluster(num_nodes=2) + @matrix(broker_version=[str(LATEST_4_2), str(LATEST_4_3)], + metadata_quorum=[quorum.combined_kraft]) + def test_describe_new_tool_old_broker(self, broker_version, metadata_quorum): + """ + The 4.4 kafka-streams-groups.sh must describe a group hosted on a 4.2/4.3 broker. Plain + --describe stays within describe v0 and must succeed. Review Comment: What about "assignor name" -- older brokers does not support it, so we need a negative check (and check that we don't crash?) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
