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::{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    /// Like [`Self::run`] but participates in a caller-supplied transaction.
37    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        // Soft delete is an UPDATE under the hood, so its RETURNING
56        // row is the post-tombstone state — the pre-delete "before"
57        // snapshot has to come from a separate row-locked read taken
58        // ahead of the mutation.
59        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}