use reliar_core::Serializer;
use reliar_outbox::{DeadLetterPage, DeadQuery, MessageRef, OutboxDeadLetters, PoisonedRow};
use crate::error::PostgresStoreError;
use crate::records::{RawRow, decode_row};
use super::PostgresOutboxStore;
const MAX_LIST_DEAD_LIMIT: u32 = 1000;
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)
}
}
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())
}