cratestack_sqlx/query/batch/
create.rs1use cratestack_core::{BatchResponse, CratestackContext, CratestackError, ModelEventKind};
6
7use crate::audit::{dispatch_audit_sink, ensure_audit_table};
8use crate::descriptor::ensure_event_outbox_table;
9use crate::{CreateModelInput, ModelDescriptor, SqlxRuntime, sqlx};
10
11use super::create_item::run_create_item;
12use super::validate::validate_batch_size;
13
14#[derive(Debug, Clone)]
15pub struct BatchCreate<'a, M: 'static, PK: 'static, I> {
16 pub(crate) runtime: &'a SqlxRuntime,
17 pub(crate) descriptor: &'static ModelDescriptor<M, PK>,
18 pub(crate) inputs: Vec<I>,
19}
20
21impl<'a, M: 'static, PK: 'static, I> BatchCreate<'a, M, PK, I>
22where
23 I: CreateModelInput<M> + Send,
24{
25 pub async fn run(self, ctx: &CratestackContext) -> Result<BatchResponse<M>, CratestackError>
26 where
27 for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
28 {
29 validate_batch_size(self.inputs.len())?;
30 if self.inputs.is_empty() {
37 return Ok(BatchResponse::from_results(vec![]));
38 }
39
40 let emits_event = self.descriptor.emits(ModelEventKind::Created);
41 let audit_enabled = self.descriptor.audit_enabled;
42
43 let (per_item, audit_events) = crate::bound::in_write_tx!(self.runtime, |tx| async {
46 if emits_event {
47 ensure_event_outbox_table(&mut **tx).await?;
48 }
49 if audit_enabled {
50 ensure_audit_table(self.runtime, &mut **tx).await?;
51 }
52
53 let mut per_item: Vec<Result<M, CratestackError>> =
54 Vec::with_capacity(self.inputs.len());
55 let mut audit_events = Vec::new();
56 for input in self.inputs {
57 let (outcome, audit_event) =
58 run_create_item(tx, self.descriptor, input, ctx, emits_event, audit_enabled)
59 .await?;
60 per_item.push(outcome);
61 audit_events.extend(audit_event);
62 }
63
64 Ok((per_item, audit_events))
65 })?;
66
67 if emits_event {
68 let _ = self.runtime.drain_event_outbox().await;
69 }
70 dispatch_audit_sink(self.runtime, &audit_events).await;
71
72 Ok(BatchResponse::from_results(per_item))
73 }
74}