From 3069596d7f5d7b8da003868befa5aaee8156e0ae Mon Sep 17 00:00:00 2001 From: Tony Le Date: Fri, 28 Aug 2026 15:53:50 -0400 Subject: [PATCH] fix: return Kafka configuration errors from constructors --- rust-arroyo/examples/transform_and_produce.rs | 2 + rust-arroyo/src/backends/kafka/mod.rs | 2 + rust-arroyo/src/backends/kafka/producer.rs | 58 +++++++++++++++---- .../src/processing/strategies/produce.rs | 4 +- rust-arroyo/src/testutils.rs | 2 + 5 files changed, 57 insertions(+), 11 deletions(-) diff --git a/rust-arroyo/examples/transform_and_produce.rs b/rust-arroyo/examples/transform_and_produce.rs index f7447bab..95a289b7 100644 --- a/rust-arroyo/examples/transform_and_produce.rs +++ b/rust-arroyo/examples/transform_and_produce.rs @@ -49,6 +49,8 @@ async fn main() { impl ProcessingStrategyFactory for ReverseStringAndProduceStrategyFactory { fn create(&self) -> Box> { let producer = KafkaProducer::new(self.config.clone()); + assert!(producer.is_ok()); + let producer = producer.unwrap(); let topic = TopicOrPartition::Topic(self.topic); let reverse_string_and_produce_strategy = RunTask::new( reverse_string, diff --git a/rust-arroyo/src/backends/kafka/mod.rs b/rust-arroyo/src/backends/kafka/mod.rs index 295ae8fc..55d21c77 100644 --- a/rust-arroyo/src/backends/kafka/mod.rs +++ b/rust-arroyo/src/backends/kafka/mod.rs @@ -679,6 +679,8 @@ mod tests { ); let producer = KafkaProducer::new(producer_configuration); + assert!(producer.is_ok()); + let producer = producer.unwrap(); let payload = KafkaPayload::new(None, None, Some("asdf".as_bytes().to_vec())); producer diff --git a/rust-arroyo/src/backends/kafka/producer.rs b/rust-arroyo/src/backends/kafka/producer.rs index 6475a142..c0d3eff7 100644 --- a/rust-arroyo/src/backends/kafka/producer.rs +++ b/rust-arroyo/src/backends/kafka/producer.rs @@ -139,7 +139,7 @@ pub struct KafkaProducer { } impl KafkaProducer { - pub fn new(config: KafkaConfig) -> Self { + pub fn new(config: KafkaConfig) -> Result { // Extract client.id from config for metrics, default to "unknown" let producer_name = config .get_config_value("client.id") @@ -147,12 +147,11 @@ impl KafkaProducer { .unwrap_or_else(|| "unknown".to_string()); let context = ProducerContext::new(producer_name.clone()); let config_obj: ClientConfig = config.into(); - let threaded_producer: ThreadedProducer<_> = - config_obj.create_with_context(context).unwrap(); + let threaded_producer: ThreadedProducer<_> = config_obj.create_with_context(context)?; - Self { + Ok(Self { producer: threaded_producer, - } + }) } } @@ -178,7 +177,7 @@ pub struct AsyncKafkaProducer { } impl AsyncKafkaProducer { - pub fn new(config: KafkaConfig) -> Self { + pub fn new(config: KafkaConfig) -> Result { // Extract client.id from config for metrics, default to "unknown" let producer_name = config .get_config_value("client.id") @@ -186,12 +185,12 @@ impl AsyncKafkaProducer { .unwrap_or_else(|| "unknown".to_string()); let context = ProducerContext::new(producer_name.clone()); let config_obj: ClientConfig = config.into(); - let future_producer: FutureProducer<_> = config_obj.create_with_context(context).unwrap(); + let future_producer: FutureProducer<_> = config_obj.create_with_context(context)?; - Self { + Ok(Self { producer: future_producer, producer_name, - } + }) } } @@ -268,6 +267,7 @@ mod tests { use crate::backends::{AsyncProducer, Producer}; use crate::types::{Topic, TopicOrPartition}; use rdkafka::client::ClientContext; + use rdkafka::error::KafkaError; use rdkafka::statistics::{Broker, Statistics, Window}; use std::collections::HashMap; @@ -389,11 +389,13 @@ mod tests { KafkaConfig::new_producer_config(vec!["127.0.0.1:9092".to_string()], None); let producer = KafkaProducer::new(configuration); + assert!(producer.is_ok()); + let producer = producer.unwrap(); let payload = KafkaPayload::new(None, None, Some("asdf".as_bytes().to_vec())); producer .produce(&destination, payload) - .expect("Message produced") + .expect("Message produced"); } #[tokio::test] @@ -404,6 +406,8 @@ mod tests { KafkaConfig::new_producer_config(vec!["127.0.0.1:9092".to_string()], None); let producer = AsyncKafkaProducer::new(configuration); + assert!(producer.is_ok()); + let producer = producer.unwrap(); let payload = KafkaPayload::new(None, None, Some("asdf".as_bytes().to_vec())); let result = producer.produce(&destination, payload).await; @@ -423,6 +427,8 @@ mod tests { ); let producer = AsyncKafkaProducer::new(configuration); + assert!(producer.is_ok()); + let producer = producer.unwrap(); let payload = KafkaPayload::new(None, None, Some("asdf".as_bytes().to_vec())); let result = producer.produce(&destination, payload).await; @@ -431,4 +437,36 @@ mod tests { "Message should not be produced successfully" ); } + + #[test] + fn test_invalid_producer_configuration_returns_error() { + let configuration = KafkaConfig::new_producer_config( + vec!["127.0.0.1:9092".to_string()], + Some(HashMap::from([( + "message.timeout.ms".to_string(), + "invalid".to_string(), + )])), + ); + + assert!(matches!( + KafkaProducer::new(configuration), + Err(KafkaError::ClientConfig(..)) + )); + } + + #[test] + fn test_invalid_async_producer_configuration_returns_error() { + let configuration = KafkaConfig::new_producer_config( + vec!["127.0.0.1:9092".to_string()], + Some(HashMap::from([( + "message.timeout.ms".to_string(), + "invalid".to_string(), + )])), + ); + + assert!(matches!( + AsyncKafkaProducer::new(configuration), + Err(KafkaError::ClientConfig(..)) + )); + } } diff --git a/rust-arroyo/src/processing/strategies/produce.rs b/rust-arroyo/src/processing/strategies/produce.rs index 7dc6406a..2c0703c7 100644 --- a/rust-arroyo/src/processing/strategies/produce.rs +++ b/rust-arroyo/src/processing/strategies/produce.rs @@ -157,7 +157,9 @@ mod tests { let partition = Partition::new(Topic::new("test"), 0); - let producer: KafkaProducer = KafkaProducer::new(config); + let producer = KafkaProducer::new(config); + assert!(producer.is_ok()); + let producer = producer.unwrap(); let concurrency = ConcurrencyConfig::new(10); let mut strategy = Produce::new( Noop {}, diff --git a/rust-arroyo/src/testutils.rs b/rust-arroyo/src/testutils.rs index 2433c30a..89bbbfe3 100644 --- a/rust-arroyo/src/testutils.rs +++ b/rust-arroyo/src/testutils.rs @@ -107,6 +107,8 @@ impl TestTopic { KafkaConfig::new_producer_config(vec![get_default_broker()], None); let producer = KafkaProducer::new(producer_configuration); + assert!(producer.is_ok()); + let producer = producer.unwrap(); producer .produce(&crate::types::TopicOrPartition::Topic(self.topic), payload)