use std::any::Any;
use std::collections::HashMap;
use std::time::Instant;
#[derive(Debug, Clone)]
pub struct MessageMetadata {
pub source: String, pub timestamp: Instant,
pub key: Option<String>, pub backend_specific: Option<Box<dyn BackendMetadata>>,
}
pub trait BackendMetadata: Send + Sync + std::fmt::Debug {
fn as_any(&self) -> &dyn Any;
fn clone_box(&self) -> Box<dyn BackendMetadata>;
}
impl Clone for Box<dyn BackendMetadata> {
fn clone(&self) -> Self {
self.clone_box()
}
}
#[derive(Debug, Clone)]
pub struct KafkaMetadata {
pub topic: String,
pub partition: i32,
pub offset: i64,
pub consumer_group: Option<String>,
pub manual_commit: bool,
pub headers: HashMap<String, String>,
}
impl BackendMetadata for KafkaMetadata {
fn as_any(&self) -> &dyn Any {
self
}
fn clone_box(&self) -> Box<dyn BackendMetadata> {
Box::new(self.clone())
}
}
#[derive(Debug, Clone)]
pub struct RedisMetadata {
pub stream: String,
pub entry_id: String,
pub consumer_group: Option<String>,
pub manual_ack: bool,
}
impl BackendMetadata for RedisMetadata {
fn as_any(&self) -> &dyn Any {
self
}
fn clone_box(&self) -> Box<dyn BackendMetadata> {
Box::new(self.clone())
}
}
impl MessageMetadata {
pub fn new(
source: String,
timestamp: Instant,
key: Option<String>,
backend_specific: Option<Box<dyn BackendMetadata>>,
) -> Self {
Self {
source,
timestamp,
key,
backend_specific,
}
}
pub fn kafka_metadata(&self) -> Option<&KafkaMetadata> {
self.backend_specific
.as_ref()?
.as_any()
.downcast_ref::<KafkaMetadata>()
}
pub fn redis_metadata(&self) -> Option<&RedisMetadata> {
self.backend_specific
.as_ref()?
.as_any()
.downcast_ref::<RedisMetadata>()
}
}