diff --git a/CHANGELOG.md b/CHANGELOG.md index 3b60257..6029d2f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,12 @@ The format is based on [Keep a Changelog](http://keepachangelog.com/) and this p - upgrade template - clarify Producer and Consumer task documentation (ports, entity serialization, dataset advanced-options reference) and fix minor typos in parameter descriptions +### Fixed + +- Consumer and Producer tasks no longer convert non-ASCII characters to unicode escape + sequences when wrapping consumed Kafka messages into JSON, or when producing messages + from a JSON dataset or from entities + ## [3.5.0] 2026-08-24 ### Changed diff --git a/cmem_plugin_kafka/kafka_handlers.py b/cmem_plugin_kafka/kafka_handlers.py index 2e71767..3455353 100644 --- a/cmem_plugin_kafka/kafka_handlers.py +++ b/cmem_plugin_kafka/kafka_handlers.py @@ -181,7 +181,7 @@ def _split_data(self, data: Response) -> Generator[KafkaMessage]: content = None if not tombstone: _content = _message.get("content") - content = json.dumps(_content) if _content else None + content = json.dumps(_content, ensure_ascii=False) if _content else None yield KafkaMessage(key=key, value=content, headers=headers, tombstone=tombstone) def _aggregate_data(self) -> Generator: @@ -361,7 +361,7 @@ def _split_data(self, data: Entities) -> Generator[KafkaMessage]: for i, path in enumerate(paths): values[path.path] = list(entity.values[i]) result["entity"] = {"uri": entity.uri, "values": values} - kafka_payload = json.dumps(result, indent=4) + kafka_payload = json.dumps(result, indent=4, ensure_ascii=False) yield KafkaMessage(key=None, value=kafka_payload) def _aggregate_data(self) -> Entities: diff --git a/cmem_plugin_kafka/utils.py b/cmem_plugin_kafka/utils.py index 7375d50..a18aec2 100644 --- a/cmem_plugin_kafka/utils.py +++ b/cmem_plugin_kafka/utils.py @@ -337,7 +337,7 @@ def get_message_with_json_wrapper(message: KafkaMessage) -> str: msg_with_wrapper = {"message": {"key": message.key, "content": json.loads(message.value)}} if message.headers: msg_with_wrapper["message"]["headers"] = dict(message.headers.items()) - return json.dumps(msg_with_wrapper, cls=BytesEncoder) + return json.dumps(msg_with_wrapper, cls=BytesEncoder, ensure_ascii=False) def get_kafka_statistics(json_data: str) -> dict: diff --git a/tests/test_split_data_unicode.py b/tests/test_split_data_unicode.py new file mode 100644 index 0000000..cb708e0 --- /dev/null +++ b/tests/test_split_data_unicode.py @@ -0,0 +1,41 @@ +"""Regression tests for non-ASCII character handling in _split_data(). + +These exercise KafkaJSONDataHandler and KafkaEntitiesDataHandler directly, +without a Kafka broker or a Corporate Memory connection, since _split_data() +only transforms in-memory data into KafkaMessage objects. +""" + +import json + +import httpx +from cmem_plugin_base.dataintegration.entity import Entities, Entity, EntityPath, EntitySchema +from cmem_plugin_base.dataintegration.plugins import PluginLogger + +from cmem_plugin_kafka.kafka_handlers import KafkaEntitiesDataHandler, KafkaJSONDataHandler + + +def test_json_data_handler_split_data_keeps_unicode_characters() -> None: + """Test that non-ASCII characters from the source dataset are not escaped""" + payload = [{"message": {"key": "1", "content": {"city": "Köln", "name": "Müller"}}}] + data = httpx.Response(200, content=json.dumps(payload).encode("utf-8")) + handler = KafkaJSONDataHandler(context=None, plugin_logger=PluginLogger()) # type: ignore[arg-type] + + messages = list(handler._split_data(data)) # noqa: SLF001 + + assert len(messages) == 1 + assert "\\u00f6" not in messages[0].value + assert "\\u00fc" not in messages[0].value + assert json.loads(messages[0].value) == {"city": "Köln", "name": "Müller"} + + +def test_entities_data_handler_split_data_keeps_unicode_characters() -> None: + """Test that non-ASCII characters in entity values are not escaped""" + schema = EntitySchema(type_uri="urn:x-test", paths=[EntityPath(path="city")]) + entities = Entities(entities=[Entity(uri="urn:x-1", values=[["Köln"]])], schema=schema) + handler = KafkaEntitiesDataHandler(context=None, plugin_logger=PluginLogger()) # type: ignore[arg-type] + + messages = list(handler._split_data(entities)) # noqa: SLF001 + + assert len(messages) == 1 + assert "\\u00f6" not in messages[0].value + assert json.loads(messages[0].value)["entity"]["values"]["city"] == ["Köln"] diff --git a/tests/test_utils.py b/tests/test_utils.py index 74898e2..2cd4720 100644 --- a/tests/test_utils.py +++ b/tests/test_utils.py @@ -2,7 +2,20 @@ import json -from cmem_plugin_kafka.utils import get_kafka_statistics +from cmem_plugin_kafka.utils import ( + KafkaMessage, + get_kafka_statistics, + get_message_with_json_wrapper, +) + + +def test_get_message_with_json_wrapper_keeps_unicode_characters() -> None: + """Test that non-ASCII characters in the message value are not escaped""" + message = KafkaMessage(key="1", value='{"city": "Köln", "name": "Müller"}') + wrapped = get_message_with_json_wrapper(message) + assert "\\u00f6" not in wrapped + assert "\\u00fc" not in wrapped + assert json.loads(wrapped)["message"]["content"] == {"city": "Köln", "name": "Müller"} def test_get_kafka_statistics() -> None: