exchange.rs (13414B)
1 /* 2 * This file is part of LibEuFin. 3 * Copyright (C) 2026 Taler Systems S.A. 4 5 * LibEuFin is free software; you can redistribute it and/or modify 6 * it under the terms of the GNU Affero General Public License as 7 * published by the Free Software Foundation; either version 3, or 8 * (at your option) any later version. 9 10 * LibEuFin is distributed in the hope that it will be useful, but 11 * WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY 12 * or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Affero General 13 * Public License for more details. 14 15 * You should have received a copy of the GNU Affero General Public 16 * License along with LibEuFin; see the file COPYING. If not, see 17 * <http://www.gnu.org/licenses/> 18 */ 19 20 //! Data access logic for exchange specific logic 21 22 use jiff::Timestamp; 23 use sqlx::{PgPool, QueryBuilder, Row, postgres::PgRow}; 24 use taler_api::{ 25 db::{BindHelper, TypeHelper, history, page}, 26 notification::NotificationChannel, 27 serialized, 28 subject::{IncomingKey, fmt_out_subject}, 29 }; 30 use taler_common::{ 31 api::{ 32 params::{History, Page}, 33 wire::{ 34 IncomingBankTransaction, OutgoingBankTransaction, TransferListStatus, TransferRequest, 35 TransferState, TransferStatus, 36 }, 37 }, 38 db::IncomingType, 39 types::amount::{Amount, Currency}, 40 }; 41 42 use crate::payto::{BankPayto, PaytoCtx, sql_bank_payto, sql_bank_simple_payto}; 43 44 /** Result of taler transfer transaction creation */ 45 pub enum TransferResult { 46 /** Transaction [id] and wire transfer [timestamp] */ 47 Success { 48 id: u64, 49 timestamp: Timestamp, 50 }, 51 NotAnExchange, 52 UnknownExchange, 53 BothPartyAreExchange, 54 BalanceInsufficient, 55 ReserveUidReuse, 56 WtidReuse, 57 AdminCreditor, 58 } 59 60 /** Perform a Taler transfer */ 61 pub async fn transfer( 62 db: &PgPool, 63 username: &str, 64 req: &TransferRequest, 65 payto: &BankPayto, 66 timestamp: &Timestamp, 67 conversion: bool, 68 ) -> sqlx::Result<TransferResult> { 69 let subject = fmt_out_subject(&req.wtid, &req.exchange_base_url, req.metadata.as_deref()); 70 serialized!( 71 sqlx::query( 72 " 73 SELECT 74 out_debtor_not_found, 75 out_debtor_not_exchange, 76 out_both_exchanges, 77 out_request_uid_reuse, 78 out_wtid_reuse, 79 out_exchange_balance_insufficient, 80 out_creditor_admin, 81 out_tx_row_id, 82 out_timestamp 83 FROM taler_transfer($1,$2,$3,$4,$5,$6,$7,$8,$9,$10) 84 ", 85 ) 86 .bind(req.request_uid) 87 .bind(req.wtid) 88 .bind(&subject) 89 .bind(req.amount) 90 .bind(req.exchange_base_url.as_str()) 91 .bind(&req.metadata) 92 .bind(payto.canonical()) 93 .bind(username) 94 .bind_timestamp(timestamp) 95 .bind(conversion) 96 .try_map(|r: PgRow| { 97 Ok(if r.try_get_flag("out_debtor_not_found")? { 98 TransferResult::UnknownExchange 99 } else if r.try_get_flag("out_debtor_not_exchange")? { 100 TransferResult::NotAnExchange 101 } else if r.try_get_flag("out_both_exchanges")? { 102 TransferResult::BothPartyAreExchange 103 } else if r.try_get_flag("out_exchange_balance_insufficient")? { 104 TransferResult::BalanceInsufficient 105 } else if r.try_get_flag("out_request_uid_reuse")? { 106 TransferResult::ReserveUidReuse 107 } else if r.try_get_flag("out_wtid_reuse")? { 108 TransferResult::WtidReuse 109 } else if r.try_get_flag("out_creditor_admin")? { 110 TransferResult::AdminCreditor 111 } else { 112 TransferResult::Success { 113 id: r.try_get_u64("out_tx_row_id")?, 114 timestamp: r.try_get_timestamp("out_timestamp")?, 115 } 116 }) 117 }) 118 .fetch_one(db) 119 ) 120 } 121 122 /** Get status of transfer [txId] of account [exchangeId] */ 123 pub async fn transfer_by_id( 124 db: &PgPool, 125 ctx: &PaytoCtx, 126 currency: &Currency, 127 exchange_id: u64, 128 tx_id: u64, 129 ) -> sqlx::Result<Option<TransferStatus>> { 130 serialized!( 131 sqlx::query( 132 " 133 SELECT 134 wtid, 135 exchange_base_url, 136 metadata, 137 transfer_date, 138 amount, 139 creditor_payto, 140 status, 141 status_msg 142 FROM transfer_operations 143 WHERE transfer_operation_id=$1 AND exchange_id=$2 144 ", 145 ) 146 .bind(tx_id as i64) 147 .bind(exchange_id as i64) 148 .try_map(|r: PgRow| { 149 Ok(TransferStatus { 150 status: r.try_get("status")?, 151 status_msg: r.try_get("status_msg")?, 152 amount: r.try_get_amount("amount", currency)?, 153 exchange_base_url: r.try_get("exchange_base_url")?, 154 metadata: r.try_get("metadata")?, 155 wtid: r.try_get("wtid")?, 156 credit_account: sql_bank_simple_payto(&r, ctx, "creditor_payto")?.as_uri(), 157 timestamp: r.try_get_timestamp("transfer_date")?.into(), 158 }) 159 }) 160 .fetch_optional(db) 161 ) 162 } 163 164 /** Get a page of transfers status of account [exchangeId] */ 165 pub async fn page_transfer( 166 db: &PgPool, 167 ctx: &PaytoCtx, 168 currency: &Currency, 169 params: &Page, 170 exchange_id: u64, 171 status: Option<TransferState>, 172 ) -> sqlx::Result<Vec<TransferListStatus>> { 173 page( 174 db, 175 params, 176 "transfer_operation_id", 177 || { 178 let mut query = QueryBuilder::new( 179 " 180 SELECT 181 transfer_operation_id 182 ,transfer_date 183 ,amount 184 ,creditor_payto 185 ,status 186 FROM transfer_operations 187 WHERE exchange_id=", 188 ); 189 query.push_bind(exchange_id as i64).push(" AND "); 190 if let Some(status) = status { 191 query.push("status=").push_bind(status).push(" AND "); 192 } 193 query 194 }, 195 |r| { 196 Ok(TransferListStatus { 197 row_id: r.try_get_u64("transfer_operation_id")?, 198 status: r.try_get("status")?, 199 amount: r.try_get_amount("amount", currency)?, 200 credit_account: sql_bank_simple_payto(&r, ctx, "creditor_payto")?.as_uri(), 201 timestamp: r.try_get_timestamp("transfer_date")?.into(), 202 }) 203 }, 204 ) 205 .await 206 } 207 208 /** Result of taler add incoming transaction creation */ 209 pub enum AddIncomingResult { 210 /** Transaction [id] and wire transfer [timestamp] */ 211 Success { 212 id: u64, 213 timestamp: Timestamp, 214 pending: bool, 215 }, 216 NotAnExchange, 217 UnknownExchange, 218 UnknownDebtor, 219 BothPartyAreExchange, 220 ReservePubReuse, 221 UnknownMapping, 222 MappingReuse, 223 BalanceInsufficient, 224 } 225 226 /** Add a new taler incoming transaction */ 227 pub async fn add_incoming( 228 db: &PgPool, 229 amount: &Amount, 230 debtor: &BankPayto, 231 subject: &str, 232 username: &str, 233 timestamp: &Timestamp, 234 metadata: &IncomingKey, 235 ) -> sqlx::Result<AddIncomingResult> { 236 serialized!( 237 sqlx::query( 238 " 239 SELECT 240 out_creditor_not_found 241 ,out_creditor_not_exchange 242 ,out_debtor_not_found 243 ,out_both_exchanges 244 ,out_reserve_pub_reuse 245 ,out_mapping_reuse 246 ,out_unknown_mapping 247 ,out_debitor_balance_insufficient 248 ,out_tx_row_id 249 ,out_pending 250 FROM taler_add_incoming($1,$2,$3,$4,$5,$6,$7::taler_incoming_type) 251 ", 252 ) 253 .bind(metadata.key) 254 .bind(subject) 255 .bind(amount) 256 .bind(debtor.canonical()) 257 .bind(username) 258 .bind_timestamp(timestamp) 259 .bind(metadata.ty.as_ref()) 260 .try_map(|r: PgRow| { 261 Ok(if r.try_get_flag("out_creditor_not_found")? { 262 AddIncomingResult::UnknownExchange 263 } else if r.try_get_flag("out_creditor_not_exchange")? { 264 AddIncomingResult::NotAnExchange 265 } else if r.try_get_flag("out_debtor_not_found")? { 266 AddIncomingResult::UnknownDebtor 267 } else if r.try_get_flag("out_both_exchanges")? { 268 AddIncomingResult::BothPartyAreExchange 269 } else if r.try_get_flag("out_debitor_balance_insufficient")? { 270 AddIncomingResult::BalanceInsufficient 271 } else if r.try_get_flag("out_reserve_pub_reuse")? { 272 AddIncomingResult::ReservePubReuse 273 } else if r.try_get_flag("out_mapping_reuse")? { 274 AddIncomingResult::MappingReuse 275 } else if r.try_get_flag("out_unknown_mapping")? { 276 AddIncomingResult::UnknownMapping 277 } else { 278 AddIncomingResult::Success { 279 id: r.try_get_u64("out_tx_row_id")?, 280 timestamp: *timestamp, 281 pending: r.try_get_flag("out_pending")?, 282 } 283 }) 284 }) 285 .fetch_one(db) 286 ) 287 } 288 289 /** Query [exchangeId] history of taler incoming transactions */ 290 pub async fn incoming_history( 291 db: &PgPool, 292 ctx: &PaytoCtx, 293 channel: &NotificationChannel<u64, i64>, 294 currency: &Currency, 295 params: &History, 296 exchange_id: u64, 297 ) -> sqlx::Result<Vec<IncomingBankTransaction>> { 298 history( 299 db, 300 "exchange_incoming_id", 301 params, 302 || channel.subscribe(exchange_id), 303 || { 304 let mut query = QueryBuilder::new( 305 "SELECT 306 exchange_incoming_id 307 ,transaction_date 308 ,amount 309 ,debtor_payto 310 ,debtor_name 311 ,type::text 312 ,metadata 313 ,authorization_pub 314 ,authorization_sig 315 FROM taler_exchange_incoming AS tfr 316 JOIN bank_account_transactions AS txs 317 ON bank_transaction=txs.bank_transaction_id 318 WHERE bank_account_id=", 319 ); 320 query.push_bind(exchange_id as i64).push(" AND "); 321 query 322 }, 323 |r| { 324 Ok(match r.try_get_parse("type")? { 325 IncomingType::reserve => IncomingBankTransaction::Reserve { 326 row_id: r.try_get_u64("exchange_incoming_id")?, 327 date: r.try_get_timestamp("transaction_date")?.into(), 328 amount: r.try_get_amount("amount", currency)?, 329 credit_fee: None, 330 debit_account: sql_bank_payto(&r, ctx, "debtor_payto", "debtor_name")?.as_uri(), 331 reserve_pub: r.try_get("metadata")?, 332 authorization_pub: r.try_get("authorization_pub")?, 333 authorization_sig: r.try_get("authorization_sig")?, 334 }, 335 IncomingType::kyc => IncomingBankTransaction::Kyc { 336 row_id: r.try_get_u64("exchange_incoming_id")?, 337 date: r.try_get_timestamp("transaction_date")?.into(), 338 amount: r.try_get_amount("amount", currency)?, 339 credit_fee: None, 340 debit_account: sql_bank_payto(&r, ctx, "debtor_payto", "debtor_name")?.as_uri(), 341 account_pub: r.try_get("metadata")?, 342 authorization_pub: r.try_get("authorization_pub")?, 343 authorization_sig: r.try_get("authorization_sig")?, 344 }, 345 IncomingType::map => unreachable!(), 346 }) 347 }, 348 ) 349 .await 350 } 351 352 /** Query [exchangeId] history of taler incoming transactions */ 353 pub async fn outgoing_history( 354 db: &PgPool, 355 ctx: &PaytoCtx, 356 channel: &NotificationChannel<u64, i64>, 357 currency: &Currency, 358 params: &History, 359 exchange_id: u64, 360 ) -> sqlx::Result<Vec<OutgoingBankTransaction>> { 361 history( 362 db, 363 "bank_transaction_id", 364 params, 365 || channel.subscribe(exchange_id), 366 || { 367 let mut query = QueryBuilder::new( 368 "SELECT 369 bank_transaction_id 370 ,transaction_date 371 ,txs.amount 372 ,txs.creditor_payto 373 ,txs.creditor_name 374 ,wtid 375 ,exchange_base_url 376 ,transfer_operations.metadata 377 FROM taler_exchange_outgoing AS tfr 378 JOIN transfer_operations USING (exchange_outgoing_id) 379 JOIN bank_account_transactions AS txs 380 ON bank_transaction=txs.bank_transaction_id 381 WHERE bank_account_id=", 382 ); 383 query.push_bind(exchange_id as i64).push(" AND "); 384 query 385 }, 386 |r| { 387 Ok(OutgoingBankTransaction { 388 row_id: r.try_get_u64("bank_transaction_id")?, 389 date: r.try_get_timestamp("transaction_date")?.into(), 390 amount: r.try_get_amount("amount", currency)?, 391 debit_fee: None, 392 credit_account: sql_bank_payto(&r, ctx, "creditor_payto", "creditor_name")? 393 .as_uri(), 394 wtid: r.try_get("wtid")?, 395 metadata: r.try_get("metadata")?, 396 exchange_base_url: r.try_get_parse("exchange_base_url")?, 397 }) 398 }, 399 ) 400 .await 401 }