use std::sync::Arc;
use super::{DeliveryHandle, Producer, Record, RecordMetadata};
use crate::error::Result;
use crate::serdes::Serializer;
pub struct TypedProducer<K: ?Sized, V: ?Sized> {
producer: Producer,
key: Arc<dyn Serializer<K>>,
value: Arc<dyn Serializer<V>>,
}
impl<K: ?Sized, V: ?Sized> std::fmt::Debug for TypedProducer<K, V> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TypedProducer")
.field("producer", &self.producer)
.finish_non_exhaustive()
}
}
impl<K: ?Sized, V: ?Sized> TypedProducer<K, V> {
pub fn new(
producer: Producer,
key: impl Serializer<K> + 'static,
value: impl Serializer<V> + 'static,
) -> Self {
Self {
producer,
key: Arc::new(key),
value: Arc::new(value),
}
}
pub async fn send(
&self,
topic: &str,
key: Option<&K>,
value: Option<&V>,
) -> Result<RecordMetadata> {
self.enqueue(topic, key, value).await?.await
}
pub async fn enqueue(
&self,
topic: &str,
key: Option<&K>,
value: Option<&V>,
) -> Result<DeliveryHandle> {
let record = self.serialize(topic, key, value)?;
self.producer.enqueue(record).await
}
fn serialize(&self, topic: &str, key: Option<&K>, value: Option<&V>) -> Result<Record> {
let mut record = Record::new(topic, bytes::Bytes::new());
record.key = key
.map(|key| self.key.serialize(topic, &mut record.headers, key))
.transpose()?;
record.value = value
.map(|value| self.value.serialize(topic, &mut record.headers, value))
.transpose()?;
Ok(record)
}
pub fn producer(&self) -> &Producer {
&self.producer
}
pub fn into_inner(self) -> Producer {
self.producer
}
pub async fn close(&self) -> Result<()> {
self.producer.close().await
}
}