Skip to main content

reliar_store_postgres/store/
dead_letters.rs

1//! [`PostgresOutboxStore`]'s [`OutboxDeadLetters`] implementation: `list_dead`/`retry_dead`/
2//! `purge_dead`.
3
4use reliar_core::Serializer;
5use reliar_outbox::{DeadLetterPage, DeadQuery, MessageRef, OutboxDeadLetters, PoisonedRow};
6
7use crate::error::PostgresStoreError;
8use crate::records::{RawRow, decode_row};
9
10use super::PostgresOutboxStore;
11
12/// The largest `DeadQuery::limit` [`OutboxDeadLetters::list_dead`] honours — a caller-supplied
13/// value above this is silently capped, never sent to the database: this store is
14/// provider-capped, with a default of 100.
15const MAX_LIST_DEAD_LIMIT: u32 = 1000;
16
17impl<Ser: Serializer + Send + Sync + 'static> OutboxDeadLetters for PostgresOutboxStore<Ser> {
18    type Error = PostgresStoreError;
19
20    /// **`ORDER BY sequence ASC` is normative**: `after_sequence` is a keyset
21    /// cursor over `sequence`, the column `ix_outbox_dead` orders by; `message_type`/
22    /// `tenant_id`/`dead_before` are filters only, expressed as `($n::type IS NULL OR ...)` so
23    /// one static statement serves every combination. The cursor returned is the largest
24    /// `sequence` **scanned**, poisoned rows included, so a poisoned tail cannot loop the
25    /// caller forever.
26    async fn list_dead(&self, query: DeadQuery) -> Result<DeadLetterPage, Self::Error> {
27        // Provider-capped, default 100: a caller-supplied
28        // limit above this never reaches the database, regardless of what `DeadQuery` carries.
29        let capped_limit = query.limit.min(MAX_LIST_DEAD_LIMIT);
30        let limit = i64::from(capped_limit);
31
32        let rows = if self.settings.statement_timeout.is_zero() {
33            list_dead_rows(&self.pool, &query, limit)
34                .await
35                .map_err(|e| self.map_err(e))?
36        } else {
37            let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
38
39            self.set_local_timeout(&mut tx).await?;
40            let rows = list_dead_rows(&mut *tx, &query, limit)
41                .await
42                .map_err(|e| self.map_err(e))?;
43            tx.commit().await.map_err(|e| self.map_err(e))?;
44
45            rows
46        };
47
48        let scanned = rows.len();
49        let mut records = Vec::with_capacity(scanned);
50        let mut poisoned = Vec::new();
51        let mut max_sequence: Option<i64> = None;
52
53        for raw in rows {
54            max_sequence = Some(max_sequence.map_or(raw.sequence, |m| m.max(raw.sequence)));
55
56            match decode_row(raw) {
57                Ok(record) => records.push(record),
58                Err(err) => poisoned.push(PoisonedRow::new(err.id, err.sequence, err.detail)),
59            }
60        }
61
62        // "Full" is scanned == limit, poisoned rows included — they occupy a row in the scan,
63        // so counting only decoded records would stop pagination early on a poisoned tail.
64        let next_after_sequence = if scanned == capped_limit as usize {
65            max_sequence
66        } else {
67            None
68        };
69
70        Ok(DeadLetterPage::new(records, poisoned, next_after_sequence))
71    }
72
73    /// Returns dead rows to pending: clears the lease that already isn't there, resets
74    /// `attempts` to 0 (the **only** operation that does), keeps `last_error` for audit. Not
75    /// worker-guarded — a dead row holds no lease, so there is no owner to check against.
76    async fn retry_dead(&self, refs: &[MessageRef]) -> Result<u64, Self::Error> {
77        if refs.is_empty() {
78            return Ok(0);
79        }
80
81        let ids: Vec<uuid::Uuid> = refs.iter().map(|r| r.id.as_uuid()).collect();
82        let affected = if self.settings.statement_timeout.is_zero() {
83            retry_dead_rows(&self.pool, &ids)
84                .await
85                .map_err(|e| self.map_err(e))?
86        } else {
87            let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
88
89            self.set_local_timeout(&mut tx).await?;
90            let affected = retry_dead_rows(&mut *tx, &ids)
91                .await
92                .map_err(|e| self.map_err(e))?;
93            tx.commit().await.map_err(|e| self.map_err(e))?;
94
95            affected
96        };
97
98        Ok(affected)
99    }
100
101    /// Deletes dead rows by reference, regardless of
102    /// [`PurgeRequest::dead_retention`](reliar_outbox::PurgeRequest::dead_retention).
103    async fn purge_dead(&self, refs: &[MessageRef]) -> Result<u64, Self::Error> {
104        if refs.is_empty() {
105            return Ok(0);
106        }
107
108        let ids: Vec<uuid::Uuid> = refs.iter().map(|r| r.id.as_uuid()).collect();
109        let affected = if self.settings.statement_timeout.is_zero() {
110            purge_dead_rows(&self.pool, &ids)
111                .await
112                .map_err(|e| self.map_err(e))?
113        } else {
114            let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
115
116            self.set_local_timeout(&mut tx).await?;
117            let affected = purge_dead_rows(&mut *tx, &ids)
118                .await
119                .map_err(|e| self.map_err(e))?;
120            tx.commit().await.map_err(|e| self.map_err(e))?;
121
122            affected
123        };
124
125        Ok(affected)
126    }
127}
128
129/// `list_dead`'s query, shared by the plain-pool and `statement_timeout`-wrapped-transaction
130/// call sites. Named against [`RawRow`] via `query_as!`, same as the claim — the `SELECT`
131/// list matches its field order exactly.
132async fn list_dead_rows<'e>(
133    executor: impl sqlx::PgExecutor<'e>,
134    query: &DeadQuery,
135    limit: i64,
136) -> Result<Vec<RawRow>, sqlx::Error> {
137    sqlx::query_as!(
138        RawRow,
139        r#"SELECT id, sequence, message_type, message_version,
140                  correlation_id, conversation_id, causation_id, request_id,
141                  content_type, payload, tenant_id, expires_at, ordering_key,
142                  metadata, headers, metadata_version,
143                  created_at, available_at,
144                  attempts, locked_by, locked_until,
145                  published_at, dead_at, dead_reason, last_error
146             FROM outbox
147            WHERE dead_at IS NOT NULL
148              AND ($1::text IS NULL OR message_type = $1)
149              AND ($2::text IS NULL OR tenant_id = $2)
150              AND ($3::timestamptz IS NULL OR dead_at < $3)
151              AND ($4::bigint IS NULL OR sequence > $4)
152            ORDER BY sequence ASC
153            LIMIT $5"#,
154        query.message_type,
155        query.tenant_id,
156        query.dead_before,
157        query.after_sequence,
158        limit,
159    )
160    .fetch_all(executor)
161    .await
162}
163
164/// `retry_dead`'s query, shared by the plain-pool and `statement_timeout`-wrapped-transaction
165/// call sites. Not worker-guarded — a dead row holds no lease, so there is no owner to check
166/// against.
167async fn retry_dead_rows<'e>(
168    executor: impl sqlx::PgExecutor<'e>,
169    ids: &[uuid::Uuid],
170) -> Result<u64, sqlx::Error> {
171    let result = sqlx::query!(
172        r#"UPDATE outbox
173              SET dead_at      = NULL,
174                  dead_reason  = NULL,
175                  available_at = now(),
176                  attempts     = 0,
177                  locked_by    = NULL,
178                  locked_until = NULL,
179                  updated_at   = now()
180            WHERE id = ANY($1) AND dead_at IS NOT NULL"#,
181        ids,
182    )
183    .execute(executor)
184    .await?;
185
186    Ok(result.rows_affected())
187}
188
189/// `purge_dead`'s query, shared by the plain-pool and `statement_timeout`-wrapped-transaction
190/// call sites.
191async fn purge_dead_rows<'e>(
192    executor: impl sqlx::PgExecutor<'e>,
193    ids: &[uuid::Uuid],
194) -> Result<u64, sqlx::Error> {
195    let result = sqlx::query!(
196        "DELETE FROM outbox WHERE id = ANY($1) AND dead_at IS NOT NULL",
197        ids,
198    )
199    .execute(executor)
200    .await?;
201
202    Ok(result.rows_affected())
203}