use bytes::Bytes;
use reliar_core::{ContentType, Message, MessageId, Serializer};
use reliar_outbox::OutboxEnqueue;
use sqlx::{Postgres, Transaction};
use tracing::Instrument as _;
use super::enqueue::{InsertOutboxRowParams, insert_outbox_row};
use super::error::{EnqueueError, map_enqueue_error};
use super::PostgresOutboxStore;
async fn insert_row<T>(
tx: &mut Transaction<'_, Postgres>,
envelope: &reliar_core::Envelope<T>,
payload: &Bytes,
content_type: &ContentType,
) -> Result<(), sqlx::Error> {
let corr = &envelope.metadata.correlation;
let sent_at_ms = envelope
.metadata
.delivery
.sent_at
.map(crate::records::encode_epoch_millis);
let rest = crate::records::MetadataRest {
trace: crate::records::TraceRest {
traceparent: envelope.metadata.trace.traceparent.clone(),
tracestate: envelope.metadata.trace.tracestate.clone(),
},
routing: crate::records::RoutingRest {
source: envelope
.metadata
.routing
.source
.as_ref()
.map(|v| v.as_str().to_owned()),
destination: envelope
.metadata
.routing
.destination
.as_ref()
.map(|v| v.as_str().to_owned()),
reply_to: envelope
.metadata
.routing
.reply_to
.as_ref()
.map(|v| v.as_str().to_owned()),
},
delivery: crate::records::DeliveryRest {
sent_at_ms,
deduplication_id: envelope.metadata.delivery.deduplication_id.clone(),
},
};
let metadata_json = if rest.trace.traceparent.is_none()
&& rest.trace.tracestate.is_none()
&& rest.routing.source.is_none()
&& rest.routing.destination.is_none()
&& rest.routing.reply_to.is_none()
&& rest.delivery.sent_at_ms.is_none()
&& rest.delivery.deduplication_id.is_none()
{
None
} else {
serde_json::to_value(&rest).ok()
};
let headers_json = envelope.headers().filter(|h| !h.is_empty()).map(|h| {
let map: serde_json::Map<String, serde_json::Value> = h
.iter()
.map(|(k, v)| (k.to_owned(), serde_json::Value::String(v.to_owned())))
.collect();
serde_json::Value::Object(map)
});
insert_outbox_row(
&mut *tx,
InsertOutboxRowParams {
message_id: envelope.id.as_uuid(),
message_type: envelope.message_type.name(),
message_version: i32::from(envelope.message_type.version()),
correlation_id: corr
.correlation_id
.as_ref()
.map(reliar_core::CorrelationId::as_str),
conversation_id: corr.conversation_id.as_uuid(),
causation_id: corr.causation_id.map(|id| id.as_uuid()),
request_id: corr.request_id.map(|id| id.as_uuid()),
content_type: content_type.as_str(),
payload: &payload[..],
tenant_id: envelope.metadata.tenant_id.as_deref(),
expires_at: envelope.metadata.delivery.expires_at,
ordering_key: None,
metadata: metadata_json,
headers: headers_json,
},
)
.await
}
impl<'c, Ser> OutboxEnqueue<Transaction<'c, Postgres>> for PostgresOutboxStore<Ser>
where
Ser: Serializer + Send + Sync + 'static,
{
type Error = EnqueueError<Ser::Error>;
fn enqueue_envelope<T: Message + Sync>(
&self,
tx: &mut Transaction<'c, Postgres>,
typed_envelope: reliar_core::Envelope<T>,
) -> impl Future<Output = Result<MessageId, Self::Error>> + Send {
let span = tracing::debug_span!(
"reliar.outbox.enqueue",
message.id = %typed_envelope.id,
message.type = %typed_envelope.message_type,
);
let payload = {
let _guard = span.enter();
self.serializer
.serialize(&typed_envelope.body)
.map_err(|source| EnqueueError::Serialize { source })
};
let envelope = typed_envelope.map_body(|_| ());
async move {
let payload = payload?;
insert_row(tx, &envelope, &payload, self.content_type())
.await
.map_err(|source| map_enqueue_error(envelope.id, source))?;
Ok(envelope.id)
}
.instrument(span)
}
}