Skip to main content

es_entity/operation/batch/
mod.rs

1//! Running a batch of items in one transaction while isolating each failure.
2//!
3//! Two shapes, both built on [`SavepointOperation::with_savepoint`] and both
4//! available on every [`AtomicOperation`](super::AtomicOperation):
5//!
6//! - [`run_isolated`](BatchIsolation::run_isolated) — one savepoint per item.
7//!   Use it when each item needs its own logic and its own error isolation.
8//! - [`run_bisected`](BatchIsolation::run_bisected) — one probe over the whole
9//!   slice, splitting only on failure. Use it when the closure handles a whole
10//!   slice in set-based statements (`create_all` / `update_all` and friends),
11//!   so the happy path costs one probe.
12//!
13//! # The closure
14//!
15//! `AsyncFnOnce + Clone + Sync`, cloned once per probe so that each probe calls
16//! its own clone exactly once.
17//!
18//! Each such call has a single opaque future type, which auto-trait inference
19//! resolves at the call site. A closure that borrows `&self` therefore composes
20//! inside an `#[async_trait]` runner or a `tokio::spawn`.
21//!
22//! Cloning the closure clones its captures, so state shared *across* probes
23//! belongs behind something whose clone is the same underlying value: `Arc<…>`,
24//! or a borrowed `&Mutex<…>`. `Sync` holds callers to that — `Cell`, `RefCell`
25//! and `Rc` are not `Sync`, so the bound rejects them.
26//!
27//! # What `f` must tolerate
28//!
29//! A bisect calls `f` repeatedly over **arbitrary contiguous sub-slices, in
30//! non-positional order**, and may re-probe the same range after a transient
31//! failure. `f` must therefore be a function of the set it is handed, and of
32//! the database state at that moment: only `*_in_op` work against the probe's
33//! operation, so that a rollback undoes all of it.
34//!
35//! # Hooks
36//!
37//! Commit hooks registered inside a probe are staged on its savepoint: folded
38//! outward on success, dropped on failure. No hook callback runs at a savepoint
39//! boundary, so a rolled-back probe contributes exactly zero hook state to
40//! match its zero database state. See [`SavepointOp`].
41//!
42//! # Observability
43//!
44//! Tracing is left to the caller, since es-entity's `tracing` dependency is
45//! optional. [`BisectOutcomes`] carries `probes_used` and `transient_retries`
46//! for reporting under the caller's own target and field names.
47
48mod search;
49mod transient;
50
51use std::future::Future;
52
53use super::{SavepointOp, SavepointOperation};
54use crate::errlanes::{Fatal, FatalKind, Fault, lanes};
55
56pub use search::*;
57pub use transient::*;
58
59/// Batch isolation for every [`AtomicOperation`](super::AtomicOperation).
60///
61/// Blanket-implemented, like [`SavepointOperation`], so `DbOp`, `SavepointOp`
62/// (nesting), `HookOperation`, and operation types defined outside this crate
63/// all get it without naming a concrete type — and so a function generic over
64/// `impl AtomicOperation` can use it.
65pub trait BatchIsolation: SavepointOperation {
66    /// Runs `f` once per item, each inside its own `SAVEPOINT`, in item order.
67    ///
68    /// A failing item unwinds only its own writes and staged hooks; the
69    /// transaction stays usable and the loop continues, so its healthy
70    /// batch-mates still commit. Outcomes are returned positionally: one entry
71    /// per input, `f`'s own `Ok`/`Err` preserved.
72    ///
73    /// The outer `Err` means the savepoint machinery itself failed, leaving
74    /// the enclosing transaction in an indeterminate state: abandon it. The
75    /// `sqlx::Error` is classified by errlanes' sqlx lane table, so a lost
76    /// connection or a deadlock on the savepoint statement itself arrives as
77    /// `Transient` — retry the whole transaction — and anything else as
78    /// `Fatal`.
79    fn run_isolated<'a, T, V, E, F>(
80        &'a mut self,
81        items: &'a [T],
82        f: F,
83    ) -> impl Future<Output = Result<Vec<Result<V, E>>, Fault<lanes!(Transient, Fatal)>>> + 'a
84    where
85        T: 'a,
86        V: 'a,
87        E: 'a,
88        F: AsyncFnOnce(&mut SavepointOp<'_>, &T) -> Result<V, E> + Clone + Sync + 'a,
89    {
90        async move {
91            let mut outcomes = Vec::with_capacity(items.len());
92            for item in items {
93                let f = f.clone();
94                outcomes.push(self.with_savepoint(async |sp| f(sp, item).await).await?);
95            }
96            Ok(outcomes)
97        }
98    }
99
100    /// Probes the whole slice at once, bisecting only on failure.
101    ///
102    /// The happy path costs **one** probe. On failure the slice splits and
103    /// pending ranges are probed largest-first (earliest start breaking ties)
104    /// until `budget` is spent, so clean siblings are salvaged and a culprit
105    /// resolves from its own single-item probe, where its error is
106    /// attributable to it.
107    ///
108    /// Deadlock victims and serialization failures re-probe the same range
109    /// unsplit and are refunded to the budget — see [`TransientPolicy`]. Use
110    /// [`run_bisected_with`](Self::run_bisected_with) to widen that class.
111    fn run_bisected<'a, T, E, F>(
112        &'a mut self,
113        items: &'a [T],
114        budget: BisectBudget,
115        f: F,
116    ) -> impl Future<Output = Result<BisectOutcomes<E>, Fault<lanes!(Transient, Fatal)>>> + 'a
117    where
118        T: 'a,
119        E: std::error::Error + 'static,
120        F: AsyncFnOnce(&mut SavepointOp<'_>, &[T]) -> Result<(), E> + Clone + Sync + 'a,
121    {
122        self.run_bisected_with(
123            items,
124            budget,
125            TransientPolicy::new(is_contention::<E> as fn(&E) -> bool),
126            f,
127        )
128    }
129
130    /// [`run_bisected`](Self::run_bisected) with a caller-supplied notion of
131    /// which failures are transient.
132    ///
133    /// Classification being the caller's, the error bound here is just
134    /// [`Display`](std::fmt::Display).
135    ///
136    /// The outer `Err` is the search's own failure, never one of `f`'s:
137    /// the savepoint machinery failed (classified as in
138    /// [`run_isolated`](Self::run_isolated)), or the transient allowance ran
139    /// out before anything could be attributed to a range —
140    /// `Fatal(Exhausted)`, since the retries have already been spent here;
141    /// the caller re-runs the whole batch later rather than automatically.
142    fn run_bisected_with<'a, T, E, F, P>(
143        &'a mut self,
144        items: &'a [T],
145        budget: BisectBudget,
146        policy: TransientPolicy<P>,
147        f: F,
148    ) -> impl Future<Output = Result<BisectOutcomes<E>, Fault<lanes!(Transient, Fatal)>>> + 'a
149    where
150        T: 'a,
151        E: std::fmt::Display + 'a,
152        P: Fn(&E) -> bool + 'a,
153        F: AsyncFnOnce(&mut SavepointOp<'_>, &[T]) -> Result<(), E> + Clone + Sync + 'a,
154    {
155        async move {
156            let mut search = BisectSearch::new(items.len(), budget)
157                .with_max_transient_retries(policy.max_retries);
158
159            while let Some(range) = search.next_range() {
160                let f = f.clone();
161                let slice = &items[range.clone()];
162
163                let verdict = match self.with_savepoint(async |sp| f(sp, slice).await).await? {
164                    Ok(()) => ProbeVerdict::Clean,
165                    Err(error) if (policy.is_transient)(&error) => ProbeVerdict::Transient(error),
166                    Err(error) => ProbeVerdict::Failed(error),
167                };
168
169                if let Err(limit) = search.report(range, verdict) {
170                    // The search learned nothing about the items, so there are
171                    // no per-item verdicts to return — the caller re-runs the
172                    // whole batch.
173                    return Err(Fatal::from_error(FatalKind::Exhausted, limit).into());
174                }
175            }
176
177            Ok(search.into_outcomes())
178        }
179    }
180}
181
182impl<T: SavepointOperation + ?Sized> BatchIsolation for T {}