use std::sync::Arc;
use std::time::Duration;
use reliar_core::{ContentType, Serializer};
use reliar_outbox::{
AcquireRequest, AcquiredBatch, FailedRecord, FailureOutcome, OutboxStats, OutboxStore,
PoisonedRow, PurgeReport, PurgeRequest, RecordRef, WorkerId,
};
use sqlx::{PgConnection, PgPool};
use crate::connection::session::Session;
use crate::records::{RawRow, decode_row};
use crate::settings::PostgresOutboxSettings;
#[cfg(feature = "json")]
use reliar_core::JsonSerializer;
use super::claim as claim_repo;
use super::error::PostgresOutboxError;
use super::outcomes as outcomes_repo;
use super::purge as purge_repo;
#[cfg_attr(not(feature = "json"), doc = "```ignore")]
#[cfg_attr(feature = "json", doc = "```no_run")]
#[non_exhaustive]
pub struct PostgresOutboxStore<
#[cfg(feature = "json")] Ser = JsonSerializer,
#[cfg(not(feature = "json"))] Ser,
> {
pub(super) session: Session,
pub(super) serializer: Arc<Ser>,
}
impl<Ser> Clone for PostgresOutboxStore<Ser> {
fn clone(&self) -> Self {
Self {
session: self.session.clone(),
serializer: Arc::clone(&self.serializer),
}
}
}
impl<Ser> std::fmt::Debug for PostgresOutboxStore<Ser> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PostgresOutboxStore")
.field("session", &self.session)
.finish_non_exhaustive()
}
}
impl<Ser: Serializer + Send + Sync + 'static> PostgresOutboxStore<Ser> {
#[cfg_attr(not(feature = "json"), doc = "```ignore")]
#[cfg_attr(feature = "json", doc = "```no_run")]
#[must_use]
#[allow(
clippy::needless_pass_by_value,
reason = "the public signature takes settings by value (ADR 0047 §1); only \
statement_timeout is read today, but PostgresOutboxSettings is #[non_exhaustive] \
and may grow a field this constructor needs to own or move out of later"
)]
pub fn with_serializer(
pool: PgPool,
settings: PostgresOutboxSettings,
serializer: Ser,
) -> Self {
let session = Session::new(pool, settings.statement_timeout);
Self {
session,
serializer: Arc::new(serializer),
}
}
#[cfg_attr(not(feature = "json"), doc = "```ignore")]
#[cfg_attr(feature = "json", doc = "```no_run")]
#[must_use]
pub fn content_type(&self) -> &ContentType {
self.serializer.content_type()
}
}
#[cfg(feature = "json")]
#[cfg_attr(docsrs, doc(cfg(feature = "json")))]
impl PostgresOutboxStore<JsonSerializer> {
#[must_use]
pub fn new(pool: PgPool) -> Self {
Self::with_serializer(pool, PostgresOutboxSettings::default(), JsonSerializer)
}
#[must_use]
pub fn with_settings(pool: PgPool, settings: PostgresOutboxSettings) -> Self {
Self::with_serializer(pool, settings, JsonSerializer)
}
}
impl<Ser: Serializer + Send + Sync + 'static> OutboxStore for PostgresOutboxStore<Ser> {
type Error = PostgresOutboxError;
async fn acquire(&self, request: AcquireRequest) -> Result<AcquiredBatch, Self::Error> {
let session = &self.session;
let batch_size = i64::from(request.batch_size);
let lease_ms = i64::try_from(request.lease.as_millis()).unwrap_or(i64::MAX);
let worker = request.worker.as_str();
let rows: Vec<RawRow> = session
.run(async |conn: &mut PgConnection| {
claim_repo::claim_rows(
&mut *conn,
claim_repo::ClaimRowsParams {
batch_size,
worker,
lease_ms,
},
)
.await
})
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))?;
let mut records = Vec::with_capacity(rows.len());
let mut poisoned = Vec::new();
let mut poisoned_ids = Vec::new();
let mut poisoned_tokens = Vec::new();
let mut poisoned_errors = Vec::new();
for raw in rows {
let claim_token = raw.claim_token;
match decode_row(raw) {
Ok(record) => records.push(record),
Err(err) => {
poisoned_ids.push(err.id.as_uuid());
poisoned_tokens.push(claim_token);
poisoned_errors.push(crate::records::truncate_last_error(err.detail.clone()));
poisoned.push(PoisonedRow::new(err.id, err.message_id, err.detail));
}
}
}
if !poisoned_ids.is_empty() {
let undecodable =
crate::records::encode_dead_reason(reliar_outbox::DeadReason::Undecodable);
let sweep_result = session
.run(async |conn: &mut PgConnection| {
claim_repo::poison_sweep_rows(
&mut *conn,
claim_repo::PoisonSweepRowsParams {
ids: &poisoned_ids,
tokens: &poisoned_tokens,
errors: &poisoned_errors,
dead_reason: undecodable,
},
)
.await
})
.await;
if let Err(err) = sweep_result {
tracing::warn!(
target: "reliar.outbox.acquire",
worker_id = %worker,
poisoned_count = poisoned_ids.len(),
error = %session.map_err::<PostgresOutboxError>(err),
"poison sweep failed; the claimed batch is returned and the undecodable rows \
stay leased until their lease lapses"
);
}
}
Ok(AcquiredBatch::new(records, poisoned))
}
async fn complete(&self, worker: &WorkerId, items: &[RecordRef]) -> Result<u64, Self::Error> {
let session = &self.session;
if items.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
let tokens: Vec<Option<uuid::Uuid>> = items
.iter()
.map(|i| i.claim_token.map(|t| t.as_uuid()))
.collect();
let applied = session
.run(async |conn: &mut PgConnection| {
outcomes_repo::complete_rows(
&mut *conn,
outcomes_repo::CompleteRowsParams {
ids: &ids,
tokens: &tokens,
},
)
.await
})
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))?;
if applied.len() < ids.len() {
tracing::warn!(
target: "reliar.outbox.complete",
requested = ids.len(),
worker.id = %worker,
applied = applied.len(),
fenced_ids = ?fenced_ids(&ids, &applied),
"fewer rows completed than requested — the fenced rows belong to a superseded claim"
);
}
Ok(applied.len() as u64)
}
async fn fail(&self, worker: &WorkerId, items: &[FailedRecord]) -> Result<u64, Self::Error> {
let session = &self.session;
if items.is_empty() {
return Ok(0);
}
let batches = classify_failures(items);
let requested = batches.retry_ids.len() + batches.dead_ids.len();
let (applied_retry, applied_dead) = apply_fail_batches(session, &batches).await?;
let applied = applied_retry.len() + applied_dead.len();
if applied < requested {
let ids: Vec<uuid::Uuid> = batches
.retry_ids
.iter()
.chain(batches.dead_ids.iter())
.copied()
.collect();
let applied_ids: Vec<uuid::Uuid> = applied_retry
.iter()
.chain(applied_dead.iter())
.copied()
.collect();
tracing::warn!(
target: "reliar.outbox.fail",
requested,
worker.id = %worker,
applied,
fenced_ids = ?fenced_ids(&ids, &applied_ids),
"fewer rows failed than requested — the fenced rows belong to a superseded claim"
);
}
Ok(applied as u64)
}
async fn release(&self, worker: &WorkerId, items: &[RecordRef]) -> Result<u64, Self::Error> {
let session = &self.session;
if items.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
let tokens: Vec<Option<uuid::Uuid>> = items
.iter()
.map(|i| i.claim_token.map(|t| t.as_uuid()))
.collect();
let applied = session
.run(async |conn: &mut PgConnection| {
outcomes_repo::release_rows(
&mut *conn,
outcomes_repo::ReleaseRowsParams {
ids: &ids,
tokens: &tokens,
},
)
.await
})
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))?;
if applied.len() < ids.len() {
tracing::warn!(
target: "reliar.outbox.release",
requested = ids.len(),
worker.id = %worker,
applied = applied.len(),
fenced_ids = ?fenced_ids(&ids, &applied),
"fewer rows released than requested — the fenced rows belong to a superseded claim"
);
}
Ok(applied.len() as u64)
}
async fn extend_lease(
&self,
worker: &WorkerId,
items: &[RecordRef],
lease: Duration,
) -> Result<u64, Self::Error> {
let session = &self.session;
if items.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
let tokens: Vec<Option<uuid::Uuid>> = items
.iter()
.map(|i| i.claim_token.map(|t| t.as_uuid()))
.collect();
let lease_ms = i64::try_from(lease.as_millis()).unwrap_or(i64::MAX);
let applied = session
.run(async |conn: &mut PgConnection| {
outcomes_repo::extend_lease_rows(
&mut *conn,
outcomes_repo::ExtendLeaseRowsParams {
ids: &ids,
tokens: &tokens,
lease_ms,
},
)
.await
})
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))?;
if applied.len() < ids.len() {
tracing::warn!(
target: "reliar.outbox.extend_lease",
requested = ids.len(),
worker.id = %worker,
applied = applied.len(),
fenced_ids = ?fenced_ids(&ids, &applied),
"fewer leases renewed than requested — the fenced rows belong to a superseded claim"
);
}
Ok(applied.len() as u64)
}
async fn purge(&self, request: PurgeRequest) -> Result<PurgeReport, Self::Error> {
let session = &self.session;
let batch_size = i64::from(request.batch_size);
let expired_reason = crate::records::encode_dead_reason(reliar_outbox::DeadReason::Expired);
let published_ms = request.published_retention.map(to_millis);
let dead_ms = request.dead_retention.map(to_millis);
let (published_deleted, dead_deleted, expired_to_dead) = session
.run(async |conn: &mut PgConnection| {
let published_deleted = match published_ms {
Some(retention_ms) => {
purge_repo::purge_published_rows(
&mut *conn,
purge_repo::PurgePublishedRowsParams {
retention_ms,
batch_size,
},
)
.await?
}
None => 0,
};
let dead_deleted = match dead_ms {
Some(retention_ms) => {
purge_repo::purge_dead_retention_rows(
&mut *conn,
purge_repo::PurgeDeadRetentionRowsParams {
retention_ms,
batch_size,
},
)
.await?
}
None => 0,
};
let expired_to_dead = purge_repo::purge_expired_sweep_rows(
&mut *conn,
purge_repo::PurgeExpiredSweepRowsParams {
batch_size,
dead_reason: expired_reason,
},
)
.await?;
Ok((published_deleted, dead_deleted, expired_to_dead))
})
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))?;
Ok(PurgeReport::new(
published_deleted,
dead_deleted,
expired_to_dead,
))
}
async fn stats(&self) -> Result<OutboxStats, Self::Error> {
let session = &self.session;
let row = session
.run(async |conn: &mut PgConnection| purge_repo::stats_row(&mut *conn).await)
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))?;
Ok(OutboxStats::new(
u64::try_from(row.pending).unwrap_or(0),
u64::try_from(row.dead).unwrap_or(0),
u64::try_from(row.expired_pending).unwrap_or(0),
row.oldest_pending_available_at,
row.as_of,
))
}
}
fn fenced_ids(requested: &[uuid::Uuid], applied: &[uuid::Uuid]) -> Vec<uuid::Uuid> {
let applied: std::collections::HashSet<&uuid::Uuid> = applied.iter().collect();
requested
.iter()
.filter(|id| !applied.contains(id))
.copied()
.collect()
}
struct FailBatches {
retry_ids: Vec<uuid::Uuid>,
retry_tokens: Vec<Option<uuid::Uuid>>,
retry_errors: Vec<String>,
retry_delays: Vec<i64>,
dead_ids: Vec<uuid::Uuid>,
dead_tokens: Vec<Option<uuid::Uuid>>,
dead_errors: Vec<String>,
dead_reasons: Vec<&'static str>,
}
fn classify_failures(items: &[FailedRecord]) -> FailBatches {
let mut batches = FailBatches {
retry_ids: Vec::new(),
retry_tokens: Vec::new(),
retry_errors: Vec::new(),
retry_delays: Vec::new(),
dead_ids: Vec::new(),
dead_tokens: Vec::new(),
dead_errors: Vec::new(),
dead_reasons: Vec::new(),
};
for item in items {
match item.outcome {
FailureOutcome::Retry { delay } => {
batches.retry_ids.push(item.record.id.as_uuid());
batches
.retry_tokens
.push(item.record.claim_token.map(|t| t.as_uuid()));
batches.retry_errors.push(item.error.clone());
batches
.retry_delays
.push(i64::try_from(delay.as_millis()).unwrap_or(i64::MAX));
}
FailureOutcome::Dead { reason } => {
batches.dead_ids.push(item.record.id.as_uuid());
batches
.dead_tokens
.push(item.record.claim_token.map(|t| t.as_uuid()));
batches.dead_errors.push(item.error.clone());
batches
.dead_reasons
.push(crate::records::encode_dead_reason(reason));
}
_ => tracing::error!(
id = %item.record.id,
"unrecognised FailureOutcome variant; row left as-is"
),
}
}
batches
}
async fn apply_fail_batches(
session: &Session,
batches: &FailBatches,
) -> Result<(Vec<uuid::Uuid>, Vec<uuid::Uuid>), PostgresOutboxError> {
session
.run(async |conn: &mut PgConnection| {
let applied_retry = if batches.retry_ids.is_empty() {
Vec::new()
} else {
outcomes_repo::fail_retry_rows(
&mut *conn,
outcomes_repo::FailRetryRowsParams {
ids: &batches.retry_ids,
tokens: &batches.retry_tokens,
errors: &batches.retry_errors,
delays_ms: &batches.retry_delays,
},
)
.await?
};
let applied_dead = if batches.dead_ids.is_empty() {
Vec::new()
} else {
outcomes_repo::fail_dead_rows(
&mut *conn,
outcomes_repo::FailDeadRowsParams {
ids: &batches.dead_ids,
tokens: &batches.dead_tokens,
errors: &batches.dead_errors,
reasons: &batches.dead_reasons,
},
)
.await?
};
Ok((applied_retry, applied_dead))
})
.await
.map_err(|e| session.map_err::<PostgresOutboxError>(e))
}
fn to_millis(duration: Duration) -> i64 {
i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)
}