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 }