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 ¶ms,
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 ¶ms,
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 ¶ms,
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 ¶ms,
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 ¶ms,
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 ¶ms,
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 ¶ms,
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, ¶ms, &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}