libeufin

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

lib.rs (21547B)


      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 #![allow(clippy::too_many_arguments)]
     21 
     22 use std::{
     23     fmt::Write as _,
     24     io::{Cursor, Read as _},
     25     sync::Arc,
     26 };
     27 
     28 use anyhow::{anyhow, bail};
     29 use axum::{body::Bytes, http::HeaderValue};
     30 use compact_str::CompactStringExt;
     31 use http_client::Client;
     32 use jiff::{SignedDuration, Timestamp, Zoned, civil::Date, tz::TimeZone};
     33 use libeufin_ebics::{
     34     cli::{EbicsArgs, EbicsLogs},
     35     db::{get_task_status, update_task_status},
     36     ebics::{
     37         EbicsClient, EbicsErrKind,
     38         administrative::{AccountInfo, HKD, OrderInfo},
     39         ebics_code::EbicsReturnCode,
     40         order::{Order, OrderDoc},
     41     },
     42     iso20022::{HacAction, hac::parse_hac},
     43     keys::{BankKeys, ClientKeys},
     44     ws::listen_for_notification,
     45 };
     46 use sqlx::PgPool;
     47 use taler_api::api::TalerRouter;
     48 use taler_build::long_version;
     49 use taler_common::{CommonArgs, cli::ConfigCmd, config::Config, types::utils::date_to_utc_ts};
     50 use tokio::{time::timeout, try_join};
     51 use tracing::{debug, error, info, trace, warn};
     52 
     53 use crate::{
     54     api::{EbisyncState, sync_api},
     55     azure::AzureBlobStorage,
     56     config::{Destination, EbisyncCfg, Source, parse_db_cfg},
     57     constants::{CHECKPOINT_KEY, FETCH_TASK_KEY},
     58     db::{dbinit, pool},
     59 };
     60 
     61 pub mod api;
     62 pub mod azure;
     63 pub mod config;
     64 pub mod constants;
     65 pub mod db;
     66 
     67 /// Testing helper commands
     68 #[derive(clap::Subcommand, Debug)]
     69 pub enum TestingCmd {
     70     // Create the container
     71     CreateContainer,
     72 }
     73 
     74 #[derive(clap::Subcommand, Debug)]
     75 pub enum Cmd {
     76     /// Initialize libeufin-ebisync database
     77     Dbinit {
     78         /// Reset database (DANGEROUS: All existing data is lost)
     79         #[arg(short, long)]
     80         reset: bool,
     81     },
     82     /// Set up the EBICS subscriber
     83     Setup {
     84         #[command(flatten)]
     85         ebics_logs: EbicsLogs,
     86 
     87         /// Resubmits all the keys to the bank
     88         #[arg(long)]
     89         force_keys_resubmission: bool,
     90 
     91         /// Accepts the bank keys without interactively asking the user
     92         #[arg(long)]
     93         auto_accept_keys: bool,
     94 
     95         /// Generates the PDF with the client public keys to send to the bank
     96         #[arg(long)]
     97         generate_registration_pdf: bool,
     98     },
     99     /// Downloads EBICS files from the bank and store them in the configured destination
    100     Fetch {
    101         #[clap(flatten)]
    102         ebics: EbicsArgs,
    103 
    104         /// Only supported in --transient mode, this option lets specify the earliest timestamp of the downloaded documents
    105         #[arg(long, value_name = "YYYY-MM-DD", requires = "transient")]
    106         pinned_start: Option<Date>,
    107 
    108         /// Only supported in --transient mode, do not consume fetched documents
    109         #[arg(long, requires = "transient")]
    110         peek: bool,
    111 
    112         /// Only supported in --transient mode, run a checkpoint
    113         #[arg(long, requires = "transient")]
    114         checkpoint: bool,
    115         /// Check whether a destination is configured. Exit with 0 if at destination is configured, otherwise 1
    116         #[arg(long)]
    117         check: bool,
    118     },
    119     /// Run libeufin-ebisync HTTP server
    120     Serve {
    121         /// Check whether an API is in use (if it's useful to start the HTTP server). Exit with 0 if at least one API is enabled, otherwise 1
    122         #[arg(long)]
    123         check: bool,
    124     },
    125     #[command(subcommand)]
    126     Config(ConfigCmd),
    127     #[command(subcommand)]
    128     Testing(TestingCmd),
    129 }
    130 
    131 #[derive(clap::Parser, Debug)]
    132 #[command(long_version = long_version(), about, long_about = None)]
    133 pub struct Args {
    134     #[clap(flatten)]
    135     pub common: CommonArgs,
    136 
    137     #[command(subcommand)]
    138     pub cmd: Cmd,
    139 }
    140 
    141 pub async fn ebics_setup(
    142     ebics: &EbicsClient<'_>,
    143     cfg: &EbisyncCfg<'_>,
    144     db: &PgPool,
    145     force_keys_resubmission: bool,
    146     generate_registration_pdf: bool,
    147     auto_accept_keys: bool,
    148 ) -> anyhow::Result<()> {
    149     let (client, bank) = libeufin_ebics::setup::ebics_setup(
    150         ebics,
    151         &cfg.ebics_setup()?,
    152         force_keys_resubmission,
    153         generate_registration_pdf,
    154         auto_accept_keys,
    155     )
    156     .await?;
    157 
    158     // Check account information
    159     info!(target: "setup", "Doing administrative request HKD");
    160     let HKD { partner, .. } = ebics.hkd(db, &client, &bank, false).await?;
    161     // Debug logging
    162     let fmt = std::fmt::from_fn(|f| {
    163         if partner.name.is_some() || !partner.accounts.is_empty() {
    164             f.write_str("Partner Info: ")?;
    165             if let Some(name) = &partner.name {
    166                 write!(f, "'{name}'")?;
    167             }
    168             for AccountInfo {
    169                 currency,
    170                 iban,
    171                 bic,
    172             } in &partner.accounts
    173             {
    174                 write!(f, "{currency}-{iban}{bic}")?;
    175             }
    176             f.write_char('\n')?;
    177         }
    178         writeln!(f, "Supported orders:")?;
    179         for OrderInfo { order, description } in &partner.orders {
    180             writeln!(f, "- {order}: {description}")?;
    181         }
    182         Ok(())
    183     });
    184     debug!(target: "setup", "{fmt}");
    185     info!(target: "setup", "EBICS ready");
    186 
    187     info!(target: "setup", "Check fetch destination setup");
    188     if let Some(dest) =
    189         DestinationClient::prepare(&cfg.fetch()?.destination, &http_client::client())?
    190     {
    191         match dest {
    192             DestinationClient::AzureBlobStorage { client, container } => {
    193                 client.container_metadata(container).await?;
    194             }
    195         }
    196     } else {
    197         warn!(target: "setup", "No destination configured");
    198     }
    199     info!(target: "setup", "Fetch destination ready");
    200 
    201     println!("setup ready");
    202     Ok(())
    203 }
    204 
    205 pub enum DestinationClient<'a> {
    206     AzureBlobStorage {
    207         client: AzureBlobStorage<'a>,
    208         container: &'a str,
    209     },
    210 }
    211 
    212 impl<'a> DestinationClient<'a> {
    213     pub fn prepare(dest: &'a Destination, client: &'a Client) -> anyhow::Result<Option<Self>> {
    214         Ok(match dest {
    215             Destination::None => None,
    216             Destination::AzureBlobStorage {
    217                 api,
    218                 name: account,
    219                 key,
    220                 container,
    221             } => Some(DestinationClient::AzureBlobStorage {
    222                 client: AzureBlobStorage::new(api.as_str(), account, key, client)?,
    223                 container,
    224             }),
    225         })
    226     }
    227 
    228     pub async fn upload(&self, name: &str, xml: Bytes) -> anyhow::Result<()> {
    229         match self {
    230             DestinationClient::AzureBlobStorage { client, container } => {
    231                 client
    232                     .put_blob(
    233                         container,
    234                         name,
    235                         xml,
    236                         HeaderValue::from_static("application/xml"),
    237                     )
    238                     .await?;
    239             }
    240         }
    241         Ok(())
    242     }
    243 }
    244 
    245 pub async fn ebics_fetch(
    246     ebics: &EbicsClient<'_>,
    247     cfg: &EbisyncCfg<'_>,
    248     client: &ClientKeys,
    249     bank: &BankKeys,
    250     db: &PgPool,
    251     pinned_start: &Option<Timestamp>,
    252     peek: bool,
    253     check: bool,
    254     transient: bool,
    255     transient_checkpoint: bool,
    256 ) -> anyhow::Result<()> {
    257     let fetch_cfg = cfg.fetch()?;
    258     let (sender, mut receiver) = tokio::sync::mpsc::channel::<Vec<Order>>(10);
    259     let http = http_client::client();
    260     let dest = DestinationClient::prepare(&fetch_cfg.destination, &http)?;
    261 
    262     if check {
    263         if dest.is_none() {
    264             info!("No destination configured, not starting the fetcher");
    265             std::process::exit(1);
    266         }
    267     } else if let Some(dest) = dest {
    268         let upload = async |orders: &[Order], since: Option<Timestamp>| -> anyhow::Result<()> {
    269             for order in orders {
    270                 if let Err(e) = ebics
    271                     .download(
    272                         db,
    273                         client,
    274                         bank,
    275                         order,
    276                         &since.map(|it| (it, Timestamp::now())),
    277                         transient && peek,
    278                         async |content| {
    279                             if order.doc() == Some(OrderDoc::acknowledgement) {
    280                                 for ack in parse_hac(&content)? {
    281                                     debug!(target: "fetch", "{ack}");
    282                                     if let Some(order_id) = &ack.order_id {
    283                                         match ack.action {
    284                                             HacAction::ORDER_HAC_FINAL_POS => {
    285                                                 info!(target: "fetch", "Order {order_id} accepted at {}", ack.timestamp);
    286                                             }
    287                                             HacAction::ORDER_HAC_FINAL_NEG => {
    288                                                 info!(target: "fetch", "Order {order_id} refused at {}", ack.timestamp);
    289                                             }
    290                                             _ => {}
    291                                         }
    292                                     }
    293                                 }
    294                             } else {
    295                                 let mut z = zip::ZipArchive::new(Cursor::new(content))?;
    296                                 for i in 0..z.len() {
    297                                     let mut file = z.by_index(i)?;
    298                                     trace!(target: "fetch", "upload {}", file.name());
    299                                     let mut buf = Vec::new();
    300                                     file.read_to_end(&mut buf)?;
    301                                     dest.upload(file.name(), buf.into())
    302                                         .await
    303                                         .map_err(|e| EbicsErrKind::Custom(e.to_string().into()))?;
    304                                 }
    305                             }
    306                             Ok(())
    307                         },
    308                     )
    309                     .await
    310                 {
    311                     if let EbicsErrKind::Code {
    312                         bank: EbicsReturnCode::EBICS_NO_DOWNLOAD_DATA_AVAILABLE,
    313                         ..
    314                     } = e.kind
    315                     {
    316                         continue;
    317                     }
    318                     return Err(e.into());
    319                 }
    320             }
    321             Ok(())
    322         };
    323         let fetch = async {
    324             if transient {
    325                 info!(target: "fetch", "Transient mode: fetching once and returning");
    326             } else {
    327                 info!(target: "fetch", "Running with a frequency of {}", fetch_cfg.frequency_raw);
    328             }
    329 
    330             let mut last_fetch = Timestamp::UNIX_EPOCH;
    331             loop {
    332                 let now = Timestamp::now();
    333                 let checkpoint = get_task_status(db, CHECKPOINT_KEY)
    334                     .await?
    335                     .unwrap_or_default();
    336                 let next_fetch = last_fetch + fetch_cfg.frequency;
    337                 let next_checkpoint = {
    338                     if let Some(last_trial) = checkpoint.last_trial {
    339                         // We run today at checkpoint_time
    340                         let checkpoint_date = Zoned::new(now, TimeZone::UTC)
    341                             .with()
    342                             .time(fetch_cfg.checkpoint_time)
    343                             .build()
    344                             .unwrap();
    345                         // If we already ran today we ran tomorrow
    346                         if last_trial > checkpoint_date.timestamp() {
    347                             checkpoint_date.tomorrow().unwrap().timestamp()
    348                         } else {
    349                             checkpoint_date.timestamp()
    350                         }
    351                     } else {
    352                         // We never ran, we must checkpoint now
    353                         Timestamp::UNIX_EPOCH
    354                     }
    355                 };
    356 
    357                 let mut success = true;
    358                 if
    359                 // Run transient checkpoint at request
    360                 (transient && transient_checkpoint)
    361                     // Or run recurrent checkpoint
    362                     || (!transient && now > next_checkpoint)
    363                 {
    364                     info!(target: "fetch", "Running checkpoint");
    365 
    366                     let since = if let Some(pinned_start) = pinned_start
    367                         && transient
    368                         && checkpoint
    369                             .last_successfull
    370                             .map(|it| *pinned_start <= it)
    371                             .unwrap_or(true)
    372                     {
    373                         Some(*pinned_start)
    374                     } else {
    375                         checkpoint.last_successfull
    376                     };
    377                     let res = async {
    378                         // We fetch HKD to only fetch supported EBICS orders and get the document versions
    379                         let hkd = ebics.hkd(db, client, bank, false).await?;
    380                         let mut supported_orders = hkd
    381                             .partner
    382                             .orders
    383                             .into_iter()
    384                             .map(|it| it.order)
    385                             .collect::<Vec<_>>();
    386                         debug!(target: "fetch", "HKD: {}", supported_orders.iter().map(|it| it.to_string()).join_compact(", "));
    387                         supported_orders
    388                             .retain(|it| it.is_downloadable());
    389                         upload( &supported_orders, since).await
    390                     }
    391                     .await;
    392                     if let Err(e) = res {
    393                         success = false;
    394                         error!(target: "fetch", "{e}");
    395                     }
    396                     try_join!(
    397                         update_task_status(db, CHECKPOINT_KEY, &now, success),
    398                         update_task_status(db, FETCH_TASK_KEY, &now, success)
    399                     )?;
    400                     last_fetch = now;
    401                 } else if transient || now > next_fetch {
    402                     if !transient {
    403                         info!(target: "fetch", "Running at frequency");
    404                     }
    405                     let res = async {
    406                         // We fetch HAA to only fetch pending & supported EBICS orders and get the document versions
    407                         let mut haa = ebics.haa(db, client, bank, false).await?;
    408                         debug!(target: "fetch", "HAA: {}", haa.orders.iter().map(|it| it.to_string()).join_compact(", "));
    409                         haa.orders
    410                             .retain(|it| it.is_downloadable());
    411                         upload( &haa.orders, *pinned_start).await
    412                     }
    413                     .await;
    414                     if let Err(e) = res {
    415                         success = false;
    416                         error!(target: "fetch", "{e}");
    417                     }
    418                     update_task_status(db, FETCH_TASK_KEY, &now, success).await?;
    419                     last_fetch = now;
    420                 }
    421 
    422                 if transient {
    423                     if success {
    424                         return anyhow::Ok(());
    425                     } else {
    426                         return Err(anyhow!("fetch failed"));
    427                     }
    428                 }
    429 
    430                 let delay = now
    431                     .duration_until(next_fetch.min(next_checkpoint))
    432                     .max(SignedDuration::ZERO)
    433                     .unsigned_abs();
    434                 let tx = timeout(delay, receiver.recv()).await;
    435                 if let Ok(Some(notification)) = tx {
    436                     info!(target: "fetch", "Running at real-time notifications reception");
    437                     if let Err(e) = upload(&notification, None).await {
    438                         error!(target: "fetch", "{e}");
    439                     }
    440                 }
    441             }
    442         };
    443         if transient {
    444             fetch.await?;
    445         } else {
    446             tokio::try_join!(fetch, async {
    447                 listen_for_notification(ebics, db, client, bank, sender).await;
    448                 Ok(())
    449             })?;
    450         }
    451     }
    452 
    453     Ok(())
    454 }
    455 
    456 pub async fn run(cfg: &Config, cmd: Cmd) -> anyhow::Result<()> {
    457     match cmd {
    458         Cmd::Dbinit { reset } => {
    459             let cfg = parse_db_cfg(cfg)?;
    460             dbinit(&cfg, reset).await?;
    461         }
    462         Cmd::Setup {
    463             ebics_logs,
    464             force_keys_resubmission,
    465             auto_accept_keys,
    466             generate_registration_pdf,
    467         } => {
    468             let cfg = EbisyncCfg::parse(cfg)?;
    469             let pool = pool(&cfg.db_cfg).await?;
    470             let ebics = EbicsClient::new(cfg.host.ebics(), ebics_logs)?;
    471             ebics_setup(
    472                 &ebics,
    473                 &cfg,
    474                 &pool,
    475                 force_keys_resubmission,
    476                 generate_registration_pdf,
    477                 auto_accept_keys,
    478             )
    479             .await?;
    480         }
    481         Cmd::Fetch {
    482             pinned_start,
    483             peek,
    484             checkpoint,
    485             check,
    486             ebics: EbicsArgs { logs, transient },
    487         } => {
    488             let cfg = EbisyncCfg::parse(cfg)?;
    489             let pool = pool(&cfg.db_cfg).await?;
    490             let ebics = EbicsClient::new(cfg.host.ebics(), logs)?;
    491             let (client, bank) = cfg.expect_full_keys()?;
    492             ebics_fetch(
    493                 &ebics,
    494                 &cfg,
    495                 &client,
    496                 &bank,
    497                 &pool,
    498                 &pinned_start.map(|it| date_to_utc_ts(&it)),
    499                 peek,
    500                 check,
    501                 transient,
    502                 transient && checkpoint,
    503             )
    504             .await?
    505         }
    506         Cmd::Serve { check } => {
    507             let cfg = EbisyncCfg::parse(cfg)?;
    508             let auth = match &cfg.submit()?.source {
    509                 Source::None => None,
    510                 Source::SyncAPI(auth_cfg) => Some(auth_cfg.method()),
    511             };
    512             if check {
    513                 if auth.is_none() {
    514                     info!("No source api, not starting the server");
    515                     std::process::exit(1);
    516                 }
    517             } else if let Some(auth) = auth {
    518                 let (client, bank) = cfg.expect_full_keys()?;
    519                 let tmp = EbisyncState {
    520                     db: pool(&cfg.db_cfg).await?,
    521                     cfg: cfg.host.clone(),
    522                     client,
    523                     bank,
    524                 };
    525                 let cfg = cfg.serve()?;
    526 
    527                 sync_api(Arc::new(tmp), &cfg.spa_path, auth)
    528                     .serve(&cfg.serve, None)
    529                     .await?;
    530             }
    531         }
    532         Cmd::Config(cmd) => cmd.run(cfg)?,
    533         Cmd::Testing(cmd) => match cmd {
    534             TestingCmd::CreateContainer => {
    535                 let cfg = EbisyncCfg::parse(cfg)?;
    536                 let http = http_client::client();
    537                 let dest = DestinationClient::prepare(&cfg.fetch()?.destination, &http)?;
    538                 if let Some(DestinationClient::AzureBlobStorage { client, container }) = dest {
    539                     client.create_container(container).await?;
    540                     info!("Created container {container}");
    541                 } else {
    542                     bail!("Destination is not azure-blob-storage")
    543                 }
    544             }
    545         },
    546     }
    547     Ok(())
    548 }
    549 
    550 #[cfg(test)]
    551 mod ebics {
    552     use clap::Parser as _;
    553     use libeufin_ebics::test::{EbicsState, TestBank};
    554     use sqlx::PgPool;
    555     use taler_common::config::Config;
    556     use taler_macros::db_test;
    557 
    558     use crate::{Args, constants::CONFIG_SOURCE, run};
    559 
    560     pub async fn ebisync_cmd(cfg: &Config, cmd: &str) -> anyhow::Result<()> {
    561         let parts = shlex::split(cmd).unwrap();
    562         let args = std::iter::once("libeufin-ebisync").chain(parts.iter().map(|it| it.as_str()));
    563 
    564         let cmd = Args::try_parse_from(args).unwrap();
    565         run(cfg, cmd.cmd).await
    566     }
    567 
    568     async fn test_setup(db: PgPool) -> (TestBank, Config, PgPool) {
    569         let test = TestBank::new().await;
    570         let cfg = Config::from_mem_with_env(
    571             CONFIG_SOURCE,
    572             &format!(
    573                 "
    574                 [paths]
    575                 EBISYNC_HOME = {:?}
    576 
    577                 [ebisync]
    578                 HOST_BASE_URL = http://bank.example.com/
    579                 UNIXPATH = {}
    580                 HOST_ID = PFEBICS
    581                 USER_ID = PFC00563
    582                 PARTNER_ID = PFC00563
    583 
    584                 [ebisyncdb-postgres]
    585                 CONFIG = postgresql:///{}
    586                 ",
    587                 test.dir.path(),
    588                 test.sock_path,
    589                 db.connect_options().get_database().unwrap()
    590             ),
    591         )
    592         .unwrap();
    593         test.sequences(&[
    594             EbicsState::hev,
    595             EbicsState::ini,
    596             EbicsState::hia,
    597             EbicsState::hpb,
    598             EbicsState::hkd,
    599             EbicsState::receipt_ok,
    600         ]);
    601         ebisync_cmd(&cfg, "setup --auto-accept-keys").await.unwrap();
    602 
    603         (test, cfg, db)
    604     }
    605 
    606     #[db_test]
    607     async fn setup(db: PgPool) {
    608         test_setup(db).await;
    609     }
    610 }