Skip to content

Commit 956cf7e

Browse files
[kafka_actions] Base64-encode binary headers and support header filtering (DataDog#23106)
* [kafka_actions] Base64-encode binary headers and support header filtering Binary Kafka headers were displayed as unhelpful `<binary, N bytes>` placeholders. Now they are base64-encoded with a `<base64>` prefix so the actual data is visible. Header filtering is now supported in both filter systems: - jq-style: `.headers.dbmOrgId == "2"` - MessageFilter: `field: "header.dbmOrgId", operator: "eq", value: "2"` Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * Remove unnecessary MessageFilter header filtering MessageFilter is not used anywhere in the codebase. Only the jq-style filter in check.py is used, so header filtering only needs to be there. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * Use lazy dict for filter context to preserve partial deserialization The filter context was eagerly accessing deserialized_msg.key and deserialized_msg.value, defeating the lazy deserialization. Now uses _LazyMessageDict which only triggers deserialization when the filter actually accesses a field (e.g., filtering on .headers won't deserialize key/value). Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * Add changelog entry for header improvements 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 7c68895 commit 956cf7e

5 files changed

Lines changed: 128 additions & 17 deletions

File tree

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Base64-encode binary headers and support header filtering in jq-style filters

kafka_actions/datadog_checks/kafka_actions/check.py

Lines changed: 36 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,38 @@
1313
from .schema_registry import SchemaRegistryClient
1414

1515

16+
class _LazyMessageDict(dict):
17+
"""Dict-like wrapper around DeserializedMessage that only deserializes fields on access."""
18+
19+
_FIELD_MAP = {
20+
'key': lambda m: m.key,
21+
'value': lambda m: m.value,
22+
'headers': lambda m: m.headers,
23+
'topic': lambda m: m.topic,
24+
'partition': lambda m: m.partition,
25+
'offset': lambda m: m.offset,
26+
}
27+
28+
def __init__(self, msg: DeserializedMessage):
29+
super().__init__()
30+
self._msg = msg
31+
32+
def __contains__(self, key):
33+
return key in self._FIELD_MAP
34+
35+
def __getitem__(self, key):
36+
getter = self._FIELD_MAP.get(key)
37+
if getter is None:
38+
raise KeyError(key)
39+
return getter(self._msg)
40+
41+
def get(self, key, default=None):
42+
getter = self._FIELD_MAP.get(key)
43+
if getter is None:
44+
return default
45+
return getter(self._msg)
46+
47+
1648
class KafkaActionsCheck(AgentCheck):
1749
"""
1850
Kafka Actions Check - Performs one-time actions on Kafka clusters.
@@ -337,18 +369,10 @@ def _evaluate_filter(self, filter_expression: str, deserialized_msg: Deserialize
337369
True if message matches filter, False otherwise
338370
"""
339371
try:
340-
msg_dict = {
341-
'key': deserialized_msg.key,
342-
'value': deserialized_msg.value,
343-
'topic': deserialized_msg.topic,
344-
'partition': deserialized_msg.partition,
345-
'offset': deserialized_msg.offset,
346-
}
347-
348-
# Simple jq expression parser
349-
# Supports: .field.subfield == "value", .field > 100, .field contains "text"
350-
# Supports: and, or operators
351-
result = self._evaluate_jq_expression(filter_expression, msg_dict)
372+
# Use a lazy dict so key/value are only deserialized when the filter accesses them
373+
context = _LazyMessageDict(deserialized_msg)
374+
375+
result = self._evaluate_jq_expression(filter_expression, context)
352376
self.log.debug("Filter '%s' evaluated to: %s", filter_expression, result)
353377
return result
354378

kafka_actions/datadog_checks/kafka_actions/message_deserializer.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -510,7 +510,7 @@ def headers(self) -> dict:
510510
try:
511511
headers[key] = value.decode('utf-8') if value else None
512512
except UnicodeDecodeError:
513-
headers[key] = f"<binary, {len(value)} bytes>"
513+
headers[key] = f"<base64>{base64.b64encode(value).decode('ascii')}"
514514
return headers
515515

516516
@staticmethod

kafka_actions/tests/test_message_deserializer.py

Lines changed: 65 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +17,13 @@
1717
class MockKafkaMessage:
1818
"""Mock confluent_kafka.Message for testing."""
1919

20-
def __init__(self, key, value, topic='test-topic', partition=0, offset=0):
20+
def __init__(self, key, value, topic='test-topic', partition=0, offset=0, headers=None):
2121
self._key = key
2222
self._value = value
2323
self._topic = topic
2424
self._partition = partition
2525
self._offset = offset
26+
self._headers = headers
2627

2728
def key(self):
2829
return self._key
@@ -43,7 +44,7 @@ def timestamp(self):
4344
return (1, 1732128000000) # (timestamp_type, timestamp_ms)
4445

4546
def headers(self):
46-
return None
47+
return self._headers
4748

4849

4950
class TestMessageDeserializer:
@@ -684,5 +685,67 @@ def test_deserialized_message_with_registry(self):
684685
assert msg.value_schema_id == 42
685686

686687

688+
class TestHeaderSerialization:
689+
"""Test header serialization in DeserializedMessage."""
690+
691+
def test_utf8_headers(self):
692+
"""Test that UTF-8 headers are decoded as strings."""
693+
log = MagicMock()
694+
deserializer = MessageDeserializer(log)
695+
kafka_msg = MockKafkaMessage(
696+
key=b'key',
697+
value=b'{}',
698+
headers=[('dbmOrgId', b'2'), ('traceparent', b'00-abc-def-00')],
699+
)
700+
config = {'key_format': 'string', 'value_format': 'json', 'value_uses_schema_registry': False}
701+
msg = DeserializedMessage(kafka_msg, deserializer, config)
702+
assert msg.headers == {'dbmOrgId': '2', 'traceparent': '00-abc-def-00'}
703+
704+
def test_binary_headers_base64_encoded(self):
705+
"""Test that binary headers are base64-encoded with prefix."""
706+
log = MagicMock()
707+
deserializer = MessageDeserializer(log)
708+
binary_value = bytes([0x80, 0xFF, 0x00, 0x01, 0xDE, 0xAD, 0xBE, 0xEF])
709+
kafka_msg = MockKafkaMessage(
710+
key=b'key',
711+
value=b'{}',
712+
headers=[('koutrisngId', binary_value)],
713+
)
714+
config = {'key_format': 'string', 'value_format': 'json', 'value_uses_schema_registry': False}
715+
msg = DeserializedMessage(kafka_msg, deserializer, config)
716+
import base64
717+
718+
expected = f"<base64>{base64.b64encode(binary_value).decode('ascii')}"
719+
assert msg.headers['koutrisngId'] == expected
720+
721+
def test_mixed_headers(self):
722+
"""Test mix of UTF-8 and binary headers."""
723+
log = MagicMock()
724+
deserializer = MessageDeserializer(log)
725+
binary_value = bytes([0x80, 0xFF])
726+
kafka_msg = MockKafkaMessage(
727+
key=b'key',
728+
value=b'{}',
729+
headers=[('text', b'hello'), ('binary', binary_value)],
730+
)
731+
config = {'key_format': 'string', 'value_format': 'json', 'value_uses_schema_registry': False}
732+
msg = DeserializedMessage(kafka_msg, deserializer, config)
733+
assert msg.headers['text'] == 'hello'
734+
assert msg.headers['binary'].startswith('<base64>')
735+
736+
def test_null_header_value(self):
737+
"""Test that null header values are preserved as None."""
738+
log = MagicMock()
739+
deserializer = MessageDeserializer(log)
740+
kafka_msg = MockKafkaMessage(
741+
key=b'key',
742+
value=b'{}',
743+
headers=[('empty', None)],
744+
)
745+
config = {'key_format': 'string', 'value_format': 'json', 'value_uses_schema_registry': False}
746+
msg = DeserializedMessage(kafka_msg, deserializer, config)
747+
assert msg.headers['empty'] is None
748+
749+
687750
if __name__ == '__main__':
688751
pytest.main([__file__, '-v'])

kafka_actions/tests/test_message_filtering.py

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +17,13 @@
1717
class MockKafkaMessage:
1818
"""Mock confluent_kafka.Message for testing."""
1919

20-
def __init__(self, key, value, topic='test-topic', partition=0, offset=0):
20+
def __init__(self, key, value, topic='test-topic', partition=0, offset=0, headers=None):
2121
self._key = key
2222
self._value = value
2323
self._topic = topic
2424
self._partition = partition
2525
self._offset = offset
26+
self._headers = headers
2627

2728
def key(self):
2829
return self._key
@@ -43,7 +44,7 @@ def timestamp(self):
4344
return (1, 1732128000000)
4445

4546
def headers(self):
46-
return None
47+
return self._headers
4748

4849

4950
class TestFilteringOperators:
@@ -534,6 +535,28 @@ def test_null_literal(self):
534535
msg = self.create_message({'optional_field': None})
535536
assert check._evaluate_filter('.value.optional_field == null', msg) is True
536537

538+
def test_jq_filter_on_header(self):
539+
"""Test jq-style .headers.key filter."""
540+
instance = {
541+
'remote_config_id': 'test',
542+
'kafka_connect_str': 'localhost:9092',
543+
'read_messages': {'cluster': 'test', 'topic': 'test'},
544+
}
545+
check = KafkaActionsCheck('kafka_actions', {}, [instance])
546+
547+
kafka_msg = MockKafkaMessage(
548+
key=b'test-key',
549+
value=json.dumps({'status': 'active'}).encode('utf-8'),
550+
headers=[('dbmOrgId', b'2'), ('traceparent', b'00-abc-def-00')],
551+
)
552+
config = {'key_format': 'string', 'value_format': 'json', 'value_uses_schema_registry': False}
553+
msg = DeserializedMessage(kafka_msg, self.deserializer, config)
554+
555+
assert check._evaluate_filter('.headers.dbmOrgId == "2"', msg) is True
556+
assert check._evaluate_filter('.headers.dbmOrgId == "3"', msg) is False
557+
assert check._evaluate_filter('.headers.traceparent contains "abc"', msg) is True
558+
assert check._evaluate_filter('.headers.nonexistent == "x"', msg) is False
559+
537560

538561
if __name__ == '__main__':
539562
pytest.main([__file__, '-vv'])

0 commit comments

Comments
 (0)