Skip to main content

es_entity/operation/
mod.rs

1//! Handle execution of database operations and transactions.
2
3pub mod hooks;
4mod with_time;
5
6use sqlx::{Acquire, Transaction};
7
8use crate::{clock::ClockHandle, db, one_time_executor::OneTimeExecutor};
9
10pub use with_time::*;
11
12/// Default return type of the derived EsRepo::begin_op().
13///
14/// Used as a wrapper of a [`sqlx::Transaction`] but can also cache the time at which the
15/// transaction is taking place.
16///
17/// When a manual clock is provided, the transaction will automatically cache that
18/// clock's time, enabling deterministic testing. This cached time will be used in all
19/// time-dependent operations.
20pub struct DbOp<'c> {
21    tx: Transaction<'c, db::Db>,
22    clock: ClockHandle,
23    now: Option<chrono::DateTime<chrono::Utc>>,
24    commit_hooks: Option<hooks::CommitHooks>,
25}
26
27impl<'c> DbOp<'c> {
28    fn new(
29        tx: Transaction<'c, db::Db>,
30        clock: ClockHandle,
31        time: Option<chrono::DateTime<chrono::Utc>>,
32    ) -> Self {
33        Self {
34            tx,
35            clock,
36            now: time,
37            commit_hooks: Some(hooks::CommitHooks::new()),
38        }
39    }
40
41    /// Initializes a transaction using the global clock.
42    ///
43    /// Delegates to [`init_with_clock`](Self::init_with_clock) using the global clock handle.
44    pub async fn init(pool: &db::Pool) -> Result<DbOp<'static>, sqlx::Error> {
45        Self::init_with_clock(pool, crate::clock::Clock::handle()).await
46    }
47
48    /// Initializes a transaction with the specified clock.
49    ///
50    /// If the clock is manual, its current time will be cached in the transaction.
51    pub async fn init_with_clock(
52        pool: &db::Pool,
53        clock: &ClockHandle,
54    ) -> Result<DbOp<'static>, sqlx::Error> {
55        let tx = pool.begin().await?;
56
57        // If a manual clock is provided, cache its time for consistent
58        // timestamps within the transaction.
59        let time = clock.manual_now();
60
61        Ok(DbOp::new(tx, clock.clone(), time))
62    }
63
64    /// Transitions to a [`DbOpWithTime`] with the given time cached.
65    pub fn with_time(self, time: chrono::DateTime<chrono::Utc>) -> DbOpWithTime<'c> {
66        DbOpWithTime::new(self, time)
67    }
68
69    /// Transitions to a [`DbOpWithTime`] using the clock.
70    ///
71    /// Uses cached time if present, otherwise uses the clock's current time.
72    pub fn with_clock_time(self) -> DbOpWithTime<'c> {
73        let time = self.now.unwrap_or_else(|| self.clock.now());
74        DbOpWithTime::new(self, time)
75    }
76
77    /// Transitions to a [`DbOpWithTime`] using the database time.
78    ///
79    /// Priority order:
80    /// 1. Cached time if present
81    /// 2. Manual clock time if the clock is manual
82    /// 3. Database time via `SELECT NOW()`
83    pub async fn with_db_time(mut self) -> Result<DbOpWithTime<'c>, sqlx::Error> {
84        let time = if let Some(time) = self.now {
85            time
86        } else if let Some(manual_time) = self.clock.manual_now() {
87            manual_time
88        } else {
89            db::database_now(&mut *self.tx).await?
90        };
91
92        Ok(DbOpWithTime::new(self, time))
93    }
94
95    /// Returns the optionally cached [`chrono::DateTime`]
96    pub fn maybe_now(&self) -> Option<chrono::DateTime<chrono::Utc>> {
97        self.now
98    }
99
100    /// Begins a nested transaction.
101    pub async fn begin(&mut self) -> Result<DbOp<'_>, sqlx::Error> {
102        Ok(DbOp::new(
103            self.tx.begin().await?,
104            self.clock.clone(),
105            self.now,
106        ))
107    }
108
109    /// Commits the inner transaction.
110    pub async fn commit(mut self) -> Result<(), sqlx::Error> {
111        let commit_hooks = self.commit_hooks.take().expect("no hooks");
112        let post_hooks = commit_hooks.execute_pre(&mut self).await?;
113        self.tx.commit().await?;
114        post_hooks.execute();
115        Ok(())
116    }
117
118    /// Gets a mutable handle to the inner transaction
119    pub fn tx_mut(&mut self) -> &mut Transaction<'c, db::Db> {
120        &mut self.tx
121    }
122}
123
124impl<'o> AtomicOperation for DbOp<'o> {
125    fn maybe_now(&self) -> Option<chrono::DateTime<chrono::Utc>> {
126        self.maybe_now()
127    }
128
129    fn clock(&self) -> &ClockHandle {
130        &self.clock
131    }
132
133    fn connection(&mut self) -> &mut db::Connection {
134        self.tx.connection()
135    }
136
137    fn add_commit_hook<H: hooks::CommitHook>(&mut self, hook: H) -> Result<(), H> {
138        self.commit_hooks.as_mut().expect("no hooks").add(hook);
139        Ok(())
140    }
141
142    fn commit_hook<H: hooks::CommitHook>(&self) -> Option<&H> {
143        self.commit_hooks.as_ref()?.get_last::<H>()
144    }
145
146    fn supports_hooks(&self) -> bool {
147        true
148    }
149}
150
151/// Equivileant of [`DbOp`] just that the time is guaranteed to be cached.
152///
153/// Used as a wrapper of a [`sqlx::Transaction`] with cached time of the transaction.
154pub struct DbOpWithTime<'c> {
155    inner: DbOp<'c>,
156    now: chrono::DateTime<chrono::Utc>,
157}
158
159impl<'c> DbOpWithTime<'c> {
160    fn new(mut inner: DbOp<'c>, time: chrono::DateTime<chrono::Utc>) -> Self {
161        inner.now = Some(time);
162        Self { inner, now: time }
163    }
164
165    /// The cached [`chrono::DateTime`]
166    pub fn now(&self) -> chrono::DateTime<chrono::Utc> {
167        self.now
168    }
169
170    /// Begins a nested transaction.
171    pub async fn begin(&mut self) -> Result<DbOpWithTime<'_>, sqlx::Error> {
172        Ok(DbOpWithTime::new(self.inner.begin().await?, self.now))
173    }
174
175    /// Commits the inner transaction.
176    pub async fn commit(self) -> Result<(), sqlx::Error> {
177        self.inner.commit().await
178    }
179
180    /// Gets a mutable handle to the inner transaction
181    pub fn tx_mut(&mut self) -> &mut Transaction<'c, db::Db> {
182        self.inner.tx_mut()
183    }
184}
185
186impl<'o> AtomicOperation for DbOpWithTime<'o> {
187    fn maybe_now(&self) -> Option<chrono::DateTime<chrono::Utc>> {
188        Some(self.now())
189    }
190
191    fn clock(&self) -> &ClockHandle {
192        self.inner.clock()
193    }
194
195    fn connection(&mut self) -> &mut db::Connection {
196        self.inner.connection()
197    }
198
199    fn add_commit_hook<H: hooks::CommitHook>(&mut self, hook: H) -> Result<(), H> {
200        self.inner.add_commit_hook(hook)
201    }
202
203    fn commit_hook<H: hooks::CommitHook>(&self) -> Option<&H> {
204        self.inner.commit_hook::<H>()
205    }
206
207    fn supports_hooks(&self) -> bool {
208        self.inner.supports_hooks()
209    }
210}
211
212impl<'o> AtomicOperationWithTime for DbOpWithTime<'o> {
213    fn now(&self) -> chrono::DateTime<chrono::Utc> {
214        self.now
215    }
216}
217
218/// Trait to signify we can make multiple consistent database roundtrips.
219///
220/// Its a stand in for [`&mut sqlx::Transaction<'_, DB>`](`sqlx::Transaction`).
221/// The reason for having a trait is to support custom types that wrap the inner
222/// transaction while providing additional functionality.
223///
224/// See [`DbOp`] or [`DbOpWithTime`].
225pub trait AtomicOperation: Send {
226    /// Function for querying when the operation is taking place - if it is cached.
227    fn maybe_now(&self) -> Option<chrono::DateTime<chrono::Utc>> {
228        None
229    }
230
231    /// Returns the clock handle for time operations.
232    ///
233    /// Default implementation returns the global clock handle.
234    fn clock(&self) -> &ClockHandle {
235        crate::clock::Clock::handle()
236    }
237
238    /// Returns the raw underlying connection.
239    /// The desired way to represent this would actually be as a GAT:
240    /// ```rust
241    /// trait AtomicOperation {
242    ///     type Executor<'c>: sqlx::PgExecutor<'c>
243    ///         where Self: 'c;
244    ///
245    ///     fn connection<'c>(&'c mut self) -> Self::Executor<'c>;
246    /// }
247    /// ```
248    ///
249    /// But GATs don't play well with `async_trait::async_trait` due to lifetime constraints
250    /// so we return the concrete [`&mut db::Connection`](`crate::db::Connection`) instead as a work around.
251    ///
252    /// Since this trait is generally applied to types that wrap a [`sqlx::Transaction`]
253    /// there is no variance in the return type - so its fine.
254    ///
255    /// Statements executed directly on the returned connection are **not**
256    /// annotated with trace context — use [`as_executor`](Self::as_executor)
257    /// unless raw connection access is required.
258    fn connection(&mut self) -> &mut db::Connection;
259
260    /// Returns the [`sqlx::Executor`] implementation that statements should be
261    /// executed through.
262    ///
263    /// The returned [`OneTimeExecutor`] annotates every statement with the
264    /// current span's `traceparent` SQL comment when the `tracing-context`
265    /// feature is enabled and a *sampled* span is active (see
266    /// [`crate::sql_commenter`]). Otherwise statements pass through untouched.
267    ///
268    /// Trade-off: the trace context makes annotated statement text unique, so
269    /// annotated statements bypass sqlx's per-connection prepared statement
270    /// cache (`persistent(false)`) — costing a server-side parse + plan per
271    /// execution. Un-annotated traffic keeps full prepared-statement reuse.
272    fn as_executor(&mut self) -> OneTimeExecutor<'_, &mut db::Connection> {
273        let now = self.maybe_now();
274        OneTimeExecutor::new(self.connection(), now)
275    }
276
277    /// Registers a commit hook that will run pre_commit before and post_commit after the transaction commits.
278    /// Returns Ok(()) if the hook was registered, Err(hook) if hooks are not supported.
279    fn add_commit_hook<H: hooks::CommitHook>(&mut self, hook: H) -> Result<(), H> {
280        Err(hook)
281    }
282
283    /// Typed shared access to the currently-accumulating commit hook of type `H`,
284    /// if this operation supports commit hooks and one is registered.
285    /// Returns the hook a subsequent `add_commit_hook::<H>` call would merge into.
286    fn commit_hook<H: hooks::CommitHook>(&self) -> Option<&H> {
287        None
288    }
289
290    /// Whether this operation supports commit hooks.
291    ///
292    /// `true` iff [`add_commit_hook`](Self::add_commit_hook) can register a hook
293    /// (i.e. the operation is backed by a [`DbOp`]-style commit-hook buffer, not
294    /// a bare [`sqlx::Transaction`]). Unlike [`commit_hook`](Self::commit_hook) —
295    /// whose `None` is ambiguous between "hooks unsupported" and "supported but
296    /// none registered yet" — this reports support directly, with no registration
297    /// attempt and no `&mut` access.
298    fn supports_hooks(&self) -> bool {
299        false
300    }
301}
302
303impl<'c> AtomicOperation for sqlx::Transaction<'c, db::Db> {
304    fn connection(&mut self) -> &mut db::Connection {
305        &mut *self
306    }
307}