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(¬ification, 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 }