use crate::models::{AmqpMessage, AmqpSimpleValue, AmqpValue, MessageId};
use azure_core::fmt::SafeDebug;
use azure_core_amqp::message::{AmqpAnnotationKey, AmqpMessageBody, AmqpMessageProperties};
use std::{
collections::HashMap,
fmt::{Debug, Formatter},
sync::OnceLock,
time::SystemTime,
};
#[derive(Default, PartialEq, Clone, SafeDebug)]
#[safe(true)]
pub struct EventData {
#[safe(false)]
body: Option<Vec<u8>>,
content_type: Option<String>,
correlation_id: Option<MessageId>,
message_id: Option<MessageId>,
#[safe(false)]
properties: Option<HashMap<String, AmqpSimpleValue>>,
}
impl EventData {
pub fn builder() -> builders::EventDataBuilder {
builders::EventDataBuilder::new()
}
pub fn properties(&self) -> Option<&HashMap<String, AmqpSimpleValue>> {
self.properties.as_ref()
}
pub fn body(&self) -> Option<&[u8]> {
self.body.as_deref()
}
pub fn content_type(&self) -> Option<&str> {
self.content_type.as_deref()
}
pub fn correlation_id(&self) -> Option<&MessageId> {
self.correlation_id.as_ref()
}
pub fn message_id(&self) -> Option<&MessageId> {
self.message_id.as_ref()
}
fn from_message(message: &AmqpMessage) -> Self {
let mut event_data_builder = EventData::builder();
if let AmqpMessageBody::Binary(binary) = &message.body {
if binary.len() == 1 {
event_data_builder = event_data_builder.with_body(binary[0].clone());
}
}
if let Some(properties) = &message.properties {
if let Some(content_type) = &properties.content_type {
event_data_builder = event_data_builder.with_content_type(content_type.into());
}
if let Some(correlation_id) = &properties.correlation_id {
event_data_builder = event_data_builder.with_correlation_id(correlation_id.clone());
}
if let Some(message_id) = &properties.message_id {
event_data_builder = event_data_builder.with_message_id(message_id.clone());
}
}
if let Some(application_properties) = &message.application_properties {
for (key, value) in application_properties.0.clone() {
event_data_builder = event_data_builder.add_property(key, value);
}
}
event_data_builder.build()
}
}
impl<T> From<T> for EventData
where
T: Into<Vec<u8>>,
{
fn from(body: T) -> Self {
Self {
body: Some(body.into()),
..Default::default()
}
}
}
impl From<EventData> for AmqpMessage {
fn from(event_data: EventData) -> Self {
let mut message_builder = AmqpMessage::builder();
if event_data.content_type.is_some()
|| event_data.correlation_id.is_some()
|| event_data.message_id.is_some()
{
let mut message_properties = AmqpMessageProperties::default();
if let Some(content_type) = event_data.content_type {
message_properties.content_type = Some(content_type.into());
}
if let Some(correlation_id) = event_data.correlation_id {
message_properties.correlation_id = Some(correlation_id.into());
}
if let Some(message_id) = event_data.message_id {
message_properties.message_id = Some(message_id.into());
}
message_builder = message_builder.with_properties(message_properties);
}
if let Some(properties) = event_data.properties {
for (key, value) in properties {
message_builder = message_builder.add_application_property(key, value);
}
}
if let Some(event_body) = event_data.body {
message_builder =
message_builder.with_body(AmqpMessageBody::Binary(vec![event_body.to_vec()]));
}
message_builder.build()
}
}
pub struct ReceivedEventData {
message: AmqpMessage,
event_data: OnceLock<EventData>,
enqueued_time: OnceLock<Option<SystemTime>>,
offset: OnceLock<Option<String>>,
sequence_number: OnceLock<Option<i64>>,
partition_key: OnceLock<Option<String>>,
system_properties: OnceLock<HashMap<String, AmqpValue>>,
}
impl Debug for ReceivedEventData {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ReceivedEventData")
.field("message", self.raw_amqp_message())
.finish()
}
}
const ENQUEUED_TIME_UTC: &str = "x-opt-enqueued-time";
const OFFSET: &str = "x-opt-offset";
const SEQUENCE_NUMBER: &str = "x-opt-sequence-number";
const PARTITION_KEY: &str = "x-opt-partition-key";
impl ReceivedEventData {
pub fn raw_amqp_message(&self) -> &AmqpMessage {
&self.message
}
pub fn event_data(&self) -> &EventData {
self.event_data
.get_or_init(|| EventData::from_message(&self.message))
}
pub fn enqueued_time(&self) -> Option<SystemTime> {
*self.enqueued_time.get_or_init(|| {
let annotations = self.message.message_annotations.as_ref()?;
for (key, value) in annotations.0.iter() {
if let AmqpAnnotationKey::Symbol(symbol) = key {
if *symbol == ENQUEUED_TIME_UTC {
if let AmqpValue::TimeStamp(timestamp) = value {
return timestamp.0;
}
}
}
}
None
})
}
pub fn offset(&self) -> &Option<String> {
self.offset.get_or_init(|| {
let annotations = self.message.message_annotations.as_ref()?;
for (key, value) in annotations.0.iter() {
if let AmqpAnnotationKey::Symbol(symbol) = key {
if *symbol == OFFSET {
if let AmqpValue::String(offset_value) = value {
return Some(offset_value.clone());
}
}
}
}
None
})
}
pub fn sequence_number(&self) -> Option<i64> {
*self.sequence_number.get_or_init(|| {
let annotations = self.message.message_annotations.as_ref()?;
for (key, value) in annotations.0.iter() {
if let AmqpAnnotationKey::Symbol(symbol) = key {
if *symbol == SEQUENCE_NUMBER {
if let AmqpValue::Long(sequence_number_value) = value {
return Some(*sequence_number_value);
}
}
}
}
None
})
}
pub fn partition_key(&self) -> &Option<String> {
self.partition_key.get_or_init(|| {
let annotations = self.message.message_annotations.as_ref()?;
for (key, value) in annotations.0.iter() {
if let AmqpAnnotationKey::Symbol(symbol) = key {
if *symbol == PARTITION_KEY {
if let AmqpValue::String(partition_key_value) = value {
return Some(partition_key_value.clone());
}
}
}
}
None
})
}
pub fn system_properties(&self) -> &HashMap<String, AmqpValue> {
self.system_properties.get_or_init(|| {
let mut system_properties = HashMap::new();
if let Some(annotations) = self.message.message_annotations.as_ref() {
for (key, value) in annotations.0.iter() {
if let AmqpAnnotationKey::Symbol(symbol) = key {
if *symbol != ENQUEUED_TIME_UTC
&& *symbol != OFFSET
&& *symbol != SEQUENCE_NUMBER
&& *symbol != PARTITION_KEY
{
system_properties.insert(symbol.0.clone(), value.clone());
}
}
}
}
system_properties
})
}
}
impl From<AmqpMessage> for ReceivedEventData {
fn from(message: AmqpMessage) -> Self {
Self {
message,
event_data: OnceLock::new(),
enqueued_time: OnceLock::new(),
offset: OnceLock::new(),
sequence_number: OnceLock::new(),
partition_key: OnceLock::new(),
system_properties: OnceLock::new(),
}
}
}
pub mod builders {
use super::*;
#[derive(Default)]
pub struct EventDataBuilder {
event_data: EventData,
}
impl EventDataBuilder {
pub(super) fn new() -> Self {
Self {
event_data: Default::default(),
}
}
pub fn with_body<T>(mut self, body: T) -> Self
where
T: Into<Vec<u8>>,
{
self.event_data.body = Some(body.into());
self
}
pub fn with_content_type(mut self, content_type: String) -> Self {
self.event_data.content_type = Some(content_type);
self
}
pub fn with_correlation_id(mut self, correlation_id: impl Into<MessageId>) -> Self {
self.event_data.correlation_id = Some(correlation_id.into());
self
}
pub fn with_message_id(mut self, message_id: impl Into<MessageId>) -> Self {
self.event_data.message_id = Some(message_id.into());
self
}
pub fn add_property(mut self, key: String, value: impl Into<AmqpSimpleValue>) -> Self {
if let Some(mut properties) = self.event_data.properties {
properties.insert(key, value.into());
self.event_data.properties = Some(properties);
} else {
let mut properties = HashMap::new();
properties.insert(key, value.into());
self.event_data.properties = Some(properties);
}
self
}
pub fn build(self) -> EventData {
self.event_data
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_event_data_builder_with_body() {
let body = vec![1, 2, 3];
let event_data = EventData::builder().with_body(body.clone()).build();
assert_eq!(event_data.body().unwrap(), &body);
}
#[test]
fn test_event_data_builder_with_content_type() {
let content_type = "application/json".to_string();
let event_data = EventData::builder()
.with_content_type(content_type.clone())
.build();
assert_eq!(event_data.content_type(), Some(content_type.as_str()));
}
#[test]
fn test_event_data_builder_with_correlation_id() {
let correlation_id = MessageId::String("correlation-id".to_string());
let event_data = EventData::builder()
.with_correlation_id(correlation_id.clone())
.build();
assert_eq!(event_data.correlation_id(), Some(&correlation_id));
}
#[test]
fn test_event_data_builder_with_message_id() {
let message_id = MessageId::String("message-id".to_string());
let event_data = EventData::builder()
.with_message_id(message_id.clone())
.build();
assert_eq!(event_data.message_id(), Some(&message_id));
}
#[test]
fn test_event_data_builder_add_property() {
let key = "key".to_string();
let value: AmqpSimpleValue = "value".into();
let event_data = EventData::builder()
.add_property(key.clone(), value.clone())
.build();
assert_eq!(event_data.properties().unwrap().get(&key), Some(&value));
}
#[test]
fn test_event_data_builder_build() {
let body = vec![1, 2, 3];
let content_type = "application/json".to_string();
let correlation_id = MessageId::String("correlation-id".to_string());
let message_id = MessageId::String("message-id".to_string());
let key = "key".to_string();
let value: AmqpSimpleValue = "value".into();
let event_data = EventData::builder()
.with_body(body.clone())
.with_content_type(content_type.clone())
.with_correlation_id(correlation_id.clone())
.with_message_id(message_id.clone())
.add_property(key.clone(), value.clone())
.build();
assert_eq!(event_data.body().unwrap(), &body);
assert_eq!(event_data.content_type(), Some(content_type.as_str()));
assert_eq!(event_data.correlation_id(), Some(&correlation_id));
assert_eq!(event_data.message_id(), Some(&message_id));
assert_eq!(event_data.properties().unwrap().get(&key), Some(&value));
}
}