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        // Whatever moved the snapshot on can equally have moved this
251        // database's enum/domain/`vector` OIDs, which the pool read once at
252        // connect time.
253        conn.pool.refresh_types().await?;
254        let fresh = schema::fetch(&conn.pool).await?;
255        *conn.schema.write().unwrap() = fresh;
256        pylon_core::query::clear_query_cache();
257        Ok(())
258    }
259
260    /// Returns a client view that injects `globals` into every query,
261    /// keyed by qualified name (`"module::name"`) — sharing the same
262    /// connection pool. Mirrors `pylon/client.py:287-304`.
263    pub fn with_globals(&self, globals: impl IntoIterator<Item = (String, DecodedValue)>) -> Client {
264        let mut merged = (*self.globals).clone();
265        merged.extend(globals);
266        Client {
267            globals: Arc::new(merged),
268            ..self.clone()
269        }
270    }
271
272    /// Returns a client view that applies `config` to every query — sharing
273    /// the same connection pool. Mirrors `pylon/client.py:306-327`.
274    pub fn with_config(&self, config: SessionConfig) -> Client {
275        Client { config, ..self.clone() }
276    }
277
278    /// Escape hatch for hand-written SQL outside PyQL — mirrors
279    /// `pylon/client.py:567-579`. Connects if this client hasn't yet.
280    pub async fn raw_connection(&self) -> Result<&pylon_pgcon::PgPool> {
281        Ok(&self.connected().await?.pool)
282    }
283
284    /// The pool, but only if this client is already connected — `None`
285    /// rather than connecting. For an observer that wants to report on
286    /// whatever connections a process happens to be holding (pool-status
287    /// metrics, say) without a metrics scrape being the thing that opens
288    /// them.
289    pub fn pool_if_connected(&self) -> Option<&pylon_pgcon::PgPool> {
290        self.connected.get().map(|conn| &conn.pool)
291    }
292
293    /// A clone of the currently-loaded schema — for callers that need to
294    /// introspect it directly (e.g. a schema-browser endpoint), not just
295    /// compile queries against it. Connects (and so fetches the snapshot) if
296    /// this client hasn't yet. Clones out from behind the lock rather than
297    /// returning a guard, same reasoning as every query method here.
298    pub async fn schema(&self) -> Result<SchemaDescriptor> {
299        Ok(self.connected().await?.schema.read().unwrap().clone())
300    }
301
302    /// The `Arc<pylon_cache::Cache>` this client reads/writes through, if
303    /// `Builder::cache` was configured — for a caller (`pylon-server`'s
304    /// worker-wiring startup) that needs to attach a `CacheInvalidationWorker`
305    /// to the exact same LMDB handle this client's own read-through caching
306    /// uses, rather than opening a second one (LMDB refuses a second
307    /// `Env::open` on the same path within one process).
308    pub fn cache_handle(&self) -> Option<Arc<pylon_cache::Cache>> {
309        self.cache.clone()
310    }
311
312    pub async fn query<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Vec<R>> {
313        let params = args.to_params();
314        let conn = self.connected().await?;
315        let schema = conn.schema.read().unwrap().clone();
316        let values = exec::query(
317            &conn.pool,
318            pyql,
319            &params,
320            &schema,
321            &self.config,
322            &self.globals,
323            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
324        )
325        .await?;
326        decode_rows(values)
327    }
328
329    pub async fn query_single<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Option<R>> {
330        let params = args.to_params();
331        let conn = self.connected().await?;
332        let schema = conn.schema.read().unwrap().clone();
333        let values = exec::query_single(
334            &conn.pool,
335            pyql,
336            &params,
337            &schema,
338            &self.config,
339            &self.globals,
340            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
341        )
342        .await?;
343        decode_optional_row(values)
344    }
345
346    pub async fn query_required_single<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<R> {
347        let params = args.to_params();
348        let conn = self.connected().await?;
349        let schema = conn.schema.read().unwrap().clone();
350        let values = exec::query_required_single(
351            &conn.pool,
352            pyql,
353            &params,
354            &schema,
355            &self.config,
356            &self.globals,
357            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
358        )
359        .await?;
360        decode_row(values)
361    }
362
363    pub async fn execute<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<()> {
364        let params = args.to_params();
365        let conn = self.connected().await?;
366        let schema = conn.schema.read().unwrap().clone();
367        exec::execute(
368            &conn.pool,
369            pyql,
370            &params,
371            &schema,
372            &self.config,
373            &self.globals,
374            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
375        )
376        .await
377    }
378
379    /// Subscribes to a schema-declared [`Channel`](pylon_core::schema::ChannelDescriptor)
380    /// (bare or `module::name` reference — the same string a schema author
381    /// already writes inside a PyQL `notify(...)` call) and returns a
382    /// [`ChannelListener`] whose `recv()` yields decoded payloads matching
383    /// that Channel's own declared shape: `Value::Uuid` for a Type-shaped
384    /// channel (the changed row's `id`, not a fetched object — see
385    /// `docs/schema/channels.md`), the matching `Value` variant for a
386    /// Scalar-shaped channel, or `Value::Object` for an Object-shaped
387    /// channel. A payload that doesn't match the declared shape comes back
388    /// as `Err(Error::MalformedPayload(_))` from that `recv()` call rather
389    /// than being silently dropped.
390    ///
391    /// Opens its own dedicated (non-pooled) connection, held for the
392    /// returned `ChannelListener`'s lifetime — `LISTEN` is per-session, so
393    /// running it on a pooled connection would leak the subscription onto
394    /// whatever unrelated query later borrows that connection back out of
395    /// the pool. The connection (and the server-side subscription with it)
396    /// closes once the `ChannelListener` is dropped.
397    ///
398    /// Mirrors `pylon/client.py`'s own `Client.listen()` — there, a typed
399    /// async generator; here, a `recv()`-based handle instead, since this
400    /// crate has no `Stream`/async-generator precedent to build on.
401    pub async fn listen(&self, channel: &str) -> Result<crate::ChannelListener> {
402        let conn = self.connected().await?;
403        let schema = conn.schema.read().unwrap().clone();
404        crate::listen::listen(&self.dsn, &schema, channel).await
405    }
406
407    pub async fn query_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
408        let params = args.to_params();
409        let conn = self.connected().await?;
410        let schema = conn.schema.read().unwrap().clone();
411        exec::query_json(
412            &conn.pool,
413            pyql,
414            &params,
415            &schema,
416            &self.config,
417            &self.globals,
418            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
419        )
420        .await
421    }
422
423    pub async fn query_single_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Option<String>> {
424        let params = args.to_params();
425        let conn = self.connected().await?;
426        let schema = conn.schema.read().unwrap().clone();
427        exec::query_single_json(
428            &conn.pool,
429            pyql,
430            &params,
431            &schema,
432            &self.config,
433            &self.globals,
434            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
435        )
436        .await
437    }
438
439    pub async fn query_required_single_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
440        let params = args.to_params();
441        let conn = self.connected().await?;
442        let schema = conn.schema.read().unwrap().clone();
443        exec::query_required_single_json(
444            &conn.pool,
445            pyql,
446            &params,
447            &schema,
448            &self.config,
449            &self.globals,
450            crate::cache::CacheAccess::read_write(self.cache.as_deref()),
451        )
452        .await
453    }
454
455    /// Current cache size, or `None` if `Builder::cache` wasn't configured
456    /// — mirrors `pylon.cache.stat()`.
457    pub fn cache_stat(&self) -> Result<Option<pylon_cache::CacheStats>> {
458        match &self.cache {
459            None => Ok(None),
460            Some(cache) => cache.stat().map(Some).map_err(|e| Error::Cache(e.to_string())),
461        }
462    }
463
464    /// Evicts every cache entry — a no-op if `Builder::cache` wasn't
465    /// configured. Mirrors `pylon.cache.clear()`.
466    pub fn cache_clear(&self) -> Result<()> {
467        match &self.cache {
468            None => Ok(()),
469            Some(cache) => cache.clear().map_err(|e| Error::Cache(e.to_string())),
470        }
471    }
472
473    /// Runs `pyql` through Postgres's `EXPLAIN (ANALYZE, FORMAT JSON)` and
474    /// returns a query plan grouped by the query's own shape instead of raw
475    /// SQL relation names. `pyql` doesn't need the leading `analyze`
476    /// keyword already written. Mirrors `pylon/client.py:461-480`.
477    pub async fn analyze<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
478        let params = args.to_params();
479        let conn = self.connected().await?;
480        let schema = conn.schema.read().unwrap().clone();
481        exec::analyze(&conn.pool, pyql, &params, &schema, &self.config, &self.globals).await
482    }
483
484    /// Runs a retrying transaction with the default isolation level
485    /// (`Serializable`) and attempt budget (3) — see
486    /// [`Client::transaction_with_attempts`] for full control.
487    ///
488    /// `body` is re-run once per attempt against a fresh [`Transaction`];
489    /// it commits automatically when `body` returns `Ok`, and rolls back
490    /// and retries (with a `0ms, 100ms, 200ms, …` back-off) when `body`
491    /// returns a serialization-failure/deadlock error, up to the attempt
492    /// budget. Any other error rolls back and propagates immediately.
493    ///
494    /// A body that returns [`Error::Rollback`] rolls back and is never
495    /// retried — a decision, not a failure. It still propagates here (there
496    /// is no `T` to return); [`Client::transaction_opt`] is the same call
497    /// with that sentinel folded into `Ok(None)`.
498    ///
499    /// ```no_run
500    /// # use pylon_client::DecodedValue;
501    /// # async fn go(client: pylon_client::Client) -> pylon_client::Result<()> {
502    /// client.transaction(pylon_client::Isolation::Serializable, |tx| Box::pin(async move {
503    ///     tx.execute("insert Person { name := <str>$name }", &[("name", DecodedValue::Str("Bob".into()))]).await
504    /// })).await?;
505    /// # Ok(()) }
506    /// ```
507    pub async fn transaction<T, F>(&self, isolation: Isolation, body: F) -> Result<T>
508    where
509        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
510    {
511        self.transaction_with_attempts(isolation, 3, body).await
512    }
513
514    /// [`Client::transaction`], but a body that deliberately rolls back is
515    /// an outcome rather than an error: `Ok(Some(value))` when it committed,
516    /// `Ok(None)` when it returned [`Error::Rollback`]. Real failures still
517    /// propagate as `Err`.
518    ///
519    /// This is the closest Rust gets to `pylon.Rollback` in the Python
520    /// client, where the exception is simply swallowed and the loop ends.
521    /// Everything the body wrote is visible to the body's own queries and
522    /// to nothing else — the point being a test or dry run that needs real
523    /// writes without leaving rows behind.
524    ///
525    /// ```no_run
526    /// # use pylon_client::{DecodedValue, Error};
527    /// # async fn go(client: pylon_client::Client) -> pylon_client::Result<()> {
528    /// let committed: Option<()> = client
529    ///     .transaction_opt(pylon_client::Isolation::Serializable, |tx| Box::pin(async move {
530    ///         tx.execute("insert Person { name := <str>$name }", &[("name", DecodedValue::Str("Bob".into()))]).await?;
531    ///         // ... assert on what the transaction can see, then discard it.
532    ///         Err(Error::Rollback)
533    ///     }))
534    ///     .await?;
535    /// assert!(committed.is_none());
536    /// # Ok(()) }
537    /// ```
538    pub async fn transaction_opt<T, F>(&self, isolation: Isolation, body: F) -> Result<Option<T>>
539    where
540        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
541    {
542        self.transaction_opt_with_attempts(isolation, 3, body).await
543    }
544
545    /// [`Client::transaction_opt`] with an explicit attempt budget — the
546    /// `Ok(None)`-on-rollback counterpart to
547    /// [`Client::transaction_with_attempts`].
548    pub async fn transaction_opt_with_attempts<T, F>(
549        &self,
550        isolation: Isolation,
551        max_attempts: u32,
552        body: F,
553    ) -> Result<Option<T>>
554    where
555        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
556    {
557        match self.transaction_with_attempts(isolation, max_attempts, body).await {
558            Ok(value) => Ok(Some(value)),
559            Err(Error::Rollback) => Ok(None),
560            Err(e) => Err(e),
561        }
562    }
563
564    pub async fn transaction_with_attempts<T, F>(
565        &self,
566        isolation: Isolation,
567        max_attempts: u32,
568        mut body: F,
569    ) -> Result<T>
570    where
571        F: for<'a> FnMut(&'a Transaction) -> TxFuture<'a, T>,
572    {
573        let conn = self.connected().await?;
574        let mut attempt = 0u32;
575        loop {
576            attempt += 1;
577            if attempt > 1 {
578                tokio::time::sleep(Duration::from_millis(100 * u64::from(attempt - 1))).await;
579            }
580            let pg_tx = conn.pool.begin(isolation.as_str()).await.map_err(Error::Db)?;
581            let tx = Transaction {
582                inner: pg_tx,
583                schema: conn.schema.clone(),
584                config: self.config.clone(),
585                globals: self.globals.clone(),
586                cache: self.cache.clone(),
587            };
588            let result = body(&tx).await;
589            match result {
590                Ok(value) => {
591                    tx.inner.commit().await.map_err(Error::Db)?;
592                    return Ok(value);
593                }
594                Err(e) if e.is_retriable() && attempt < max_attempts => {
595                    let _ = tx.inner.rollback().await;
596                }
597                // `Error::Rollback` is not retriable, so a body that asks to
598                // be discarded lands here and is rolled back once rather
599                // than re-run — re-running a body that already said "don't
600                // keep this" would just repeat its writes to throw them away
601                // again.
602                Err(e) => {
603                    let _ = tx.inner.rollback().await;
604                    return Err(e);
605                }
606            }
607        }
608    }
609}
610
611#[cfg(test)]
612mod tests {
613    use super::*;
614
615    /// A DSN nothing is listening on, so a connection attempt fails fast
616    /// without needing a live Postgres — port 1 is reserved and never bound.
617    const UNREACHABLE_DSN: &str = "postgresql://nobody@127.0.0.1:1/nothing";
618
619    #[test]
620    fn build_does_not_connect() {
621        let client = Client::builder(UNREACHABLE_DSN).build().unwrap();
622        assert!(client.pool_if_connected().is_none());
623    }
624
625    /// A failed attempt must leave the cell empty rather than storing the
626    /// failure, or a client built before its database finished starting
627    /// would stay broken for the life of the process.
628    #[tokio::test]
629    async fn a_failed_connect_is_retried_rather_than_remembered() {
630        let client = Client::builder(UNREACHABLE_DSN).build().unwrap();
631
632        assert!(client.ensure_connected().await.is_err());
633        assert!(client.pool_if_connected().is_none());
634        assert!(client.ensure_connected().await.is_err());
635    }
636
637    /// Views share the one connection slot, the way `with_globals` siblings
638    /// share `pylon/client.py`'s `_PoolRef`.
639    #[test]
640    fn views_share_the_connection_slot() {
641        let client = Client::builder(UNREACHABLE_DSN).build().unwrap();
642        let view = client.with_globals([]);
643
644        assert!(Arc::ptr_eq(&client.connected, &view.connected));
645    }
646
647    /// A client that reloaded only the schema would compile against the new
648    /// one and then fail to decode its enum columns.
649    #[tokio::test]
650    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
651    async fn reload_schema_also_refreshes_the_type_registry() {
652        let dsn = std::env::var("PYLON_PGCON_TEST_DSN").expect("PYLON_PGCON_TEST_DSN must be set for live tests");
653        let setup = pylon_pgcon::PgPool::connect(&dsn, 2).await.unwrap();
654        setup.batch_execute("CREATE SCHEMA IF NOT EXISTS _pylon").await.unwrap();
655        pylon_core::migrate::ensure_internal_schema(&setup).await.unwrap();
656        // One shared row for this whole DSN — put back what was there.
657        let previous = pylon_core::migrate::read_schema_snapshot(&setup).await.unwrap();
658        let placeholder = serde_json::to_string(&SchemaDescriptor::default()).unwrap();
659        pylon_core::migrate::write_schema_snapshot(&setup, &placeholder)
660            .await
661            .unwrap();
662        setup
663            .batch_execute("DROP TYPE IF EXISTS pylon_client_reload_enum")
664            .await
665            .unwrap();
666
667        let client = Client::builder(&dsn).build().unwrap();
668        client.ensure_connected().await.unwrap();
669
670        // Created after the client's pool connected.
671        setup
672            .batch_execute("CREATE TYPE pylon_client_reload_enum AS ENUM ('a')")
673            .await
674            .unwrap();
675        let oid = match setup
676            .query_typed(
677                "SELECT (oid::int8) AS result FROM pg_type WHERE typname = 'pylon_client_reload_enum'",
678                &[],
679                &pylon_pgcon::ExtensionOids::default(),
680            )
681            .await
682            .unwrap()
683            .first()
684        {
685            Some(pylon_value::DecodedValue::I64(oid)) => *oid as u32,
686            other => panic!("expected the new enum's oid, got {other:?}"),
687        };
688        assert!(
689            !client.pool_if_connected().unwrap().types().enums.contains(&oid),
690            "the connect-time snapshot must not know a type created after it"
691        );
692
693        client.reload_schema().await.unwrap();
694
695        assert!(
696            client.pool_if_connected().unwrap().types().enums.contains(&oid),
697            "reload_schema must re-discover type OIDs, not just re-fetch the schema"
698        );
699
700        setup.batch_execute("DROP TYPE pylon_client_reload_enum").await.unwrap();
701        if let Some(previous) = previous {
702            pylon_core::migrate::write_schema_snapshot(&setup, &previous)
703                .await
704                .unwrap();
705        }
706    }
707}