Skip to content

Commit fb1f14e

Browse files
fix(kafka_actions): accept list format for kafka_connect_str (DataDog#23115)
* fix(kafka_actions): accept list format for kafka_connect_str The kafka_actions check receives kafka_connect_str from kafka_consumer config via autodiscovery, which may be a list of broker strings. Normalize lists to comma-separated strings (matching kafka_consumer behavior) in both the Pydantic model validator and the config class. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * Add changelog entry for kafka_actions list fix Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent 73d460c commit fb1f14e

4 files changed

Lines changed: 38 additions & 10 deletions

File tree

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Accept list format for `kafka_connect_str` when copied from kafka_consumer config via autodiscovery

kafka_actions/datadog_checks/kafka_actions/config.py

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,10 @@ def __init__(self, instance, log):
1414
self.log = log
1515

1616
self.remote_config_id = instance.get('remote_config_id')
17-
self.kafka_connect_str = instance.get('kafka_connect_str')
17+
kafka_connect_str = instance.get('kafka_connect_str')
18+
if isinstance(kafka_connect_str, list):
19+
kafka_connect_str = ','.join(str(s) for s in kafka_connect_str)
20+
self.kafka_connect_str = kafka_connect_str
1821
self.tags = instance.get('tags', [])
1922

2023
# Authentication fields (same pattern as kafka_consumer)

kafka_actions/datadog_checks/kafka_actions/config_models/validators.py

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,12 @@
22
# All rights reserved
33
# Licensed under a 3-clause BSD style license (see LICENSE)
44

5-
# Here you can include additional config validators or transformers
6-
#
7-
# def initialize_instance(values, **kwargs):
8-
# if 'my_option' not in values and 'my_legacy_option' in values:
9-
# values['my_option'] = values['my_legacy_option']
10-
# if values.get('my_number') > 10:
11-
# raise ValueError('my_number max value is 10, got %s' % str(values.get('my_number')))
12-
#
13-
# return values
5+
6+
def initialize_instance(values, **kwargs):
7+
# kafka_connect_str may be passed as a list of broker strings (e.g. from kafka_consumer config
8+
# via autodiscovery). Normalize to a comma-separated string for librdkafka's bootstrap.servers.
9+
kafka_connect_str = values.get('kafka_connect_str')
10+
if isinstance(kafka_connect_str, list):
11+
values['kafka_connect_str'] = ','.join(str(s) for s in kafka_connect_str)
12+
13+
return values

kafka_actions/tests/test_config_validation.py

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,30 @@ def test_no_action_configured(self, dd_run_check):
4848
config = KafkaActionsConfig(instance, None)
4949
config.validate_config()
5050

51+
def test_kafka_connect_str_list(self, dd_run_check):
52+
"""Test that kafka_connect_str accepts a list of brokers and joins them."""
53+
instance = {
54+
'remote_config_id': 'test-id',
55+
'kafka_connect_str': ['broker1:9092', 'broker2:9092'],
56+
'produce_message': {'cluster': 'test', 'topic': 'test', 'value': 'test'},
57+
}
58+
59+
config = KafkaActionsConfig(instance, None)
60+
config.validate_config()
61+
assert config.kafka_connect_str == 'broker1:9092,broker2:9092'
62+
63+
def test_kafka_connect_str_string(self, dd_run_check):
64+
"""Test that kafka_connect_str works as a plain string."""
65+
instance = {
66+
'remote_config_id': 'test-id',
67+
'kafka_connect_str': 'broker1:9092,broker2:9092',
68+
'produce_message': {'cluster': 'test', 'topic': 'test', 'value': 'test'},
69+
}
70+
71+
config = KafkaActionsConfig(instance, None)
72+
config.validate_config()
73+
assert config.kafka_connect_str == 'broker1:9092,broker2:9092'
74+
5175

5276
class TestReadMessagesValidation:
5377
"""Test read_messages action validation."""

0 commit comments

Comments
 (0)