db.rs (5547B)
1 /* 2 This file is part of TALER 3 Copyright (C) 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 sqlx::PgPool; 18 use taler_api::config::DbCfg; 19 use tokio::sync::watch::Sender; 20 21 pub mod exchange; 22 pub mod initiated; 23 pub mod list; 24 pub mod payment; 25 pub mod transfer; 26 27 const SCHEMA: &str = "libeufin_nexus"; 28 pub const UNSETTLED: &str = "status NOT IN ('success', 'permanent_failure', 'late_failure')"; 29 pub const PENDING: &str = "status IN ('unsubmitted', 'pending')"; 30 31 pub async fn pool(cfg: &DbCfg) -> anyhow::Result<PgPool> { 32 let pool = taler_common::db::pool(cfg.cfg.clone(), SCHEMA).await?; 33 Ok(pool) 34 } 35 36 pub async fn dbinit(cfg: &DbCfg, reset: bool) -> anyhow::Result<PgPool> { 37 let pool = taler_common::db::pool(cfg.cfg.clone(), SCHEMA).await?; 38 let mut db = pool.acquire().await?; 39 taler_common::db::dbinit(&mut db, cfg.sql_dir.as_ref(), "libeufin-nexus", reset).await?; 40 Ok(pool) 41 } 42 43 pub async fn notification_listener( 44 pool: PgPool, 45 in_channel: Sender<i64>, 46 taler_in_channel: Sender<i64>, 47 taler_out_channel: Sender<i64>, 48 ) -> sqlx::Result<()> { 49 taler_api::notification::notification_listener!(&pool, 50 "nexus_revenue_tx" => (row_id: i64) { 51 in_channel.send_replace(row_id); 52 }, 53 "nexus_incoming_tx" => (row_id: i64) { 54 taler_in_channel.send_replace(row_id); 55 }, 56 "nexus_outgoing_tx" => (row_id: i64) { 57 taler_out_channel.send_replace(row_id); 58 } 59 ) 60 } 61 62 #[cfg(test)] 63 pub mod test { 64 use libeufin_ebics::db::test::ebics_routine; 65 use sqlx::{PgPool, Row, postgres::PgRow}; 66 use taler_api::db::TypeHelper; 67 use taler_common::db::IncomingType; 68 use taler_macros::db_test; 69 use taler_test_utils::routine::Status; 70 71 pub async fn check_count(db: &PgPool, nb_tx: usize, nb_bounce: usize) { 72 assert_eq!( 73 sqlx::query_as::<_, (i64, i64)>( 74 "SELECT (SELECT count(*) FROM incoming_transactions) + (SELECT count(*) FROM outgoing_transactions), (SELECT count(*) FROM bounced_transactions)", 75 ) 76 .fetch_one(db) 77 .await 78 .unwrap(), 79 (nb_tx as i64, nb_bounce as i64) 80 ); 81 } 82 83 pub async fn check_in_count( 84 db: &PgPool, 85 nb_incoming: usize, 86 nb_bounce: usize, 87 nb_talerable: usize, 88 ) { 89 assert_eq!( 90 sqlx::query_as::<_, (i64, i64, i64)>( 91 " 92 SELECT (SELECT count(*) FROM incoming_transactions), 93 (SELECT count(*) FROM bounced_transactions), 94 (SELECT count(*) FROM talerable_incoming_transactions) 95 ", 96 ) 97 .fetch_one(db) 98 .await 99 .unwrap(), 100 (nb_incoming as i64, nb_bounce as i64, nb_talerable as i64) 101 ); 102 } 103 104 pub async fn check_out_count(db: &PgPool, nb_outgoing: u64, nb_talerable: u64) { 105 assert_eq!( 106 sqlx::query_as::<_, (i64, i64)>( 107 " 108 SELECT (SELECT count(*) FROM outgoing_transactions), 109 (SELECT count(*) FROM talerable_outgoing_transactions) 110 ", 111 ) 112 .fetch_one(db) 113 .await 114 .unwrap(), 115 (nb_outgoing as i64, nb_talerable as i64) 116 ); 117 } 118 119 pub async fn check_in(db: &PgPool) -> Vec<Status> { 120 sqlx::query( 121 " 122 SELECT pending_recurrent_incoming_transactions.authorization_pub IS NOT NULL, initiated_outgoing_transaction_id IS NOT NULL, debit_payto IS NULL OR subject IS NULL, type::text, metadata 123 FROM incoming_transactions 124 LEFT JOIN talerable_incoming_transactions USING (incoming_transaction_id) 125 LEFT JOIN pending_recurrent_incoming_transactions USING (incoming_transaction_id) 126 LEFT JOIN bounced_transactions USING (incoming_transaction_id) 127 ORDER BY incoming_transaction_id 128 ", 129 ) 130 .try_map(|r: PgRow| { 131 Ok( 132 if r.try_get_flag(0)? { 133 Status::Pending 134 } else if r.try_get_flag(1)? { 135 Status::Bounced 136 } else if r.try_get_flag(2)? { 137 Status::Incomplete 138 } else { 139 match r.try_get_opt_parse(3)? { 140 None => Status::Simple, 141 Some(IncomingType::reserve) => Status::Reserve(r.try_get(4)?), 142 Some(IncomingType::kyc) => Status::Kyc(r.try_get(4)?), 143 Some(e) => unreachable!("{e:?}") 144 } 145 } 146 ) 147 }) 148 .fetch_all(db) 149 .await 150 .unwrap() 151 } 152 153 pub async fn check_in_state(db: &PgPool, state: &[Status]) { 154 assert_eq!(state, check_in(db).await); 155 } 156 157 #[db_test] 158 pub async fn ebics(db: PgPool) { 159 ebics_routine(&db).await; 160 } 161 }