1use std::sync::Arc;
5
6use bytes::Bytes;
7use reliar_core::{ContentType, Message, MessageId, Serializer};
8use reliar_outbox::{
9 AcquiredBatch, CompletedMessage, DeadLetterPage, DeadQuery, FailedMessage, FailureOutcome,
10 MessageRef, OutboxDeadLetters, OutboxStats, OutboxStore, PoisonedRow, PurgeReport,
11 PurgeRequest, WorkerId,
12};
13use sqlx::{PgPool, Postgres, Transaction};
14
15use crate::error::{
16 EnqueueError, PostgresStoreError, is_undefined_table, map_enqueue_error, map_operational_error,
17};
18use crate::records::{RawRow, decode_row};
19use crate::settings::PostgresOutboxSettings;
20
21#[cfg(feature = "json")]
22use reliar_core::JsonSerializer;
23
24const MAX_LIST_DEAD_LIMIT: u32 = 1000;
28
29#[derive(Clone, Debug, Default)]
32#[non_exhaustive]
33pub struct EnqueueOptions<'a> {
34 pub ordering_key: Option<&'a str>,
37}
38
39impl<'a> EnqueueOptions<'a> {
40 #[must_use]
43 pub const fn ordering_key(mut self, key: &'a str) -> Self {
44 self.ordering_key = Some(key);
45 self
46 }
47}
48
49#[non_exhaustive]
57pub struct PostgresOutboxStore<
58 #[cfg(feature = "json")] Ser = JsonSerializer,
59 #[cfg(not(feature = "json"))] Ser,
60> {
61 pool: PgPool,
62 settings: PostgresOutboxSettings,
63 serializer: Arc<Ser>,
64}
65
66impl<Ser> Clone for PostgresOutboxStore<Ser> {
70 fn clone(&self) -> Self {
71 Self {
72 pool: self.pool.clone(),
73 settings: self.settings.clone(),
74 serializer: Arc::clone(&self.serializer),
75 }
76 }
77}
78
79impl<Ser> std::fmt::Debug for PostgresOutboxStore<Ser> {
80 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
81 f.debug_struct("PostgresOutboxStore")
82 .field("settings", &self.settings)
83 .finish_non_exhaustive()
84 }
85}
86
87struct SchemaCheck {
89 resolved_schema: Option<String>,
90 configured_exists: bool,
91 search_path: String,
92}
93
94async fn verify_schema(pool: &PgPool, schema: &str) -> Result<SchemaCheck, PostgresStoreError> {
95 let qualified = format!("{schema}.outbox");
96 let row = sqlx::query!(
97 r#"SELECT
98 current_setting('search_path') AS "search_path!",
99 (SELECT n.nspname
100 FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
101 WHERE c.oid = to_regclass('outbox')) AS resolved_schema,
102 (to_regclass($1) IS NOT NULL) AS "configured_exists!""#,
103 qualified,
104 )
105 .fetch_one(pool)
106 .await
107 .map_err(|err| {
108 if is_undefined_table(&err) {
109 PostgresStoreError::NotMigrated {
110 schema: schema.to_owned(),
111 }
112 } else {
113 PostgresStoreError::from(err)
114 }
115 })?;
116
117 Ok(SchemaCheck {
118 resolved_schema: row.resolved_schema,
119 configured_exists: row.configured_exists,
120 search_path: row.search_path,
121 })
122}
123
124async fn other_outbox_schemas(
127 pool: &PgPool,
128 schema: &str,
129) -> Result<Vec<String>, PostgresStoreError> {
130 let schemas = sqlx::query_scalar!(
131 r#"SELECT n.nspname
132 FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
133 WHERE c.relname = 'outbox' AND n.nspname <> $1"#,
134 schema,
135 )
136 .fetch_all(pool)
137 .await?;
138 Ok(schemas)
139}
140
141impl<Ser: Serializer + Send + Sync + 'static> PostgresOutboxStore<Ser> {
142 pub async fn connect(
154 pool: PgPool,
155 settings: PostgresOutboxSettings,
156 serializer: Ser,
157 ) -> Result<Self, PostgresStoreError> {
158 if !crate::error::is_valid_schema_name(&settings.schema) {
159 return Err(PostgresStoreError::InvalidSchema {
160 schema: settings.schema,
161 });
162 }
163
164 let check = verify_schema(&pool, &settings.schema).await?;
165
166 let resolved_here = check.resolved_schema.as_deref() == Some(settings.schema.as_str());
167 if !resolved_here {
168 if !check.configured_exists {
169 return Err(PostgresStoreError::NotMigrated {
170 schema: settings.schema,
171 });
172 }
173 return Err(PostgresStoreError::SchemaResolution {
174 configured: settings.schema,
175 observed: check.search_path,
176 });
177 }
178
179 let others = other_outbox_schemas(&pool, &settings.schema).await?;
180 if !others.is_empty() {
181 tracing::warn!(
182 configured_schema = %settings.schema,
183 other_schemas = ?others,
184 "a table named `outbox` also exists outside the configured schema; \
185 an unqualified reference from another session could resolve to it"
186 );
187 }
188
189 Ok(Self {
190 pool,
191 settings,
192 serializer: Arc::new(serializer),
193 })
194 }
195
196 #[must_use]
201 pub fn content_type(&self) -> &ContentType {
202 self.serializer.content_type()
203 }
204
205 fn map_err(&self, err: sqlx::Error) -> PostgresStoreError {
209 map_operational_error(&self.settings.schema, err)
210 }
211
212 async fn set_local_timeout(
215 &self,
216 tx: &mut Transaction<'_, Postgres>,
217 ) -> Result<(), PostgresStoreError> {
218 let timeout_ms = i64::try_from(self.settings.statement_timeout.as_millis())
219 .unwrap_or(i64::MAX)
220 .to_string();
221 sqlx::query_scalar!(
222 "SELECT set_config('statement_timeout', $1, true)",
223 timeout_ms
224 )
225 .fetch_one(&mut **tx)
226 .await
227 .map_err(|e| self.map_err(e))?;
228 Ok(())
229 }
230
231 pub async fn enqueue<T: Message>(
242 &self,
243 tx: &mut Transaction<'_, Postgres>,
244 envelope: &reliar_core::Envelope<T>,
245 ) -> Result<MessageId, EnqueueError<Ser::Error>> {
246 self.enqueue_with(tx, envelope, EnqueueOptions::default())
247 .await
248 }
249
250 pub async fn enqueue_with<T: Message>(
257 &self,
258 tx: &mut Transaction<'_, Postgres>,
259 envelope: &reliar_core::Envelope<T>,
260 options: EnqueueOptions<'_>,
261 ) -> Result<MessageId, EnqueueError<Ser::Error>> {
262 let payload = self
263 .serializer
264 .serialize(&envelope.body)
265 .map_err(|source| EnqueueError::Serialize { source })?;
266
267 let restore = if self.settings.enqueue_sets_search_path {
268 Some(set_search_path(tx, &self.settings.schema).await?)
269 } else {
270 None
271 };
272
273 let result = insert_row(tx, envelope, &payload, self.content_type(), options).await;
274
275 if result.is_ok()
281 && let Some(previous) = restore
282 {
283 restore_search_path(tx, &previous).await?;
284 }
285
286 result.map_err(|source| map_enqueue_error(envelope.id, source))?;
287 Ok(envelope.id)
288 }
289}
290
291async fn set_search_path<E>(
295 tx: &mut Transaction<'_, Postgres>,
296 schema: &str,
297) -> Result<String, EnqueueError<E>> {
298 let previous: String = sqlx::query_scalar!("SELECT current_setting('search_path')")
299 .fetch_one(&mut **tx)
300 .await
301 .map_err(|source| EnqueueError::Database { source })?
302 .unwrap_or_default();
303 let wanted = format!("{schema},public");
304 sqlx::query_scalar!("SELECT set_config('search_path', $1, true)", wanted)
305 .fetch_one(&mut **tx)
306 .await
307 .map_err(|source| EnqueueError::Database { source })?;
308 Ok(previous)
309}
310
311async fn restore_search_path<E>(
312 tx: &mut Transaction<'_, Postgres>,
313 previous: &str,
314) -> Result<(), EnqueueError<E>> {
315 sqlx::query_scalar!("SELECT set_config('search_path', $1, true)", previous)
316 .fetch_one(&mut **tx)
317 .await
318 .map_err(|source| EnqueueError::Database { source })?;
319 Ok(())
320}
321
322async fn insert_row<T: Message>(
323 tx: &mut Transaction<'_, Postgres>,
324 envelope: &reliar_core::Envelope<T>,
325 payload: &Bytes,
326 content_type: &ContentType,
327 options: EnqueueOptions<'_>,
328) -> Result<(), sqlx::Error> {
329 let corr = &envelope.metadata.correlation;
330 let sent_at_ms = envelope
331 .metadata
332 .delivery
333 .sent_at
334 .map(crate::records::encode_epoch_millis);
335 let rest = crate::records::MetadataRest {
336 trace: crate::records::TraceRest {
337 traceparent: envelope.metadata.trace.traceparent.clone(),
338 tracestate: envelope.metadata.trace.tracestate.clone(),
339 },
340 routing: crate::records::RoutingRest {
341 source: envelope
342 .metadata
343 .routing
344 .source
345 .as_ref()
346 .map(|v| v.as_str().to_owned()),
347 destination: envelope
348 .metadata
349 .routing
350 .destination
351 .as_ref()
352 .map(|v| v.as_str().to_owned()),
353 reply_to: envelope
354 .metadata
355 .routing
356 .reply_to
357 .as_ref()
358 .map(|v| v.as_str().to_owned()),
359 },
360 delivery: crate::records::DeliveryRest {
361 sent_at_ms,
362 deduplication_id: envelope.metadata.delivery.deduplication_id.clone(),
363 },
364 };
365 let metadata_json = if rest.trace.traceparent.is_none()
367 && rest.trace.tracestate.is_none()
368 && rest.routing.source.is_none()
369 && rest.routing.destination.is_none()
370 && rest.routing.reply_to.is_none()
371 && rest.delivery.sent_at_ms.is_none()
372 && rest.delivery.deduplication_id.is_none()
373 {
374 None
375 } else {
376 serde_json::to_value(&rest).ok()
384 };
385
386 let headers_json = envelope.headers().filter(|h| !h.is_empty()).map(|h| {
387 let map: serde_json::Map<String, serde_json::Value> = h
388 .iter()
389 .map(|(k, v)| (k.to_owned(), serde_json::Value::String(v.to_owned())))
390 .collect();
391 serde_json::Value::Object(map)
392 });
393
394 sqlx::query!(
395 r#"INSERT INTO outbox (
396 id, message_type, message_version,
397 correlation_id, conversation_id, causation_id, request_id,
398 content_type, payload, tenant_id, expires_at, ordering_key,
399 metadata, headers, available_at
400 ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14, now())"#,
401 envelope.id.as_uuid(),
402 T::TYPE,
403 i32::from(T::VERSION),
404 corr.correlation_id
405 .as_ref()
406 .map(reliar_core::CorrelationId::as_str),
407 corr.conversation_id.as_uuid(),
408 corr.causation_id.map(|id| id.as_uuid()),
409 corr.request_id.map(|id| id.as_uuid()),
410 content_type.as_str(),
411 &payload[..],
412 envelope.metadata.tenant_id.as_deref(),
413 envelope.metadata.delivery.expires_at,
414 options.ordering_key,
415 metadata_json,
416 headers_json,
417 )
418 .execute(&mut **tx)
419 .await?;
420 Ok(())
421}
422
423#[cfg(feature = "json")]
424impl PostgresOutboxStore<JsonSerializer> {
425 pub async fn new(pool: PgPool) -> Result<Self, PostgresStoreError> {
431 Self::connect(pool, PostgresOutboxSettings::default(), JsonSerializer).await
432 }
433
434 pub async fn with_settings(
441 pool: PgPool,
442 settings: PostgresOutboxSettings,
443 ) -> Result<Self, PostgresStoreError> {
444 Self::connect(pool, settings, JsonSerializer).await
445 }
446}
447
448async fn poison_sweep_rows<'e>(
453 executor: impl sqlx::PgExecutor<'e>,
454 poisoned_ids: &[uuid::Uuid],
455 poisoned_errors: &[String],
456 worker: &str,
457 undecodable: &str,
458) -> Result<(), sqlx::Error> {
459 sqlx::query!(
460 r#"UPDATE outbox o
461 SET dead_at = now(),
462 dead_reason = $4,
463 last_error = f.err,
464 locked_by = NULL,
465 locked_until = NULL,
466 updated_at = now()
467 FROM UNNEST($1::uuid[], $2::text[]) AS f(id, err)
468 WHERE o.id = f.id AND o.locked_by = $3"#,
469 poisoned_ids,
470 poisoned_errors,
471 worker,
472 undecodable,
473 )
474 .execute(executor)
475 .await?;
476 Ok(())
477}
478
479async fn purge_published_rows<'e>(
487 executor: impl sqlx::PgExecutor<'e>,
488 retention_ms: i64,
489 batch_size: i64,
490) -> Result<u64, sqlx::Error> {
491 let result = sqlx::query!(
492 r#"DELETE FROM outbox WHERE id IN (
493 SELECT id FROM outbox
494 WHERE published_at IS NOT NULL
495 AND published_at < now() - ($1::bigint * interval '1 millisecond')
496 LIMIT $2
497 )
498 AND published_at IS NOT NULL
499 AND published_at < now() - ($1::bigint * interval '1 millisecond')"#,
500 retention_ms,
501 batch_size,
502 )
503 .execute(executor)
504 .await?;
505 Ok(result.rows_affected())
506}
507
508async fn purge_dead_retention_rows<'e>(
514 executor: impl sqlx::PgExecutor<'e>,
515 retention_ms: i64,
516 batch_size: i64,
517) -> Result<u64, sqlx::Error> {
518 let result = sqlx::query!(
519 r#"DELETE FROM outbox WHERE id IN (
520 SELECT id FROM outbox
521 WHERE dead_at IS NOT NULL
522 AND dead_at < now() - ($1::bigint * interval '1 millisecond')
523 LIMIT $2
524 )
525 AND dead_at IS NOT NULL
526 AND dead_at < now() - ($1::bigint * interval '1 millisecond')"#,
527 retention_ms,
528 batch_size,
529 )
530 .execute(executor)
531 .await?;
532 Ok(result.rows_affected())
533}
534
535async fn purge_expired_sweep_rows<'e>(
541 executor: impl sqlx::PgExecutor<'e>,
542 batch_size: i64,
543 expired_reason: &str,
544) -> Result<u64, sqlx::Error> {
545 let result = sqlx::query!(
546 r#"UPDATE outbox
547 SET dead_at = now(),
548 dead_reason = $2,
549 last_error = 'reliar: expired before publication',
550 locked_by = NULL,
551 locked_until = NULL,
552 updated_at = now()
553 WHERE id IN (
554 SELECT id FROM outbox
555 WHERE expires_at IS NOT NULL AND expires_at < now()
556 AND published_at IS NULL AND dead_at IS NULL
557 AND (locked_until IS NULL OR locked_until < now())
558 LIMIT $1
559 )
560 AND published_at IS NULL AND dead_at IS NULL
561 AND (locked_until IS NULL OR locked_until < now())"#,
562 batch_size,
563 expired_reason,
564 )
565 .execute(executor)
566 .await?;
567 Ok(result.rows_affected())
568}
569
570async fn complete_rows<'e>(
574 executor: impl sqlx::PgExecutor<'e>,
575 ids: &[uuid::Uuid],
576 worker: &str,
577) -> Result<u64, sqlx::Error> {
578 let result = sqlx::query!(
579 r#"UPDATE outbox
580 SET published_at = now(),
581 attempts = attempts + 1,
582 locked_by = NULL,
583 locked_until = NULL,
584 updated_at = now()
585 WHERE id = ANY($1) AND locked_by = $2"#,
586 ids,
587 worker,
588 )
589 .execute(executor)
590 .await?;
591 Ok(result.rows_affected())
592}
593
594async fn release_rows<'e>(
595 executor: impl sqlx::PgExecutor<'e>,
596 ids: &[uuid::Uuid],
597 worker: &str,
598) -> Result<u64, sqlx::Error> {
599 let result = sqlx::query!(
600 r#"UPDATE outbox
601 SET locked_by = NULL,
602 locked_until = NULL,
603 updated_at = now()
604 WHERE id = ANY($1) AND locked_by = $2"#,
605 ids,
606 worker,
607 )
608 .execute(executor)
609 .await?;
610 Ok(result.rows_affected())
611}
612
613async fn extend_lease_rows<'e>(
614 executor: impl sqlx::PgExecutor<'e>,
615 ids: &[uuid::Uuid],
616 lease_ms: i64,
617 worker: &str,
618) -> Result<u64, sqlx::Error> {
619 let result = sqlx::query!(
620 r#"UPDATE outbox
621 SET locked_until = now() + ($2::bigint * interval '1 millisecond'),
622 updated_at = now()
623 WHERE id = ANY($1) AND locked_by = $3"#,
624 ids,
625 lease_ms,
626 worker,
627 )
628 .execute(executor)
629 .await?;
630 Ok(result.rows_affected())
631}
632
633async fn fail_retry_rows<'e>(
634 executor: impl sqlx::PgExecutor<'e>,
635 ids: &[uuid::Uuid],
636 errors: &[String],
637 delays_ms: &[i64],
638 worker: &str,
639) -> Result<u64, sqlx::Error> {
640 let result = sqlx::query!(
641 r#"UPDATE outbox o
642 SET attempts = o.attempts + 1,
643 last_error = f.err,
644 locked_by = NULL,
645 locked_until = NULL,
646 available_at = now() + (f.delay_ms * interval '1 millisecond'),
647 updated_at = now()
648 FROM UNNEST($1::uuid[], $2::text[], $3::bigint[]) AS f(id, err, delay_ms)
649 WHERE o.id = f.id AND o.locked_by = $4"#,
650 ids,
651 errors,
652 delays_ms,
653 worker,
654 )
655 .execute(executor)
656 .await?;
657 Ok(result.rows_affected())
658}
659
660async fn fail_dead_rows<'e>(
661 executor: impl sqlx::PgExecutor<'e>,
662 ids: &[uuid::Uuid],
663 errors: &[String],
664 reasons: &[&str],
665 worker: &str,
666) -> Result<u64, sqlx::Error> {
667 let result = sqlx::query!(
668 r#"UPDATE outbox o
669 SET attempts = o.attempts + 1,
670 last_error = f.err,
671 dead_at = now(),
672 dead_reason = f.reason,
673 locked_by = NULL,
674 locked_until = NULL,
675 updated_at = now()
676 FROM UNNEST($1::uuid[], $2::text[], $3::text[]) AS f(id, err, reason)
677 WHERE o.id = f.id AND o.locked_by = $4"#,
678 ids,
679 errors,
680 reasons as &[&str],
681 worker,
682 )
683 .execute(executor)
684 .await?;
685 Ok(result.rows_affected())
686}
687
688async fn claim_rows<'e>(
696 executor: impl sqlx::PgExecutor<'e>,
697 batch_size: i64,
698 worker: &str,
699 lease_ms: i64,
700) -> Result<Vec<RawRow>, sqlx::Error> {
701 sqlx::query_as!(
702 RawRow,
703 r#"WITH claimed AS (
704 SELECT id FROM outbox
705 WHERE published_at IS NULL AND dead_at IS NULL
706 AND available_at <= now()
707 AND (locked_until IS NULL OR locked_until < now())
708 AND (expires_at IS NULL OR expires_at > now())
709 ORDER BY available_at, sequence
710 LIMIT $1
711 FOR UPDATE SKIP LOCKED
712 )
713 UPDATE outbox o
714 SET locked_by = $2,
715 locked_until = now() + ($3::bigint * interval '1 millisecond'),
716 updated_at = now()
717 FROM claimed
718 WHERE o.id = claimed.id
719 RETURNING o.id, o.sequence, o.message_type, o.message_version,
720 o.correlation_id, o.conversation_id, o.causation_id, o.request_id,
721 o.content_type, o.payload, o.tenant_id, o.expires_at, o.ordering_key,
722 o.metadata, o.headers, o.metadata_version,
723 o.created_at, o.available_at,
724 o.attempts, o.locked_by, o.locked_until,
725 o.published_at, o.dead_at, o.dead_reason, o.last_error"#,
726 batch_size,
727 worker,
728 lease_ms,
729 )
730 .fetch_all(executor)
731 .await
732}
733
734async fn list_dead_rows<'e>(
738 executor: impl sqlx::PgExecutor<'e>,
739 query: &DeadQuery,
740 limit: i64,
741) -> Result<Vec<RawRow>, sqlx::Error> {
742 sqlx::query_as!(
743 RawRow,
744 r#"SELECT id, sequence, message_type, message_version,
745 correlation_id, conversation_id, causation_id, request_id,
746 content_type, payload, tenant_id, expires_at, ordering_key,
747 metadata, headers, metadata_version,
748 created_at, available_at,
749 attempts, locked_by, locked_until,
750 published_at, dead_at, dead_reason, last_error
751 FROM outbox
752 WHERE dead_at IS NOT NULL
753 AND ($1::text IS NULL OR message_type = $1)
754 AND ($2::text IS NULL OR tenant_id = $2)
755 AND ($3::timestamptz IS NULL OR dead_at < $3)
756 AND ($4::bigint IS NULL OR sequence > $4)
757 ORDER BY sequence ASC
758 LIMIT $5"#,
759 query.message_type,
760 query.tenant_id,
761 query.dead_before,
762 query.after_sequence,
763 limit,
764 )
765 .fetch_all(executor)
766 .await
767}
768
769async fn retry_dead_rows<'e>(
772 executor: impl sqlx::PgExecutor<'e>,
773 ids: &[uuid::Uuid],
774) -> Result<u64, sqlx::Error> {
775 let result = sqlx::query!(
776 r#"UPDATE outbox
777 SET dead_at = NULL,
778 dead_reason = NULL,
779 available_at = now(),
780 attempts = 0,
781 locked_by = NULL,
782 locked_until = NULL,
783 updated_at = now()
784 WHERE id = ANY($1) AND dead_at IS NOT NULL"#,
785 ids,
786 )
787 .execute(executor)
788 .await?;
789 Ok(result.rows_affected())
790}
791
792async fn purge_dead_rows<'e>(
795 executor: impl sqlx::PgExecutor<'e>,
796 ids: &[uuid::Uuid],
797) -> Result<u64, sqlx::Error> {
798 let result = sqlx::query!(
799 "DELETE FROM outbox WHERE id = ANY($1) AND dead_at IS NOT NULL",
800 ids,
801 )
802 .execute(executor)
803 .await?;
804 Ok(result.rows_affected())
805}
806
807impl<Ser: Serializer + Send + Sync + 'static> OutboxStore for PostgresOutboxStore<Ser> {
808 type Error = PostgresStoreError;
809
810 async fn acquire(
819 &self,
820 request: reliar_outbox::AcquireRequest,
821 ) -> Result<AcquiredBatch, Self::Error> {
822 let batch_size = i64::from(request.batch_size);
823 let lease_ms = i64::try_from(request.lease.as_millis()).unwrap_or(i64::MAX);
824 let worker = request.worker.as_str();
825
826 let rows = if self.settings.statement_timeout.is_zero() {
831 claim_rows(&self.pool, batch_size, worker, lease_ms)
832 .await
833 .map_err(|e| self.map_err(e))?
834 } else {
835 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
836 let timeout_ms = i64::try_from(self.settings.statement_timeout.as_millis())
837 .unwrap_or(i64::MAX)
838 .to_string();
839 sqlx::query_scalar!(
840 "SELECT set_config('statement_timeout', $1, true)",
841 timeout_ms
842 )
843 .fetch_one(&mut *tx)
844 .await
845 .map_err(|e| self.map_err(e))?;
846 let rows = claim_rows(&mut *tx, batch_size, worker, lease_ms)
847 .await
848 .map_err(|e| self.map_err(e))?;
849 tx.commit().await.map_err(|e| self.map_err(e))?;
850 rows
851 };
852
853 let mut records = Vec::with_capacity(rows.len());
854 let mut poisoned = Vec::new();
855 let mut poisoned_ids = Vec::new();
856 let mut poisoned_errors = Vec::new();
857
858 for raw in rows {
859 match decode_row(raw) {
860 Ok(record) => records.push(record),
861 Err(err) => {
862 poisoned_ids.push(err.id.as_uuid());
863 poisoned_errors.push(crate::records::truncate_last_error(err.detail.clone()));
864 poisoned.push(PoisonedRow::new(err.id, err.sequence, err.detail));
865 }
866 }
867 }
868
869 if !poisoned_ids.is_empty() {
870 let undecodable =
877 crate::records::encode_dead_reason(reliar_outbox::DeadReason::Undecodable);
878 if self.settings.statement_timeout.is_zero() {
879 poison_sweep_rows(
880 &self.pool,
881 &poisoned_ids,
882 &poisoned_errors,
883 worker,
884 undecodable,
885 )
886 .await
887 .map_err(|e| self.map_err(e))?;
888 } else {
889 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
890 self.set_local_timeout(&mut tx).await?;
891 poison_sweep_rows(
892 &mut *tx,
893 &poisoned_ids,
894 &poisoned_errors,
895 worker,
896 undecodable,
897 )
898 .await
899 .map_err(|e| self.map_err(e))?;
900 tx.commit().await.map_err(|e| self.map_err(e))?;
901 }
902 }
903
904 Ok(AcquiredBatch::new(records, poisoned))
905 }
906
907 async fn complete(
911 &self,
912 worker: &WorkerId,
913 items: &[CompletedMessage],
914 ) -> Result<u64, Self::Error> {
915 if items.is_empty() {
916 return Ok(0);
917 }
918 let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.message.id.as_uuid()).collect();
919 let affected = if self.settings.statement_timeout.is_zero() {
920 complete_rows(&self.pool, &ids, worker.as_str())
921 .await
922 .map_err(|e| self.map_err(e))?
923 } else {
924 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
925 self.set_local_timeout(&mut tx).await?;
926 let affected = complete_rows(&mut *tx, &ids, worker.as_str())
927 .await
928 .map_err(|e| self.map_err(e))?;
929 tx.commit().await.map_err(|e| self.map_err(e))?;
930 affected
931 };
932 log_shortfall("complete", items.len(), affected);
933 Ok(affected)
934 }
935
936 async fn fail(&self, worker: &WorkerId, items: &[FailedMessage]) -> Result<u64, Self::Error> {
941 if items.is_empty() {
942 return Ok(0);
943 }
944
945 let mut retry_ids = Vec::new();
946 let mut retry_errors = Vec::new();
947 let mut retry_delays = Vec::new();
948 let mut dead_ids = Vec::new();
949 let mut dead_errors = Vec::new();
950 let mut dead_reasons = Vec::new();
951
952 for item in items {
953 match item.outcome {
954 FailureOutcome::Retry { delay } => {
955 retry_ids.push(item.message.id.as_uuid());
956 retry_errors.push(item.error.clone());
957 retry_delays.push(i64::try_from(delay.as_millis()).unwrap_or(i64::MAX));
958 }
959 FailureOutcome::Dead { reason } => {
960 dead_ids.push(item.message.id.as_uuid());
961 dead_errors.push(item.error.clone());
962 dead_reasons.push(crate::records::encode_dead_reason(reason));
963 }
964 _ => tracing::error!(
969 id = %item.message.id,
970 "unrecognised FailureOutcome variant; row left as-is"
971 ),
972 }
973 }
974
975 let affected = if self.settings.statement_timeout.is_zero() {
976 let mut affected = 0u64;
977 if !retry_ids.is_empty() {
978 affected += fail_retry_rows(
979 &self.pool,
980 &retry_ids,
981 &retry_errors,
982 &retry_delays,
983 worker.as_str(),
984 )
985 .await
986 .map_err(|e| self.map_err(e))?;
987 }
988 if !dead_ids.is_empty() {
989 affected += fail_dead_rows(
990 &self.pool,
991 &dead_ids,
992 &dead_errors,
993 &dead_reasons,
994 worker.as_str(),
995 )
996 .await
997 .map_err(|e| self.map_err(e))?;
998 }
999 affected
1000 } else {
1001 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1002 self.set_local_timeout(&mut tx).await?;
1003 let mut affected = 0u64;
1004 if !retry_ids.is_empty() {
1005 affected += fail_retry_rows(
1006 &mut *tx,
1007 &retry_ids,
1008 &retry_errors,
1009 &retry_delays,
1010 worker.as_str(),
1011 )
1012 .await
1013 .map_err(|e| self.map_err(e))?;
1014 }
1015 if !dead_ids.is_empty() {
1016 affected += fail_dead_rows(
1017 &mut *tx,
1018 &dead_ids,
1019 &dead_errors,
1020 &dead_reasons,
1021 worker.as_str(),
1022 )
1023 .await
1024 .map_err(|e| self.map_err(e))?;
1025 }
1026 tx.commit().await.map_err(|e| self.map_err(e))?;
1027 affected
1028 };
1029 log_shortfall("fail", items.len(), affected);
1030 Ok(affected)
1031 }
1032
1033 async fn release(&self, worker: &WorkerId, items: &[MessageRef]) -> Result<u64, Self::Error> {
1036 if items.is_empty() {
1037 return Ok(0);
1038 }
1039 let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
1040 let affected = if self.settings.statement_timeout.is_zero() {
1041 release_rows(&self.pool, &ids, worker.as_str())
1042 .await
1043 .map_err(|e| self.map_err(e))?
1044 } else {
1045 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1046 self.set_local_timeout(&mut tx).await?;
1047 let affected = release_rows(&mut *tx, &ids, worker.as_str())
1048 .await
1049 .map_err(|e| self.map_err(e))?;
1050 tx.commit().await.map_err(|e| self.map_err(e))?;
1051 affected
1052 };
1053 log_shortfall("release", items.len(), affected);
1054 Ok(affected)
1055 }
1056
1057 async fn extend_lease(
1060 &self,
1061 worker: &WorkerId,
1062 items: &[MessageRef],
1063 lease: std::time::Duration,
1064 ) -> Result<u64, Self::Error> {
1065 if items.is_empty() {
1066 return Ok(0);
1067 }
1068 let ids: Vec<uuid::Uuid> = items.iter().map(|i| i.id.as_uuid()).collect();
1069 let lease_ms = i64::try_from(lease.as_millis()).unwrap_or(i64::MAX);
1070 let affected = if self.settings.statement_timeout.is_zero() {
1071 extend_lease_rows(&self.pool, &ids, lease_ms, worker.as_str())
1072 .await
1073 .map_err(|e| self.map_err(e))?
1074 } else {
1075 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1076 self.set_local_timeout(&mut tx).await?;
1077 let affected = extend_lease_rows(&mut *tx, &ids, lease_ms, worker.as_str())
1078 .await
1079 .map_err(|e| self.map_err(e))?;
1080 tx.commit().await.map_err(|e| self.map_err(e))?;
1081 affected
1082 };
1083 log_shortfall("extend_lease", items.len(), affected);
1084 Ok(affected)
1085 }
1086
1087 async fn purge(&self, request: PurgeRequest) -> Result<PurgeReport, Self::Error> {
1094 let batch_size = i64::from(request.batch_size);
1095 let expired_reason = crate::records::encode_dead_reason(reliar_outbox::DeadReason::Expired);
1096
1097 let (published_deleted, dead_deleted, expired_to_dead) =
1098 if self.settings.statement_timeout.is_zero() {
1099 let published_deleted = if let Some(retention) = request.published_retention {
1100 let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
1101 purge_published_rows(&self.pool, retention_ms, batch_size)
1102 .await
1103 .map_err(|e| self.map_err(e))?
1104 } else {
1105 0
1106 };
1107
1108 let dead_deleted = if let Some(retention) = request.dead_retention {
1109 let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
1110 purge_dead_retention_rows(&self.pool, retention_ms, batch_size)
1111 .await
1112 .map_err(|e| self.map_err(e))?
1113 } else {
1114 0
1115 };
1116
1117 let expired_to_dead =
1118 purge_expired_sweep_rows(&self.pool, batch_size, expired_reason)
1119 .await
1120 .map_err(|e| self.map_err(e))?;
1121
1122 (published_deleted, dead_deleted, expired_to_dead)
1123 } else {
1124 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1130 self.set_local_timeout(&mut tx).await?;
1131
1132 let published_deleted = if let Some(retention) = request.published_retention {
1133 let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
1134 purge_published_rows(&mut *tx, retention_ms, batch_size)
1135 .await
1136 .map_err(|e| self.map_err(e))?
1137 } else {
1138 0
1139 };
1140
1141 let dead_deleted = if let Some(retention) = request.dead_retention {
1142 let retention_ms = i64::try_from(retention.as_millis()).unwrap_or(i64::MAX);
1143 purge_dead_retention_rows(&mut *tx, retention_ms, batch_size)
1144 .await
1145 .map_err(|e| self.map_err(e))?
1146 } else {
1147 0
1148 };
1149
1150 let expired_to_dead =
1151 purge_expired_sweep_rows(&mut *tx, batch_size, expired_reason)
1152 .await
1153 .map_err(|e| self.map_err(e))?;
1154
1155 tx.commit().await.map_err(|e| self.map_err(e))?;
1156 (published_deleted, dead_deleted, expired_to_dead)
1157 };
1158
1159 Ok(PurgeReport::new(
1160 published_deleted,
1161 dead_deleted,
1162 expired_to_dead,
1163 ))
1164 }
1165 async fn stats(&self) -> Result<OutboxStats, Self::Error> {
1179 if self.settings.statement_timeout.is_zero() {
1180 let row = sqlx::query!(
1181 r#"SELECT
1182 count(*) FILTER (
1183 WHERE published_at IS NULL AND dead_at IS NULL
1184 AND available_at <= now()
1185 AND (locked_until IS NULL OR locked_until < now())
1186 AND (expires_at IS NULL OR expires_at > now())
1187 ) AS "pending!",
1188 count(*) FILTER (WHERE dead_at IS NOT NULL) AS "dead!",
1189 count(*) FILTER (
1190 WHERE published_at IS NULL AND dead_at IS NULL
1191 AND expires_at IS NOT NULL AND expires_at < now()
1192 ) AS "expired_pending!",
1193 min(available_at) FILTER (
1194 WHERE published_at IS NULL AND dead_at IS NULL
1195 AND available_at <= now()
1196 AND (locked_until IS NULL OR locked_until < now())
1197 AND (expires_at IS NULL OR expires_at > now())
1198 ) AS oldest_pending_available_at,
1199 now() AS "as_of!"
1200 FROM outbox"#
1201 )
1202 .fetch_one(&self.pool)
1203 .await
1204 .map_err(|e| self.map_err(e))?;
1205
1206 return Ok(OutboxStats::new(
1207 u64::try_from(row.pending).unwrap_or(0),
1208 u64::try_from(row.dead).unwrap_or(0),
1209 u64::try_from(row.expired_pending).unwrap_or(0),
1210 row.oldest_pending_available_at,
1211 row.as_of,
1212 ));
1213 }
1214
1215 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1217 self.set_local_timeout(&mut tx).await?;
1218
1219 let row = sqlx::query!(
1220 r#"SELECT
1221 count(*) FILTER (
1222 WHERE published_at IS NULL AND dead_at IS NULL
1223 AND available_at <= now()
1224 AND (locked_until IS NULL OR locked_until < now())
1225 AND (expires_at IS NULL OR expires_at > now())
1226 ) AS "pending!",
1227 count(*) FILTER (WHERE dead_at IS NOT NULL) AS "dead!",
1228 count(*) FILTER (
1229 WHERE published_at IS NULL AND dead_at IS NULL
1230 AND expires_at IS NOT NULL AND expires_at < now()
1231 ) AS "expired_pending!",
1232 min(available_at) FILTER (
1233 WHERE published_at IS NULL AND dead_at IS NULL
1234 AND available_at <= now()
1235 AND (locked_until IS NULL OR locked_until < now())
1236 AND (expires_at IS NULL OR expires_at > now())
1237 ) AS oldest_pending_available_at,
1238 now() AS "as_of!"
1239 FROM outbox"#
1240 )
1241 .fetch_one(&mut *tx)
1242 .await
1243 .map_err(|e| self.map_err(e))?;
1244
1245 tx.commit().await.map_err(|e| self.map_err(e))?;
1246
1247 Ok(OutboxStats::new(
1248 u64::try_from(row.pending).unwrap_or(0),
1249 u64::try_from(row.dead).unwrap_or(0),
1250 u64::try_from(row.expired_pending).unwrap_or(0),
1251 row.oldest_pending_available_at,
1252 row.as_of,
1253 ))
1254 }
1255}
1256
1257impl<Ser: Serializer + Send + Sync + 'static> OutboxDeadLetters for PostgresOutboxStore<Ser> {
1258 type Error = PostgresStoreError;
1259
1260 async fn list_dead(&self, query: DeadQuery) -> Result<DeadLetterPage, Self::Error> {
1267 let capped_limit = query.limit.min(MAX_LIST_DEAD_LIMIT);
1270 let limit = i64::from(capped_limit);
1271
1272 let rows = if self.settings.statement_timeout.is_zero() {
1273 list_dead_rows(&self.pool, &query, limit)
1274 .await
1275 .map_err(|e| self.map_err(e))?
1276 } else {
1277 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1278 self.set_local_timeout(&mut tx).await?;
1279 let rows = list_dead_rows(&mut *tx, &query, limit)
1280 .await
1281 .map_err(|e| self.map_err(e))?;
1282 tx.commit().await.map_err(|e| self.map_err(e))?;
1283 rows
1284 };
1285
1286 let scanned = rows.len();
1287 let mut records = Vec::with_capacity(scanned);
1288 let mut poisoned = Vec::new();
1289 let mut max_sequence: Option<i64> = None;
1290
1291 for raw in rows {
1292 max_sequence = Some(max_sequence.map_or(raw.sequence, |m| m.max(raw.sequence)));
1293 match decode_row(raw) {
1294 Ok(record) => records.push(record),
1295 Err(err) => poisoned.push(PoisonedRow::new(err.id, err.sequence, err.detail)),
1296 }
1297 }
1298
1299 let next_after_sequence = if scanned == capped_limit as usize {
1302 max_sequence
1303 } else {
1304 None
1305 };
1306
1307 Ok(DeadLetterPage::new(records, poisoned, next_after_sequence))
1308 }
1309
1310 async fn retry_dead(&self, refs: &[MessageRef]) -> Result<u64, Self::Error> {
1314 if refs.is_empty() {
1315 return Ok(0);
1316 }
1317 let ids: Vec<uuid::Uuid> = refs.iter().map(|r| r.id.as_uuid()).collect();
1318 let affected = if self.settings.statement_timeout.is_zero() {
1319 retry_dead_rows(&self.pool, &ids)
1320 .await
1321 .map_err(|e| self.map_err(e))?
1322 } else {
1323 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1324 self.set_local_timeout(&mut tx).await?;
1325 let affected = retry_dead_rows(&mut *tx, &ids)
1326 .await
1327 .map_err(|e| self.map_err(e))?;
1328 tx.commit().await.map_err(|e| self.map_err(e))?;
1329 affected
1330 };
1331 Ok(affected)
1332 }
1333
1334 async fn purge_dead(&self, refs: &[MessageRef]) -> Result<u64, Self::Error> {
1336 if refs.is_empty() {
1337 return Ok(0);
1338 }
1339 let ids: Vec<uuid::Uuid> = refs.iter().map(|r| r.id.as_uuid()).collect();
1340 let affected = if self.settings.statement_timeout.is_zero() {
1341 purge_dead_rows(&self.pool, &ids)
1342 .await
1343 .map_err(|e| self.map_err(e))?
1344 } else {
1345 let mut tx = self.pool.begin().await.map_err(|e| self.map_err(e))?;
1346 self.set_local_timeout(&mut tx).await?;
1347 let affected = purge_dead_rows(&mut *tx, &ids)
1348 .await
1349 .map_err(|e| self.map_err(e))?;
1350 tx.commit().await.map_err(|e| self.map_err(e))?;
1351 affected
1352 };
1353 Ok(affected)
1354 }
1355}
1356
1357fn log_shortfall(operation: &'static str, claimed: usize, affected: u64) {
1360 let claimed = claimed as u64;
1361 if affected < claimed {
1362 tracing::debug!(
1363 operation,
1364 claimed,
1365 affected,
1366 "fewer rows affected than claimed"
1367 );
1368 }
1369}