libeufin

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

api.rs (17451B)


      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 use const_format::formatcp;
     21 use jiff::Timestamp;
     22 use libeufin_ebics::{
     23     ebics::{TaskStatus, rand_ebics_id},
     24     iso20022::model::{InId, InTx},
     25 };
     26 use prometheus_client::{
     27     encoding::EncodeLabelSet,
     28     metrics::{family::Family, gauge::Gauge},
     29     registry::{Registry, Unit},
     30 };
     31 use sqlx::{PgPool, types::Json};
     32 use taler_api::{
     33     api::{
     34         TalerApi,
     35         observability::Observability,
     36         prepared::{PreparedTransfer, simple_subject},
     37         revenue::Revenue,
     38         wire::WireGateway,
     39     },
     40     error::{ApiResult, failure_code},
     41     subject::{IncomingKey, fmt_in_subject, subject_fmt_qr_bill},
     42 };
     43 use taler_common::{
     44     api::{
     45         params::{History, Page},
     46         prepared::{
     47             RegistrationRequest, RegistrationResponse, SubjectFormat, TransferSubject,
     48             Unregistration,
     49         },
     50         revenue::RevenueIncomingHistory,
     51         wire::{
     52             AddIncomingRequest, AddIncomingResponse, AddKycauthRequest, AddMappedRequest,
     53             IncomingHistory, OutgoingHistory, TransferList, TransferRequest, TransferResponse,
     54             TransferState, TransferStatus,
     55         },
     56     },
     57     error_code::ErrorCode,
     58     types::{
     59         amount::{Amount, Currency},
     60         iban::IBAN,
     61         payto::{FullIbanPayto, IbanPayto, PaytoURI},
     62         time::TalerTimestamp,
     63     },
     64 };
     65 use tokio::sync::watch::Sender;
     66 
     67 use crate::{
     68     constants::{FETCH_TASK_KEY, SUBMIT_TASK_KEY},
     69     db::{
     70         self,
     71         exchange::{
     72             TransferResult, incoming_history, outgoing_history, revenue_history, transfer,
     73             transfer_by_id, transfer_page,
     74         },
     75         payment::{IncomingRegistrationResult, register_in_talerable},
     76         transfer::{RegistrationResult, transfer_register, transfer_unregister},
     77     },
     78 };
     79 
     80 pub struct NexusApi {
     81     pub pool: sqlx::PgPool,
     82     pub currency: Currency,
     83     pub payto: FullIbanPayto,
     84     pub qr_iban: Option<IBAN>,
     85     pub in_channel: Sender<i64>,
     86     pub taler_in_channel: Sender<i64>,
     87     pub taler_out_channel: Sender<i64>,
     88     metrics: Metrics,
     89     registry: Registry,
     90 }
     91 
     92 #[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)]
     93 struct TaskLabel {
     94     name: &'static str,
     95 }
     96 
     97 #[derive(Default)]
     98 pub struct Metrics {
     99     db_access: Gauge,
    100     task_execution: Family<TaskLabel, Gauge>,
    101     task_success: Family<TaskLabel, Gauge>,
    102 }
    103 
    104 impl Metrics {
    105     pub fn registry(&self) -> Registry {
    106         let mut registry = Registry::default();
    107 
    108         registry.register(
    109             "db_access",
    110             "Whether the last database metrics refresh succeeded",
    111             self.db_access.clone(),
    112         );
    113 
    114         registry.register_with_unit(
    115             "task_execution_timestamp_seconds",
    116             "Unix timestamp of the last task execution",
    117             Unit::Seconds,
    118             self.task_execution.clone(),
    119         );
    120 
    121         registry.register_with_unit(
    122             "task_success_timestamp_seconds",
    123             "Unix timestamp of the last successful task execution",
    124             Unit::Seconds,
    125             self.task_success.clone(),
    126         );
    127 
    128         registry
    129     }
    130 
    131     pub async fn sync(&self, db: &PgPool) {
    132         let success =
    133             match sqlx::query_as::<_, (Option<Json<TaskStatus>>, Option<Json<TaskStatus>>)>(
    134                 formatcp!(
    135                     "
    136             SELECT
    137             (SELECT value FROM kv WHERE key='{SUBMIT_TASK_KEY}'),
    138             (SELECT value FROM kv WHERE key='{FETCH_TASK_KEY}')
    139         "
    140                 ),
    141             )
    142             .fetch_one(db)
    143             .await
    144             {
    145                 Ok((submit, fetch)) => {
    146                     for (name, status) in [("submit", submit), ("fetch", fetch)] {
    147                         if let Some(status) = status {
    148                             if let Some(time) = status.last_trial {
    149                                 self.task_execution
    150                                     .get_or_create(&TaskLabel { name })
    151                                     .set(time.as_second());
    152                             }
    153                             if let Some(time) = status.last_successfull {
    154                                 self.task_success
    155                                     .get_or_create(&TaskLabel { name })
    156                                     .set(time.as_second());
    157                             }
    158                         }
    159                     }
    160                     true
    161                 }
    162                 Err(_) => false,
    163             };
    164         self.db_access.set(if success { 1 } else { 0 });
    165     }
    166 }
    167 
    168 impl NexusApi {
    169     pub async fn start(
    170         pool: sqlx::PgPool,
    171         currency: Currency,
    172         payto: FullIbanPayto,
    173         qr_iban: Option<IBAN>,
    174     ) -> Self {
    175         let in_channel = Sender::new(0);
    176         let taler_in_channel = Sender::new(0);
    177         let taler_out_channel = Sender::new(0);
    178 
    179         let metrics = Metrics::default();
    180 
    181         let tmp = Self {
    182             pool: pool.clone(),
    183             payto,
    184             qr_iban,
    185             currency,
    186             in_channel: in_channel.clone(),
    187             taler_in_channel: taler_in_channel.clone(),
    188             taler_out_channel: taler_out_channel.clone(),
    189             registry: metrics.registry(),
    190             metrics,
    191         };
    192         tokio::spawn(db::notification_listener(
    193             pool,
    194             in_channel,
    195             taler_in_channel,
    196             taler_out_channel,
    197         ));
    198         tmp
    199     }
    200 }
    201 
    202 impl TalerApi for NexusApi {
    203     fn currency(&self) -> Currency {
    204         self.currency
    205     }
    206 
    207     fn implementation(&self) -> &'static str {
    208         "urn:net:taler:specs:libeufin-nexus:taler-rust"
    209     }
    210 }
    211 
    212 async fn add_incoming(
    213     db: &PgPool,
    214     subject: &IncomingKey,
    215     amount: Amount,
    216     debit_account: PaytoURI,
    217 ) -> ApiResult<AddIncomingResponse> {
    218     FullIbanPayto::try_from(&debit_account)?;
    219     let now = Timestamp::now();
    220     match register_in_talerable(
    221         db,
    222         &InTx {
    223             id: InId {
    224                 uetr: None,
    225                 tx_id: Some(rand_ebics_id()),
    226                 sref: None,
    227             },
    228             amount,
    229             credit_fee: Amount::zero(&amount.currency),
    230             subject: Some(format!(
    231                 "Manual incoming {}",
    232                 fmt_in_subject(subject.ty, &subject.key)
    233             )),
    234             execution_time: now,
    235             debtor: Some(debit_account),
    236         },
    237         subject,
    238     )
    239     .await?
    240     {
    241         IncomingRegistrationResult::Success(in_result) => Ok(AddIncomingResponse {
    242             row_id: in_result.id,
    243             timestamp: now.into(),
    244         }),
    245         IncomingRegistrationResult::ReservePubReuse => {
    246             Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT))
    247         }
    248         IncomingRegistrationResult::MappingReuse => {
    249             Err(failure_code(ErrorCode::BANK_TRANSFER_MAPPING_REUSED))
    250         }
    251         IncomingRegistrationResult::UnknownMapping => {
    252             Err(failure_code(ErrorCode::BANK_TRANSFER_MAPPING_UNKNOWN))
    253         }
    254     }
    255 }
    256 
    257 impl WireGateway for NexusApi {
    258     async fn transfer(&self, req: TransferRequest) -> ApiResult<TransferResponse> {
    259         FullIbanPayto::try_from(&req.credit_account)?;
    260         let result = transfer(&self.pool, &req, &rand_ebics_id(), &Timestamp::now()).await?;
    261         match result {
    262             TransferResult::Success { id, timestamp } => Ok(TransferResponse {
    263                 timestamp: timestamp.into(),
    264                 row_id: id,
    265             }),
    266             TransferResult::RequestUidReuse => {
    267                 Err(failure_code(ErrorCode::BANK_TRANSFER_REQUEST_UID_REUSED))
    268             }
    269             TransferResult::WtidReuse => Err(failure_code(ErrorCode::BANK_TRANSFER_WTID_REUSED)),
    270         }
    271     }
    272 
    273     async fn transfer_page(
    274         &self,
    275         page: Page,
    276         status: Option<TransferState>,
    277     ) -> ApiResult<TransferList> {
    278         Ok(TransferList {
    279             transfers: transfer_page(&self.pool, &self.currency, &page, &status).await?,
    280             debit_account: self.payto.as_uri(),
    281         })
    282     }
    283 
    284     async fn transfer_by_id(&self, id: u64) -> ApiResult<Option<TransferStatus>> {
    285         Ok(transfer_by_id(&self.pool, &self.currency, id).await?)
    286     }
    287 
    288     async fn outgoing_history(&self, params: History) -> ApiResult<OutgoingHistory> {
    289         Ok(OutgoingHistory {
    290             outgoing_transactions: outgoing_history(&self.pool, &self.currency, &params, || {
    291                 self.taler_out_channel.subscribe()
    292             })
    293             .await?,
    294             debit_account: self.payto.as_uri(),
    295         })
    296     }
    297 
    298     async fn incoming_history(&self, params: History) -> ApiResult<IncomingHistory> {
    299         Ok(IncomingHistory {
    300             incoming_transactions: incoming_history(&self.pool, &self.currency, &params, || {
    301                 self.taler_in_channel.subscribe()
    302             })
    303             .await?,
    304             credit_account: self.payto.as_uri(),
    305         })
    306     }
    307 
    308     async fn add_incoming_reserve(
    309         &self,
    310         req: AddIncomingRequest,
    311     ) -> ApiResult<AddIncomingResponse> {
    312         add_incoming(
    313             &self.pool,
    314             &IncomingKey::reserve(req.reserve_pub),
    315             req.amount,
    316             req.debit_account,
    317         )
    318         .await
    319     }
    320 
    321     async fn add_incoming_kyc(&self, req: AddKycauthRequest) -> ApiResult<AddIncomingResponse> {
    322         add_incoming(
    323             &self.pool,
    324             &IncomingKey::kyc(req.account_pub),
    325             req.amount,
    326             req.debit_account,
    327         )
    328         .await
    329     }
    330 
    331     async fn add_incoming_mapped(&self, req: AddMappedRequest) -> ApiResult<AddIncomingResponse> {
    332         add_incoming(
    333             &self.pool,
    334             &IncomingKey::map(req.authorization_pub),
    335             req.amount,
    336             req.debit_account,
    337         )
    338         .await
    339     }
    340 
    341     fn support_account_check(&self) -> bool {
    342         false
    343     }
    344 }
    345 
    346 impl Revenue for NexusApi {
    347     async fn history(&self, params: History) -> ApiResult<RevenueIncomingHistory> {
    348         Ok(RevenueIncomingHistory {
    349             incoming_transactions: revenue_history(&self.pool, &self.currency, &params, || {
    350                 self.in_channel.subscribe()
    351             })
    352             .await?,
    353             credit_account: self.payto.as_uri(),
    354         })
    355     }
    356 }
    357 
    358 impl PreparedTransfer for NexusApi {
    359     fn supported_formats(&self) -> &[SubjectFormat] {
    360         if self.qr_iban.is_some() {
    361             &[SubjectFormat::SIMPLE, SubjectFormat::CH_QR_BILL]
    362         } else {
    363             &[SubjectFormat::SIMPLE]
    364         }
    365     }
    366 
    367     async fn registration(&self, req: RegistrationRequest) -> ApiResult<RegistrationResponse> {
    368         let creditor = IbanPayto::try_from(&req.credit_account)?;
    369         let reference_number = if creditor.iban == self.payto.iban {
    370             None
    371         } else if Some(creditor.iban) == self.qr_iban {
    372             Some(subject_fmt_qr_bill(req.authorization_pub.as_ref()))
    373         } else {
    374             return Err(failure_code(ErrorCode::BANK_UNKNOWN_CREDITOR));
    375         };
    376         match transfer_register(
    377             &self.pool,
    378             req.r#type.into(),
    379             &req.account_pub,
    380             &req.authorization_pub,
    381             &req.authorization_sig,
    382             req.recurrent,
    383             reference_number.as_deref(),
    384             &Timestamp::now(),
    385         )
    386         .await?
    387         {
    388             RegistrationResult::Success => ApiResult::Ok(RegistrationResponse {
    389                 subjects: vec![if let Some(qr_reference_number) = reference_number {
    390                     TransferSubject::QrBill {
    391                         credit_amount: req.credit_amount,
    392                         qr_reference_number,
    393                     }
    394                 } else {
    395                     simple_subject(req)
    396                 }],
    397                 expiration: TalerTimestamp::Never,
    398             }),
    399             RegistrationResult::ReservePubReuse => {
    400                 ApiResult::Err(failure_code(ErrorCode::BANK_DUPLICATE_RESERVE_PUB_SUBJECT))
    401             }
    402             RegistrationResult::SubjectReuse => {
    403                 ApiResult::Err(failure_code(ErrorCode::BANK_DERIVATION_REUSE))
    404             }
    405         }
    406     }
    407 
    408     async fn unregistration(&self, req: Unregistration) -> ApiResult<bool> {
    409         Ok(transfer_unregister(&self.pool, &req.authorization_pub, &Timestamp::now()).await?)
    410     }
    411 }
    412 
    413 impl Observability for NexusApi {
    414     async fn metrics(&self) -> ApiResult<&Registry> {
    415         self.metrics.sync(&self.pool).await;
    416         Ok(&self.registry)
    417     }
    418 }
    419 
    420 #[cfg(test)]
    421 pub mod test {
    422     use std::sync::Arc;
    423 
    424     use sqlx::PgPool;
    425     use taler_api::{api::TalerRouter as _, auth::AuthMethod, subject::OutgoingSubject};
    426     use taler_common::api::{
    427         prepared::PreparedTransferConfig,
    428         revenue::RevenueConfig,
    429         wire::{TransferState, WireConfig},
    430     };
    431     use taler_macros::db_test;
    432     use taler_test_utils::{
    433         Router,
    434         routine::{
    435             admin_add_incoming_routine, in_history_routine, out_history_routine,
    436             registration_routine, revenue_routine, transfer_routine,
    437         },
    438         server::TestServer as _,
    439         tasks,
    440     };
    441 
    442     use crate::{
    443         api::NexusApi,
    444         db::{payment::register_out_tx, test::check_in},
    445         test::{CLIENT, CURR, EXCHANGE, UNKNOWN, gen_out_pay},
    446     };
    447 
    448     pub async fn api_setup(db: &PgPool) -> Router {
    449         let api = Arc::new(NexusApi::start(db.clone(), CURR, EXCHANGE.clone(), None).await);
    450         Router::new()
    451             .wire_gateway(api.clone(), AuthMethod::None)
    452             .prepared_transfer(api.clone())
    453             .revenue(api, AuthMethod::None)
    454             .finalize()
    455     }
    456 
    457     #[db_test]
    458     async fn config(db: PgPool) {
    459         let server = api_setup(&db).await;
    460         server
    461             .get("/taler-wire-gateway/config")
    462             .await
    463             .assert_ok_json::<WireConfig>();
    464         server
    465             .get("/taler-prepared-transfer/config")
    466             .await
    467             .assert_ok_json::<PreparedTransferConfig>();
    468         server
    469             .get("/taler-revenue/config")
    470             .await
    471             .assert_ok_json::<RevenueConfig>();
    472     }
    473 
    474     #[db_test]
    475     async fn transfer(db: PgPool) {
    476         let server = api_setup(&db).await;
    477         transfer_routine(
    478             &server.prefix("/taler-wire-gateway"),
    479             TransferState::pending,
    480             &EXCHANGE.as_uri(),
    481         )
    482         .await;
    483         // TODO
    484         /*db.initiated.batchSubmissionSuccess(1, Instant.now(), "ORDER1")
    485         db.initiated.batchSubmissionFailure(2, Instant.now(), "Failure")
    486         db.initiated.batchSubmissionFailure(3, Instant.now(), "Failure")
    487         client.getA("/taler-wire-gateway/transfers?status=transient_failure").assertOkJson<TransferList> {
    488             assertEquals(2, it.transfers.size)
    489         }
    490         client.getA("/taler-wire-gateway/transfers?status=pending").assertOkJson<TransferList> {
    491             assertEquals(4, it.transfers.size)
    492         }*/
    493     }
    494 
    495     #[db_test]
    496     async fn outgoing_history(db: PgPool) {
    497         let server = api_setup(&db).await;
    498         out_history_routine(
    499             &server.prefix("/taler-wire-gateway"),
    500             tasks!({
    501                 register_out_tx(&db, &gen_out_pay("subject"), Some(&OutgoingSubject::rand()))
    502                     .await
    503                     .unwrap();
    504             }),
    505             tasks!(),
    506         )
    507         .await;
    508     }
    509 
    510     #[db_test]
    511     async fn incoming_history(db: PgPool) {
    512         let server = api_setup(&db).await;
    513         in_history_routine(
    514             &server.prefix("/taler-wire-gateway"),
    515             &server.prefix("/taler-prepared-transfer"),
    516             &CLIENT.as_uri(),
    517             &EXCHANGE.as_uri(),
    518             tasks!(),
    519             tasks!(),
    520         )
    521         .await;
    522     }
    523 
    524     #[db_test]
    525     async fn admin_add_incoming(db: PgPool) {
    526         let server = api_setup(&db).await;
    527         admin_add_incoming_routine(
    528             &server.prefix("/taler-wire-gateway"),
    529             &server.prefix("/taler-prepared-transfer"),
    530             &CLIENT.as_uri(),
    531             &EXCHANGE.as_uri(),
    532         )
    533         .await;
    534     }
    535 
    536     #[db_test]
    537     async fn revenue(db: PgPool) {
    538         let server = api_setup(&db).await;
    539         revenue_routine(
    540             &server.prefix("/taler-wire-gateway"),
    541             &server.prefix("/taler-revenue"),
    542             &CLIENT.as_uri(),
    543             tasks!(),
    544             tasks!(),
    545         )
    546         .await;
    547     }
    548 
    549     #[db_test]
    550     async fn registration(db: PgPool) {
    551         let server = api_setup(&db).await;
    552         registration_routine(
    553             &server.prefix("/taler-wire-gateway"),
    554             &server.prefix("/taler-prepared-transfer"),
    555             &CLIENT.as_uri(),
    556             &EXCHANGE.as_uri(),
    557             &UNKNOWN,
    558             || check_in(&db),
    559         )
    560         .await;
    561     }
    562 }