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