use bytes::Bytes;
use reliar_core::{ContentType, Message, MessageId, Serializer};
use reliar_outbox::OutboxEnqueue;
use sqlx::{Postgres, Transaction};
use tracing::Instrument as _;
use super::error::{EnqueueError, map_enqueue_error};
use crate::settings::PostgresOutboxSettings;
use super::PostgresOutboxStore;
use crate::connection::schema::{restore_search_path, set_search_path};
async fn insert_enqueued<T>(
tx: &mut Transaction<'_, Postgres>,
settings: &PostgresOutboxSettings,
envelope: &reliar_core::Envelope<T>,
payload: &Bytes,
content_type: &ContentType,
) -> Result<(), sqlx::Error> {
let restore = if settings.enqueue_sets_search_path {
Some(set_search_path(tx, &settings.schema).await?)
} else {
None
};
let result = insert_row(tx, envelope, payload, content_type).await;
if result.is_ok()
&& let Some(previous) = restore
{
restore_search_path(tx, &previous).await?;
}
result
}
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)
});
sqlx::query!(
r#"INSERT INTO outbox (
message_id, message_type, message_version,
correlation_id, conversation_id, causation_id, request_id,
content_type, payload, tenant_id, expires_at, ordering_key,
metadata, headers, available_at
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14, now())"#,
envelope.id.as_uuid(),
envelope.message_type.name(),
i32::from(envelope.message_type.version()),
corr.correlation_id
.as_ref()
.map(reliar_core::CorrelationId::as_str),
corr.conversation_id.as_uuid(),
corr.causation_id.map(|id| id.as_uuid()),
corr.request_id.map(|id| id.as_uuid()),
content_type.as_str(),
&payload[..],
envelope.metadata.tenant_id.as_deref(),
envelope.metadata.delivery.expires_at,
None::<&str>,
metadata_json,
headers_json,
)
.execute(&mut **tx)
.await?;
Ok(())
}
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_enqueued(tx, &self.settings, &envelope, &payload, self.content_type())
.await
.map_err(|source| map_enqueue_error(envelope.id, source))?;
Ok(envelope.id)
}
.instrument(span)
}
}