taler-rust

GNU Taler code in Rust. Largely core banking integrations.
Log | Files | Refs | Submodules | README | LICENSE

db.rs (6341B)


      1 /*
      2   This file is part of TALER
      3   Copyright (C) 2025-2026 Taler Systems SA
      4 
      5   TALER is free software; you can redistribute it and/or modify it under the
      6   terms of the GNU Affero General Public License as published by the Free Software
      7   Foundation; either version 3, or (at your option) any later version.
      8 
      9   TALER is distributed in the hope that it will be useful, but WITHOUT ANY
     10   WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR
     11   A PARTICULAR PURPOSE.  See the GNU Affero General Public License for more details.
     12 
     13   You should have received a copy of the GNU Affero General Public License along with
     14   TALER; see the file COPYING.  If not, see <http://www.gnu.org/licenses/>
     15 */
     16 
     17 use std::{path::Path, str::FromStr as _};
     18 
     19 use sqlx::{Connection, PgConnection, PgPool, postgres::PgConnectOptions};
     20 use taler_api::config::DbCfg;
     21 use taler_common::{
     22     config::{Config, parser::ConfigSource},
     23     db::{dbinit, pool},
     24 };
     25 use tracing::info;
     26 
     27 use crate::setup_tracing;
     28 
     29 /// An exclusive reservation of a reusable test database.
     30 ///
     31 /// The administrative connection holds a session advisory lock. Dropping this
     32 /// reservation closes that connection and releases the slot. Keep it alive
     33 /// until every task and connection using the database has been dropped.
     34 pub struct TestDb {
     35     _lock: PgConnection,
     36     options: PgConnectOptions,
     37 }
     38 
     39 impl TestDb {
     40     /// Connection options for the reserved database.
     41     pub fn connect_options(&self) -> PgConnectOptions {
     42         self.options.clone()
     43     }
     44 
     45     /// Reset the component schema and create a pool using its local config.
     46     pub async fn setup(&self, src: ConfigSource) -> PgPool {
     47         let cfg = Config::load(src, None::<&str>).unwrap();
     48         let name = format!("{}db-postgres", src.component_name);
     49         let db_cfg = DbCfg::parse(cfg.section(&name)).unwrap();
     50         self.setup_manual(db_cfg.sql_dir.as_ref(), src.component_name)
     51             .await
     52     }
     53 
     54     /// Reset the component schema using an explicit SQL directory.
     55     pub async fn setup_manual(&self, sql_dir: &Path, component_name: &str) -> PgPool {
     56         let pool = pool(self.connect_options(), &component_name.replace('-', "_"))
     57             .await
     58             .unwrap();
     59         let mut conn = pool.acquire().await.unwrap();
     60         dbinit(&mut conn, sql_dir, component_name, true)
     61             .await
     62             .unwrap();
     63         drop(conn);
     64         pool
     65     }
     66 }
     67 
     68 /// Run a database test using a component's local configuration.
     69 ///
     70 /// The test gets its own current-thread Tokio runtime. On return or panic the
     71 /// runtime drops its remaining tasks before the database reservation is released.
     72 /// Return values must not retain database resources; use `()` or `Result<(), E>`.
     73 /// Keep all work using the database on this runtime so teardown can stop it.
     74 pub fn run_db_test(src: ConfigSource, test_name: &str, test: impl AsyncFnOnce(PgPool)) {
     75     run(test_name, async |db| test(db.setup(src).await).await)
     76 }
     77 
     78 /// Run a database test using an explicit SQL directory and component name.
     79 pub fn run_db_test_manual(
     80     sql_dir: &Path,
     81     component_name: &str,
     82     test_name: &str,
     83     test: impl AsyncFnOnce(PgPool),
     84 ) {
     85     run(test_name, async |db| {
     86         test(db.setup_manual(sql_dir, component_name).await).await
     87     })
     88 }
     89 
     90 /// Run a test with reserved connection options and no schema initialization.
     91 ///
     92 /// Use this for migration tests, custom fixtures, and tests sharing a database
     93 /// across components. As with `run_db_test`, the reservation outlives the test
     94 /// runtime, and all database work must remain on that runtime.
     95 pub fn run_db_test_raw(test_name: &str, test: impl AsyncFnOnce(PgConnectOptions)) {
     96     run(test_name, async |db| test(db.connect_options()).await)
     97 }
     98 
     99 fn run(test_name: &str, test: impl AsyncFnOnce(&TestDb)) {
    100     setup_tracing();
    101     // Locals drop in reverse declaration order, including during unwinding.
    102     // The reservation MUST outlive the runtime: detached tasks may own pools,
    103     // checked-out connections, or listeners until runtime shutdown drops them.
    104     let mut reservation = None;
    105     let runtime = tokio::runtime::Builder::new_current_thread()
    106         .enable_all()
    107         .build()
    108         .expect("create database test runtime");
    109     runtime.block_on(async {
    110         reservation = Some(
    111             reserve_db(test_name)
    112                 .await
    113                 .unwrap_or_else(|e| panic!("{test_name}: reserve test database: {e}")),
    114         );
    115         let db = reservation.as_ref().unwrap();
    116         info!(target: "test", test_name, database = ?db.options.get_database(), "Reserved test database");
    117         test(db).await
    118     })
    119 }
    120 
    121 /// Reserve a reusable test database through a session advisory lock.
    122 ///
    123 /// Connection defaults are `postgres:///taler_rust_check`. PostgreSQL's usual
    124 /// `PGHOST`, `PGPORT`, and `PGUSER` settings can select another test server.
    125 pub async fn test_db() -> sqlx::Result<TestDb> {
    126     reserve_db("taler-test-utils").await
    127 }
    128 
    129 async fn reserve_db(test_name: &str) -> sqlx::Result<TestDb> {
    130     let admin =
    131         PgConnectOptions::from_str("postgres:///taler_rust_check")?.application_name(test_name);
    132     let mut conn = PgConnection::connect_with(&admin).await?;
    133 
    134     // Find a free slot via Advisory Locks.
    135     let slot: Option<i32> = sqlx::query_scalar(
    136         "SELECT id FROM generate_series(0, 1000) AS id
    137          WHERE pg_try_advisory_lock(id) LIMIT 1",
    138     )
    139     .fetch_optional(&mut conn)
    140     .await?;
    141 
    142     let Some(slot) = slot else {
    143         return Err(sqlx::Error::Protocol(
    144             "Could not find a free database slot after 1001 attempts.".into(),
    145         ));
    146     };
    147     let name = match slot {
    148         0 => "taler_rust_test".to_owned(),
    149         id => format!("taler_rust_test_{id}"),
    150     };
    151     // Check in a fresh statement after acquiring the lock. A previous holder
    152     // may have created this database after the allocation query's snapshot.
    153     let exists: bool = sqlx::query_scalar(
    154         "SELECT EXISTS (SELECT 1 FROM pg_catalog.pg_database WHERE datname = $1)",
    155     )
    156     .bind(&name)
    157     .fetch_one(&mut conn)
    158     .await?;
    159     if !exists {
    160         sqlx::raw_sql(&format!("CREATE DATABASE {name}"))
    161             .execute(&mut conn)
    162             .await?;
    163     }
    164     let options = admin.database(&name);
    165     Ok(TestDb {
    166         _lock: conn,
    167         options,
    168     })
    169 }