use crate::kafka;
use crate::kafka::{MessageHandler, StreamError};
use futures::TryStreamExt;
use rdkafka::ClientConfig;
use rdkafka::consumer::{CommitMode, Consumer as ConsumerExt, StreamConsumer};
use rdkafka::error::KafkaError;
use rdkafka::message::BorrowedMessage;
use std::error::Error;
use std::marker::PhantomData;
use tokio_stream::StreamExt;
pub struct Consumer<M, KM, DE, AE>
where
M: MessageHandler<AE, Message = KM> + Send + Sync,
AE: Error + Send + Sync,
DE: Error + Send + Sync,
KM: for<'a> TryFrom<&'a BorrowedMessage<'a>, Error = DE>,
{
consumer: StreamConsumer,
application_error: PhantomData<AE>,
handler: M,
topics: &'static [&'static str],
}
impl<M, KM, DE, AE> Consumer<M, KM, DE, AE>
where
M: MessageHandler<AE, Message = KM> + Send + Sync,
AE: Error + Send + Sync,
DE: Error + Send + Sync,
KM: for<'a> TryFrom<&'a BorrowedMessage<'a>, Error = DE> + Send + Sync,
{
pub fn new(
config: &kafka::Config,
topics: &'static [&'static str],
handler: M,
) -> Result<Self, KafkaError> {
let config = config.clone();
let mut cfg = ClientConfig::new();
cfg.extend(config.properties.into_iter().map(|(k, v)| (k, v.into())));
cfg.extend(config.env_properties);
let rdkafka_consumer: StreamConsumer = cfg.create()?;
Ok(Self {
consumer: rdkafka_consumer,
application_error: PhantomData,
topics,
handler,
})
}
pub async fn start(&self) -> Result<(), StreamError<DE, AE>> {
self.consumer.subscribe(self.topics)?;
self.consumer
.stream()
.then(async |m| -> Result<BorrowedMessage, StreamError<DE, AE>> {
let bm = m?;
let message = KM::try_from(&bm).map_err(StreamError::Decode)?;
self.handler
.handle(message)
.await
.map_err(StreamError::Application)?;
Ok(bm)
})
.try_for_each(async |bm| {
self.consumer
.commit_message(&bm, CommitMode::Async)
.map_err(StreamError::Kafka)
})
.await
}
}