use std::fmt;
use bytes::Bytes;
use crate::consumer::ConsumerRecord;
use crate::producer::Record;
pub const HEADER_ORIGINAL_TOPIC: &str = "__krafka.dlq.original.topic";
pub const HEADER_ORIGINAL_PARTITION: &str = "__krafka.dlq.original.partition";
pub const HEADER_ORIGINAL_OFFSET: &str = "__krafka.dlq.original.offset";
pub const HEADER_EXCEPTION_MESSAGE: &str = "__krafka.dlq.exception.message";
pub fn record_for(dlq_topic: &str, original: &ConsumerRecord, error: &dyn fmt::Display) -> Record {
let mut headers = original.headers.clone();
headers.push((
HEADER_ORIGINAL_TOPIC.to_string(),
Some(Bytes::copy_from_slice(original.topic.as_bytes())),
));
headers.push((
HEADER_ORIGINAL_PARTITION.to_string(),
Some(Bytes::from(original.partition.to_string())),
));
headers.push((
HEADER_ORIGINAL_OFFSET.to_string(),
Some(Bytes::from(original.offset.to_string())),
));
headers.push((
HEADER_EXCEPTION_MESSAGE.to_string(),
Some(Bytes::from(error.to_string())),
));
Record {
topic: dlq_topic.to_string(),
partition: None,
key: original.key.clone(),
value: original.value.clone(),
timestamp: None,
headers,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_record_for_provenance_headers() {
let original = ConsumerRecord::new(
"source-topic",
2,
42,
Some(Bytes::from("key")),
Some(Bytes::from("value")),
);
let record = record_for("source-topic.DLQ", &original, &"decode error");
assert_eq!(record.topic, "source-topic.DLQ");
assert_eq!(record.key, Some(Bytes::from("key")));
assert_eq!(record.value, Some(Bytes::from("value")));
let hdr = |name: &str| -> Option<Bytes> {
record
.headers
.iter()
.find(|(k, _)| k == name)
.and_then(|(_, v)| v.clone())
};
assert_eq!(
hdr("__krafka.dlq.original.topic"),
Some(Bytes::from("source-topic"))
);
assert_eq!(
hdr("__krafka.dlq.original.partition"),
Some(Bytes::from("2"))
);
assert_eq!(hdr("__krafka.dlq.original.offset"), Some(Bytes::from("42")));
assert_eq!(
hdr("__krafka.dlq.exception.message"),
Some(Bytes::from("decode error"))
);
}
#[test]
fn test_record_for_original_headers_preserved() {
let mut original = ConsumerRecord::new("t", 0, 0, None, Some(Bytes::from("v")));
original
.headers
.push(("x-trace-id".to_string(), Some(Bytes::from("abc123"))));
let record = record_for("t.DLQ", &original, &"error");
assert_eq!(record.headers[0].0, "x-trace-id");
assert_eq!(record.headers[0].1, Some(Bytes::from("abc123")));
assert!(
record
.headers
.iter()
.any(|(k, _)| k == "__krafka.dlq.original.topic")
);
}
#[test]
fn test_record_for_preserves_tombstone() {
let tombstone = ConsumerRecord::new("t", 0, 0, Some(Bytes::from("k")), None);
let from_tombstone = record_for("t.DLQ", &tombstone, &"tombstone");
assert_eq!(from_tombstone.value, None);
assert!(from_tombstone.is_tombstone());
let empty = ConsumerRecord::new("t", 0, 0, Some(Bytes::from("k")), Some(Bytes::new()));
let from_empty = record_for("t.DLQ", &empty, &"tombstone");
assert_eq!(from_empty.value, Some(Bytes::new()));
assert!(!from_empty.is_tombstone());
assert_ne!(from_tombstone.value, from_empty.value);
}
#[test]
fn test_record_for_preserves_null_header_value() {
let mut original = ConsumerRecord::new("t", 0, 0, None, Some(Bytes::from("v")));
original.headers.push(("null-hdr".to_string(), None));
original
.headers
.push(("empty-hdr".to_string(), Some(Bytes::new())));
let record = record_for("t.DLQ", &original, &"error");
assert_eq!(record.headers[0].0, "null-hdr");
assert_eq!(record.headers[0].1, None);
assert_eq!(record.headers[1].0, "empty-hdr");
assert_eq!(record.headers[1].1, Some(Bytes::new()));
assert_ne!(record.headers[0].1, record.headers[1].1);
}
}