reliar_store_postgres/store/
dead_letters.rs1use 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
12const MAX_LIST_DEAD_LIMIT: u32 = 1000;
16
17impl<Ser: Serializer + Send + Sync + 'static> OutboxDeadLetters for PostgresOutboxStore<Ser> {
18 type Error = PostgresStoreError;
19
20 async fn list_dead(&self, query: DeadQuery) -> Result<DeadLetterPage, Self::Error> {
27 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 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 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 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
129async 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
164async 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
189async 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}