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 =
100 delete_record_in_conn(tx, self.descriptor, self.id, ctx, self.if_match).await?;
101 if emits_event {
102 enqueue_event_outbox(
103 &mut **tx,
104 self.descriptor.schema_name,
105 ModelEventKind::Deleted,
106 &record,
107 )
108 .await?;
109 }
110 let mut audit_event = None;
111 if audit_enabled {
112 let (before, after) = if soft_delete {
113 (before_snapshot, serde_json::to_value(&record).ok())
114 } else {
115 (serde_json::to_value(&record).ok(), None)
116 };
117 let event =
118 build_audit_event(self.descriptor, AuditOperation::Delete, before, after, ctx);
119 enqueue_audit_event(&mut **tx, &event).await?;
120 audit_event = Some(event);
121 }
122 Ok(RunInTxOutcome::new(
123 record,
124 audit_event.into_iter().collect(),
125 ))
126 }
127
128 pub async fn run(self, ctx: &CratestackContext) -> Result<M, CratestackError>
129 where
130 for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
131 PK: Send + Clone + sqlx::Type<sqlx::Postgres> + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
132 {
133 if let Some(bound) = self.runtime.bound() {
137 let emits = self
138 .descriptor
139 .emits(cratestack_core::ModelEventKind::Deleted);
140 let outcome = crate::bound::in_bound_savepoint!(bound, |sp| self.run_in_tx(sp, ctx))?;
141 return Ok(bound.settle(outcome, emits));
142 }
143 if self.descriptor.version_column.is_some() && self.if_match.is_none() {
144 return Err(CratestackError::PreconditionFailed(
145 "If-Match header required for versioned model".to_owned(),
146 ));
147 }
148 let emits_event = self.descriptor.emits(ModelEventKind::Deleted);
149 let audit_enabled = self.descriptor.audit_enabled;
150 let soft_delete = self.descriptor.soft_delete_column.is_some();
151 let needs_tx = emits_event || audit_enabled;
152 let mut audit_event = None;
153 let record = if needs_tx {
154 let mut tx = self
155 .runtime
156 .pool()
157 .begin()
158 .await
159 .map_err(cratestack_error_from_sqlx)?;
160 if emits_event {
161 ensure_event_outbox_table(&mut *tx).await?;
162 }
163 if audit_enabled {
164 ensure_audit_table(self.runtime, &mut *tx).await?;
165 }
166
167 let before_record = if audit_enabled && soft_delete {
168 fetch_for_audit(&mut *tx, self.descriptor, self.id.clone()).await?
169 } else {
170 None
171 };
172 let before_snapshot = before_record
173 .as_ref()
174 .and_then(|m| serde_json::to_value(m).ok());
175 let record =
176 delete_record_in_conn(&mut tx, self.descriptor, self.id, ctx, self.if_match)
177 .await?;
178 if emits_event {
179 enqueue_event_outbox(
180 &mut *tx,
181 self.descriptor.schema_name,
182 ModelEventKind::Deleted,
183 &record,
184 )
185 .await?;
186 }
187 if audit_enabled {
188 let (before, after) = if soft_delete {
189 (before_snapshot, serde_json::to_value(&record).ok())
190 } else {
191 (serde_json::to_value(&record).ok(), None)
192 };
193 let event =
194 build_audit_event(self.descriptor, AuditOperation::Delete, before, after, ctx);
195 enqueue_audit_event(&mut *tx, &event).await?;
196 audit_event = Some(event);
197 }
198 tx.commit().await.map_err(cratestack_error_from_sqlx)?;
199 record
200 } else {
201 delete_returning_record(
202 self.runtime.pool(),
203 self.runtime.pool(),
204 self.descriptor,
205 self.id,
206 ctx,
207 self.if_match,
208 )
209 .await?
210 };
211
212 if emits_event {
213 let _ = self.runtime.drain_event_outbox().await;
214 }
215 if let Some(event) = &audit_event {
216 dispatch_audit_sink(self.runtime, std::slice::from_ref(event)).await;
217 }
218
219 Ok(record)
220 }
221}