Skip to main content

cratestack_sqlx/query/batch/
update.rs

1//! `batch_update` driver — opens the outer tx, fans the items out to
2//! [`super::update_item::run_update_item`] (one savepoint per item),
3//! commits, drains the outbox.
4
5use std::hash::Hash;
6
7use cratestack_core::{BatchResponse, CratestackContext, CratestackError, ModelEventKind};
8
9use crate::audit::{dispatch_audit_sink, ensure_audit_table};
10use crate::descriptor::ensure_event_outbox_table;
11use crate::{ModelDescriptor, SqlxRuntime, UpdateModelInput, sqlx};
12
13use super::update_item::run_update_item;
14use super::validate::{reject_duplicate_pks, validate_batch_size};
15
16/// One per-item update: `(id, patch, optional expected version)`.
17pub type BatchUpdateItem<PK, I> = (PK, I, Option<i64>);
18
19#[derive(Debug, Clone)]
20pub struct BatchUpdate<'a, M: 'static, PK: 'static, I> {
21    pub(crate) runtime: &'a SqlxRuntime,
22    pub(crate) descriptor: &'static ModelDescriptor<M, PK>,
23    pub(crate) items: Vec<BatchUpdateItem<PK, I>>,
24}
25
26impl<'a, M: 'static, PK: 'static, I> BatchUpdate<'a, M, PK, I>
27where
28    I: UpdateModelInput<M> + Send,
29{
30    pub async fn run(self, ctx: &CratestackContext) -> Result<BatchResponse<M>, CratestackError>
31    where
32        for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
33        PK: Clone
34            + Eq
35            + Hash
36            + Send
37            + sqlx::Type<sqlx::Postgres>
38            + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
39    {
40        validate_batch_size(self.items.len())?;
41        let ids: Vec<PK> = self.items.iter().map(|(id, _, _)| id.clone()).collect();
42        reject_duplicate_pks(&ids)?;
43        if self.items.is_empty() {
44            return Ok(BatchResponse::from_results(vec![]));
45        }
46
47        let emits_event = self.descriptor.emits(ModelEventKind::Updated);
48        let audit_enabled = self.descriptor.audit_enabled;
49
50        // A savepoint of the `@isolation` transaction when there is one
51        // (docs/design/procedure-isolation.md §4).
52        let (per_item, audit_events) = crate::bound::in_write_tx!(self.runtime, |tx| async {
53            if emits_event {
54                ensure_event_outbox_table(&mut **tx).await?;
55            }
56            if audit_enabled {
57                ensure_audit_table(self.runtime, &mut **tx).await?;
58            }
59
60            let mut per_item: Vec<Result<M, CratestackError>> =
61                Vec::with_capacity(self.items.len());
62            let mut audit_events = Vec::new();
63            for item in self.items {
64                let (outcome, audit_event) =
65                    run_update_item(tx, self.descriptor, item, ctx, emits_event, audit_enabled)
66                        .await?;
67                per_item.push(outcome);
68                audit_events.extend(audit_event);
69            }
70
71            Ok((per_item, audit_events))
72        })?;
73
74        if emits_event {
75            let _ = self.runtime.drain_event_outbox().await;
76        }
77        dispatch_audit_sink(self.runtime, &audit_events).await;
78
79        Ok(BatchResponse::from_results(per_item))
80    }
81}