From c0f89f7d7e92cc231891b909862089ee6ac43800 Mon Sep 17 00:00:00 2001 From: StandingMan Date: Tue, 14 Jul 2026 14:35:17 +0800 Subject: [PATCH] fix(connectors): preserve transformed payload schema across FFI Sink transforms can change the Payload variant after the input stream has been decoded, but the runtime forwarded decoder.schema() to sink plugins. Plugins then reconstructed the transformed bytes using the stale input schema. Derive the schema from transformed payloads and use it for both FFI metadata containers. Reject mixed-schema messages within a batch and cover Raw-to-Text conversion at the FFI boundary. Signed-off-by: StandingMan --- core/connectors/runtime/src/sink.rs | 133 +++++++++++++++++++++++++++- core/connectors/sdk/src/lib.rs | 11 +++ 2 files changed, 142 insertions(+), 2 deletions(-) diff --git a/core/connectors/runtime/src/sink.rs b/core/connectors/runtime/src/sink.rs index 7a17724510..810806e9d8 100644 --- a/core/connectors/runtime/src/sink.rs +++ b/core/connectors/runtime/src/sink.rs @@ -552,7 +552,7 @@ pub(crate) async fn setup_sink_consumers( #[allow(clippy::too_many_arguments)] async fn process_messages( plugin_id: u32, - messages_metadata: MessagesMetadata, + mut messages_metadata: MessagesMetadata, topic_metadata: &TopicMetadata, messages: Vec, consume: &ConsumeCallback, @@ -603,6 +603,7 @@ async fn process_messages( let decode_elapsed = decode_start.elapsed(); let mut messages = Vec::with_capacity(decoded.len()); + let mut output_schema = None; for message in decoded { let mut current_message = Some(message); for transform in transforms.iter() { @@ -673,6 +674,7 @@ async fn process_messages( continue; }; + let payload_schema = message.payload.schema(); let Ok(payload) = message.payload.try_into_vec() else { error!( "Failed to get message payload for message with ID: {id}, offset: {offset} for sink connector with ID: {plugin_id}" @@ -695,6 +697,15 @@ async fn process_messages( None => vec![], }; + if output_schema.is_some_and(|schema| schema != payload_schema) { + error!( + "Transformed payload schema {payload_schema} does not match the other messages in the batch for message with ID: {id}, offset: {offset}, sink connector with ID: {plugin_id}" + ); + error_count += 1; + continue; + } + output_schema = Some(payload_schema); + messages.push(RawMessage { id, offset, @@ -712,6 +723,8 @@ async fn process_messages( } let processed_count = messages.len(); + let output_schema = output_schema.unwrap_or(messages_metadata.schema); + messages_metadata.schema = output_schema; let topic_meta = postcard::to_allocvec(topic_metadata).map_err(|error| { error!( @@ -728,7 +741,7 @@ async fn process_messages( })?; let messages = postcard::to_allocvec(&RawMessages { - schema: decoder.schema(), + schema: output_schema, messages, }) .map_err(|error| { @@ -760,3 +773,119 @@ struct SinkBatchTiming { decode_elapsed: Duration, ffi_elapsed: Duration, } + +#[cfg(test)] +mod tests { + use std::sync::Mutex; + + use iggy_connector_sdk::{Error, Payload, transforms::TransformType}; + + use super::*; + + static CAPTURED_BATCH: Mutex> = Mutex::new(None); + + struct CapturedBatch { + metadata_schema: Schema, + batch_schema: Schema, + payloads: Vec>, + } + + struct RawToText; + + impl Transform for RawToText { + fn r#type(&self) -> TransformType { + TransformType::ProtoConvert + } + + fn transform( + &self, + _metadata: &TopicMetadata, + message: DecodedMessage, + ) -> Result, Error> { + let Payload::Raw(payload) = message.payload else { + return Ok(Some(message)); + }; + let text = String::from_utf8(payload).map_err(|_| Error::InvalidTextPayload)?; + Ok(Some(DecodedMessage { + payload: Payload::Text(text), + ..message + })) + } + } + + extern "C" fn capture_schemas( + _plugin_id: u32, + _topic_meta_ptr: *const u8, + _topic_meta_len: usize, + messages_meta_ptr: *const u8, + messages_meta_len: usize, + messages_ptr: *const u8, + messages_len: usize, + ) -> i32 { + unsafe { + let messages_metadata = postcard::from_bytes::( + std::slice::from_raw_parts(messages_meta_ptr, messages_meta_len), + ) + .expect("messages metadata should deserialize"); + let raw_messages = postcard::from_bytes::(std::slice::from_raw_parts( + messages_ptr, + messages_len, + )) + .expect("raw messages should deserialize"); + *CAPTURED_BATCH.lock().expect("capture lock should succeed") = Some(CapturedBatch { + metadata_schema: messages_metadata.schema, + batch_schema: raw_messages.schema, + payloads: raw_messages + .messages + .into_iter() + .map(|message| message.payload) + .collect(), + }); + } + 0 + } + + #[tokio::test] + async fn given_transform_changes_schema_when_crossing_ffi_should_use_output_schema() { + let metrics = Arc::new(Metrics::init()); + let labels = SinkLabels::new("schema-aware"); + let decoder: Arc = Schema::Raw.decoder(); + let transforms: Vec> = vec![Arc::new(RawToText)]; + let consume: ConsumeCallback = capture_schemas; + let message = IggyMessage::builder() + .payload("transformed".into()) + .build() + .expect("message should build"); + + let timing = process_messages( + 1, + MessagesMetadata { + partition_id: 1, + current_offset: 0, + schema: Schema::Raw, + }, + &TopicMetadata { + stream: "stream".to_string(), + topic: "topic".to_string(), + }, + vec![message], + &consume, + &transforms, + &decoder, + &metrics, + &labels, + ) + .await + .expect("message processing should succeed"); + + let captured = CAPTURED_BATCH + .lock() + .expect("capture lock should succeed") + .take() + .expect("FFI callback should capture schemas"); + assert_eq!(timing.processed_count, 1); + assert_eq!(captured.metadata_schema, Schema::Text); + assert_eq!(captured.batch_schema, Schema::Text); + assert_eq!(captured.payloads, vec![b"transformed".to_vec()]); + } +} diff --git a/core/connectors/sdk/src/lib.rs b/core/connectors/sdk/src/lib.rs index c8ed2ff94d..90f2446cea 100644 --- a/core/connectors/sdk/src/lib.rs +++ b/core/connectors/sdk/src/lib.rs @@ -145,6 +145,17 @@ pub enum Payload { } impl Payload { + pub const fn schema(&self) -> Schema { + match self { + Self::Json(_) => Schema::Json, + Self::Raw(_) => Schema::Raw, + Self::Text(_) => Schema::Text, + Self::Proto(_) => Schema::Proto, + Self::FlatBuffer(_) => Schema::FlatBuffer, + Self::Avro(_) => Schema::Avro, + } + } + /// Consuming conversion — transfers ownership of inner buffers. pub fn try_into_vec(self) -> Result, Error> { match self {