use std::fmt;
use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicU64, Ordering};
use bytes::Bytes;
use crate::consumer::ConsumerRecord;
use crate::producer::{Producer, ProducerRecord};
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 trait DeadLetterQueue: Send + Sync + fmt::Debug {
fn send(
&self,
record: ProducerRecord,
error: String,
) -> Pin<Box<dyn Future<Output = ()> + Send + '_>>;
}
pub struct KafkaDeadLetterQueue {
producer: Producer,
topic: String,
routed: AtomicU64,
failures: AtomicU64,
}
impl KafkaDeadLetterQueue {
#[must_use]
pub fn new(producer: Producer, topic: impl Into<String>) -> Self {
Self {
producer,
topic: topic.into(),
routed: AtomicU64::new(0),
failures: AtomicU64::new(0),
}
}
#[inline]
#[must_use]
pub fn topic(&self) -> &str {
&self.topic
}
#[inline]
#[must_use]
pub fn routed(&self) -> u64 {
self.routed.load(Ordering::Relaxed)
}
#[inline]
#[must_use]
pub fn failures(&self) -> u64 {
self.failures.load(Ordering::Relaxed)
}
}
impl fmt::Debug for KafkaDeadLetterQueue {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("KafkaDeadLetterQueue")
.field("topic", &self.topic)
.field("routed", &self.routed())
.field("failures", &self.failures())
.finish_non_exhaustive()
}
}
impl DeadLetterQueue for KafkaDeadLetterQueue {
fn send(
&self,
mut record: ProducerRecord,
error: String,
) -> Pin<Box<dyn Future<Output = ()> + Send + '_>> {
let original_topic = std::mem::replace(&mut record.topic, self.topic.clone());
record.headers.push((
HEADER_ORIGINAL_TOPIC.to_string(),
Some(Bytes::from(original_topic)),
));
record.headers.push((
HEADER_EXCEPTION_MESSAGE.to_string(),
Some(Bytes::from(error)),
));
record.partition = None;
Box::pin(async move {
match self.producer.send_record(record).await {
Ok(_) => {
self.routed.fetch_add(1, Ordering::Relaxed);
}
Err(e) => {
self.failures.fetch_add(1, Ordering::Relaxed);
tracing::error!(
error = %e,
topic = %self.topic,
"failed to route a record to the dead-letter topic; the record is lost"
);
}
}
})
}
}
pub fn build_dlq_record(
dlq_topic: &str,
original: &ConsumerRecord,
error: &dyn fmt::Display,
) -> ProducerRecord {
let mut headers: Vec<(String, Option<Bytes>)> = original
.headers
.iter()
.map(|(k, v)| {
(
match std::str::from_utf8(k) {
Ok(s) => s.to_owned(),
Err(_) => {
use std::fmt::Write;
let mut s = String::with_capacity(4 + k.len() * 2);
s.push_str("hex:");
for byte in k.iter() {
let _ = write!(s, "{byte:02x}");
}
s
}
},
v.clone(),
)
})
.collect();
headers.push((
HEADER_ORIGINAL_TOPIC.to_string(),
Some(Bytes::from(original.topic.clone())),
));
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())),
));
ProducerRecord {
topic: dlq_topic.to_string(),
partition: None,
key: original.key.clone(),
value: original.value.clone(),
timestamp: None,
headers,
record_name: None,
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod kafka_dlq_tests {
use super::*;
fn record() -> ProducerRecord {
ProducerRecord {
topic: "orders".to_string(),
partition: Some(7),
key: Some(Bytes::from_static(b"k")),
value: Some(Bytes::from_static(b"v")),
timestamp: None,
headers: vec![("app".to_string(), Some(Bytes::from_static(b"x")))],
record_name: None,
}
}
#[test]
fn retargeting_preserves_provenance_and_drops_the_source_partition() {
let mut record = record();
let original_topic = std::mem::replace(&mut record.topic, "orders.DLQ".to_string());
record.headers.push((
HEADER_ORIGINAL_TOPIC.to_string(),
Some(Bytes::from(original_topic)),
));
record.headers.push((
HEADER_EXCEPTION_MESSAGE.to_string(),
Some(Bytes::from("broker said no")),
));
record.partition = None;
assert_eq!(record.topic, "orders.DLQ");
assert_eq!(
record.partition, None,
"partition 7 may not exist on the dead-letter topic"
);
assert_eq!(record.headers[0].0, "app");
assert_eq!(
record
.headers
.iter()
.find(|(k, _)| k == HEADER_ORIGINAL_TOPIC)
.and_then(|(_, v)| v.clone()),
Some(Bytes::from_static(b"orders"))
);
assert_eq!(
record
.headers
.iter()
.find(|(k, _)| k == HEADER_EXCEPTION_MESSAGE)
.and_then(|(_, v)| v.clone()),
Some(Bytes::from_static(b"broker said no"))
);
}
#[test]
fn consumer_helper_uses_the_same_header_names() {
let original = ConsumerRecord::new("source", 2, 42, None, Some(Bytes::from_static(b"v")));
let built = build_dlq_record("source.DLQ", &original, &"boom");
let names: Vec<&str> = built.headers.iter().map(|(k, _)| k.as_str()).collect();
assert!(names.contains(&HEADER_ORIGINAL_TOPIC));
assert!(names.contains(&HEADER_ORIGINAL_PARTITION));
assert!(names.contains(&HEADER_ORIGINAL_OFFSET));
assert!(names.contains(&HEADER_EXCEPTION_MESSAGE));
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::consumer::ConsumerRecord;
#[test]
fn test_build_dlq_record_provenance_headers() {
let original = ConsumerRecord::new(
"source-topic",
2,
42,
Some(Bytes::from("key")),
Some(Bytes::from("value")),
);
let record = build_dlq_record("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_build_dlq_record_original_headers_preserved() {
let mut original = ConsumerRecord::new("t", 0, 0, None, Some(Bytes::from("v")));
original
.headers
.push((Bytes::from("x-trace-id"), Some(Bytes::from("abc123"))));
let record = build_dlq_record("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_build_dlq_record_preserves_tombstone() {
let tombstone = ConsumerRecord::new("t", 0, 0, Some(Bytes::from("k")), None);
let from_tombstone = build_dlq_record("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 = build_dlq_record("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_build_dlq_record_preserves_null_header_value() {
let mut original = ConsumerRecord::new("t", 0, 0, None, Some(Bytes::from("v")));
original.headers.push((Bytes::from("null-hdr"), None));
original
.headers
.push((Bytes::from("empty-hdr"), Some(Bytes::new())));
let record = build_dlq_record("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);
}
}