cratestack_sqlx/
descriptor.rs1mod event_outbox;
2
3use std::sync::Arc;
4use std::sync::atomic::AtomicBool;
5
6use crate::sqlx;
7
8use cratestack_core::{
9 AuditSink, CratestackError, CratestackEventBus, CratestackEventEnvelope, CratestackEventFuture,
10 ModelEventKind, NoopAuditSink, SubscriptionHandle,
11};
12
13use crate::error::cratestack_error_from_sqlx;
14use event_outbox::EventOutboxRow;
15
16pub use event_outbox::{enqueue_event_outbox, ensure_event_outbox_table};
17
18#[derive(Clone)]
19pub struct SqlxRuntime {
20 pool: sqlx::PgPool,
21 events: CratestackEventBus,
22 audit_table_ensured: Arc<AtomicBool>,
29 audit_sink: Arc<dyn AuditSink>,
38 pub(crate) bound: Option<Arc<crate::bound::BoundTx>>,
42 pub(crate) isolation_max_retries: u32,
43}
44
45impl std::fmt::Debug for SqlxRuntime {
50 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
51 f.debug_struct("SqlxRuntime")
52 .field("pool", &self.pool)
53 .field("events", &self.events)
54 .field("bound", &self.bound.is_some())
55 .field("isolation_max_retries", &self.isolation_max_retries)
56 .field(
57 "audit_table_ensured",
58 &self
59 .audit_table_ensured
60 .load(std::sync::atomic::Ordering::Relaxed),
61 )
62 .finish_non_exhaustive()
63 }
64}
65
66impl SqlxRuntime {
67 pub fn new(pool: sqlx::PgPool) -> Self {
68 Self {
69 pool,
70 events: CratestackEventBus::default(),
71 audit_table_ensured: Arc::new(AtomicBool::new(false)),
72 audit_sink: Arc::new(NoopAuditSink),
73 bound: None,
74 isolation_max_retries: crate::isolation::MAX_RETRIES_DEFAULT,
75 }
76 }
77
78 pub fn pool(&self) -> &sqlx::PgPool {
79 &self.pool
80 }
81
82 pub(crate) fn audit_table_ensured(&self) -> &AtomicBool {
83 &self.audit_table_ensured
84 }
85
86 pub fn with_audit_sink(mut self, sink: Arc<dyn AuditSink>) -> Self {
93 self.audit_sink = sink;
94 self
95 }
96
97 pub(crate) fn audit_sink(&self) -> &Arc<dyn AuditSink> {
98 &self.audit_sink
99 }
100
101 #[doc(hidden)]
102 pub fn subscribe<F>(
103 &self,
104 model: &'static str,
105 operation: ModelEventKind,
106 handler: F,
107 ) -> SubscriptionHandle
108 where
109 F: Fn(CratestackEventEnvelope) -> CratestackEventFuture + Send + Sync + 'static,
110 {
111 self.events.subscribe(model, operation, handler)
112 }
113
114 #[doc(hidden)]
119 pub fn events_bus(&self) -> CratestackEventBus {
120 self.events.clone()
121 }
122
123 #[doc(hidden)]
124 pub async fn drain_event_outbox(&self) -> Result<usize, CratestackError> {
125 if let Some(bound) = self.bound() {
128 bound.request_drain();
129 return Ok(0);
130 }
131 ensure_event_outbox_table(&self.pool).await?;
132
133 let rows = sqlx::query_as::<_, EventOutboxRow>(
134 "SELECT event_id, model, operation, occurred_at, payload, attempts, last_error \
135 FROM cratestack_event_outbox \
136 WHERE delivered_at IS NULL \
137 ORDER BY occurred_at ASC, event_id ASC",
138 )
139 .fetch_all(&self.pool)
140 .await
141 .map_err(cratestack_error_from_sqlx)?;
142
143 let mut delivered = 0usize;
144 for row in rows {
145 let event_id = row.event_id;
146 let envelope = row.try_into_envelope()?;
147 match self.events.emit(envelope).await {
148 Ok(()) => {
149 sqlx::query(
150 "UPDATE cratestack_event_outbox \
151 SET delivered_at = NOW(), last_error = NULL, attempts = attempts + 1 \
152 WHERE event_id = $1",
153 )
154 .bind(event_id)
155 .execute(&self.pool)
156 .await
157 .map_err(cratestack_error_from_sqlx)?;
158 delivered += 1;
159 }
160 Err(error) => {
161 sqlx::query(
162 "UPDATE cratestack_event_outbox \
163 SET attempts = attempts + 1, last_error = $2 \
164 WHERE event_id = $1",
165 )
166 .bind(event_id)
167 .bind(error.to_string())
168 .execute(&self.pool)
169 .await
170 .map_err(cratestack_error_from_sqlx)?;
171 }
172 }
173 }
174
175 Ok(delivered)
176 }
177}