Skip to main content

cratestack_sqlx/
descriptor.rs

1mod 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    // Shared (not per-clone) so every handle onto the same logical
23    // runtime agrees on whether `cratestack_audit` has been
24    // bootstrapped. See `crate::audit::ensure_audit_table` — this is
25    // what lets it skip re-issuing `CREATE INDEX IF NOT EXISTS` after
26    // the first call, which is what self-deadlocked chained
27    // `run_in_tx` audit writes in a caller-managed transaction.
28    audit_table_ensured: Arc<AtomicBool>,
29    // Installation point for cratestack#473: defaults to `NoopAuditSink`
30    // so existing callers of `new()` see no behavior change. Installed
31    // via `with_audit_sink` (mirrors `IdempotencyLayer::new`/
32    // `with_principal_fingerprint`'s builder shape). The DB write in
33    // `crate::audit::enqueue_audit_event` remains the sole source of
34    // truth; this is a best-effort downstream projection dispatched
35    // from `crate::audit::dispatch_audit_sink` after the owning
36    // transaction commits.
37    audit_sink: Arc<dyn AuditSink>,
38    // `Some` only on the per-attempt runtime `run_isolated` builds for an
39    // `@isolation` procedure: every executing path then runs on this
40    // transaction instead of the pool (docs/design/procedure-isolation.md §4).
41    pub(crate) bound: Option<Arc<crate::bound::BoundTx>>,
42    pub(crate) isolation_max_retries: u32,
43}
44
45// `dyn AuditSink` has no `Debug` bound (matching `IdempotencyStore` /
46// `RateLimitStore`, neither of which require one either), so this can't
47// be `#[derive(Debug)]` — same reason `CratestackEventBus` hand-rolls its own
48// `Debug` impl instead of deriving one.
49impl 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    /// Install a custom [`AuditSink`] that every `@@audit` mutation on
87    /// this runtime fans out to, in addition to the in-database
88    /// `cratestack_audit` table row `enqueue_audit_event` always
89    /// writes. Composable via [`cratestack_core::MulticastAuditSink`]
90    /// for more than one downstream (Kafka, Redis pubsub, a webhook —
91    /// this crate ships none of them; see `AuditSink`'s doc comment).
92    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    /// An owned, cheaply-cloneable handle onto the underlying
115    /// `CratestackEventBus` — needed by callers (e.g. `@@subscribe` SSE
116    /// dispatch, cratestack#390) that outlive the `&SqlxRuntime` borrow
117    /// `subscribe`/`unsubscribe` would otherwise require.
118    #[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        // Inside an `@isolation` procedure the outbox rows are still
126        // uncommitted: remember to drain once the attempt commits.
127        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}