use core::fmt;
use bytes::Bytes;
use crate::{
ConversationId, CorrelationId, CorrelationMetadata, HeaderError, Headers, Message, MessageId,
MessageType, Metadata,
};
#[non_exhaustive]
pub struct Envelope<T> {
pub id: MessageId,
pub message_type: MessageType,
pub body: T,
pub metadata: Metadata,
pub(crate) headers: Option<Headers>,
}
pub type SerializedEnvelope = Envelope<Bytes>;
impl<T> Envelope<T> {
#[must_use]
pub fn headers(&self) -> Option<&Headers> {
self.headers.as_ref()
}
pub fn headers_mut(&mut self) -> &mut Headers {
self.headers.get_or_insert_with(Headers::default)
}
pub fn set_headers(&mut self, headers: Option<Headers>) {
self.headers = headers;
}
#[must_use]
pub fn map_body<U>(self, f: impl FnOnce(T) -> U) -> Envelope<U> {
Envelope {
id: self.id,
message_type: self.message_type,
body: f(self.body),
metadata: self.metadata,
headers: self.headers,
}
}
pub fn try_map_body<U, E>(self, f: impl FnOnce(T) -> Result<U, E>) -> Result<Envelope<U>, E> {
Ok(Envelope {
id: self.id,
message_type: self.message_type,
body: f(self.body)?,
metadata: self.metadata,
headers: self.headers,
})
}
}
impl<T: Message> Envelope<T> {
pub fn builder(body: T) -> EnvelopeBuilder<T> {
EnvelopeBuilder::new(body)
}
}
impl SerializedEnvelope {
#[must_use]
pub fn from_parts(
id: MessageId,
message_type: MessageType,
body: Bytes,
metadata: Metadata,
headers: Option<Headers>,
) -> Self {
Self {
id,
message_type,
body,
metadata,
headers,
}
}
}
impl<T> fmt::Debug for Envelope<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Envelope")
.field("id", &self.id)
.field("message_type", &self.message_type)
.field("body", &"<elided>")
.field("metadata", &self.metadata)
.field("headers", &self.headers)
.finish()
}
}
impl<T: PartialEq> PartialEq for Envelope<T> {
fn eq(&self, other: &Self) -> bool {
self.id == other.id
&& self.message_type == other.message_type
&& self.body == other.body
&& self.metadata == other.metadata
&& self.headers == other.headers
}
}
impl<T: Clone> Clone for Envelope<T> {
fn clone(&self) -> Self {
Self {
id: self.id,
message_type: self.message_type.clone(),
body: self.body.clone(),
metadata: self.metadata.clone(),
headers: self.headers.clone(),
}
}
}
#[must_use]
pub struct EnvelopeBuilder<T> {
id: Option<MessageId>,
body: T,
metadata: Metadata,
headers: Option<Headers>,
}
impl<T> fmt::Debug for EnvelopeBuilder<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("EnvelopeBuilder")
.field("id", &self.id)
.field("body", &"<elided>")
.field("metadata", &self.metadata)
.field("headers", &self.headers)
.finish()
}
}
impl<T: Message> EnvelopeBuilder<T> {
fn new(body: T) -> Self {
Self {
id: None,
body,
metadata: Metadata::default(),
headers: None,
}
}
pub fn id(mut self, id: MessageId) -> Self {
self.id = Some(id);
self
}
pub fn metadata(mut self, metadata: Metadata) -> Self {
self.metadata = metadata;
self
}
pub fn correlation(mut self, correlation: CorrelationMetadata) -> Self {
self.metadata.correlation = correlation;
self
}
pub fn correlation_id(mut self, id: CorrelationId) -> Self {
self.metadata.correlation.correlation_id = Some(id);
self
}
pub fn conversation(mut self, id: ConversationId) -> Self {
self.metadata.correlation.conversation_id = id;
self
}
pub fn causation(mut self, parent: MessageId) -> Self {
self.metadata.correlation.causation_id = Some(parent);
self
}
pub fn tenant(mut self, tenant_id: impl Into<String>) -> Self {
self.metadata.tenant_id = Some(tenant_id.into());
self
}
pub fn expires_at(mut self, at: time::OffsetDateTime) -> Self {
self.metadata.delivery.expires_at = Some(at);
self
}
pub fn trace(mut self, traceparent: impl Into<String>, tracestate: Option<String>) -> Self {
self.metadata.trace.traceparent = Some(traceparent.into());
self.metadata.trace.tracestate = tracestate;
self
}
pub fn header(
mut self,
k: impl Into<String>,
v: impl Into<String>,
) -> Result<Self, HeaderError> {
self.headers
.get_or_insert_with(Headers::default)
.insert(k, v)?;
Ok(self)
}
#[must_use]
pub fn build(mut self) -> Envelope<T> {
let id = self.id.unwrap_or_default();
if self.metadata.correlation.conversation_id.is_unset() {
self.metadata.correlation.conversation_id = ConversationId::from_uuid(id.as_uuid());
}
Envelope {
id,
message_type: MessageType::of::<T>(),
body: self.body,
metadata: self.metadata,
headers: self.headers,
}
}
}