Skip to main content

pylon_client/
client.rs

1//
2// This source file is part of the Pylon open source project.
3//
4// Copyright (c) 2026 Jaldis B.V.
5//
6// Licensed under the MIT OR Apache-2.0 license (the "License");
7// you may not use this file except in compliance with the License.
8// You may obtain a copy of the License at
9//
10//     https://opensource.org/licenses/MIT
11//     https://www.apache.org/licenses/LICENSE-2.0
12//
13// Unless required by applicable law or agreed to in writing, software
14// distributed under the License is distributed on an "AS IS" BASIS,
15// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16// See the License for the specific language governing permissions and
17// limitations under the License.
18//
19
20//! [`Client`] construction, connection, globals/config, and query methods —
21//! see `crate::transaction` for the closure-based retrying transaction API.
22
23use std::collections::HashMap;
24use std::future::Future;
25use std::path::PathBuf;
26use std::pin::Pin;
27use std::sync::{Arc, RwLock};
28use std::time::Duration;
29
30use pylon_core::ir::SessionConfig;
31use pylon_core::schema::SchemaDescriptor;
32use pylon_value::DecodedValue;
33use tokio::sync::OnceCell;
34
35use crate::error::{Error, Result};
36use crate::exec;
37use crate::query_arg::QueryArgs;
38use crate::queryable::{Queryable, decode_optional_row, decode_row, decode_rows};
39use crate::schema;
40use crate::transaction::{Isolation, Transaction};
41
42/// A transaction attempt body's return type — boxed since a plain generic
43/// `Fut: Future` can't express "this future borrows the `&Transaction` it
44/// was handed" without higher-ranked lifetimes on the future type itself;
45/// boxing sidesteps that and matches how most async-closure-taking APIs in
46/// the ecosystem handle exactly this shape.
47pub type TxFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T>> + Send + 'a>>;
48
49/// Either a `(path, max_size_mb)` this `Builder` should open itself, or an
50/// already-open handle a caller wants shared in as-is — see
51/// `Builder::cache`/`Builder::cache_handle`.
52enum CacheSource {
53    Open { path: PathBuf, max_size_mb: usize },
54    Shared(Arc<pylon_cache::Cache>),
55}
56
57/// Builds a [`Client`]. `dsn` is required up front (this crate never
58/// parses `pylon.toml` — see the crate-level docs); everything else has a
59/// sensible default.
60pub struct Builder {
61    dsn: String,
62    max_pool_size: usize,
63    cache: Option<CacheSource>,
64}
65
66impl Builder {
67    pub fn new(dsn: impl Into<String>) -> Self {
68        Self {
69            dsn: dsn.into(),
70            max_pool_size: 10,
71            cache: None,
72        }
73    }
74
75    pub fn max_pool_size(mut self, max_pool_size: usize) -> Self {
76        self.max_pool_size = max_pool_size;
77        self
78    }
79
80    /// Opts into read-through result caching at an LMDB-backed directory —
81    /// mirrors `pylon.toml`'s `[cache]` section, minus per-type
82    /// (`[cache.sets.<Name>]`) overrides: caching here is a single global
83    /// on/off. Omit this entirely for no caching (today's default
84    /// behavior). Only `Client`'s own query methods read/write the cache —
85    /// `Transaction` never populates it, since rows read inside a
86    /// transaction aren't committed, matching `pylon/client.py`'s
87    /// `AsyncTransaction`.
88    ///
89    /// This client's *own* writes evict the tags they touch, from inside a
90    /// transaction too, so a program that writes and then re-reads sees its
91    /// own change. Writes from *other* processes still need a listener on
92    /// the same directory (e.g. `pylon worker start`) to invalidate.
93    ///
94    /// Opens its own LMDB handle for `path` — `heed` (the LMDB binding this
95    /// crate uses) refuses a second `Env::open` on the same canonicalized
96    /// path while an earlier handle onto it is still alive within the same
97    /// process, so a caller building more than one `Client` that should
98    /// share one cache directory (e.g. one per named `pylon.toml`
99    /// connection) must use `Builder::cache_handle` with one already-open
100    /// `Arc<pylon_cache::Cache>` instead of calling this per client.
101    pub fn cache(mut self, path: impl Into<PathBuf>, max_size_mb: usize) -> Self {
102        self.cache = Some(CacheSource::Open {
103            path: path.into(),
104            max_size_mb,
105        });
106        self
107    }
108
109    /// Like `Builder::cache`, but attaches to an already-open handle instead
110    /// of opening a new one — see that method's doc comment for why this
111    /// exists.
112    pub fn cache_handle(mut self, cache: Arc<pylon_cache::Cache>) -> Self {
113        self.cache = Some(CacheSource::Shared(cache));
114        self
115    }
116
117    /// Builds the client without touching the database: the connection pool
118    /// and the schema snapshot are opened together on first use, or on an
119    /// explicit [`Client::ensure_connected`]. A `Client` can therefore be
120    /// constructed outside an async context, and a configured-but-never-queried
121    /// connection costs nothing.
122    ///
123    /// The cache, when [`Builder::cache`] configured one, *is* opened here —
124    /// it's a local LMDB directory rather than a network resource, and a bad
125    /// path is worth failing on at construction rather than on whichever
126    /// query happens to run first.
127    ///
128    /// A process that wants a bad DSN/host/credentials — or a database where
129    /// neither `pylon migration apply` nor `pylon migration watch` has ever
130    /// run, so there is no schema snapshot to fetch — to fail at startup
131    /// instead of on its first query should call [`Client::ensure_connected`]
132    /// right after this.
133    pub fn build(self) -> Result<Client> {
134        let cache = match self.cache {
135            None => None,
136            Some(CacheSource::Open { path, max_size_mb }) => Some(Arc::new(
137                pylon_cache::Cache::open(&path, max_size_mb).map_err(|e| Error::Cache(e.to_string()))?,
138            )),
139            Some(CacheSource::Shared(cache)) => Some(cache),
140        };
141        Ok(Client {
142            dsn: Arc::new(self.dsn),
143            max_pool_size: self.max_pool_size,
144            connected: Arc::new(OnceCell::new()),
145            globals: Arc::new(HashMap::new()),
146            config: SessionConfig::default(),
147            cache,
148        })
149    }
150}
151
152/// The half of a [`Client`] that only exists once it has actually reached
153/// the database — held behind a `OnceCell` so construction stays cheap and
154/// synchronous. The pool and the schema snapshot are initialised together
155/// because a client with one and not the other can't serve a query anyway.
156struct Connected {
157    pool: pylon_pgcon::PgPool,
158    /// An `Arc` rather than a plain `RwLock` so `Client::transaction` can
159    /// hand a `Transaction` its own handle on the same schema slot without
160    /// borrowing from the `OnceCell` for the transaction's whole lifetime.
161    schema: Arc<RwLock<SchemaDescriptor>>,
162}
163
164/// An async Pylon client — a connection pool plus a compiled schema,
165/// shared cheaply across every [`Client::with_globals`]/[`Client::with_config`]
166/// view of it (mirrors `pylon/client.py`'s own shared-pool-ref pattern).
167///
168/// Connecting is lazy: [`Builder::build`] reaches nothing over the network,
169/// and the first query (or an explicit [`Client::ensure_connected`]) opens
170/// the pool and fetches the schema. Clones — including every
171/// `with_globals`/`with_config` view — share one connection, so connecting
172/// through any of them connects all of them, exactly as
173/// `pylon/client.py`'s `_PoolRef` is shared across its own views.
174#[derive(Clone)]
175pub struct Client {
176    /// Used to open the pool on first use, and on every `Client::listen()`
177    /// call — a `LISTEN` subscription needs its own dedicated (non-pooled)
178    /// connection, opened fresh from this DSN each time.
179    dsn: Arc<String>,
180    max_pool_size: usize,
181    /// Shared across clones so that all of them see one pool and one schema
182    /// slot. `tokio::sync::OnceCell` (not `std`'s) because initialising it
183    /// has to await; it serialises concurrent initialisers, so N requests
184    /// racing to be the first make one connection attempt between them, and
185    /// it stays empty when an attempt fails, so a database that is merely
186    /// not up yet is retried by the next query rather than poisoning the
187    /// client for good.
188    connected: Arc<OnceCell<Connected>>,
189    globals: Arc<HashMap<String, DecodedValue>>,
190    config: SessionConfig,
191    /// `None` unless `Builder::cache` was called — read-through caching is
192    /// opt-in. Shared across `with_globals`/`with_config` clones, same as
193    /// `connected`. Opened eagerly by `Builder::build`, since it's local.
194    cache: Option<Arc<pylon_cache::Cache>>,
195}
196
197impl Client {
198    pub fn builder(dsn: impl Into<String>) -> Builder {
199        Builder::new(dsn)
200    }
201
202    /// The pool and schema, opening them on first use.
203    ///
204    /// A client is reached from request handlers, background workers and CLI
205    /// commands alike, which share no startup between them to connect from,
206    /// so every query method goes through here rather than making callers
207    /// remember an explicit connect step — the same reasoning as
208    /// `pylon/client.py`'s `_connected_pool`.
209    ///
210    /// Note that every caller clones the schema out from behind the lock
211    /// rather than holding the guard: a `std::sync::RwLockReadGuard` isn't
212    /// `Send`, and holding one across an `.await` point would make the
213    /// resulting future non-`Send` (fatal for `Client::transaction`'s boxed
214    /// futures, and a footgun on a multi-threaded runtime generally).
215    async fn connected(&self) -> Result<&Connected> {
216        self.connected
217            .get_or_try_init(|| async {
218                let pool = pylon_pgcon::PgPool::connect(&self.dsn, self.max_pool_size)
219                    .await
220                    .map_err(Error::Db)?;
221                let schema = schema::fetch(&pool).await?;
222                Ok(Connected {
223                    pool,
224                    schema: Arc::new(RwLock::new(schema)),
225                })
226            })
227            .await
228    }
229
230    /// Opens the connection pool and fetches the schema snapshot if that
231    /// hasn't happened yet. Safe to call repeatedly; after the first success
232    /// it costs one atomic load.
233    ///
234    /// Queries connect on their own, so this is never required — it exists
235    /// for a process that would rather learn about an unreachable database,
236    /// bad credentials or a database with no schema snapshot
237    /// ([`Error::NoSchemaSnapshot`]) at startup than on whichever request
238    /// arrives first. Mirrors `pylon/client.py`'s `Client.ensure_connected`.
239    pub async fn ensure_connected(&self) -> Result<()> {
240        self.connected().await.map(|_| ())
241    }
242
243    /// Re-fetches the schema snapshot from `_pylon."Schema"`. Visible to
244    /// every clone sharing this client's pool (`with_globals`/`with_config`
245    /// views included) — there's only one schema slot per underlying
246    /// connection pool, matching `pylon/client.py`'s single process-level
247    /// singleton.
248    pub async fn reload_schema(&self) -> Result<()> {
249        let conn = self.connected().await?;
250        let fresh = schema::fetch(&conn.pool).await?;
251        *conn.schema.write().unwrap() = fresh;
252        pylon_core::query::clear_query_cache();
253        Ok(())
254    }
255
256    /// Returns a client view that injects `globals` into every query,
257    /// keyed by qualified name (`"module::name"`) — sharing the same
258    /// connection pool. Mirrors `pylon/client.py:287-304`.
259    pub fn with_globals(&self, globals: impl IntoIterator<Item = (String, DecodedValue)>) -> Client {
260        let mut merged = (*self.globals).clone();
261        merged.extend(globals);
262        Client {
263            globals: Arc::new(merged),
264            ..self.clone()
265        }
266    }
267
268    /// Returns a client view that applies `config` to every query — sharing
269    /// the same connection pool. Mirrors `pylon/client.py:306-327`.
270    pub fn with_config(&self, config: SessionConfig) -> Client {
271        Client { config, ..self.clone() }
272    }
273
274    /// Escape hatch for hand-written SQL outside PyQL — mirrors
275    /// `pylon/client.py:567-579`. Connects if this client hasn't yet.
276    pub async fn raw_connection(&self) -> Result<&pylon_pgcon::PgPool> {
277        Ok(&self.connected().await?.pool)
278    }
279
280    /// The pool, but only if this client is already connected — `None`
281    /// rather than connecting. For an observer that wants to report on
282    /// whatever connections a process happens to be holding (pool-status
283    /// metrics, say) without a metrics scrape being the thing that opens
284    /// them.
285    pub fn pool_if_connected(&self) -> Option<&pylon_pgcon::PgPool> {
286        self.connected.get().map(|conn| &conn.pool)
287    }
288
289    /// A clone of the currently-loaded schema — for callers that need to
290    /// introspect it directly (e.g. a schema-browser endpoint), not just
291    /// compile queries against it. Connects (and so fetches the snapshot) if
292    /// this client hasn't yet. Clones out from behind the lock rather than
293    /// returning a guard, same reasoning as every query method here.
294    pub async fn schema(&self) -> Result<SchemaDescriptor> {
295        Ok(self.connected().await?.schema.read().unwrap().clone())
296    }
297
298    /// The `Arc<pylon_cache::Cache>` this client reads/writes through, if
299    /// `Builder::cache` was configured — for a caller (`pylon-server`'s
300    /// worker-wiring startup) that needs to attach a `CacheInvalidationWorker`
301    /// to the exact same LMDB handle this client's own read-through caching
302    /// uses, rather than opening a second one (LMDB refuses a second
303    /// `Env::open` on the same path within one process).
304    pub fn cache_handle(&self) -> Option<Arc<pylon_cache::Cache>> {
305        self.cache.clone()
306    }
307
308    pub async fn query<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Vec<R>> {
309        let params = args.to_params();
310        let conn = self.connected().await?;
311        let schema = conn.schema.read().unwrap().clone();
312        let values = exec::query(
313            &conn.pool,
314            pyql,
315            &params,
316            &schema,
317            &self.config,
318            &self.globals,
319            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
320        )
321        .await?;
322        decode_rows(values)
323    }
324
325    pub async fn query_single<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Option<R>> {
326        let params = args.to_params();
327        let conn = self.connected().await?;
328        let schema = conn.schema.read().unwrap().clone();
329        let values = exec::query_single(
330            &conn.pool,
331            pyql,
332            &params,
333            &schema,
334            &self.config,
335            &self.globals,
336            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
337        )
338        .await?;
339        decode_optional_row(values)
340    }
341
342    pub async fn query_required_single<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<R> {
343        let params = args.to_params();
344        let conn = self.connected().await?;
345        let schema = conn.schema.read().unwrap().clone();
346        let values = exec::query_required_single(
347            &conn.pool,
348            pyql,
349            &params,
350            &schema,
351            &self.config,
352            &self.globals,
353            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
354        )
355        .await?;
356        decode_row(values)
357    }
358
359    pub async fn execute<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<()> {
360        let params = args.to_params();
361        let conn = self.connected().await?;
362        let schema = conn.schema.read().unwrap().clone();
363        exec::execute(
364            &conn.pool,
365            pyql,
366            &params,
367            &schema,
368            &self.config,
369            &self.globals,
370            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
371        )
372        .await
373    }
374
375    /// Subscribes to a schema-declared [`Channel`](pylon_core::schema::ChannelDescriptor)
376    /// (bare or `module::name` reference — the same string a schema author
377    /// already writes inside a PyQL `notify(...)` call) and returns a
378    /// [`ChannelListener`] whose `recv()` yields decoded payloads matching
379    /// that Channel's own declared shape: `Value::Uuid` for a Type-shaped
380    /// channel (the changed row's `id`, not a fetched object — see
381    /// `docs/schema/channels.md`), the matching `Value` variant for a
382    /// Scalar-shaped channel, or `Value::Object` for an Object-shaped
383    /// channel. A payload that doesn't match the declared shape comes back
384    /// as `Err(Error::MalformedPayload(_))` from that `recv()` call rather
385    /// than being silently dropped.
386    ///
387    /// Opens its own dedicated (non-pooled) connection, held for the
388    /// returned `ChannelListener`'s lifetime — `LISTEN` is per-session, so
389    /// running it on a pooled connection would leak the subscription onto
390    /// whatever unrelated query later borrows that connection back out of
391    /// the pool. The connection (and the server-side subscription with it)
392    /// closes once the `ChannelListener` is dropped.
393    ///
394    /// Mirrors `pylon/client.py`'s own `Client.listen()` — there, a typed
395    /// async generator; here, a `recv()`-based handle instead, since this
396    /// crate has no `Stream`/async-generator precedent to build on.
397    pub async fn listen(&self, channel: &str) -> Result<crate::ChannelListener> {
398        let conn = self.connected().await?;
399        let schema = conn.schema.read().unwrap().clone();
400        crate::listen::listen(&self.dsn, &schema, channel).await
401    }
402
403    pub async fn query_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
404        let params = args.to_params();
405        let conn = self.connected().await?;
406        let schema = conn.schema.read().unwrap().clone();
407        exec::query_json(
408            &conn.pool,
409            pyql,
410            &params,
411            &schema,
412            &self.config,
413            &self.globals,
414            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
415        )
416        .await
417    }
418
419    pub async fn query_single_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Option<String>> {
420        let params = args.to_params();
421        let conn = self.connected().await?;
422        let schema = conn.schema.read().unwrap().clone();
423        exec::query_single_json(
424            &conn.pool,
425            pyql,
426            &params,
427            &schema,
428            &self.config,
429            &self.globals,
430            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
431        )
432        .await
433    }
434
435    pub async fn query_required_single_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
436        let params = args.to_params();
437        let conn = self.connected().await?;
438        let schema = conn.schema.read().unwrap().clone();
439        exec::query_required_single_json(
440            &conn.pool,
441            pyql,
442            &params,
443            &schema,
444            &self.config,
445            &self.globals,
446            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
447        )
448        .await
449    }
450
451    /// Current cache size, or `None` if `Builder::cache` wasn't configured
452    /// — mirrors `pylon.cache.stat()`.
453    pub fn cache_stat(&self) -> Result<Option<pylon_cache::CacheStats>> {
454        match &self.cache {
455            None => Ok(None),
456            Some(cache) => cache.stat().map(Some).map_err(|e| Error::Cache(e.to_string())),
457        }
458    }
459
460    /// Evicts every cache entry — a no-op if `Builder::cache` wasn't
461    /// configured. Mirrors `pylon.cache.clear()`.
462    pub fn cache_clear(&self) -> Result<()> {
463        match &self.cache {
464            None => Ok(()),
465            Some(cache) => cache.clear().map_err(|e| Error::Cache(e.to_string())),
466        }
467    }
468
469    /// Runs `pyql` through Postgres's `EXPLAIN (ANALYZE, FORMAT JSON)` and
470    /// returns a query plan grouped by the query's own shape instead of raw
471    /// SQL relation names. `pyql` doesn't need the leading `analyze`
472    /// keyword already written. Mirrors `pylon/client.py:461-480`.
473    pub async fn analyze<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
474        let params = args.to_params();
475        let conn = self.connected().await?;
476        let schema = conn.schema.read().unwrap().clone();
477        exec::analyze(&conn.pool, pyql, &params, &schema, &self.config, &self.globals).await
478    }
479
480    /// Runs a retrying transaction with the default isolation level
481    /// (`Serializable`) and attempt budget (3) — see
482    /// [`Client::transaction_with_attempts`] for full control.
483    ///
484    /// `body` is re-run once per attempt against a fresh [`Transaction`];
485    /// it commits automatically when `body` returns `Ok`, and rolls back
486    /// and retries (with a `0ms, 100ms, 200ms, …` back-off) when `body`
487    /// returns a serialization-failure/deadlock error, up to the attempt
488    /// budget. Any other error rolls back and propagates immediately.
489    ///
490    /// A body that returns [`Error::Rollback`] rolls back and is never
491    /// retried — a decision, not a failure. It still propagates here (there
492    /// is no `T` to return); [`Client::transaction_opt`] is the same call
493    /// with that sentinel folded into `Ok(None)`.
494    ///
495    /// ```no_run
496    /// # use pylon_client::DecodedValue;
497    /// # async fn go(client: pylon_client::Client) -> pylon_client::Result<()> {
498    /// client.transaction(pylon_client::Isolation::Serializable, |tx| Box::pin(async move {
499    ///     tx.execute("insert Person { name := <str>$name }", &[("name", DecodedValue::Str("Bob".into()))]).await
500    /// })).await?;
501    /// # Ok(()) }
502    /// ```
503    pub async fn transaction<T, F>(&self, isolation: Isolation, body: F) -> Result<T>
504    where
505        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
506    {
507        self.transaction_with_attempts(isolation, 3, body).await
508    }
509
510    /// [`Client::transaction`], but a body that deliberately rolls back is
511    /// an outcome rather than an error: `Ok(Some(value))` when it committed,
512    /// `Ok(None)` when it returned [`Error::Rollback`]. Real failures still
513    /// propagate as `Err`.
514    ///
515    /// This is the closest Rust gets to `pylon.Rollback` in the Python
516    /// client, where the exception is simply swallowed and the loop ends.
517    /// Everything the body wrote is visible to the body's own queries and
518    /// to nothing else — the point being a test or dry run that needs real
519    /// writes without leaving rows behind.
520    ///
521    /// ```no_run
522    /// # use pylon_client::{DecodedValue, Error};
523    /// # async fn go(client: pylon_client::Client) -> pylon_client::Result<()> {
524    /// let committed: Option<()> = client
525    ///     .transaction_opt(pylon_client::Isolation::Serializable, |tx| Box::pin(async move {
526    ///         tx.execute("insert Person { name := <str>$name }", &[("name", DecodedValue::Str("Bob".into()))]).await?;
527    ///         // ... assert on what the transaction can see, then discard it.
528    ///         Err(Error::Rollback)
529    ///     }))
530    ///     .await?;
531    /// assert!(committed.is_none());
532    /// # Ok(()) }
533    /// ```
534    pub async fn transaction_opt<T, F>(&self, isolation: Isolation, body: F) -> Result<Option<T>>
535    where
536        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
537    {
538        self.transaction_opt_with_attempts(isolation, 3, body).await
539    }
540
541    /// [`Client::transaction_opt`] with an explicit attempt budget — the
542    /// `Ok(None)`-on-rollback counterpart to
543    /// [`Client::transaction_with_attempts`].
544    pub async fn transaction_opt_with_attempts<T, F>(
545        &self,
546        isolation: Isolation,
547        max_attempts: u32,
548        body: F,
549    ) -> Result<Option<T>>
550    where
551        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
552    {
553        match self.transaction_with_attempts(isolation, max_attempts, body).await {
554            Ok(value) => Ok(Some(value)),
555            Err(Error::Rollback) => Ok(None),
556            Err(e) => Err(e),
557        }
558    }
559
560    pub async fn transaction_with_attempts<T, F>(
561        &self,
562        isolation: Isolation,
563        max_attempts: u32,
564        mut body: F,
565    ) -> Result<T>
566    where
567        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
568    {
569        let conn = self.connected().await?;
570        let mut attempt = 0u32;
571        loop {
572            attempt += 1;
573            if attempt > 1 {
574                tokio::time::sleep(Duration::from_millis(100 * u64::from(attempt - 1))).await;
575            }
576            let pg_tx = conn.pool.begin(isolation.as_str()).await.map_err(Error::Db)?;
577            let tx = Transaction {
578                inner: pg_tx,
579                schema: conn.schema.clone(),
580                config: self.config.clone(),
581                globals: self.globals.clone(),
582                cache: self.cache.clone(),
583            };
584            let result = body(&tx).await;
585            match result {
586                Ok(value) => {
587                    tx.inner.commit().await.map_err(Error::Db)?;
588                    return Ok(value);
589                }
590                Err(e) if e.is_retriable() && attempt < max_attempts => {
591                    let _ = tx.inner.rollback().await;
592                }
593                // `Error::Rollback` is not retriable, so a body that asks to
594                // be discarded lands here and is rolled back once rather
595                // than re-run — re-running a body that already said "don't
596                // keep this" would just repeat its writes to throw them away
597                // again.
598                Err(e) => {
599                    let _ = tx.inner.rollback().await;
600                    return Err(e);
601                }
602            }
603        }
604    }
605}
606
607#[cfg(test)]
608mod tests {
609    use super::*;
610
611    /// A DSN nothing is listening on, so a connection attempt fails fast
612    /// without needing a live Postgres — port 1 is reserved and never bound.
613    const UNREACHABLE_DSN: &str = "postgresql://nobody@127.0.0.1:1/nothing";
614
615    #[test]
616    fn build_does_not_connect() {
617        let client = Client::builder(UNREACHABLE_DSN).build().unwrap();
618        assert!(client.pool_if_connected().is_none());
619    }
620
621    /// A failed attempt must leave the cell empty rather than storing the
622    /// failure, or a client built before its database finished starting
623    /// would stay broken for the life of the process.
624    #[tokio::test]
625    async fn a_failed_connect_is_retried_rather_than_remembered() {
626        let client = Client::builder(UNREACHABLE_DSN).build().unwrap();
627
628        assert!(client.ensure_connected().await.is_err());
629        assert!(client.pool_if_connected().is_none());
630        assert!(client.ensure_connected().await.is_err());
631    }
632
633    /// Views share the one connection slot, the way `with_globals` siblings
634    /// share `pylon/client.py`'s `_PoolRef`.
635    #[test]
636    fn views_share_the_connection_slot() {
637        let client = Client::builder(UNREACHABLE_DSN).build().unwrap();
638        let view = client.with_globals([]);
639
640        assert!(Arc::ptr_eq(&client.connected, &view.connected));
641    }
642}