use std::{collections::BTreeMap, sync::Arc};
use bytes::Bytes;
pub use kacrab_protocol::record::RecordHeader;
use crate::common::TopicPartition;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct OffsetAndTimestamp {
pub offset: i64,
pub timestamp: i64,
pub leader_epoch: Option<i32>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TimestampType {
CreateTime,
LogAppendTime,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerRecord {
pub topic: Arc<str>,
pub partition: i32,
pub offset: i64,
pub timestamp: i64,
pub timestamp_type: TimestampType,
pub key: Option<Bytes>,
pub value: Option<Bytes>,
pub headers: Vec<RecordHeader>,
pub leader_epoch: Option<i32>,
}
impl ConsumerRecord {
#[must_use]
pub fn topic_partition(&self) -> TopicPartition {
TopicPartition::new(self.topic.as_ref(), self.partition)
}
pub fn deserialized<K, V>(
&self,
key: &impl super::ConsumerDeserializer<K>,
value: &impl super::ConsumerDeserializer<V>,
) -> super::Result<(Option<K>, Option<V>)> {
Ok((
key.deserialize(&self.topic, self.key.as_ref())?,
value.deserialize(&self.topic, self.value.as_ref())?,
))
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ConsumerRecords {
by_partition: BTreeMap<(String, i32), Vec<ConsumerRecord>>,
count: usize,
}
impl ConsumerRecords {
#[must_use]
pub fn empty() -> Self {
Self::default()
}
pub(crate) fn push_partition(
&mut self,
topic: String,
partition: i32,
records: Vec<ConsumerRecord>,
) {
if records.is_empty() {
return;
}
self.count = self.count.saturating_add(records.len());
match self.by_partition.entry((topic, partition)) {
std::collections::btree_map::Entry::Vacant(entry) => {
let _records = entry.insert(records);
},
std::collections::btree_map::Entry::Occupied(mut entry) => {
entry.get_mut().extend(records);
},
}
}
#[must_use]
pub const fn count(&self) -> usize {
self.count
}
#[must_use]
pub const fn is_empty(&self) -> bool {
self.count == 0
}
#[must_use]
pub fn partitions(&self) -> Vec<TopicPartition> {
self.by_partition
.keys()
.map(|(topic, partition)| TopicPartition::new(topic.clone(), *partition))
.collect()
}
#[must_use]
pub fn records(&self, partition: &TopicPartition) -> &[ConsumerRecord] {
self.by_partition
.get(&(partition.topic.clone(), partition.partition))
.map_or(&[], Vec::as_slice)
}
pub fn iter(&self) -> impl Iterator<Item = &ConsumerRecord> {
self.by_partition.values().flatten()
}
}
impl<'a> IntoIterator for &'a ConsumerRecords {
type Item = &'a ConsumerRecord;
type IntoIter = std::iter::Flatten<
std::collections::btree_map::Values<'a, (String, i32), Vec<ConsumerRecord>>,
>;
fn into_iter(self) -> Self::IntoIter {
self.by_partition.values().flatten()
}
}
impl IntoIterator for ConsumerRecords {
type Item = ConsumerRecord;
type IntoIter = std::iter::Flatten<
std::collections::btree_map::IntoValues<(String, i32), Vec<ConsumerRecord>>,
>;
fn into_iter(self) -> Self::IntoIter {
self.by_partition.into_values().flatten()
}
}
#[cfg(test)]
mod tests {
use bytes::Bytes;
use super::*;
fn record(topic: &str, partition: i32, offset: i64) -> ConsumerRecord {
ConsumerRecord {
topic: Arc::from(topic),
partition,
offset,
timestamp: offset,
timestamp_type: TimestampType::CreateTime,
key: None,
value: Some(Bytes::from(format!("v{offset}"))),
headers: Vec::new(),
leader_epoch: Some(4),
}
}
#[test]
fn records_group_by_partition_and_iterate_in_order() {
let mut records = ConsumerRecords::empty();
assert!(records.is_empty());
records.push_partition("t".to_owned(), 0, Vec::new());
assert!(records.is_empty());
records.push_partition("t".to_owned(), 1, vec![record("t", 1, 5)]);
records.push_partition(
"t".to_owned(),
0,
vec![record("t", 0, 0), record("t", 0, 1)],
);
assert_eq!(records.count(), 3);
assert!(!records.is_empty());
assert_eq!(
records.partitions(),
vec![TopicPartition::new("t", 0), TopicPartition::new("t", 1)]
);
assert_eq!(records.records(&TopicPartition::new("t", 0)).len(), 2);
assert!(records.records(&TopicPartition::new("t", 9)).is_empty());
let offsets: Vec<i64> = records.iter().map(|record| record.offset).collect();
assert_eq!(offsets, vec![0, 1, 5]);
let by_ref: Vec<i64> = (&records).into_iter().map(|record| record.offset).collect();
assert_eq!(by_ref, vec![0, 1, 5]);
let owned: Vec<i64> = records.into_iter().map(|record| record.offset).collect();
assert_eq!(owned, vec![0, 1, 5]);
}
#[test]
fn record_exposes_topic_partition_and_timestamp() {
let record = record("orders", 3, 9);
assert_eq!(record.topic_partition(), TopicPartition::new("orders", 3));
assert_eq!(record.leader_epoch, Some(4));
let timestamp = OffsetAndTimestamp {
offset: 9,
timestamp: 100,
leader_epoch: Some(4),
};
assert_eq!(timestamp.offset, 9);
}
}