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 {}