exchange.rs (11024B)
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 jiff::Timestamp; 18 use sqlx::{PgPool, QueryBuilder, Row as _, postgres::PgRow}; 19 use taler_api::{ 20 db::{BindHelper, TypeHelper as _, history, page}, 21 serialized, 22 subject::fmt_out_subject, 23 }; 24 use taler_common::{ 25 api::{ 26 params::{History, Page}, 27 revenue::RevenueIncomingBankTransaction, 28 wire::{ 29 IncomingBankTransaction, OutgoingBankTransaction, TransferListStatus, TransferRequest, 30 TransferState, TransferStatus, 31 }, 32 }, 33 db::IncomingType, 34 types::amount::Currency, 35 }; 36 use tokio::sync::watch::Receiver; 37 38 use crate::model::SubmissionState; 39 40 pub async fn outgoing_history( 41 db: &PgPool, 42 currency: &Currency, 43 params: &History, 44 listen: impl FnOnce() -> Receiver<i64>, 45 ) -> sqlx::Result<Vec<OutgoingBankTransaction>> { 46 history( 47 db, 48 "outgoing_transaction_id", 49 params, 50 listen, 51 || { 52 QueryBuilder::new( 53 " 54 SELECT 55 outgoing_transaction_id 56 ,execution_time 57 ,amount 58 ,debit_fee 59 ,credit_payto 60 ,wtid 61 ,exchange_base_url 62 ,metadata 63 FROM talerable_outgoing_transactions 64 JOIN outgoing_transactions USING(outgoing_transaction_id) 65 WHERE 66 ", 67 ) 68 }, 69 |r: PgRow| { 70 Ok(OutgoingBankTransaction { 71 row_id: r.try_get_u64("outgoing_transaction_id")?, 72 amount: r.try_get_amount("amount", currency)?, 73 debit_fee: r 74 .try_get_opt_amount("debit_fee", currency)? 75 .filter(|it| !it.is_zero()), 76 credit_account: r.try_get_payto("credit_payto")?, 77 date: r.try_get_timestamp("execution_time")?.into(), 78 exchange_base_url: r.try_get_url("exchange_base_url")?, 79 wtid: r.try_get("wtid")?, 80 metadata: r.try_get("metadata")?, 81 }) 82 }, 83 ) 84 .await 85 } 86 87 pub async fn incoming_history( 88 db: &PgPool, 89 currency: &Currency, 90 params: &History, 91 listen: impl FnOnce() -> Receiver<i64>, 92 ) -> sqlx::Result<Vec<IncomingBankTransaction>> { 93 history( 94 db, 95 "taler_in_id", 96 params, 97 listen, 98 || { 99 // What does the order do when it's finally complete 100 QueryBuilder::new( 101 " 102 SELECT 103 taler_in_id 104 ,execution_time 105 ,amount 106 ,credit_fee 107 ,debit_payto 108 ,type::text 109 ,metadata 110 ,authorization_pub 111 ,authorization_sig 112 FROM talerable_incoming_transactions 113 JOIN incoming_transactions USING(incoming_transaction_id) 114 WHERE 115 ", 116 ) 117 }, 118 |r: PgRow| { 119 let credit_fee = r 120 .try_get_opt_amount("credit_fee", currency)? 121 .filter(|it| !it.is_zero()); 122 Ok(match r.try_get_parse("type")? { 123 IncomingType::reserve => IncomingBankTransaction::Reserve { 124 row_id: r.try_get_u64("taler_in_id")?, 125 amount: r.try_get_amount("amount", currency)?, 126 credit_fee, 127 debit_account: r.try_get_payto("debit_payto")?, 128 date: r.try_get_timestamp("execution_time")?.into(), 129 reserve_pub: r.try_get("metadata")?, 130 authorization_pub: r.try_get("authorization_pub")?, 131 authorization_sig: r.try_get("authorization_sig")?, 132 }, 133 IncomingType::kyc => IncomingBankTransaction::Kyc { 134 row_id: r.try_get_u64("taler_in_id")?, 135 amount: r.try_get_amount("amount", currency)?, 136 credit_fee, 137 debit_account: r.try_get_payto("debit_payto")?, 138 date: r.try_get_timestamp("execution_time")?.into(), 139 account_pub: r.try_get("metadata")?, 140 authorization_pub: r.try_get("authorization_pub")?, 141 authorization_sig: r.try_get("authorization_sig")?, 142 }, 143 IncomingType::map => unimplemented!("MAP are never listed in the history"), 144 }) 145 }, 146 ) 147 .await 148 } 149 150 pub async fn revenue_history( 151 db: &PgPool, 152 currency: &Currency, 153 params: &History, 154 listen: impl FnOnce() -> Receiver<i64>, 155 ) -> sqlx::Result<Vec<RevenueIncomingBankTransaction>> { 156 history( 157 db, 158 "incoming_transaction_id", 159 params, 160 listen, 161 || { 162 QueryBuilder::new( 163 " 164 SELECT 165 incoming_transaction_id 166 ,execution_time 167 ,amount 168 ,credit_fee 169 ,debit_payto 170 ,subject 171 FROM incoming_transactions 172 WHERE debit_payto IS NOT NULL AND subject IS NOT NULL AND 173 ", 174 ) 175 }, 176 |r: PgRow| { 177 Ok(RevenueIncomingBankTransaction { 178 row_id: r.try_get_u64("incoming_transaction_id")?, 179 amount: r.try_get_amount("amount", currency)?, 180 credit_fee: r 181 .try_get_opt_amount("credit_fee", currency)? 182 .filter(|it| !it.is_zero()), 183 debit_account: r.try_get_payto("debit_payto")?, 184 date: r.try_get_timestamp("execution_time")?.into(), 185 subject: r.try_get("subject")?, 186 }) 187 }, 188 ) 189 .await 190 } 191 192 pub enum TransferResult { 193 Success { id: u64, timestamp: Timestamp }, 194 RequestUidReuse, 195 WtidReuse, 196 } 197 198 pub async fn transfer( 199 db: &PgPool, 200 req: &TransferRequest, 201 e2e_id: &str, 202 timestamp: &Timestamp, 203 ) -> sqlx::Result<TransferResult> { 204 let subject = fmt_out_subject(&req.wtid, &req.exchange_base_url, req.metadata.as_deref()); 205 serialized!( 206 sqlx::query( 207 " 208 SELECT 209 out_request_uid_reuse 210 ,out_wtid_reuse 211 ,out_tx_row_id 212 ,out_timestamp 213 FROM taler_transfer($1,$2,$3,$4,$5,$6,$7,$8,$9) 214 ", 215 ) 216 .bind(req.request_uid) 217 .bind(req.wtid) 218 .bind(&subject) 219 .bind(req.amount) 220 .bind(req.exchange_base_url.as_str()) 221 .bind(&req.metadata) 222 .bind(req.credit_account.as_ref().as_str()) 223 .bind(e2e_id) 224 .bind_timestamp(timestamp) 225 .try_map(|r: PgRow| { 226 Ok(if r.try_get_flag("out_request_uid_reuse")? { 227 TransferResult::RequestUidReuse 228 } else if r.try_get_flag("out_wtid_reuse")? { 229 TransferResult::WtidReuse 230 } else { 231 TransferResult::Success { 232 id: r.try_get_u64("out_tx_row_id")?, 233 timestamp: r.try_get_timestamp("out_timestamp")?, 234 } 235 }) 236 }) 237 .fetch_one(db) 238 ) 239 } 240 241 pub async fn transfer_by_id( 242 db: &PgPool, 243 currency: &Currency, 244 id: u64, 245 ) -> sqlx::Result<Option<TransferStatus>> { 246 serialized!( 247 sqlx::query( 248 " 249 SELECT 250 wtid 251 ,exchange_base_url 252 ,metadata 253 ,amount 254 ,credit_payto 255 ,initiation_time 256 ,status 257 ,status_msg 258 FROM transfer_operations 259 JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id) 260 WHERE initiated_outgoing_transaction_id=$1 261 ", 262 ) 263 .bind(id as i64) 264 .try_map(|r: PgRow| { 265 Ok(TransferStatus { 266 status: r 267 .try_get::<SubmissionState, _>("status")? 268 .to_transfer_status(), 269 status_msg: r.try_get("status_msg")?, 270 amount: r.try_get_amount("amount", currency)?, 271 exchange_base_url: r.try_get("exchange_base_url")?, 272 metadata: r.try_get("metadata")?, 273 wtid: r.try_get("wtid")?, 274 credit_account: r.try_get_payto("credit_payto")?, 275 timestamp: r.try_get_timestamp("initiation_time")?.into(), 276 }) 277 }) 278 .fetch_optional(db) 279 ) 280 } 281 282 pub async fn transfer_page( 283 db: &PgPool, 284 currency: &Currency, 285 params: &Page, 286 status: &Option<TransferState>, 287 ) -> sqlx::Result<Vec<TransferListStatus>> { 288 page( 289 db, 290 params, 291 "initiated_outgoing_transaction_id", 292 || { 293 let mut builder = QueryBuilder::new( 294 " 295 SELECT 296 initiated_outgoing_transaction_id 297 ,amount 298 ,status 299 ,credit_payto 300 ,initiation_time 301 FROM transfer_operations 302 JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id) 303 WHERE 304 ", 305 ); 306 if let Some(status) = status { 307 match status { 308 TransferState::pending => { 309 builder.push("( status = ").push_bind(SubmissionState::pending).push(" OR ").push(" status = ").push_bind(SubmissionState::unsubmitted).push(") AND "); 310 } 311 status => { 312 builder.push(" status = ").push_bind(SubmissionState::from(*status)).push(" AND ");} 313 } 314 } 315 builder 316 }, 317 |r: PgRow| { 318 Ok(TransferListStatus { 319 row_id: r.try_get_u64("initiated_outgoing_transaction_id")?, 320 status: r.try_get::<SubmissionState, _>("status")?.to_transfer_status(), 321 amount: r.try_get_amount("amount", currency)?, 322 credit_account: r.try_get_payto("credit_payto")?, 323 timestamp: r.try_get_timestamp("initiation_time")?.into(), 324 }) 325 }, 326 ) 327 .await 328 }