libeufin

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

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 }