cratestack_sqlx/query/write/
delete.rs1use cratestack_core::{AuditOperation, CratestackContext, CratestackError, ModelEventKind};
12
13use crate::audit::{
14 RunInTxOutcome, build_audit_event, dispatch_audit_sink, enqueue_audit_event,
15 ensure_audit_table, fetch_for_audit,
16};
17use crate::descriptor::{enqueue_event_outbox, ensure_event_outbox_table};
18use crate::{ModelDescriptor, SqlxRuntime, cratestack_error_from_sqlx, sqlx};
19
20use super::delete_exec::{delete_record_in_conn, delete_returning_record};
21
22#[derive(Debug, Clone)]
23pub struct DeleteRecord<'a, M: 'static, PK: 'static> {
24 pub(crate) runtime: &'a SqlxRuntime,
25 pub(crate) descriptor: &'static ModelDescriptor<M, PK>,
26 pub(crate) id: PK,
27 pub(crate) if_match: Option<i64>,
28}
29
30impl<'a, M: 'static, PK: 'static> DeleteRecord<'a, M, PK> {
31 pub fn if_match(mut self, expected: i64) -> Self {
37 self.if_match = Some(expected);
38 self
39 }
40
41 pub fn preview_sql(&self) -> String {
42 match self.descriptor.version_column {
43 Some(version_col) => format!(
44 "DELETE FROM {} WHERE {} = $1 AND {} = $2 RETURNING {}",
45 self.descriptor.table_name,
46 self.descriptor.primary_key,
47 version_col,
48 self.descriptor.select_projection(),
49 ),
50 None => format!(
51 "DELETE FROM {} WHERE {} = $1 RETURNING {}",
52 self.descriptor.table_name,
53 self.descriptor.primary_key,
54 self.descriptor.select_projection(),
55 ),
56 }
57 }
58
59 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 PK: Send + Clone + sqlx::Type<sqlx::Postgres> + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
72 {
73 if self.descriptor.version_column.is_some() && self.if_match.is_none() {
74 return Err(CratestackError::PreconditionFailed(
75 "If-Match header required for versioned model".to_owned(),
76 ));
77 }
78 let emits_event = self.descriptor.emits(ModelEventKind::Deleted);
79 let audit_enabled = self.descriptor.audit_enabled;
80 let soft_delete = self.descriptor.soft_delete_column.is_some();
81 if emits_event {
82 ensure_event_outbox_table(&mut **tx).await?;
83 }
84 if audit_enabled {
85 ensure_audit_table(self.runtime, &mut **tx).await?;
86 }
87 let before_record = if audit_enabled && soft_delete {
92 fetch_for_audit(&mut **tx, self.descriptor, self.id.clone()).await?
93 } else {
94 None
95 };
96 let before_snapshot = before_record
97 .as_ref()
98 .and_then(|m| serde_json::to_value(m).ok());
99 let record = delete_record_in_conn(
100 self.runtime,
101 tx,
102 self.descriptor,
103 self.id,
104 ctx,
105 self.if_match,
106 )
107 .await?;
108 if emits_event {
109 enqueue_event_outbox(
110 &mut **tx,
111 self.descriptor.schema_name,
112 ModelEventKind::Deleted,
113 &record,
114 )
115 .await?;
116 }
117 let mut audit_event = None;
118 if audit_enabled {
119 let (before, after) = if soft_delete {
120 (before_snapshot, serde_json::to_value(&record).ok())
121 } else {
122 (serde_json::to_value(&record).ok(), None)
123 };
124 let event =
125 build_audit_event(self.descriptor, AuditOperation::Delete, before, after, ctx);
126 enqueue_audit_event(&mut **tx, &event).await?;
127 audit_event = Some(event);
128 }
129 Ok(RunInTxOutcome::new(
130 record,
131 audit_event.into_iter().collect(),
132 ))
133 }
134
135 pub async fn run(self, ctx: &CratestackContext) -> Result<M, CratestackError>
136 where
137 for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
138 PK: Send + Clone + sqlx::Type<sqlx::Postgres> + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
139 {
140 if let Some(bound) = self.runtime.bound() {
144 let emits = self
145 .descriptor
146 .emits(cratestack_core::ModelEventKind::Deleted);
147 let outcome = crate::bound::in_bound_savepoint!(bound, |sp| self.run_in_tx(sp, ctx))?;
148 return Ok(bound.settle(outcome, emits));
149 }
150 if self.descriptor.version_column.is_some() && self.if_match.is_none() {
151 return Err(CratestackError::PreconditionFailed(
152 "If-Match header required for versioned model".to_owned(),
153 ));
154 }
155 let emits_event = self.descriptor.emits(ModelEventKind::Deleted);
156 let audit_enabled = self.descriptor.audit_enabled;
157 let soft_delete = self.descriptor.soft_delete_column.is_some();
158 let needs_tx = emits_event || audit_enabled;
159 let mut audit_event = None;
160 let record = if needs_tx {
161 let mut tx = self
162 .runtime
163 .pool()
164 .begin()
165 .await
166 .map_err(cratestack_error_from_sqlx)?;
167 if emits_event {
168 ensure_event_outbox_table(&mut *tx).await?;
169 }
170 if audit_enabled {
171 ensure_audit_table(self.runtime, &mut *tx).await?;
172 }
173
174 let before_record = if audit_enabled && soft_delete {
175 fetch_for_audit(&mut *tx, self.descriptor, self.id.clone()).await?
176 } else {
177 None
178 };
179 let before_snapshot = before_record
180 .as_ref()
181 .and_then(|m| serde_json::to_value(m).ok());
182 let record = delete_record_in_conn(
183 self.runtime,
184 &mut tx,
185 self.descriptor,
186 self.id,
187 ctx,
188 self.if_match,
189 )
190 .await?;
191 if emits_event {
192 enqueue_event_outbox(
193 &mut *tx,
194 self.descriptor.schema_name,
195 ModelEventKind::Deleted,
196 &record,
197 )
198 .await?;
199 }
200 if audit_enabled {
201 let (before, after) = if soft_delete {
202 (before_snapshot, serde_json::to_value(&record).ok())
203 } else {
204 (serde_json::to_value(&record).ok(), None)
205 };
206 let event =
207 build_audit_event(self.descriptor, AuditOperation::Delete, before, after, ctx);
208 enqueue_audit_event(&mut *tx, &event).await?;
209 audit_event = Some(event);
210 }
211 tx.commit().await.map_err(cratestack_error_from_sqlx)?;
212 record
213 } else {
214 delete_returning_record(
215 self.runtime.pool(),
216 self.runtime.pool(),
217 self.descriptor,
218 self.id,
219 ctx,
220 self.if_match,
221 )
222 .await?
223 };
224
225 if emits_event {
226 let _ = self.runtime.drain_event_outbox().await;
227 }
228 if let Some(event) = &audit_event {
229 dispatch_audit_sink(self.runtime, std::slice::from_ref(event)).await;
230 }
231
232 Ok(record)
233 }
234}