From 659966d32e19588cc9589c568fc83d1ac75f058e Mon Sep 17 00:00:00 2001 From: "gopikrishna.c" Date: Fri, 7 Aug 2026 17:43:03 +0530 Subject: [PATCH] fix(events): serialize connector payloads as JSON text --- crates/common/common_utils/src/events.rs | 32 ++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/crates/common/common_utils/src/events.rs b/crates/common/common_utils/src/events.rs index e9bc9c1a5f..be84f9bc84 100644 --- a/crates/common/common_utils/src/events.rs +++ b/crates/common/common_utils/src/events.rs @@ -472,6 +472,18 @@ impl EventStage { pub struct KafkaTopicConfig { pub topic: String, pub partition_key_field: String, + #[serde(default)] + pub payload_format: EventPayloadFormat, +} + +#[derive( + Debug, Clone, Copy, Default, Deserialize, Serialize, PartialEq, config_patch_derive::Patch, +)] +#[serde(rename_all = "snake_case")] +pub enum EventPayloadFormat { + #[default] + Json, + JsonString, } /// Configuration for events system @@ -614,9 +626,29 @@ pub(crate) fn process_event_with_config( } } + if config.topic_config(&event.stage).payload_format == EventPayloadFormat::JsonString { + stringify_event_payloads(&mut result); + } + Ok(result) } +fn stringify_event_payloads(result: &mut serde_json::Value) { + let Some(obj) = result.as_object_mut() else { + return; + }; + + for field in ["request_data", "response_data"] { + let Some(value) = obj.get_mut(field) else { + continue; + }; + + if !value.is_null() && !value.is_string() { + *value = serde_json::Value::String(value.to_string()); + } + } +} + pub(crate) fn extract_from_request( event_value: &serde_json::Value, extraction_path: &str,