use core::fmt;
use crate::{
ContentType, ConversationId, CorrelationId, MessageId, RequestId,
ids::{IdError, contains_control_char},
};
#[derive(Clone, Debug, Default, PartialEq)]
#[non_exhaustive]
pub struct Metadata {
pub correlation: CorrelationMetadata,
pub trace: TraceContext,
pub routing: RoutingMetadata,
pub delivery: DeliveryMetadata,
pub tenant_id: Option<String>,
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub struct CorrelationMetadata {
pub correlation_id: Option<CorrelationId>,
pub conversation_id: ConversationId,
pub causation_id: Option<MessageId>,
pub request_id: Option<RequestId>,
}
impl Default for CorrelationMetadata {
fn default() -> Self {
Self {
correlation_id: None,
conversation_id: ConversationId::UNSET,
causation_id: None,
request_id: None,
}
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[non_exhaustive]
pub struct TraceContext {
pub traceparent: Option<String>,
pub tracestate: Option<String>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct RoutingMetadata {
pub source: Option<EndpointAddress>,
pub destination: Option<EndpointAddress>,
pub reply_to: Option<EndpointAddress>,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct EndpointAddress(String);
impl EndpointAddress {
pub const MAX_LEN: usize = 256;
pub fn parse(s: impl Into<String>) -> Result<Self, IdError> {
let s = s.into();
if s.is_empty() {
return Err(IdError::Empty);
}
if contains_control_char(&s) {
return Err(IdError::ControlCharacter);
}
if s.len() > Self::MAX_LEN {
return Err(IdError::TooLong {
len: s.len(),
max: Self::MAX_LEN,
});
}
Ok(Self(s))
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
impl fmt::Display for EndpointAddress {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.0)
}
}
#[cfg(feature = "serde")]
#[cfg_attr(docsrs, doc(cfg(feature = "serde")))]
impl serde::Serialize for EndpointAddress {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
s.collect_str(&self.0)
}
}
#[cfg(feature = "serde")]
#[cfg_attr(docsrs, doc(cfg(feature = "serde")))]
impl<'de> serde::Deserialize<'de> for EndpointAddress {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
let raw = String::deserialize(d)?;
Self::parse(raw).map_err(serde::de::Error::custom)
}
}
#[derive(Clone, Debug, PartialEq)]
#[non_exhaustive]
pub struct DeliveryMetadata {
pub content_type: ContentType,
pub sent_at: Option<time::OffsetDateTime>,
pub expires_at: Option<time::OffsetDateTime>,
pub deduplication_id: Option<String>,
}
impl Default for DeliveryMetadata {
fn default() -> Self {
Self {
content_type: ContentType::JSON,
sent_at: None,
expires_at: None,
deduplication_id: None,
}
}
}
#[cfg(feature = "serde")]
#[cfg_attr(docsrs, doc(cfg(feature = "serde")))]
mod serde_impls {
use serde::{Deserialize, Serialize};
use super::{CorrelationMetadata, DeliveryMetadata, Metadata, RoutingMetadata, TraceContext};
#[derive(Serialize, Deserialize)]
#[serde(remote = "Metadata")]
struct MetadataDef {
#[serde(default)]
correlation: CorrelationMetadata,
#[serde(default)]
trace: TraceContext,
#[serde(default)]
routing: RoutingMetadata,
#[serde(default)]
delivery: DeliveryMetadata,
#[serde(default)]
tenant_id: Option<String>,
}
impl Serialize for Metadata {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
MetadataDef::serialize(self, s)
}
}
impl<'de> Deserialize<'de> for Metadata {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
MetadataDef::deserialize(d)
}
}
impl Serialize for CorrelationMetadata {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
#[derive(Serialize)]
#[allow(clippy::struct_field_names)]
struct Def<'a> {
correlation_id: &'a Option<super::CorrelationId>,
conversation_id: &'a super::ConversationId,
causation_id: &'a Option<super::MessageId>,
request_id: &'a Option<super::RequestId>,
}
Def {
correlation_id: &self.correlation_id,
conversation_id: &self.conversation_id,
causation_id: &self.causation_id,
request_id: &self.request_id,
}
.serialize(s)
}
}
impl<'de> Deserialize<'de> for CorrelationMetadata {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
fn default_conversation_id() -> super::ConversationId {
super::ConversationId::UNSET
}
#[derive(Deserialize)]
#[allow(clippy::struct_field_names)]
struct Def {
#[serde(default)]
correlation_id: Option<super::CorrelationId>,
#[serde(default = "default_conversation_id")]
conversation_id: super::ConversationId,
#[serde(default)]
causation_id: Option<super::MessageId>,
#[serde(default)]
request_id: Option<super::RequestId>,
}
let def = Def::deserialize(d)?;
Ok(CorrelationMetadata {
correlation_id: def.correlation_id,
conversation_id: def.conversation_id,
causation_id: def.causation_id,
request_id: def.request_id,
})
}
}
impl Serialize for RoutingMetadata {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
#[derive(Serialize)]
struct Def<'a> {
source: &'a Option<super::EndpointAddress>,
destination: &'a Option<super::EndpointAddress>,
reply_to: &'a Option<super::EndpointAddress>,
}
Def {
source: &self.source,
destination: &self.destination,
reply_to: &self.reply_to,
}
.serialize(s)
}
}
impl<'de> Deserialize<'de> for RoutingMetadata {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
#[derive(Deserialize)]
struct Def {
#[serde(default)]
source: Option<super::EndpointAddress>,
#[serde(default)]
destination: Option<super::EndpointAddress>,
#[serde(default)]
reply_to: Option<super::EndpointAddress>,
}
let def = Def::deserialize(d)?;
Ok(RoutingMetadata {
source: def.source,
destination: def.destination,
reply_to: def.reply_to,
})
}
}
impl Serialize for DeliveryMetadata {
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
#[derive(Serialize)]
struct Def<'a> {
content_type: &'a super::ContentType,
#[serde(with = "time::serde::rfc3339::option")]
sent_at: &'a Option<time::OffsetDateTime>,
#[serde(with = "time::serde::rfc3339::option")]
expires_at: &'a Option<time::OffsetDateTime>,
deduplication_id: &'a Option<String>,
}
Def {
content_type: &self.content_type,
sent_at: &self.sent_at,
expires_at: &self.expires_at,
deduplication_id: &self.deduplication_id,
}
.serialize(s)
}
}
impl<'de> Deserialize<'de> for DeliveryMetadata {
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
fn default_content_type() -> super::ContentType {
super::ContentType::JSON
}
#[derive(Deserialize)]
struct Def {
#[serde(default = "default_content_type")]
content_type: super::ContentType,
#[serde(default, with = "time::serde::rfc3339::option")]
sent_at: Option<time::OffsetDateTime>,
#[serde(default, with = "time::serde::rfc3339::option")]
expires_at: Option<time::OffsetDateTime>,
#[serde(default)]
deduplication_id: Option<String>,
}
let def = Def::deserialize(d)?;
Ok(DeliveryMetadata {
content_type: def.content_type,
sent_at: def.sent_at,
expires_at: def.expires_at,
deduplication_id: def.deduplication_id,
})
}
}
}