Skip to main content

reliar_store_postgres/outbox/
dead_letters.rs

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