reliar_store_postgres/outbox/
dead_letters.rs1use 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
15const MAX_LIST_DEAD_LIMIT: u32 = 1000;
19
20impl<Ser: Serializer + Send + Sync + 'static> OutboxDeadLetters for PostgresOutboxStore<Ser> {
21 type Error = PostgresOutboxError;
22
23 async fn list_dead(&self, query: DeadQuery) -> Result<DeadLetterPage, Self::Error> {
28 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 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 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 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
134async 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
174async 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
199async 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}