use std::sync::Arc;
use bytes::Bytes;
use reliar_core::{ContentType, Message, MessageId, Serializer};
use reliar_outbox::{
AcquiredBatch, CompletedMessage, DeadLetterPage, DeadQuery, FailedMessage, FailureOutcome,
MessageRef, OutboxDeadLetters, OutboxStats, OutboxStore, PoisonedRow, PurgeReport,
PurgeRequest, WorkerId,
};
use sqlx::{PgPool, Postgres, Transaction};
use crate::error::{
EnqueueError, PostgresStoreError, is_undefined_table, map_enqueue_error, map_operational_error,
};
use crate::records::{RawRow, decode_row};
use crate::settings::PostgresOutboxSettings;
#[cfg(feature = "json")]
use reliar_core::JsonSerializer;
const MAX_LIST_DEAD_LIMIT: u32 = 1000;
#[derive(Clone, Debug, Default)]
#[non_exhaustive]
pub struct EnqueueOptions<'a> {
pub ordering_key: Option<&'a str>,
}
impl<'a> EnqueueOptions<'a> {
#[must_use]
pub const fn ordering_key(mut self, key: &'a str) -> Self {
self.ordering_key = Some(key);
self
}
}
#[non_exhaustive]
pub struct PostgresOutboxStore<
#[cfg(feature = "json")] Ser = JsonSerializer,
#[cfg(not(feature = "json"))] Ser,
> {
pool: PgPool,
settings: PostgresOutboxSettings,
serializer: Arc<Ser>,
}
impl<Ser> Clone for PostgresOutboxStore<Ser> {
fn clone(&self) -> Self {
Self {
pool: self.pool.clone(),
settings: self.settings.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("settings", &self.settings)
.finish_non_exhaustive()
}
}
struct SchemaCheck {
resolved_schema: Option<String>,
configured_exists: bool,
search_path: String,
}
async fn verify_schema(pool: &PgPool, schema: &str) -> Result<SchemaCheck, PostgresStoreError> {
let qualified = format!("{schema}.outbox");
let row = sqlx::query!(
r#"SELECT
current_setting('search_path') AS "search_path!",
(SELECT n.nspname
FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.oid = to_regclass('outbox')) AS resolved_schema,
(to_regclass($1) IS NOT NULL) AS "configured_exists!""#,
qualified,
)
.fetch_one(pool)
.await
.map_err(|err| {
if is_undefined_table(&err) {
PostgresStoreError::NotMigrated {
schema: schema.to_owned(),
}
} else {
PostgresStoreError::from(err)
}
})?;
Ok(SchemaCheck {
resolved_schema: row.resolved_schema,
configured_exists: row.configured_exists,
search_path: row.search_path,
})
}
async fn other_outbox_schemas(
pool: &PgPool,
schema: &str,
) -> Result<Vec<String>, PostgresStoreError> {
let schemas = sqlx::query_scalar!(
r#"SELECT n.nspname
FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relname = 'outbox' AND n.nspname <> $1"#,
schema,
)
.fetch_all(pool)
.await?;
Ok(schemas)
}
impl<Ser: Serializer + Send + Sync + 'static> PostgresOutboxStore<Ser> {
pub async fn connect(
pool: PgPool,
settings: PostgresOutboxSettings,
serializer: Ser,
) -> Result<Self, PostgresStoreError> {
if !crate::error::is_valid_schema_name(&settings.schema) {
return Err(PostgresStoreError::InvalidSchema {
schema: settings.schema,
});
}
let check = verify_schema(&pool, &settings.schema).await?;
let resolved_here = check.resolved_schema.as_deref() == Some(settings.schema.as_str());
if !resolved_here {
if !check.configured_exists {
return Err(PostgresStoreError::NotMigrated {
schema: settings.schema,
});
}
return Err(PostgresStoreError::SchemaResolution {
configured: settings.schema,
observed: check.search_path,
});
}
let others = other_outbox_schemas(&pool, &settings.schema).await?;
if !others.is_empty() {
tracing::warn!(
configured_schema = %settings.schema,
other_schemas = ?others,
"a table named `outbox` also exists outside the configured schema; \
an unqualified reference from another session could resolve to it"
);
}
Ok(Self {
pool,
settings,
serializer: Arc::new(serializer),
})
}
#[must_use]
pub fn content_type(&self) -> &ContentType {
self.serializer.content_type()
}
fn map_err(&self, err: sqlx::Error) -> PostgresStoreError {
map_operational_error(&self.settings.schema, err)
}
async fn set_local_timeout(
&self,
tx: &mut Transaction<'_, Postgres>,
) -> Result<(), PostgresStoreError> {
let timeout_ms = i64::try_from(self.settings.statement_timeout.as_millis())
.unwrap_or(i64::MAX)
.to_string();
sqlx::query_scalar!(
"SELECT set_config('statement_timeout', $1, true)",
timeout_ms
)
.fetch_one(&mut **tx)
.await
.map_err(|e| self.map_err(e))?;
Ok(())
}
pub async fn enqueue<T: Message>(
&self,
tx: &mut Transaction<'_, Postgres>,
envelope: &reliar_core::Envelope<T>,
) -> Result<MessageId, EnqueueError<Ser::Error>> {
self.enqueue_with(tx, envelope, EnqueueOptions::default())
.await
}
pub async fn enqueue_with<T: Message>(
&self,
tx: &mut Transaction<'_, Postgres>,
envelope: &reliar_core::Envelope<T>,
options: EnqueueOptions<'_>,
) -> Result<MessageId, EnqueueError<Ser::Error>> {
let payload = self
.serializer
.serialize(&envelope.body)
.map_err(|source| EnqueueError::Serialize { source })?;
let restore = if self.settings.enqueue_sets_search_path {
Some(set_search_path(tx, &self.settings.schema).await?)
} else {
None
};
let result = insert_row(tx, envelope, &payload, self.content_type(), options).await;
if result.is_ok()
&& let Some(previous) = restore
{
restore_search_path(tx, &previous).await?;
}
result.map_err(|source| map_enqueue_error(envelope.id, source))?;
Ok(envelope.id)
}
}
async fn set_search_path<E>(
tx: &mut Transaction<'_, Postgres>,
schema: &str,
) -> Result<String, EnqueueError<E>> {
let previous: String = sqlx::query_scalar!("SELECT current_setting('search_path')")
.fetch_one(&mut **tx)
.await
.map_err(|source| EnqueueError::Database { source })?
.unwrap_or_default();
let wanted = format!("{schema},public");
sqlx::query_scalar!("SELECT set_config('search_path', $1, true)", wanted)
.fetch_one(&mut **tx)
.await
.map_err(|source| EnqueueError::Database { source })?;
Ok(previous)
}
async fn restore_search_path<E>(
tx: &mut Transaction<'_, Postgres>,
previous: &str,
) -> Result<(), EnqueueError<E>> {
sqlx::query_scalar!("SELECT set_config('search_path', $1, true)", previous)
.fetch_one(&mut **tx)
.await
.map_err(|source| EnqueueError::Database { source })?;
Ok(())
}
async fn insert_row<T: Message>(
tx: &mut Transaction<'_, Postgres>,
envelope: &reliar_core::Envelope<T>,
payload: &Bytes,
content_type: &ContentType,
options: EnqueueOptions<'_>,
) -> 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 (
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(),
T::TYPE,
i32::from(T::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,
options.ordering_key,
metadata_json,
headers_json,
)
.execute(&mut **tx)
.await?;
Ok(())
}
#[cfg(feature = "json")]
impl PostgresOutboxStore<JsonSerializer> {
pub async fn new(pool: PgPool) -> Result<Self, PostgresStoreError> {
Self::connect(pool, PostgresOutboxSettings::default(), JsonSerializer).await
}
pub async fn with_settings(
pool: PgPool,
settings: PostgresOutboxSettings,
) -> Result<Self, PostgresStoreError> {
Self::connect(pool, settings, JsonSerializer).await
}
}
async fn poison_sweep_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
poisoned_ids: &[uuid::Uuid],
poisoned_errors: &[String],
worker: &str,
undecodable: &str,
) -> Result<(), sqlx::Error> {
sqlx::query!(
r#"UPDATE outbox o
SET dead_at = now(),
dead_reason = $4,
last_error = f.err,
locked_by = NULL,
locked_until = NULL,
updated_at = now()
FROM UNNEST($1::uuid[], $2::text[]) AS f(id, err)
WHERE o.id = f.id AND o.locked_by = $3"#,
poisoned_ids,
poisoned_errors,
worker,
undecodable,
)
.execute(executor)
.await?;
Ok(())
}
async fn purge_published_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
retention_ms: i64,
batch_size: i64,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"DELETE FROM outbox WHERE id IN (
SELECT id FROM outbox
WHERE published_at IS NOT NULL
AND published_at < now() - ($1::bigint * interval '1 millisecond')
LIMIT $2
)
AND published_at IS NOT NULL
AND published_at < now() - ($1::bigint * interval '1 millisecond')"#,
retention_ms,
batch_size,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn purge_dead_retention_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
retention_ms: i64,
batch_size: i64,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"DELETE FROM outbox WHERE id IN (
SELECT id FROM outbox
WHERE dead_at IS NOT NULL
AND dead_at < now() - ($1::bigint * interval '1 millisecond')
LIMIT $2
)
AND dead_at IS NOT NULL
AND dead_at < now() - ($1::bigint * interval '1 millisecond')"#,
retention_ms,
batch_size,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn purge_expired_sweep_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
batch_size: i64,
expired_reason: &str,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox
SET dead_at = now(),
dead_reason = $2,
last_error = 'reliar: expired before publication',
locked_by = NULL,
locked_until = NULL,
updated_at = now()
WHERE id IN (
SELECT id FROM outbox
WHERE expires_at IS NOT NULL AND expires_at < now()
AND published_at IS NULL AND dead_at IS NULL
AND (locked_until IS NULL OR locked_until < now())
LIMIT $1
)
AND published_at IS NULL AND dead_at IS NULL
AND (locked_until IS NULL OR locked_until < now())"#,
batch_size,
expired_reason,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn complete_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
worker: &str,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox
SET published_at = now(),
attempts = attempts + 1,
locked_by = NULL,
locked_until = NULL,
updated_at = now()
WHERE id = ANY($1) AND locked_by = $2"#,
ids,
worker,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn release_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
worker: &str,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox
SET locked_by = NULL,
locked_until = NULL,
updated_at = now()
WHERE id = ANY($1) AND locked_by = $2"#,
ids,
worker,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn extend_lease_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
lease_ms: i64,
worker: &str,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox
SET locked_until = now() + ($2::bigint * interval '1 millisecond'),
updated_at = now()
WHERE id = ANY($1) AND locked_by = $3"#,
ids,
lease_ms,
worker,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn fail_retry_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
errors: &[String],
delays_ms: &[i64],
worker: &str,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox o
SET attempts = o.attempts + 1,
last_error = f.err,
locked_by = NULL,
locked_until = NULL,
available_at = now() + (f.delay_ms * interval '1 millisecond'),
updated_at = now()
FROM UNNEST($1::uuid[], $2::text[], $3::bigint[]) AS f(id, err, delay_ms)
WHERE o.id = f.id AND o.locked_by = $4"#,
ids,
errors,
delays_ms,
worker,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn fail_dead_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
errors: &[String],
reasons: &[&str],
worker: &str,
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox o
SET attempts = o.attempts + 1,
last_error = f.err,
dead_at = now(),
dead_reason = f.reason,
locked_by = NULL,
locked_until = NULL,
updated_at = now()
FROM UNNEST($1::uuid[], $2::text[], $3::text[]) AS f(id, err, reason)
WHERE o.id = f.id AND o.locked_by = $4"#,
ids,
errors,
reasons as &[&str],
worker,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn claim_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
batch_size: i64,
worker: &str,
lease_ms: i64,
) -> Result<Vec<RawRow>, sqlx::Error> {
sqlx::query_as!(
RawRow,
r#"WITH claimed AS (
SELECT id FROM outbox
WHERE published_at IS NULL AND dead_at IS NULL
AND available_at <= now()
AND (locked_until IS NULL OR locked_until < now())
AND (expires_at IS NULL OR expires_at > now())
ORDER BY available_at, sequence
LIMIT $1
FOR UPDATE SKIP LOCKED
)
UPDATE outbox o
SET locked_by = $2,
locked_until = now() + ($3::bigint * interval '1 millisecond'),
updated_at = now()
FROM claimed
WHERE o.id = claimed.id
RETURNING o.id, o.sequence, o.message_type, o.message_version,
o.correlation_id, o.conversation_id, o.causation_id, o.request_id,
o.content_type, o.payload, o.tenant_id, o.expires_at, o.ordering_key,
o.metadata, o.headers, o.metadata_version,
o.created_at, o.available_at,
o.attempts, o.locked_by, o.locked_until,
o.published_at, o.dead_at, o.dead_reason, o.last_error"#,
batch_size,
worker,
lease_ms,
)
.fetch_all(executor)
.await
}
async fn list_dead_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
query: &DeadQuery,
limit: i64,
) -> Result<Vec<RawRow>, sqlx::Error> {
sqlx::query_as!(
RawRow,
r#"SELECT id, sequence, message_type, message_version,
correlation_id, conversation_id, causation_id, request_id,
content_type, payload, tenant_id, expires_at, ordering_key,
metadata, headers, metadata_version,
created_at, available_at,
attempts, locked_by, locked_until,
published_at, dead_at, dead_reason, last_error
FROM outbox
WHERE dead_at IS NOT NULL
AND ($1::text IS NULL OR message_type = $1)
AND ($2::text IS NULL OR tenant_id = $2)
AND ($3::timestamptz IS NULL OR dead_at < $3)
AND ($4::bigint IS NULL OR sequence > $4)
ORDER BY sequence ASC
LIMIT $5"#,
query.message_type,
query.tenant_id,
query.dead_before,
query.after_sequence,
limit,
)
.fetch_all(executor)
.await
}
async fn retry_dead_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
r#"UPDATE outbox
SET dead_at = NULL,
dead_reason = NULL,
available_at = now(),
attempts = 0,
locked_by = NULL,
locked_until = NULL,
updated_at = now()
WHERE id = ANY($1) AND dead_at IS NOT NULL"#,
ids,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
async fn purge_dead_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
ids: &[uuid::Uuid],
) -> Result<u64, sqlx::Error> {
let result = sqlx::query!(
"DELETE FROM outbox WHERE id = ANY($1) AND dead_at IS NOT NULL",
ids,
)
.execute(executor)
.await?;
Ok(result.rows_affected())
}
impl<Ser: Serializer + Send + Sync + 'static> OutboxStore for PostgresOutboxStore<Ser> {
type Error = PostgresStoreError;
async fn acquire(
&self,
request: reliar_outbox::AcquireRequest,
) -> Result<AcquiredBatch, Self::Error> {
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 = if self.settings.statement_timeout.is_zero() {
claim_rows(&self.pool, batch_size, worker, lease_ms)
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
let timeout_ms = i64::try_from(self.settings.statement_timeout.as_millis())
.unwrap_or(i64::MAX)
.to_string();
sqlx::query_scalar!(
"SELECT set_config('statement_timeout', $1, true)",
timeout_ms
)
.fetch_one(&mut *tx)
.await
.map_err(|e| self.map_err(e))?;
let rows = claim_rows(&mut *tx, batch_size, worker, lease_ms)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
rows
};
let mut records = Vec::with_capacity(rows.len());
let mut poisoned = Vec::new();
let mut poisoned_ids = Vec::new();
let mut poisoned_errors = Vec::new();
for raw in rows {
match decode_row(raw) {
Ok(record) => records.push(record),
Err(err) => {
poisoned_ids.push(err.id.as_uuid());
poisoned_errors.push(crate::records::truncate_last_error(err.detail.clone()));
poisoned.push(PoisonedRow::new(err.id, err.sequence, err.detail));
}
}
}
if !poisoned_ids.is_empty() {
let undecodable =
crate::records::encode_dead_reason(reliar_outbox::DeadReason::Undecodable);
if self.settings.statement_timeout.is_zero() {
poison_sweep_rows(
&self.pool,
&poisoned_ids,
&poisoned_errors,
worker,
undecodable,
)
.await
.map_err(|e| self.map_err(e))?;
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
poison_sweep_rows(
&mut *tx,
&poisoned_ids,
&poisoned_errors,
worker,
undecodable,
)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
}
}
Ok(AcquiredBatch::new(records, poisoned))
}
async fn complete(
&self,
worker: &WorkerId,
items: &[CompletedMessage],
) -> Result<u64, Self::Error> {
if items.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.message.id.as_uuid()).collect();
let affected = if self.settings.statement_timeout.is_zero() {
complete_rows(&self.pool, &ids, worker.as_str())
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let affected = complete_rows(&mut *tx, &ids, worker.as_str())
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
affected
};
log_shortfall("complete", items.len(), affected);
Ok(affected)
}
async fn fail(&self, worker: &WorkerId, items: &[FailedMessage]) -> Result<u64, Self::Error> {
if items.is_empty() {
return Ok(0);
}
let mut retry_ids = Vec::new();
let mut retry_errors = Vec::new();
let mut retry_delays = Vec::new();
let mut dead_ids = Vec::new();
let mut dead_errors = Vec::new();
let mut dead_reasons = Vec::new();
for item in items {
match item.outcome {
FailureOutcome::Retry { delay } => {
retry_ids.push(item.message.id.as_uuid());
retry_errors.push(item.error.clone());
retry_delays.push(i64::try_from(delay.as_millis()).unwrap_or(i64::MAX));
}
FailureOutcome::Dead { reason } => {
dead_ids.push(item.message.id.as_uuid());
dead_errors.push(item.error.clone());
dead_reasons.push(crate::records::encode_dead_reason(reason));
}
_ => tracing::error!(
id = %item.message.id,
"unrecognised FailureOutcome variant; row left as-is"
),
}
}
let affected = if self.settings.statement_timeout.is_zero() {
let mut affected = 0u64;
if !retry_ids.is_empty() {
affected += fail_retry_rows(
&self.pool,
&retry_ids,
&retry_errors,
&retry_delays,
worker.as_str(),
)
.await
.map_err(|e| self.map_err(e))?;
}
if !dead_ids.is_empty() {
affected += fail_dead_rows(
&self.pool,
&dead_ids,
&dead_errors,
&dead_reasons,
worker.as_str(),
)
.await
.map_err(|e| self.map_err(e))?;
}
affected
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let mut affected = 0u64;
if !retry_ids.is_empty() {
affected += fail_retry_rows(
&mut *tx,
&retry_ids,
&retry_errors,
&retry_delays,
worker.as_str(),
)
.await
.map_err(|e| self.map_err(e))?;
}
if !dead_ids.is_empty() {
affected += fail_dead_rows(
&mut *tx,
&dead_ids,
&dead_errors,
&dead_reasons,
worker.as_str(),
)
.await
.map_err(|e| self.map_err(e))?;
}
tx.commit().await.map_err(|e| self.map_err(e))?;
affected
};
log_shortfall("fail", items.len(), affected);
Ok(affected)
}
async fn release(&self, worker: &WorkerId, items: &[MessageRef]) -> Result<u64, Self::Error> {
if items.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
let affected = if self.settings.statement_timeout.is_zero() {
release_rows(&self.pool, &ids, worker.as_str())
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let affected = release_rows(&mut *tx, &ids, worker.as_str())
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
affected
};
log_shortfall("release", items.len(), affected);
Ok(affected)
}
async fn extend_lease(
&self,
worker: &WorkerId,
items: &[MessageRef],
lease: std::time::Duration,
) -> Result<u64, Self::Error> {
if items.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
let lease_ms = i64::try_from(lease.as_millis()).unwrap_or(i64::MAX);
let affected = if self.settings.statement_timeout.is_zero() {
extend_lease_rows(&self.pool, &ids, lease_ms, worker.as_str())
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let affected = extend_lease_rows(&mut *tx, &ids, lease_ms, worker.as_str())
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
affected
};
log_shortfall("extend_lease", items.len(), affected);
Ok(affected)
}
async fn purge(&self, request: PurgeRequest) -> Result<PurgeReport, Self::Error> {
let batch_size = i64::from(request.batch_size);
let expired_reason = crate::records::encode_dead_reason(reliar_outbox::DeadReason::Expired);
let (published_deleted, dead_deleted, expired_to_dead) =
if self.settings.statement_timeout.is_zero() {
let published_deleted = if let Some(retention) = request.published_retention {
let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
purge_published_rows(&self.pool, retention_ms, batch_size)
.await
.map_err(|e| self.map_err(e))?
} else {
0
};
let dead_deleted = if let Some(retention) = request.dead_retention {
let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
purge_dead_retention_rows(&self.pool, retention_ms, batch_size)
.await
.map_err(|e| self.map_err(e))?
} else {
0
};
let expired_to_dead =
purge_expired_sweep_rows(&self.pool, batch_size, expired_reason)
.await
.map_err(|e| self.map_err(e))?;
(published_deleted, dead_deleted, expired_to_dead)
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let published_deleted = if let Some(retention) = request.published_retention {
let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
purge_published_rows(&mut *tx, retention_ms, batch_size)
.await
.map_err(|e| self.map_err(e))?
} else {
0
};
let dead_deleted = if let Some(retention) = request.dead_retention {
let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
purge_dead_retention_rows(&mut *tx, retention_ms, batch_size)
.await
.map_err(|e| self.map_err(e))?
} else {
0
};
let expired_to_dead =
purge_expired_sweep_rows(&mut *tx, batch_size, expired_reason)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
(published_deleted, dead_deleted, expired_to_dead)
};
Ok(PurgeReport::new(
published_deleted,
dead_deleted,
expired_to_dead,
))
}
async fn stats(&self) -> Result<OutboxStats, Self::Error> {
if self.settings.statement_timeout.is_zero() {
let row = sqlx::query!(
r#"SELECT
count(*) FILTER (
WHERE published_at IS NULL AND dead_at IS NULL
AND available_at <= now()
AND (locked_until IS NULL OR locked_until < now())
AND (expires_at IS NULL OR expires_at > now())
) AS "pending!",
count(*) FILTER (WHERE dead_at IS NOT NULL) AS "dead!",
count(*) FILTER (
WHERE published_at IS NULL AND dead_at IS NULL
AND expires_at IS NOT NULL AND expires_at < now()
) AS "expired_pending!",
min(available_at) FILTER (
WHERE published_at IS NULL AND dead_at IS NULL
AND available_at <= now()
AND (locked_until IS NULL OR locked_until < now())
AND (expires_at IS NULL OR expires_at > now())
) AS oldest_pending_available_at,
now() AS "as_of!"
FROM outbox"#
)
.fetch_one(&self.pool)
.await
.map_err(|e| self.map_err(e))?;
return 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,
));
}
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let row = sqlx::query!(
r#"SELECT
count(*) FILTER (
WHERE published_at IS NULL AND dead_at IS NULL
AND available_at <= now()
AND (locked_until IS NULL OR locked_until < now())
AND (expires_at IS NULL OR expires_at > now())
) AS "pending!",
count(*) FILTER (WHERE dead_at IS NOT NULL) AS "dead!",
count(*) FILTER (
WHERE published_at IS NULL AND dead_at IS NULL
AND expires_at IS NOT NULL AND expires_at < now()
) AS "expired_pending!",
min(available_at) FILTER (
WHERE published_at IS NULL AND dead_at IS NULL
AND available_at <= now()
AND (locked_until IS NULL OR locked_until < now())
AND (expires_at IS NULL OR expires_at > now())
) AS oldest_pending_available_at,
now() AS "as_of!"
FROM outbox"#
)
.fetch_one(&mut *tx)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(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,
))
}
}
impl<Ser: Serializer + Send + Sync + 'static> OutboxDeadLetters for PostgresOutboxStore<Ser> {
type Error = PostgresStoreError;
async fn list_dead(&self, query: DeadQuery) -> Result<DeadLetterPage, Self::Error> {
let capped_limit = query.limit.min(MAX_LIST_DEAD_LIMIT);
let limit = i64::from(capped_limit);
let rows = if self.settings.statement_timeout.is_zero() {
list_dead_rows(&self.pool, &query, limit)
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let rows = list_dead_rows(&mut *tx, &query, limit)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
rows
};
let scanned = rows.len();
let mut records = Vec::with_capacity(scanned);
let mut poisoned = Vec::new();
let mut max_sequence: Option<i64> = None;
for raw in rows {
max_sequence = Some(max_sequence.map_or(raw.sequence, |m| m.max(raw.sequence)));
match decode_row(raw) {
Ok(record) => records.push(record),
Err(err) => poisoned.push(PoisonedRow::new(err.id, err.sequence, err.detail)),
}
}
let next_after_sequence = if scanned == capped_limit as usize {
max_sequence
} else {
None
};
Ok(DeadLetterPage::new(records, poisoned, next_after_sequence))
}
async fn retry_dead(&self, refs: &[MessageRef]) -> Result<u64, Self::Error> {
if refs.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = refs.iter().map(|r| r.id.as_uuid()).collect();
let affected = if self.settings.statement_timeout.is_zero() {
retry_dead_rows(&self.pool, &ids)
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let affected = retry_dead_rows(&mut *tx, &ids)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
affected
};
Ok(affected)
}
async fn purge_dead(&self, refs: &[MessageRef]) -> Result<u64, Self::Error> {
if refs.is_empty() {
return Ok(0);
}
let ids: Vec<uuid::Uuid> = refs.iter().map(|r| r.id.as_uuid()).collect();
let affected = if self.settings.statement_timeout.is_zero() {
purge_dead_rows(&self.pool, &ids)
.await
.map_err(|e| self.map_err(e))?
} else {
let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
self.set_local_timeout(&mut tx).await?;
let affected = purge_dead_rows(&mut *tx, &ids)
.await
.map_err(|e| self.map_err(e))?;
tx.commit().await.map_err(|e| self.map_err(e))?;
affected
};
Ok(affected)
}
}
fn log_shortfall(operation: &'static str, claimed: usize, affected: u64) {
let claimed = claimed as u64;
if affected < claimed {
tracing::debug!(
operation,
claimed,
affected,
"fewer rows affected than claimed"
);
}
}