use std::sync::Arc;
use bytes::Bytes;
pub use crate::protocol::TimestampType;
use crate::{Headers, Offset, PartitionId, Timestamp};
#[non_exhaustive]
#[must_use = "contains data consumed from Kafka"]
#[derive(Debug, Clone)]
pub struct ConsumerRecord {
pub topic: Arc<str>,
pub partition: PartitionId,
pub offset: Offset,
pub timestamp: Timestamp,
pub timestamp_type: TimestampType,
pub key: Option<Bytes>,
pub value: Option<Bytes>,
pub headers: Headers,
pub leader_epoch: Option<i32>,
pub delivery_count: Option<i16>,
}
impl ConsumerRecord {
pub fn new(
topic: impl Into<Arc<str>>,
partition: PartitionId,
offset: Offset,
key: Option<Bytes>,
value: Option<Bytes>,
) -> Self {
Self {
topic: topic.into(),
partition,
offset,
timestamp: 0,
timestamp_type: TimestampType::CreateTime,
key,
value,
headers: Headers::new(),
leader_epoch: None,
delivery_count: None,
}
}
#[inline]
pub fn is_tombstone(&self) -> bool {
self.key.is_some() && self.value.is_none()
}
#[inline]
pub fn serialized_key_size(&self) -> Option<usize> {
self.key.as_ref().map(|k| k.len())
}
#[inline]
pub fn serialized_value_size(&self) -> Option<usize> {
self.value.as_ref().map(|v| v.len())
}
#[inline]
pub fn key_str(&self) -> Option<&str> {
self.key.as_ref().and_then(|k| std::str::from_utf8(k).ok())
}
#[inline]
pub fn value_str(&self) -> Option<&str> {
self.value
.as_ref()
.and_then(|v| std::str::from_utf8(v).ok())
}
#[inline]
pub fn header(&self, key: &str) -> Option<Option<&Bytes>> {
self.headers
.iter()
.find(|(k, _)| k == key)
.map(|(_, v)| v.as_ref())
}
#[inline]
pub fn header_str(&self, key: &str) -> Option<&str> {
self.header(key)
.flatten()
.and_then(|v| std::str::from_utf8(v).ok())
}
#[inline]
pub fn headers_by_key(&self, key: &str) -> Vec<Option<&Bytes>> {
self.headers
.iter()
.filter(|(k, _)| k == key)
.map(|(_, v)| v.as_ref())
.collect()
}
}
pub(crate) fn header_key(key: &[u8]) -> String {
String::from_utf8_lossy(key).into_owned()
}
pub(crate) fn headers_from_wire(headers: Vec<crate::protocol::RecordHeader>) -> Headers {
headers
.into_iter()
.map(|h| (header_key(&h.key), h.value))
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub struct TopicPartition {
pub topic: String,
pub partition: PartitionId,
}
impl TopicPartition {
pub fn new(topic: impl Into<String>, partition: PartitionId) -> Self {
Self {
topic: topic.into(),
partition,
}
}
#[inline]
pub fn topic(&self) -> &str {
&self.topic
}
#[inline]
pub fn partition(&self) -> PartitionId {
self.partition
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
#[test]
fn test_consumer_record_new() {
let record = ConsumerRecord::new(
"test-topic",
0,
42,
Some(Bytes::from("key")),
Some(Bytes::from("value")),
);
assert_eq!(&*record.topic, "test-topic");
assert_eq!(record.partition, 0);
assert_eq!(record.offset, 42);
assert_eq!(record.key_str(), Some("key"));
assert_eq!(record.value_str(), Some("value"));
assert_eq!(record.serialized_key_size(), Some(3));
assert_eq!(record.serialized_value_size(), Some(5));
}
#[test]
fn test_consumer_record_serialized_sizes_absent() {
let record = ConsumerRecord::new("topic", 0, 0, None, None);
assert_eq!(record.serialized_key_size(), None);
assert_eq!(record.serialized_value_size(), None);
}
#[test]
fn test_consumer_record_is_tombstone() {
let tombstone = ConsumerRecord::new("t", 0, 0, Some(Bytes::from("key")), None);
assert!(tombstone.is_tombstone());
let normal = ConsumerRecord::new(
"t",
0,
0,
Some(Bytes::from("key")),
Some(Bytes::from("val")),
);
assert!(!normal.is_tombstone());
let keyless = ConsumerRecord::new("t", 0, 0, None, None);
assert!(!keyless.is_tombstone());
let no_key = ConsumerRecord::new("t", 0, 0, None, Some(Bytes::from("val")));
assert!(!no_key.is_tombstone());
}
#[test]
fn test_consumer_record_duplicate_headers_preserved() {
let mut record = ConsumerRecord::new("test-topic", 0, 0, None, Some(Bytes::from("value")));
record
.headers
.push(("trace-id".to_string(), Some(Bytes::from("abc"))));
record
.headers
.push(("trace-id".to_string(), Some(Bytes::from("def"))));
record
.headers
.push(("other".to_string(), Some(Bytes::from("xyz"))));
assert_eq!(
record.headers.len(),
3,
"all headers including duplicates should be preserved"
);
assert_eq!(
record.header("trace-id"),
Some(Some(&Bytes::from("abc"))),
"header() should return the first matching header value"
);
}
#[test]
fn test_consumer_record_headers_by_key() {
let mut record = ConsumerRecord::new("test-topic", 0, 0, None, Some(Bytes::from("value")));
record
.headers
.push(("trace-id".to_string(), Some(Bytes::from("first"))));
record
.headers
.push(("trace-id".to_string(), Some(Bytes::from("second"))));
record
.headers
.push(("trace-id".to_string(), Some(Bytes::from("third"))));
record
.headers
.push(("other-key".to_string(), Some(Bytes::from("other"))));
let trace_values = record.headers_by_key("trace-id");
assert_eq!(
trace_values.len(),
3,
"headers_by_key should return all values for a duplicate key"
);
assert_eq!(trace_values[0], Some(&Bytes::from("first")));
assert_eq!(trace_values[1], Some(&Bytes::from("second")));
assert_eq!(trace_values[2], Some(&Bytes::from("third")));
let other_values = record.headers_by_key("other-key");
assert_eq!(other_values.len(), 1);
let missing_values = record.headers_by_key("nonexistent");
assert!(
missing_values.is_empty(),
"headers_by_key for missing key should return empty vec"
);
}
#[test]
fn test_consumer_record_header_with_null_value() {
let mut record = ConsumerRecord::new("t", 0, 0, None, Some(Bytes::from("v")));
record.headers.push(("x-null".to_string(), None));
record
.headers
.push(("x-present".to_string(), Some(Bytes::from("data"))));
assert_eq!(record.header("x-null"), Some(None));
assert_eq!(record.header("x-present"), Some(Some(&Bytes::from("data"))));
assert_eq!(record.header("missing"), None);
}
#[test]
fn test_consumer_record_header_str_returns_none_for_null() {
let mut record = ConsumerRecord::new("t", 0, 0, None, None);
record.headers.push(("h".to_string(), None));
record
.headers
.push(("h2".to_string(), Some(Bytes::from("text"))));
assert_eq!(record.header_str("h"), None);
assert_eq!(record.header_str("h2"), Some("text"));
}
#[test]
fn test_consumer_record_headers_by_key_with_nulls() {
let mut record = ConsumerRecord::new("t", 0, 0, None, None);
record
.headers
.push(("k".to_string(), Some(Bytes::from("a"))));
record.headers.push(("k".to_string(), None));
record
.headers
.push(("k".to_string(), Some(Bytes::from("b"))));
let vals = record.headers_by_key("k");
assert_eq!(vals.len(), 3);
assert_eq!(vals[0], Some(&Bytes::from("a")));
assert_eq!(vals[1], None);
assert_eq!(vals[2], Some(&Bytes::from("b")));
}
#[test]
fn a_non_utf8_header_key_is_decoded_lossily() {
let headers = headers_from_wire(vec![crate::protocol::RecordHeader::new(
Bytes::from_static(b"tr\xffce"),
Bytes::from_static(b"v"),
)]);
assert_eq!(headers[0].0, "tr\u{fffd}ce");
assert_eq!(headers[0].1.as_deref(), Some(&b"v"[..]));
}
}