Skip to main content

cratestack_sqlx/query/write/
create.rs

1//! `CreateRecord` — single-row INSERT with policy + audit + event
2//! fan-out. `run()` opens its own tx only when audit/event capture is
3//! enabled; otherwise it goes straight against the pool.
4
5use cratestack_core::{AuditOperation, CratestackContext, CratestackError, ModelEventKind};
6
7use crate::audit::{
8    RunInTxOutcome, build_audit_event, dispatch_audit_sink, enqueue_audit_event, ensure_audit_table,
9};
10use crate::descriptor::{enqueue_event_outbox, ensure_event_outbox_table};
11use crate::{CreateModelInput, ModelDescriptor, SqlxRuntime, cratestack_error_from_sqlx, sqlx};
12
13use super::create_exec::{create_record_in_conn, create_record_with_executor};
14
15#[derive(Debug, Clone)]
16pub struct CreateRecord<'a, M: 'static, PK: 'static, I> {
17    pub(crate) runtime: &'a SqlxRuntime,
18    pub(crate) descriptor: &'static ModelDescriptor<M, PK>,
19    pub(crate) input: I,
20}
21
22impl<'a, M: 'static, PK: 'static, I> CreateRecord<'a, M, PK, I>
23where
24    I: CreateModelInput<M>,
25{
26    pub fn preview_sql(&self) -> String {
27        let values = self.input.sql_values();
28        let placeholders = (1..=values.len())
29            .map(|index| format!("${index}"))
30            .collect::<Vec<_>>()
31            .join(", ");
32        let columns = values
33            .iter()
34            .map(|value| value.column)
35            .collect::<Vec<_>>()
36            .join(", ");
37
38        format!(
39            "INSERT INTO {} ({}) VALUES ({}) RETURNING {}",
40            self.descriptor.table_name,
41            columns,
42            placeholders,
43            self.descriptor.select_projection(),
44        )
45    }
46
47    /// Like [`Self::run`] but participates in a caller-supplied
48    /// transaction. The insert + outbox + audit writes all happen
49    /// inside `tx`; caller commits.
50    ///
51    /// Neither the event outbox nor the `AuditSink` fan-out run here —
52    /// both are post-commit, best-effort projections, and this
53    /// function has no visibility into when, or whether, the caller
54    /// commits `tx` (see [`crate::audit::dispatch_audit_sink`]'s doc
55    /// comment for the full reasoning). Unlike before cratestack#534,
56    /// though, neither is a dead end: this returns a
57    /// [`RunInTxOutcome`] carrying the `AuditEvent` this call built (if
58    /// `@@audit`), for the caller to pass to
59    /// `Cratestack::dispatch_audit_sink` once *their* commit succeeds;
60    /// and if the model `@@emit`s, the caller can call the pre-existing
61    /// `Cratestack::events().drain()` (cratestack#390) the same way — it
62    /// re-scans for any undelivered outbox row, so it needs no event
63    /// handed back to find this one.
64    pub async fn run_in_tx<'tx>(
65        self,
66        tx: &mut sqlx::Transaction<'tx, sqlx::Postgres>,
67        ctx: &CratestackContext,
68    ) -> Result<RunInTxOutcome<M>, CratestackError>
69    where
70        for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
71    {
72        let emits_event = self.descriptor.emits(ModelEventKind::Created);
73        let audit_enabled = self.descriptor.audit_enabled;
74        if emits_event {
75            ensure_event_outbox_table(&mut **tx).await?;
76        }
77        if audit_enabled {
78            ensure_audit_table(self.runtime, &mut **tx).await?;
79        }
80        let record =
81            create_record_in_conn(self.runtime, tx, self.descriptor, self.input, ctx).await?;
82        if emits_event {
83            enqueue_event_outbox(
84                &mut **tx,
85                self.descriptor.schema_name,
86                ModelEventKind::Created,
87                &record,
88            )
89            .await?;
90        }
91        let mut audit_event = None;
92        if audit_enabled {
93            let after = serde_json::to_value(&record).ok();
94            let event =
95                build_audit_event(self.descriptor, AuditOperation::Create, None, after, ctx);
96            enqueue_audit_event(&mut **tx, &event).await?;
97            audit_event = Some(event);
98        }
99        Ok(RunInTxOutcome::new(
100            record,
101            audit_event.into_iter().collect(),
102        ))
103    }
104
105    pub async fn run(self, ctx: &CratestackContext) -> Result<M, CratestackError>
106    where
107        for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
108    {
109        // Inside an `@isolation` procedure: run on its transaction and
110        // defer the post-commit fan-out (docs/design/procedure-isolation.md
111        // §4, §6).
112        if let Some(bound) = self.runtime.bound() {
113            let emits = self
114                .descriptor
115                .emits(cratestack_core::ModelEventKind::Created);
116            let outcome = crate::bound::in_bound_savepoint!(bound, |sp| self.run_in_tx(sp, ctx))?;
117            return Ok(bound.settle(outcome, emits));
118        }
119        let emits_event = self.descriptor.emits(ModelEventKind::Created);
120        let audit_enabled = self.descriptor.audit_enabled;
121        let needs_tx = emits_event || audit_enabled;
122        let mut audit_event = None;
123        let record = if needs_tx {
124            let mut tx = self
125                .runtime
126                .pool()
127                .begin()
128                .await
129                .map_err(cratestack_error_from_sqlx)?;
130            if emits_event {
131                ensure_event_outbox_table(&mut *tx).await?;
132            }
133            if audit_enabled {
134                ensure_audit_table(self.runtime, &mut *tx).await?;
135            }
136            let record =
137                create_record_in_conn(self.runtime, &mut tx, self.descriptor, self.input, ctx)
138                    .await?;
139            if emits_event {
140                enqueue_event_outbox(
141                    &mut *tx,
142                    self.descriptor.schema_name,
143                    ModelEventKind::Created,
144                    &record,
145                )
146                .await?;
147            }
148            if audit_enabled {
149                let after = serde_json::to_value(&record).ok();
150                let event =
151                    build_audit_event(self.descriptor, AuditOperation::Create, None, after, ctx);
152                enqueue_audit_event(&mut *tx, &event).await?;
153                audit_event = Some(event);
154            }
155            tx.commit().await.map_err(cratestack_error_from_sqlx)?;
156            record
157        } else {
158            create_record_with_executor(
159                self.runtime.pool(),
160                self.runtime.pool(),
161                self.descriptor,
162                self.input,
163                ctx,
164            )
165            .await?
166        };
167
168        if emits_event {
169            let _ = self.runtime.drain_event_outbox().await;
170        }
171        if let Some(event) = &audit_event {
172            dispatch_audit_sink(self.runtime, std::slice::from_ref(event)).await;
173        }
174
175        Ok(record)
176    }
177}