libeufin

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

commit 995bed3a7128a9d9b31b62c398213ed67113a192
parent 2d74f521631a6cf1c4f4797b7c46767e8b0f1d5c
Author: Antoine A <>
Date:   Wed, 27 May 2026 18:38:38 +0200

nexus: add manual cmd

Diffstat:
Mcrates/libeufin-ebics/src/iso20022/pain001.rs | 8++++----
Mcrates/libeufin-nexus/src/db/initiated.rs | 12++++++------
Mcrates/libeufin-nexus/src/fetch.rs | 2+-
Mcrates/libeufin-nexus/src/lib.rs | 72+++++++++++++++++++++++++++++++++++++++++-------------------------------
Acrates/libeufin-nexus/src/manual.rs | 164+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/libeufin-nexus/src/model.rs | 8+++++---
Mcrates/libeufin-nexus/src/test.rs | 7+++----
7 files changed, 224 insertions(+), 49 deletions(-)

diff --git a/crates/libeufin-ebics/src/iso20022/pain001.rs b/crates/libeufin-ebics/src/iso20022/pain001.rs @@ -32,7 +32,7 @@ use crate::{ /** pain.001 transaction metadata */ pub struct Pain001Tx<'a> { - pub creditor: FullIbanPayto, + pub creditor: &'a FullIbanPayto, pub amount: Amount, pub subject: &'a str, pub e2e_id: &'a str, @@ -193,7 +193,7 @@ mod test { #[test] fn pain001() { - let creditor = FullIbanPayto::new( + let creditor = &FullIbanPayto::new( BankID { iban: "CH4189144589712575493".parse().expect("invalid IBAN"), bic: None, @@ -214,13 +214,13 @@ mod test { sum: amount("CHF:47.32"), txs: vec![ Pain001Tx { - creditor: creditor.clone(), + creditor, amount: amount("CHF:42"), subject: "Test 42", e2e_id: "TX_FIRST", }, Pain001Tx { - creditor: creditor.clone(), + creditor, amount: amount("CHF:5.11"), subject: "Test 5.11", e2e_id: "TX_SECOND", diff --git a/crates/libeufin-nexus/src/db/initiated.rs b/crates/libeufin-nexus/src/db/initiated.rs @@ -26,7 +26,7 @@ use taler_api::{ }; use taler_common::types::{ amount::{Amount, Currency}, - payto::PaytoURI, + payto::FullIbanPayto, }; use crate::{ @@ -46,7 +46,7 @@ pub async fn initiate( pool: &PgPool, amount: &Amount, subject: &str, - creditor: &PaytoURI, + creditor: &FullIbanPayto, initiation_time: &Timestamp, e2e_id: &str, ) -> sqlx::Result<PaymentInitiationResult> { @@ -65,7 +65,7 @@ pub async fn initiate( ) .bind(amount) .bind(subject) - .bind(creditor.as_ref().as_str()) + .bind(creditor.as_uri().as_ref().as_str()) .bind_timestamp(initiation_time) .bind(e2e_id) .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?))) @@ -98,13 +98,13 @@ pub async fn batch_initiated( Ok(()) } -pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<()> { - serialized!( +pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<bool> { + let res = serialized!( sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1") .bind(id as i64) .execute(db) )?; - Ok(()) + Ok(res.rows_affected() > 0) } pub async fn initiated_submittable( diff --git a/crates/libeufin-nexus/src/fetch.rs b/crates/libeufin-nexus/src/fetch.rs @@ -415,7 +415,7 @@ pub async fn ebics_fetch( Ok(()) } -async fn register_camt(db: &PgPool, cfg: &NexusCfg, xml: &[u8]) -> anyhow::Result<usize> { +pub async fn register_camt(db: &PgPool, cfg: &NexusCfg, xml: &[u8]) -> anyhow::Result<usize> { let account = &cfg.ebics()?.account; let ingest_cfg = cfg.ingest()?; let mut nb_tx = 0; diff --git a/crates/libeufin-nexus/src/lib.rs b/crates/libeufin-nexus/src/lib.rs @@ -17,7 +17,7 @@ * <http://www.gnu.org/licenses/> */ -use std::{fmt::Write, str::FromStr, sync::Arc, time::Duration}; +use std::{fmt::Write, sync::Arc, time::Duration}; use anyhow::{anyhow, bail}; use compact_str::{CompactString, CompactStringExt, ToCompactString}; @@ -41,18 +41,14 @@ use taler_common::{ CommonArgs, cli::ConfigCmd, config::{Config, parser::ConfigSource}, - types::{ - amount::Amount, - payto::{FullIbanPayto, TransferIbanPayto}, - utils::date_to_utc_ts, - }, + types::{amount::Amount, payto::TransferIbanPayto, utils::date_to_utc_ts}, }; use taler_test_utils::Router; use tracing::{debug, error, info, warn}; use crate::{ api::NexusApi, - config::NexusCfg, + config::{NexusCfg, NexusEbicsConfig}, db::{ dbinit, initiated::{ @@ -62,6 +58,7 @@ use crate::{ }, fetch::ebics_fetch, list::ListCmd, + manual::ManualCmd, model::PaymentBatch, testing::TestingCmd, }; @@ -72,6 +69,7 @@ pub mod config; pub mod db; pub mod fetch; pub mod list; +pub mod manual; pub mod model; #[cfg(test)] pub mod test; @@ -165,7 +163,8 @@ pub enum Cmd { /// The credited account IBAN payto UR payto: TransferIbanPayto, }, - Manual {}, + #[command(subcommand)] + Manual(ManualCmd), #[command(subcommand)] List(ListCmd), #[command(subcommand)] @@ -184,6 +183,33 @@ pub struct Args { pub cmd: Cmd, } +pub fn batch_pain001( + batch: &PaymentBatch, + cfg: &NexusEbicsConfig, + instant: bool, +) -> Result<String, EbicsErrKind> { + let msg = Pain001Msg { + msg_id: &batch.msg_id, + timestamp: &Timestamp::now(), + debtor: &cfg.account, + sum: batch.sum, + txs: batch + .payments + .iter() + .map(|tx| { + // TODO handle missing name ? + Pain001Tx { + creditor: &tx.creditor, + amount: tx.amount, + subject: &tx.subject, + e2e_id: &tx.e2e_id, + } + }) + .collect(), + }; + create_pain001(&msg, &cfg.dialect, instant) +} + pub async fn ebics_submit( ebics: &EbicsClient<'_>, cfg: &NexusCfg, @@ -200,27 +226,7 @@ pub async fn ebics_submit( instant: bool| -> Result<CompactString, EbicsError> { let ctx = EbicsCtx::new(order); - let msg = Pain001Msg { - msg_id: &batch.msg_id, - timestamp: &Timestamp::now(), - debtor: &ebics_cfg.account, - sum: batch.sum, - txs: batch - .payments - .iter() - .map(|tx| { - let creditor = FullIbanPayto::from_str(tx.creditor.as_ref().as_str()).unwrap(); - // TODO handle missing name ? - Pain001Tx { - creditor, - amount: tx.amount, - subject: &tx.subject, - e2e_id: &tx.e2e_id, - } - }) - .collect(), - }; - let xml = create_pain001(&msg, &ebics_cfg.dialect, instant).ctx(&ctx)?; + let xml = batch_pain001(batch, ebics_cfg, instant).ctx(&ctx)?; ebics.upload(client, bank, order, &xml).await }; @@ -533,7 +539,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { &pool, amount, subject, - &payto.as_uri(), + &payto.clone().into(), &Timestamp::now(), &end_to_end_id .as_ref() @@ -577,7 +583,11 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { server.finalize().serve(&cfg.serve_cfg, None).await?; } } - Cmd::Manual {} => todo!(), + Cmd::Manual(cmd) => { + let pool = pool(&cfg).await?; + let cfg = NexusCfg::parse(cfg)?; + cmd.run(&pool, &cfg).await?; + } Cmd::List(cmd) => { let pool = pool(&cfg).await?; let cfg = NexusCfg::parse(cfg)?; diff --git a/crates/libeufin-nexus/src/manual.rs b/crates/libeufin-nexus/src/manual.rs @@ -0,0 +1,164 @@ +/* +* This file is part of LibEuFin. +* Copyright (C) 2026 Taler Systems S.A. + +* LibEuFin is free software; you can redistribute it and/or modify +* it under the terms of the GNU Affero General Public License as +* published by the Free Software Foundation; either version 3, or +* (at your option) any later version. + +* LibEuFin is distributed in the hope that it will be useful, but +* WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY +* or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Affero General +* Public License for more details. + +* You should have received a copy of the GNU Affero General Public +* License along with LibEuFin; see the file COPYING. If not, see +* <http://www.gnu.org/licenses/> +*/ + +use std::{ + fmt::Write as _, + io::{BufWriter, Read, Write, stdin, stdout}, +}; + +use anyhow::bail; +use compact_str::CompactString; +use jiff::{Timestamp, tz::TimeZone}; +use libeufin_ebics::ebics::rand_ebics_id; +use sqlx::PgPool; +use taler_macros::EnumMeta; +use tracing::{info, warn}; +use zip::write::FileOptions; + +use crate::{ + batch_pain001, + config::NexusCfg, + db::initiated::{ + batch_initiated, batch_status_update, initiated_ack, initiated_submittable, + tx_status_update, + }, + fetch::register_camt, + model::SubmissionState, +}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, EnumMeta)] +#[enum_meta(Str)] +pub enum Kind { + batch, + tx, +} + +/// Manual management commands +#[derive(clap::Subcommand, Debug)] +pub enum ManualCmd { + /// Export pending batches as pain001 messages + Export { out: CompactString }, + /// Import EBICS camt files + Import { sources: Vec<CompactString> }, + /// Change batches or transactions status + Status { + kind: Kind, + id: CompactString, + status: SubmissionState, + msg: Option<CompactString>, + }, + /// Manually acknowledge the outgoing transaction for submission + Ack { ids: Vec<u64> }, +} + +impl ManualCmd { + pub async fn run(self, db: &PgPool, cfg: &NexusCfg) -> anyhow::Result<()> { + match self { + ManualCmd::Export { out } => { + let out: &mut dyn Write = if out == "-" { + &mut stdout().lock() + } else { + &mut BufWriter::new(std::fs::File::open(out)?) + }; + + // Create and get pending batches + batch_initiated( + db, + &Timestamp::now(), + &rand_ebics_id(), + cfg.submit()?.require_ack, + ) + .await?; + let mut zip = zip::ZipWriter::new_stream(&mut *out); + let batches = initiated_submittable(&db, &cfg.currency).await?; + let mut metadata = String::new(); + write!(&mut metadata, "Exported {} pain.001 files:", batches.len()).unwrap(); + for batch in batches { + let datetime = batch.creation_date.to_zoned(TimeZone::UTC).datetime(); + zip.start_file( + format!("{datetime}-{}.xml", batch.msg_id), + FileOptions::DEFAULT, + )?; + let xml = batch_pain001(&batch, cfg.ebics()?, false)?; + zip.write_all(xml.as_bytes())?; + write!(&mut metadata, "\nbatch {}:", batch.msg_id)?; + for tx in batch.payments { + write!( + &mut metadata, + "\n- tx {} {} {} '{}'", + tx.e2e_id, tx.amount, tx.creditor.iban, tx.creditor.name + )?; + } + write!(&mut metadata, "\n")?; + } + zip.start_file("README.txt", FileOptions::DEFAULT)?; + zip.write_all(metadata.as_bytes())?; + info!("{metadata}"); + } + ManualCmd::Import { sources } => { + for source in sources { + let xml = if source == "-" { + let mut buf = Vec::with_capacity(8 * 1024); + stdin().lock().read_to_end(&mut buf)?; + buf + } else { + std::fs::read(&source)? + }; + let nb = register_camt(db, cfg, &xml).await?; + info!("Imported {nb} transactions from {source}") + } + } + ManualCmd::Status { + kind, + id, + status, + msg, + } => match kind { + Kind::batch => { + if batch_status_update(&db, &id, status, msg.as_deref().unwrap_or_default()) + .await? + { + info!("Updated batch '{id}' to {status}"); + } else { + bail!("Unknown batch {id}") + } + } + Kind::tx => { + if tx_status_update(&db, &id, "", status, msg.as_deref().unwrap_or_default()) + .await? + { + info!("Updated tx '{id}' to {status}"); + } else { + bail!("Unknown tx {id}") + } + } + }, + ManualCmd::Ack { ids } => { + for id in ids { + if initiated_ack(&db, id).await? { + info!("Mark {id} as acknowledge for submission"); + } else { + warn!("Unknown transaction {id}"); + } + } + } + } + Ok(()) + } +} diff --git a/crates/libeufin-nexus/src/model.rs b/crates/libeufin-nexus/src/model.rs @@ -21,12 +21,14 @@ use compact_str::CompactString; use jiff::Timestamp; use taler_common::{ api::wire::TransferState, - types::{amount::Amount, payto::PaytoURI}, + types::{amount::Amount, payto::FullIbanPayto}, }; +use taler_macros::EnumMeta; -#[derive(Debug, Clone, Copy, PartialEq, Eq, sqlx::Type)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, sqlx::Type, EnumMeta)] #[allow(non_camel_case_types)] #[sqlx(type_name = "submission_state")] +#[enum_meta(Str)] /** Outgoing transactions and batches submission status */ pub enum SubmissionState { // Initiated but not yet submitted @@ -82,7 +84,7 @@ pub struct Initiated { pub id: u64, pub amount: Amount, pub subject: String, - pub creditor: PaytoURI, + pub creditor: FullIbanPayto, pub initiation_time: Timestamp, pub e2e_id: CompactString, } diff --git a/crates/libeufin-nexus/src/test.rs b/crates/libeufin-nexus/src/test.rs @@ -30,7 +30,7 @@ use taler_common::{ types::{ amount::{Amount, Currency}, base32::Base32, - payto::{IbanPayto, PaytoURI, payto}, + payto::{FullIbanPayto, IbanPayto, PaytoURI, payto}, }, }; @@ -77,9 +77,8 @@ pub fn gen_init_pay( Initiated { id: 0, amount: Amount::new(&CURR, 44, 0), - creditor: IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test") - .unwrap() - .as_uri(), + creditor: FullIbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test") + .unwrap(), subject: subject.into(), initiation_time: Timestamp::now(), e2e_id: end_to_end_id.into(),