Skip to main content

pylon_client/
transaction.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//! Closure-based retrying transactions — see [`crate::Client::transaction`].
21//!
22//! Unlike `pylon/client.py`'s `RetryingTransaction` (an async iterator the
23//! caller drives with `async for tx in client.transaction(): async with
24//! tx: ...`), the body here is a closure re-run once per attempt; the
25//! transaction commits automatically when it returns `Ok`, rolls back and
26//! retries on a retriable error (serialization failure/deadlock), and
27//! rolls back and propagates on anything else. `Transaction` itself has no
28//! public `commit`/`rollback` — that decision is made for the caller by
29//! the closure's own return value.
30//!
31//! That includes a deliberate abort: returning [`crate::Error::Rollback`]
32//! rolls back without retrying, and
33//! [`Client::transaction_opt`](crate::Client::transaction_opt) reports it
34//! as `Ok(None)` instead of an error. It is the Rust counterpart of
35//! `pylon.Rollback` in `pylon/client.py`.
36
37use std::collections::HashMap;
38use std::sync::{Arc, RwLock};
39
40use pylon_core::ir::SessionConfig;
41use pylon_core::schema::SchemaDescriptor;
42use pylon_value::DecodedValue;
43
44use crate::error::Result;
45use crate::exec;
46use crate::query_arg::QueryArgs;
47use crate::queryable::{Queryable, decode_optional_row, decode_row, decode_rows};
48
49/// PostgreSQL transaction isolation level. Defaults to `Serializable`,
50/// matching `pylon/client.py`'s own default.
51#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
52pub enum Isolation {
53    ReadUncommitted,
54    ReadCommitted,
55    RepeatableRead,
56    #[default]
57    Serializable,
58}
59
60impl Isolation {
61    pub(crate) fn as_str(self) -> &'static str {
62        match self {
63            Isolation::ReadUncommitted => "read_uncommitted",
64            Isolation::ReadCommitted => "read_committed",
65            Isolation::RepeatableRead => "repeatable_read",
66            Isolation::Serializable => "serializable",
67        }
68    }
69}
70
71/// A single transaction attempt, handed to the closure passed to
72/// [`crate::Client::transaction`]. Exposes the same query methods as
73/// [`crate::Client`] itself (minus `analyze`, which the Python client also
74/// never runs inside an explicit transaction).
75pub struct Transaction {
76    pub(crate) inner: pylon_pgcon::PgTransaction,
77    pub(crate) schema: Arc<RwLock<SchemaDescriptor>>,
78    pub(crate) config: SessionConfig,
79    pub(crate) globals: Arc<HashMap<String, DecodedValue>>,
80    /// Held only to *evict* on a write — see `execute`. Never used to read
81    /// or populate: rows read inside a transaction aren't committed yet.
82    pub(crate) cache: Option<Arc<pylon_cache::Cache>>,
83}
84
85impl Transaction {
86    // Every method passes `CacheAccess::evict_only`: uncommitted rows must
87    // never populate the read-through cache (matching `pylon/client.py`'s
88    // `AsyncTransaction`), but a write still has to evict — and a write can
89    // arrive through any of these, not just `execute`. `query("insert ...")`
90    // is a normal way to insert and read the row back, and `Client.save`'s
91    // Python counterpart uses `query_single` for exactly that.
92
93    pub async fn query<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Vec<R>> {
94        let params = args.to_params();
95        let schema = self.schema.read().unwrap().clone();
96        let values = exec::query(
97            &self.inner,
98            pyql,
99            &params,
100            &schema,
101            &self.config,
102            &self.globals,
103            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
104        )
105        .await?;
106        decode_rows(values)
107    }
108
109    pub async fn query_single<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Option<R>> {
110        let params = args.to_params();
111        let schema = self.schema.read().unwrap().clone();
112        let values = exec::query_single(
113            &self.inner,
114            pyql,
115            &params,
116            &schema,
117            &self.config,
118            &self.globals,
119            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
120        )
121        .await?;
122        decode_optional_row(values)
123    }
124
125    pub async fn query_required_single<R: Queryable, A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<R> {
126        let params = args.to_params();
127        let schema = self.schema.read().unwrap().clone();
128        let values = exec::query_required_single(
129            &self.inner,
130            pyql,
131            &params,
132            &schema,
133            &self.config,
134            &self.globals,
135            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
136        )
137        .await?;
138        decode_row(values)
139    }
140
141    /// A write inside a transaction still evicts immediately: the cache is
142    /// never populated from inside a transaction, so an aborted attempt can
143    /// only over-evict — which costs a re-read and never serves stale data.
144    pub async fn execute<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<()> {
145        let params = args.to_params();
146        let schema = self.schema.read().unwrap().clone();
147        exec::execute(
148            &self.inner,
149            pyql,
150            &params,
151            &schema,
152            &self.config,
153            &self.globals,
154            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
155        )
156        .await
157    }
158
159    pub async fn query_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
160        let params = args.to_params();
161        let schema = self.schema.read().unwrap().clone();
162        exec::query_json(
163            &self.inner,
164            pyql,
165            &params,
166            &schema,
167            &self.config,
168            &self.globals,
169            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
170        )
171        .await
172    }
173
174    pub async fn query_single_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<Option<String>> {
175        let params = args.to_params();
176        let schema = self.schema.read().unwrap().clone();
177        exec::query_single_json(
178            &self.inner,
179            pyql,
180            &params,
181            &schema,
182            &self.config,
183            &self.globals,
184            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
185        )
186        .await
187    }
188
189    pub async fn query_required_single_json<A: QueryArgs + ?Sized>(&self, pyql: &str, args: &A) -> Result<String> {
190        let params = args.to_params();
191        let schema = self.schema.read().unwrap().clone();
192        exec::query_required_single_json(
193            &self.inner,
194            pyql,
195            &params,
196            &schema,
197            &self.config,
198            &self.globals,
199            crate::cache::CacheAccess::evict_only(self.cache.as_deref()),
200        )
201        .await
202    }
203}