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, 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    /// Expected version for optimistic locking. Required on models
32    /// that declare `@version`; ignored otherwise. Mirrors
33    /// [`crate::UpdateRecordSet::if_match`] — see that doc comment for
34    /// the rationale of an `Option<i64>` builder step rather than a
35    /// required constructor argument.
36    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    /// Like [`Self::run`] but participates in a caller-supplied
60    /// transaction. Neither the `AuditSink` fan-out nor the event
61    /// outbox drain happens here — see `create.rs`'s `run_in_tx` doc
62    /// comment for the full contract and how a caller opts into both
63    /// after their own commit.
64    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        // Soft delete is an UPDATE under the hood, so its RETURNING
88        // row is the post-tombstone state — the pre-delete "before"
89        // snapshot has to come from a separate row-locked read taken
90        // ahead of the mutation.
91        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        // Inside an `@isolation` procedure: run on its transaction and
134        // defer the post-commit fan-out (docs/design/procedure-isolation.md
135        // §4, §6).
136        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}