Skip to main content

cratestack_sqlx/query/write/
delete.rs

1//! `DeleteRecord` — single-row DELETE (soft or hard) with policy +
2//! audit + event fan-out. For a hard delete, the `RETURNING` row IS
3//! the pre-delete state, so it doubles as the audit "before" snapshot
4//! and "after" stays `None` (the row no longer exists). For a soft
5//! delete, `delete_returning_record` actually runs an `UPDATE ...
6//! RETURNING`, so that row is the *post*-tombstone state — it's
7//! captured as "after", and "before" comes from a separate row-locked
8//! fetch taken ahead of the mutation, mirroring how `update.rs` splits
9//! its own before/after snapshots.
10
11use 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    /// Like [`Self::run`] but participates in a caller-supplied
40    /// transaction; no `AuditSink` fan-out happens here for the same
41    /// reason the event outbox isn't drained here — see `create.rs`'s
42    /// `run_in_tx` doc comment.
43    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        // Soft delete is an UPDATE under the hood, so its RETURNING
62        // row is the post-tombstone state — the pre-delete "before"
63        // snapshot has to come from a separate row-locked read taken
64        // ahead of the mutation.
65        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}