cratestack_sqlx/query/write/
delete.rs1use cratestack_core::{AuditOperation, CoolContext, CoolError, ModelEventKind};
12
13use crate::audit::{
14 build_audit_event, dispatch_audit_sink, enqueue_audit_event, ensure_audit_table,
15 fetch_for_audit,
16};
17use crate::descriptor::{enqueue_event_outbox, ensure_event_outbox_table};
18use crate::{ModelDescriptor, SqlxRuntime, cool_error_from_sqlx, sqlx};
19
20use super::delete_exec::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}
28
29impl<'a, M: 'static, PK: 'static> DeleteRecord<'a, M, PK> {
30 pub fn preview_sql(&self) -> String {
31 format!(
32 "DELETE FROM {} WHERE {} = $1 RETURNING {}",
33 self.descriptor.table_name,
34 self.descriptor.primary_key,
35 self.descriptor.select_projection(),
36 )
37 }
38
39 pub async fn run_in_tx<'tx>(
44 self,
45 tx: &mut sqlx::Transaction<'tx, sqlx::Postgres>,
46 ctx: &CoolContext,
47 ) -> Result<M, CoolError>
48 where
49 for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
50 PK: Send + Clone + sqlx::Type<sqlx::Postgres> + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
51 {
52 let emits_event = self.descriptor.emits(ModelEventKind::Deleted);
53 let audit_enabled = self.descriptor.audit_enabled;
54 let soft_delete = self.descriptor.soft_delete_column.is_some();
55 if emits_event {
56 ensure_event_outbox_table(&mut **tx).await?;
57 }
58 if audit_enabled {
59 ensure_audit_table(self.runtime).await?;
60 }
61 let before_record = if audit_enabled && soft_delete {
66 fetch_for_audit(&mut **tx, self.descriptor, self.id.clone()).await?
67 } else {
68 None
69 };
70 let before_snapshot = before_record
71 .as_ref()
72 .and_then(|m| serde_json::to_value(m).ok());
73 let record = delete_returning_record(&mut **tx, self.descriptor, self.id, ctx).await?;
74 if emits_event {
75 enqueue_event_outbox(
76 &mut **tx,
77 self.descriptor.schema_name,
78 ModelEventKind::Deleted,
79 &record,
80 )
81 .await?;
82 }
83 if audit_enabled {
84 let (before, after) = if soft_delete {
85 (before_snapshot, serde_json::to_value(&record).ok())
86 } else {
87 (serde_json::to_value(&record).ok(), None)
88 };
89 let event =
90 build_audit_event(self.descriptor, AuditOperation::Delete, before, after, ctx);
91 enqueue_audit_event(&mut **tx, &event).await?;
92 }
93 Ok(record)
94 }
95
96 pub async fn run(self, ctx: &CoolContext) -> Result<M, CoolError>
97 where
98 for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
99 PK: Send + Clone + sqlx::Type<sqlx::Postgres> + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
100 {
101 let emits_event = self.descriptor.emits(ModelEventKind::Deleted);
102 let audit_enabled = self.descriptor.audit_enabled;
103 let soft_delete = self.descriptor.soft_delete_column.is_some();
104 let needs_tx = emits_event || audit_enabled;
105 let mut audit_event = None;
106 let record = if needs_tx {
107 let mut tx = self
108 .runtime
109 .pool()
110 .begin()
111 .await
112 .map_err(cool_error_from_sqlx)?;
113 if emits_event {
114 ensure_event_outbox_table(&mut *tx).await?;
115 }
116 if audit_enabled {
117 ensure_audit_table(self.runtime).await?;
118 }
119
120 let before_record = if audit_enabled && soft_delete {
121 fetch_for_audit(&mut *tx, self.descriptor, self.id.clone()).await?
122 } else {
123 None
124 };
125 let before_snapshot = before_record
126 .as_ref()
127 .and_then(|m| serde_json::to_value(m).ok());
128 let record = delete_returning_record(&mut *tx, self.descriptor, self.id, ctx).await?;
129 if emits_event {
130 enqueue_event_outbox(
131 &mut *tx,
132 self.descriptor.schema_name,
133 ModelEventKind::Deleted,
134 &record,
135 )
136 .await?;
137 }
138 if audit_enabled {
139 let (before, after) = if soft_delete {
140 (before_snapshot, serde_json::to_value(&record).ok())
141 } else {
142 (serde_json::to_value(&record).ok(), None)
143 };
144 let event =
145 build_audit_event(self.descriptor, AuditOperation::Delete, before, after, ctx);
146 enqueue_audit_event(&mut *tx, &event).await?;
147 audit_event = Some(event);
148 }
149 tx.commit().await.map_err(cool_error_from_sqlx)?;
150 record
151 } else {
152 delete_returning_record(self.runtime.pool(), self.descriptor, self.id, ctx).await?
153 };
154
155 if emits_event {
156 let _ = self.runtime.drain_event_outbox().await;
157 }
158 if let Some(event) = &audit_event {
159 dispatch_audit_sink(self.runtime, std::slice::from_ref(event)).await;
160 }
161
162 Ok(record)
163 }
164}