cratestack_sqlx/query/write/
create.rs1use 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 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 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}