cratestack_sqlx/query/batch/
delete.rs1use std::collections::HashMap;
6use std::hash::Hash;
7
8use cratestack_core::{
9 AuditOperation, BatchResponse, CratestackContext, CratestackError, ModelEventKind,
10};
11
12use crate::audit::{
13 build_audit_event, dispatch_audit_sink, enqueue_audit_event, ensure_audit_table,
14};
15use crate::descriptor::{enqueue_event_outbox, ensure_event_outbox_table};
16use crate::query::support::push_action_policy_query;
17use crate::{ModelDescriptor, ModelPrimaryKey, SqlxRuntime, cratestack_error_from_sqlx, sqlx};
18
19use super::validate::{reject_duplicate_pks, validate_batch_size};
20
21#[derive(Debug, Clone)]
22pub struct BatchDelete<'a, M: 'static, PK: 'static> {
23 pub(crate) runtime: &'a SqlxRuntime,
24 pub(crate) descriptor: &'static ModelDescriptor<M, PK>,
25 pub(crate) ids: Vec<PK>,
26}
27
28impl<'a, M: 'static, PK: 'static> BatchDelete<'a, M, PK> {
29 pub async fn run(self, ctx: &CratestackContext) -> Result<BatchResponse<M>, CratestackError>
30 where
31 for<'r> M: Send
32 + Unpin
33 + sqlx::FromRow<'r, sqlx::postgres::PgRow>
34 + ModelPrimaryKey<PK>
35 + serde::Serialize,
36 PK: Clone
37 + Eq
38 + Hash
39 + Send
40 + sqlx::Type<sqlx::Postgres>
41 + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
42 {
43 validate_batch_size(self.ids.len())?;
44 reject_duplicate_pks(&self.ids)?;
45 if self.ids.is_empty() {
46 return Ok(BatchResponse::from_results(vec![]));
47 }
48
49 let emits_event = self.descriptor.emits(ModelEventKind::Deleted);
50 let audit_enabled = self.descriptor.audit_enabled;
51
52 let (deleted, audit_events) = crate::bound::in_write_tx!(self.runtime, |tx| async {
55 if emits_event {
56 ensure_event_outbox_table(&mut **tx).await?;
57 }
58 if audit_enabled {
59 ensure_audit_table(self.runtime, &mut **tx).await?;
60 }
61
62 let mut query = sqlx::QueryBuilder::<sqlx::Postgres>::new("");
63 match self.descriptor.soft_delete_column {
64 Some(col) => {
65 query.push("UPDATE ").push(self.descriptor.table_name);
66 query.push(" SET ").push(col).push(" = NOW()");
67 if let Some(version_col) = self.descriptor.version_column {
68 query
69 .push(", ")
70 .push(version_col)
71 .push(" = ")
72 .push(version_col)
73 .push(" + 1");
74 }
75 query.push(" WHERE ").push(col).push(" IS NULL AND ");
76 }
77 None => {
78 query.push("DELETE FROM ").push(self.descriptor.table_name);
79 query.push(" WHERE ");
80 }
81 }
82 query.push(self.descriptor.primary_key).push(" IN (");
83 for (index, id) in self.ids.iter().enumerate() {
84 if index > 0 {
85 query.push(", ");
86 }
87 query.push_bind(id.clone());
88 }
89 query.push(") AND ");
90 push_action_policy_query(
91 &mut query,
92 self.descriptor.delete_allow_policies,
93 self.descriptor.delete_deny_policies,
94 ctx,
95 );
96 query
97 .push(" RETURNING ")
98 .push(self.descriptor.select_projection());
99
100 let deleted: Vec<M> = query
101 .build_query_as::<M>()
102 .fetch_all(&mut **tx)
103 .await
104 .map_err(cratestack_error_from_sqlx)?;
105
106 let mut audit_events = Vec::new();
109 for record in &deleted {
110 if emits_event {
111 enqueue_event_outbox(
112 &mut **tx,
113 self.descriptor.schema_name,
114 ModelEventKind::Deleted,
115 record,
116 )
117 .await?;
118 }
119 if audit_enabled {
120 let before = serde_json::to_value(record).ok();
121 let event = build_audit_event(
122 self.descriptor,
123 AuditOperation::Delete,
124 before,
125 None,
126 ctx,
127 );
128 enqueue_audit_event(&mut **tx, &event).await?;
129 audit_events.push(event);
130 }
131 }
132
133 Ok((deleted, audit_events))
134 })?;
135
136 if emits_event {
137 let _ = self.runtime.drain_event_outbox().await;
138 }
139 dispatch_audit_sink(self.runtime, &audit_events).await;
140
141 let mut by_pk: HashMap<PK, M> = deleted.into_iter().map(|m| (m.primary_key(), m)).collect();
145 let per_item: Vec<Result<M, CratestackError>> = self
146 .ids
147 .into_iter()
148 .map(|id| {
149 by_pk
150 .remove(&id)
151 .ok_or_else(|| CratestackError::NotFound("no row matched".to_owned()))
152 })
153 .collect();
154
155 Ok(BatchResponse::from_results(per_item))
156 }
157}