use bytes::Bytes;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct KafkaHeader {
pub key: String,
pub value: Option<Bytes>,
}
impl KafkaHeader {
#[must_use]
pub fn new(key: impl Into<String>, value: impl Into<Option<Bytes>>) -> Self {
Self {
key: key.into(),
value: value.into(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum KafkaTimestamp {
NotAvailable,
CreateTime(i64),
LogAppendTime(i64),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerRecord {
pub topic: String,
pub partition: i32,
pub offset: i64,
pub timestamp: KafkaTimestamp,
pub key: Option<Bytes>,
pub payload: Option<Bytes>,
pub headers: Vec<KafkaHeader>,
}
impl ConsumerRecord {
#[must_use]
pub fn topic_partition(&self) -> crate::TopicPartition {
crate::TopicPartition::new(self.topic.clone(), self.partition)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProducerRecord {
pub topic: String,
pub key: Option<Bytes>,
pub payload: Option<Bytes>,
pub partition: Option<i32>,
pub timestamp: Option<i64>,
pub headers: Vec<KafkaHeader>,
}
impl ProducerRecord {
#[must_use]
pub fn new(topic: impl Into<String>, payload: impl Into<Option<Bytes>>) -> Self {
Self {
topic: topic.into(),
key: None,
payload: payload.into(),
partition: None,
timestamp: None,
headers: Vec::new(),
}
}
#[must_use]
pub fn with_key(mut self, key: impl Into<Option<Bytes>>) -> Self {
self.key = key.into();
self
}
#[must_use]
pub fn with_partition(mut self, partition: i32) -> Self {
self.partition = Some(partition);
self
}
#[must_use]
pub fn with_timestamp(mut self, timestamp: i64) -> Self {
self.timestamp = Some(timestamp);
self
}
#[must_use]
pub fn with_header(mut self, key: impl Into<String>, value: impl Into<Option<Bytes>>) -> Self {
self.headers.push(KafkaHeader::new(key, value));
self
}
}