use sqlx::PgConnection;
use reliar_core::Serializer;
use reliar_outbox::{
DeadCursor, DeadLetterPage, DeadQuery, OutboxDeadLetters, OutboxRecordId, PoisonedRow,
RecordRef,
};
use super::dead_letters as repo;
use super::error::PostgresOutboxError;
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 = PostgresOutboxError;
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 (after_dead_at, after_id) = query.after.map_or((None, None), |cursor| {
(Some(cursor.dead_at()), Some(cursor.id().as_uuid()))
});
let rows: Vec<RawRow> = self
.session
.run(async |conn: &mut PgConnection| {
repo::list_dead_rows(
&mut *conn,
repo::ListDeadRowsParams {
message_type: query.message_type.as_deref(),
tenant_id: query.tenant_id.as_deref(),
dead_before: query.dead_before,
after_dead_at,
after_id,
limit,
},
)
.await
})
.await
.map_err(|e| self.session.map_err::<PostgresOutboxError>(e))?;
let scanned = rows.len();
let mut records = Vec::with_capacity(scanned);
let mut poisoned = Vec::new();
let mut last_cursor: Option<DeadCursor> = None;
for raw in rows {
debug_assert!(raw.dead_at.is_some(), "list_dead selects only dead rows");
if let Some(dead_at) = raw.dead_at {
last_cursor = Some(DeadCursor::new(dead_at, OutboxRecordId::from_uuid(raw.id)));
}
match decode_row(raw) {
Ok(record) => records.push(record),
Err(err) => poisoned.push(PoisonedRow::new(err.id, err.message_id, err.detail)),
}
}
let next_after = if scanned == capped_limit as usize {
last_cursor
} else {
None
};
Ok(DeadLetterPage::new(records, poisoned, next_after))
}
async fn retry_dead(&self, refs: &[RecordRef]) -> 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 = self
.session
.run(async |conn: &mut PgConnection| {
repo::retry_dead_rows(&mut *conn, repo::RetryDeadRowsParams { ids: &ids }).await
})
.await
.map_err(|e| self.session.map_err::<PostgresOutboxError>(e))?;
Ok(affected)
}
async fn purge_dead(&self, refs: &[RecordRef]) -> 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 = self
.session
.run(async |conn: &mut PgConnection| {
repo::purge_dead_rows(&mut *conn, repo::PurgeDeadRowsParams { ids: &ids }).await
})
.await
.map_err(|e| self.session.map_err::<PostgresOutboxError>(e))?;
Ok(affected)
}
}