Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions core/connectors/runtime/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@ pub enum RuntimeError {
FailedToSerializeMessagesMetadata,
#[error("Failed to serialize raw messages")]
FailedToSerializeRawMessages,
#[error("Sink plugin consume failed with status: {status} (plugin ID: {plugin_id})")]
SinkConsumeFailed { plugin_id: u32, status: i32 },
#[error("Connector SDK error")]
ConnectorSdkError(#[from] iggy_connector_sdk::Error),
#[error("Iggy client error")]
Expand Down
121 changes: 120 additions & 1 deletion core/connectors/runtime/src/sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -737,7 +737,7 @@ async fn process_messages(
})?;

let ffi_start = Instant::now();
(consume)(
let status = (consume)(
plugin_id,
topic_meta.as_ptr(),
topic_meta.len(),
Expand All @@ -748,6 +748,16 @@ async fn process_messages(
);
let ffi_elapsed = ffi_start.elapsed();

// Non-zero status is the plugin's only channel for reporting a failed write; propagate
// it so the runtime fails fast instead of counting the batch as processed. This does
// not redeliver the batch (offsets are already committed at poll time) — see #2927.
if status != 0 {
error!(
"Sink plugin consume failed with status: {status} for sink connector with ID: {plugin_id}"
);
return Err(RuntimeError::SinkConsumeFailed { plugin_id, status });
}

Ok(SinkBatchTiming {
processed_count,
decode_elapsed,
Expand All @@ -760,3 +770,112 @@ struct SinkBatchTiming {
decode_elapsed: Duration,
ffi_elapsed: Duration,
}

#[cfg(test)]
mod tests {
use super::*;
use iggy_connector_sdk::decoders::json::JsonStreamDecoder;

extern "C" fn consume_ok(
_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 {
0
}

extern "C" fn consume_fail(
_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 {
1
}

fn fixtures() -> (
TopicMetadata,
MessagesMetadata,
Vec<IggyMessage>,
Arc<dyn StreamDecoder>,
Arc<Metrics>,
SinkLabels,
) {
let topic_metadata = TopicMetadata {
stream: "test-stream".to_owned(),
topic: "test-topic".to_owned(),
};
let messages_metadata = MessagesMetadata {
partition_id: 1,
current_offset: 0,
schema: Schema::Json,
};
let messages = vec![
IggyMessage::builder()
.payload(br#"{"key":"value"}"#.to_vec().into())
.build()
.expect("test message should be valid"),
];
(
topic_metadata,
messages_metadata,
messages,
Arc::new(JsonStreamDecoder),
Arc::new(Metrics::init()),
SinkLabels::new("test-sink"),
)
}

#[tokio::test]
async fn given_zero_status_when_plugin_consumes_should_return_timing() {
let (topic_metadata, messages_metadata, messages, decoder, metrics, labels) = fixtures();
let result = process_messages(
1,
messages_metadata,
&topic_metadata,
messages,
&(consume_ok as ConsumeCallback),
&Vec::new(),
&decoder,
&metrics,
&labels,
)
.await;
match result {
Ok(timing) => assert_eq!(timing.processed_count, 1),
Err(error) => panic!("expected success, got error: {error}"),
}
}

#[tokio::test]
async fn given_nonzero_status_when_plugin_consume_fails_should_propagate_error() {
let (topic_metadata, messages_metadata, messages, decoder, metrics, labels) = fixtures();
let result = process_messages(
2,
messages_metadata,
&topic_metadata,
messages,
&(consume_fail as ConsumeCallback),
&Vec::new(),
&decoder,
&metrics,
&labels,
)
.await;
match result {
Err(RuntimeError::SinkConsumeFailed { plugin_id, status }) => {
assert_eq!(plugin_id, 2);
assert_eq!(status, 1);
}
Err(error) => panic!("expected SinkConsumeFailed, got: {error}"),
Ok(_) => panic!("expected SinkConsumeFailed, got success"),
}
}
}