libeufin

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

commit 199defca06ea8d65d0b772864aadda861c3319d5
parent ef707257df28e301cc91c74cef54e8a4dc698113
Author: Antoine A <>
Date:   Mon,  7 Sep 2026 17:38:34 +0200

nexus: many fixes

Diffstat:
MCargo.lock | 12++++++------
Mlibeufin-ebics/src/setup.rs | 2+-
Mlibeufin-ebisync/src/lib.rs | 2+-
Mlibeufin-nexus/src/db/initiated.rs | 2+-
Mlibeufin-nexus/src/fetch.rs | 45++++++++++++++++++++++++++++-----------------
Mlibeufin-nexus/src/lib.rs | 28+++++++++++++++-------------
Mlibeufin-nexus/src/list.rs | 2+-
Mlibeufin-nexus/src/manual.rs | 3++-
Mlibeufin-nexus/src/testing.rs | 22+++++++++++++++++++---
Mtestbench/src/main.rs | 9+++++++++
10 files changed, 83 insertions(+), 44 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -1361,7 +1361,7 @@ dependencies = [ [[package]] name = "http-client" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#18ef68372e0ee690dd625bb99dc277bc12e1894c" +source = "git+git://git.taler.net/taler-rust.git/#8fa0a1d6d78cd49f88509404c15b1c2032d62401" dependencies = [ "compact_str", "futures-util", @@ -3493,7 +3493,7 @@ dependencies = [ [[package]] name = "taler-api" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#18ef68372e0ee690dd625bb99dc277bc12e1894c" +source = "git+git://git.taler.net/taler-rust.git/#8fa0a1d6d78cd49f88509404c15b1c2032d62401" dependencies = [ "aws-lc-rs", "axum", @@ -3521,12 +3521,12 @@ dependencies = [ [[package]] name = "taler-build" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#18ef68372e0ee690dd625bb99dc277bc12e1894c" +source = "git+git://git.taler.net/taler-rust.git/#8fa0a1d6d78cd49f88509404c15b1c2032d62401" [[package]] name = "taler-common" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#18ef68372e0ee690dd625bb99dc277bc12e1894c" +source = "git+git://git.taler.net/taler-rust.git/#8fa0a1d6d78cd49f88509404c15b1c2032d62401" dependencies = [ "anyhow", "aws-lc-rs", @@ -3556,7 +3556,7 @@ dependencies = [ [[package]] name = "taler-macros" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#18ef68372e0ee690dd625bb99dc277bc12e1894c" +source = "git+git://git.taler.net/taler-rust.git/#8fa0a1d6d78cd49f88509404c15b1c2032d62401" dependencies = [ "proc-macro2", "quote", @@ -3566,7 +3566,7 @@ dependencies = [ [[package]] name = "taler-test-utils" version = "1.5.0" -source = "git+git://git.taler.net/taler-rust.git/#18ef68372e0ee690dd625bb99dc277bc12e1894c" +source = "git+git://git.taler.net/taler-rust.git/#8fa0a1d6d78cd49f88509404c15b1c2032d62401" dependencies = [ "aws-lc-rs", "axum", diff --git a/libeufin-ebics/src/setup.rs b/libeufin-ebics/src/setup.rs @@ -94,7 +94,7 @@ pub async fn ebics_setup( let path = "/tmp/libeufin-ebics-keys.pdf"; std::fs::write("/tmp/libeufin-ebics-keys.pdf", &pdf) .map_err(|e| anyhow!("Could not write PDF to '{path}': {}", e.kind()))?; - info!(target: "setup", "PDF file with keys created at '{path}'"); + println!("PDF file with keys created at '{path}'"); } if !client.submitted_hia || force_keys_resubmission { ebics diff --git a/libeufin-ebisync/src/lib.rs b/libeufin-ebisync/src/lib.rs @@ -201,7 +201,7 @@ pub async fn ebics_setup( } info!(target: "setup", "Fetch destination ready"); - eprintln!("setup ready"); + println!("setup ready"); Ok(()) } diff --git a/libeufin-nexus/src/db/initiated.rs b/libeufin-nexus/src/db/initiated.rs @@ -436,7 +436,7 @@ pub async fn batch_status_update( )) .bind(msg_id) .bind(state) - .bind(msg) + .bind((!msg.is_empty()).then_some(msg)) .try_map(|r: PgRow| r.try_get(0)) .fetch_one(db) ) diff --git a/libeufin-nexus/src/fetch.rs b/libeufin-nexus/src/fetch.rs @@ -48,7 +48,7 @@ use sqlx::PgPool; use taler_api::subject::{ IncomingSubject, parse_incoming_unstructured, parse_outgoing, subject_is_qr_bill, }; -use taler_common::types::amount::Currency; +use taler_common::types::{amount::Currency, payto::ParsedQuery}; use tokio::{time::timeout, try_join}; use tracing::{debug, error, info, trace, warn}; @@ -76,7 +76,7 @@ pub async fn ebics_fetch( client: &ClientKeys, bank: &BankKeys, db: &PgPool, - documents: Option<&[OrderDoc]>, + documents: &[OrderDoc], pinned_start: &Option<Timestamp>, peek: bool, transient: bool, @@ -152,11 +152,14 @@ pub async fn ebics_fetch( }; // EBICS order than should be fetched - let orders: Vec<_> = documents - .unwrap_or(OrderDoc::entries) - .iter() - .flat_map(|it| ebics_cfg.dialect.standard().downloads(it)) - .collect(); + let orders: Vec<_> = if documents.is_empty() { + OrderDoc::entries + } else { + documents + } + .iter() + .flat_map(|it| ebics_cfg.dialect.standard().downloads(it)) + .collect(); let fetch_cfg = cfg.fetch()?; @@ -192,7 +195,7 @@ pub async fn ebics_fetch( } } else { // We never ran, we must checkpoint now - now + Timestamp::UNIX_EPOCH } }; @@ -279,7 +282,9 @@ pub async fn ebics_fetch( notification.retain(|order| orders.iter().find(|it| order.eq(it)).is_some()); if !notification.is_empty() { info!(target: "fetch", "Running at real-time notifications reception"); - fetch(&notification, None).await?; + if let Err(e) = fetch(&notification, None).await { + error!(target: "fetch", "{e}"); + } } } } @@ -344,7 +349,7 @@ async fn register_file( } PaymentGroupStatus::Rejected => { error!(target: "fetch", "Batch {} failed: {msg}", msg_status.id); - SubmissionState::success + SubmissionState::permanent_failure } _ => SubmissionState::pending, }, @@ -366,7 +371,7 @@ async fn register_file( } PaymentGroupStatus::Rejected => { error!(target: "fetch", "Batch {} failed: {msg}", msg_status.id); - SubmissionState::success + SubmissionState::permanent_failure } _ => SubmissionState::pending, }, @@ -447,7 +452,7 @@ pub async fn register_incoming( let log_res = |res: InResult, kind: &str, suffix: &str| { let fmt = std::fmt::from_fn(|f| { write!(f, "{payment}")?; - if kind.is_empty() { + if !kind.is_empty() { write!(f, " {kind}")?; } if res.new { @@ -466,7 +471,7 @@ pub async fn register_incoming( } } } - if suffix.is_empty() { + if !suffix.is_empty() { write!(f, " {suffix}")?; } Ok(()) @@ -538,16 +543,22 @@ pub async fn register_incoming( }; // Check we have enough info to handle this transaction - if payment.debtor.is_none() { - // TODO payment.debtor.receiverName == null + if match &payment.debtor { + Some(payto) => match payto.query::<ParsedQuery>() { + Ok(q) => q.receiver_name.is_none(), + Err(_) => true, + }, + None => true, + } { let res = register_in(db, payment).await?; log_res(res, "incomplete", ""); return Ok(()); } - // TODO if payment.debtor.is_none() && payment.debtor.map(|it| it.rec) if let Some(regex) = &cfg.restriction_payto_regex && let Some(debtor) = &payment.debtor - && !regex.is_match(debtor.as_ref().as_str()) + && !regex + .find(debtor.as_ref().as_str()) + .is_some_and(|m| m.as_str() == debtor.as_ref().as_str()) { bounce("restricted account").await?; return Ok(()); diff --git a/libeufin-nexus/src/lib.rs b/libeufin-nexus/src/lib.rs @@ -19,7 +19,7 @@ #![allow(clippy::too_many_arguments)] -use std::{fmt::Write, sync::Arc, time::Duration}; +use std::{fmt::Write, sync::Arc}; use anyhow::{anyhow, bail}; use compact_str::{CompactString, CompactStringExt, ToCompactString}; @@ -30,7 +30,7 @@ use libeufin_ebics::{ ebics::{ EbicsClient, EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper as _, administrative::{AccountInfo, HKD, OrderInfo}, - order::Order, + order::{Order, OrderDoc}, rand_ebics_id, }, iso20022::pain001::{Pain001Msg, Pain001Tx, create_pain001}, @@ -122,7 +122,7 @@ pub enum Cmd { ebics: EbicsArgs, /// Only supported in --transient mode, this option lets specify the earliest timestamp of the downloaded documents - #[arg(long, value_name = "YYYY-MM-DD")] + #[arg(long, value_name = "YYYY-MM-DD", requires = "transient")] pinned_start: Option<Date>, /// Only supported in --transient mode, do not consume fetched documents @@ -132,6 +132,8 @@ pub enum Cmd { /// Only supported in --transient mode, run a checkpoint #[arg(long, requires = "transient")] checkpoint: bool, + + documents: Vec<OrderDoc>, }, /// Run libeufin-nexus HTTP server Serve { @@ -244,7 +246,7 @@ pub async fn ebics_submit( for batch in initiated_submittable(db, &cfg.currency).await? { debug!(target: "submit", "Submitting batch {}", batch.msg_id); let res = async { - if let Some(instant) = standard.instant_direct_debit() { + if let Some(instant) = instant_order.clone() { match submit_batch(&instant, &batch, true).await { Ok(id) => return Ok(id), Err(e) => if let EbicsErrKind::Code { .. } = e.kind { @@ -300,13 +302,12 @@ pub async fn ebics_submit( if let Err(e) = update_task_status(db, SUBMIT_TASK_KEY, &now, success).await { warn!(target: "submit", "{e}"); } - tokio::time::sleep(Duration::from_millis( - Timestamp::now() - .duration_until(now + submit_cfg.frequency) - .abs() - .as_millis() as u64, - )) - .await; + let wait = Timestamp::now().duration_until(now + submit_cfg.frequency); + if wait.is_positive() { + tokio::time::sleep(wait.unsigned_abs()).await; + } else { + tokio::task::yield_now().await; + } } } } @@ -441,7 +442,7 @@ pub async fn ebics_setup( info!(target: "setup", "Subscriber status: {}", user.status.description()) } - eprintln!("setup ready"); + println!("setup ready"); Ok(()) } @@ -475,6 +476,7 @@ pub async fn run(cfg: &Config, cmd: Cmd) -> anyhow::Result<()> { peek, checkpoint, ebics: EbicsArgs { logs, transient }, + documents, } => { let cfg = NexusCfg::parse(cfg)?; let pool = pool(&cfg.db_cfg).await?; @@ -486,7 +488,7 @@ pub async fn run(cfg: &Config, cmd: Cmd) -> anyhow::Result<()> { &client, &bank, &pool, - None, + &documents, &pinned_start.map(|it| date_to_utc_ts(&it)), peek, transient, diff --git a/libeufin-nexus/src/list.rs b/libeufin-nexus/src/list.rs @@ -91,7 +91,7 @@ impl ListCmd { writeln!(out, " subject: {subject}")?; } if let Some(talerable) = talerable { - writeln!(out, " subject: {talerable}")?; + writeln!(out, " talerable: {talerable}")?; } if let Some(bounced) = bounced { writeln!(out, " bounced: {bounced}")?; diff --git a/libeufin-nexus/src/manual.rs b/libeufin-nexus/src/manual.rs @@ -74,7 +74,7 @@ impl ManualCmd { let out: &mut dyn Write = if out == "-" { &mut stdout().lock() } else { - &mut BufWriter::new(std::fs::File::open(out)?) + &mut BufWriter::new(std::fs::File::create_new(out)?) }; // Create and get pending batches @@ -109,6 +109,7 @@ impl ManualCmd { } zip.start_file("README.txt", FileOptions::DEFAULT)?; zip.write_all(metadata.as_bytes())?; + zip.finish()?; info!("{metadata}"); } ManualCmd::Import { sources } => { diff --git a/libeufin-nexus/src/testing.rs b/libeufin-nexus/src/testing.rs @@ -17,6 +17,8 @@ * <http://www.gnu.org/licenses/> */ +use std::io::{Cursor, Read}; + use anyhow::{anyhow, bail}; use compact_str::CompactString; use jiff::{Timestamp, civil::Date, tz::TimeZone}; @@ -47,7 +49,10 @@ use crate::{ #[derive(clap::Subcommand, Debug)] pub enum IbanCmd { /// Generate fake IBANs for testing - Gen { country: Country }, + Gen { + #[arg(long)] + country: Country, + }, } impl IbanCmd { @@ -249,7 +254,7 @@ impl TestingCmd { service: name, scope, option, - container, + container: container.clone(), msg: message_name, version: message_version, }), @@ -270,7 +275,18 @@ impl TestingCmd { ) }), peek, - async |_| { + async move |content| { + if container.as_deref() == Some("ZIP") { + let mut z = zip::ZipArchive::new(Cursor::new(content))?; + for i in 0..z.len() { + let mut file = z.by_index(i)?; + let mut buf = Vec::new(); + file.read_to_end(&mut buf)?; + println!("{}\n{}", file.name(), String::from_utf8_lossy(&buf)); + } + } else { + println!("{}", String::from_utf8_lossy(&content)); + } if dry_run { Err(EbicsErrKind::Custom("dry run".into())) } else { diff --git a/testbench/src/main.rs b/testbench/src/main.rs @@ -111,6 +111,9 @@ pub enum NexusCmd { path: PathBuf, glob: Vec<Glob>, }, + Export { + path: PathBuf, + }, Compact, } @@ -365,6 +368,9 @@ async fn main() -> anyhow::Result<()> { run(&format!("manual import {}", paths.join(" "))).await; } + NexusCmd::Export { path } => { + run(&format!("manual export {}", path.to_string_lossy())).await; + } NexusCmd::Exit => break, } } else { @@ -487,6 +493,9 @@ fn compact(p: String) { let content = std::fs::read_to_string(&file).unwrap(); rm_if_similar(&file, &content, &mut content_hashes); } + if is_dir_empty(&payload_dir_path) { + remove(&request, "empty payload"); + } } if request.exists() && !payload_file_path.exists() && !payload_dir_path.exists() { remove(&request, "empty request");