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