libeufin

Integration and sandbox testing for FinTech APIs and data formats
Log | Files | Refs | Submodules | README | LICENSE

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 }