libeufin

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

commit d94adf8e06df23cd3c6dee58b13198fbbefea345
parent 9f6974a6bdec8324079045b9d3cee71833bf1c3e
Author: Antoine A <>
Date:   Fri, 24 Apr 2026 10:38:00 +0200

nexus: add websocket notification

Diffstat:
MCargo.lock | 149+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
MCargo.toml | 8++++++--
Msrc/api.rs | 11++++-------
Msrc/bin/testbench.rs | 6+++++-
Msrc/config.rs | 1+
Msrc/dialect.rs | 54+++++++++++++++++++++++++++---------------------------
Msrc/ebics/administrative.rs | 8++++----
Msrc/ebics/bts.rs | 8++++----
Msrc/ebics/key_management.rs | 449+++++++++++++++++++++++++++++++++++++++----------------------------------------
Msrc/ebics/mod.rs | 478+++++++++++++++++++++++++++++++++++++++----------------------------------------
Msrc/ebics/order.rs | 73++++++++++++++++++++++++++++++++++++++-----------------------------------
Msrc/iso20022/camt.rs | 14+++++---------
Msrc/lib.rs | 124+++++++++++++++++++++++++++++--------------------------------------------------
Msrc/list.rs | 6+++---
Msrc/model.rs | 6+++---
Asrc/testing.rs | 236+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Asrc/ws.rs | 439+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
17 files changed, 1407 insertions(+), 663 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -128,6 +128,22 @@ dependencies = [ ] [[package]] +name = "async-tungstenite" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8acc405d38be14342132609f06f02acaf825ddccfe76c4824a69281e0458ebd4" +dependencies = [ + "atomic-waker", + "futures-core", + "futures-io", + "futures-task", + "futures-util", + "log", + "pin-project-lite", + "tungstenite 0.28.0", +] + +[[package]] name = "atoi" version = "2.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -187,6 +203,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" dependencies = [ "axum-core", + "base64", "bytes", "form_urlencoded", "futures-util", @@ -205,8 +222,10 @@ dependencies = [ "serde_json", "serde_path_to_error", "serde_urlencoded", + "sha1", "sync_wrapper", "tokio", + "tokio-tungstenite", "tower", "tower-layer", "tower-service", @@ -478,9 +497,9 @@ checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" [[package]] name = "const_format" -version = "0.2.35" +version = "0.2.36" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7faa7469a93a566e9ccc1c73fe783b4a65c274c5ace346038dca9c39fe0030ad" +checksum = "4481a617ad9a412be3b97c5d403fef8ed023103368908b9c50af598ff467cc1e" dependencies = [ "const_format_proc_macros", "konst", @@ -977,6 +996,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cecba35d7ad927e23624b22ad55235f2239cfa44fd10428eecbeba6d6a717718" [[package]] +name = "futures-macro" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] name = "futures-sink" version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -996,6 +1026,7 @@ checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" dependencies = [ "futures-core", "futures-io", + "futures-macro", "futures-sink", "futures-task", "memchr", @@ -1607,12 +1638,14 @@ version = "0.1.0" dependencies = [ "anyhow", "aws-lc-rs", + "axum", "base64", "calamine", "clap", "compact_str", "const_format", "flate2", + "futures-util", "getrandom 0.4.2", "hex", "indicatif", @@ -1625,6 +1658,7 @@ dependencies = [ "reedline", "regex", "reqwest", + "reqwest-websocket", "roxmltree", "serde", "serde_json", @@ -1637,6 +1671,7 @@ dependencies = [ "taler-test-utils", "thiserror 2.0.18", "tokio", + "tokio-tungstenite", "tracing", "tracing-subscriber", "url", @@ -1819,7 +1854,7 @@ dependencies = [ "num-integer", "num-iter", "num-traits", - "rand 0.8.5", + "rand 0.8.6", "smallvec", "zeroize", ] @@ -1994,9 +2029,9 @@ checksum = "c33a9471896f1c69cecef8d20cbe2f7accd12527ce60845ff44c153bb2a21b49" [[package]] name = "portable-atomic-util" -version = "0.2.6" +version = "0.2.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "091397be61a01d4be58e7841595bd4bfedb15f1cd54977d79b8271e94ed799a3" +checksum = "c2a106d1259c23fac8e543272398ae0e3c0b8d33c88ed73d0cc71b0f1d902618" dependencies = [ "portable-atomic", ] @@ -2143,9 +2178,9 @@ checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" [[package]] name = "rand" -version = "0.8.5" +version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" +checksum = "5ca0ecfa931c29007047d1bc58e623ab12e5590e8c7cc53200d5202b69266d8a" dependencies = [ "libc", "rand_chacha 0.3.1", @@ -2337,6 +2372,24 @@ dependencies = [ ] [[package]] +name = "reqwest-websocket" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7705b649c3b66b85c4e9c304a6898b1ae3eecb880c474720ebf925e4a932ae02" +dependencies = [ + "async-tungstenite", + "bytes", + "futures-util", + "reqwest", + "thiserror 2.0.18", + "tokio", + "tokio-util", + "tracing", + "tungstenite 0.28.0", + "web-sys", +] + +[[package]] name = "ring" version = "0.17.14" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2481,9 +2534,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.12" +version = "0.103.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8279bb85272c9f10811ae6a6c547ff594d6a7f3c6c6b02ee9726d1d0dcfcdd06" +checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" dependencies = [ "aws-lc-rs", "ring", @@ -2889,7 +2942,7 @@ dependencies = [ "memchr", "once_cell", "percent-encoding", - "rand 0.8.5", + "rand 0.8.6", "rsa", "serde", "sha1", @@ -2928,7 +2981,7 @@ dependencies = [ "md-5", "memchr", "once_cell", - "rand 0.8.5", + "rand 0.8.6", "serde", "serde_json", "sha2", @@ -3291,9 +3344,9 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" [[package]] name = "tokio" -version = "1.52.0" +version = "1.52.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a91135f59b1cbf38c91e73cf3386fca9bb77915c45ce2771460c9d92f0f3d776" +checksum = "b67dee974fe86fd92cc45b7a95fdd2f99a36a6d7b0d431a231178d3d670bbcc6" dependencies = [ "bytes", "libc", @@ -3338,6 +3391,18 @@ dependencies = [ ] [[package]] +name = "tokio-tungstenite" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f72a05e828585856dacd553fba484c242c46e391fb0e58917c942ee9202915c" +dependencies = [ + "futures-util", + "log", + "tokio", + "tungstenite 0.29.0", +] + +[[package]] name = "tokio-util" version = "0.7.18" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3345,6 +3410,7 @@ checksum = "9ae9cec805b01e8fc3fd2fe289f89149a9b66dd16786abd8b19cfa7b48cb0098" dependencies = [ "bytes", "futures-core", + "futures-io", "futures-sink", "pin-project-lite", "tokio", @@ -3461,6 +3527,39 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] +name = "tungstenite" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8628dcc84e5a09eb3d8423d6cb682965dea9133204e8fb3efee74c2a0c259442" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.4", + "sha1", + "thiserror 2.0.18", + "utf-8", +] + +[[package]] +name = "tungstenite" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c01152af293afb9c7c2a57e4b559c5620b421f6d133261c60dd2d0cdb38e6b8" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.4", + "sha1", + "thiserror 2.0.18", +] + +[[package]] name = "typed-path" version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3468,9 +3567,9 @@ checksum = "8e28f89b80c87b8fb0cf04ab448d5dd0dd0ade2f8891bae878de66a75a28600e" [[package]] name = "typenum" -version = "1.19.0" +version = "1.20.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "562d481066bde0658276a35467c4af00bdc6ee726305698a55b86e61d7ad82bb" +checksum = "40ce102ab67701b8526c123c1bab5cbe42d7040ccfd0f64af1a385808d2f43de" [[package]] name = "unicase" @@ -3555,6 +3654,12 @@ dependencies = [ ] [[package]] +name = "utf-8" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" + +[[package]] name = "utf8_iter" version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3632,11 +3737,11 @@ checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" [[package]] name = "wasip2" -version = "1.0.2+wasi-0.2.9" +version = "1.0.3+wasi-0.2.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9517f9239f02c069db75e65f174b3da828fe5f5b945c4dd26bd25d89c03ebcf5" +checksum = "20064672db26d7cdc89c7798c48a0fdfac8213434a1186e5ef29fd560ae223d6" dependencies = [ - "wit-bindgen", + "wit-bindgen 0.57.1", ] [[package]] @@ -3645,7 +3750,7 @@ version = "0.4.0+wasi-0.3.0-rc-2026-01-06" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5428f8bf88ea5ddc08faddef2ac4a67e390b88186c703ce6dbd955e1c145aca5" dependencies = [ - "wit-bindgen", + "wit-bindgen 0.51.0", ] [[package]] @@ -4208,6 +4313,12 @@ dependencies = [ ] [[package]] +name = "wit-bindgen" +version = "0.57.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" + +[[package]] name = "wit-bindgen-core" version = "0.51.0" source = "registry+https://github.com/rust-lang/crates.io-index" diff --git a/Cargo.toml b/Cargo.toml @@ -55,4 +55,8 @@ calamine = "*" indicatif = "0.18.0" tracing-subscriber = "*" owo-colors = "*" -shlex = "*" -\ No newline at end of file +shlex = "*" +axum = { version = "*", features = ["ws"]} +reqwest-websocket = "*" +futures-util = "*" +tokio-tungstenite = "*" +\ No newline at end of file diff --git a/src/api.rs b/src/api.rs @@ -120,13 +120,10 @@ async fn add_incoming( }, amount, credit_fee: Amount::zero(&amount.currency), - subject: Some( - format!( - "Manual incoming {}", - fmt_in_subject(subject.ty(), subject.key()) - ) - .into_boxed_str(), - ), + subject: Some(format!( + "Manual incoming {}", + fmt_in_subject(subject.ty(), subject.key()) + )), execution_time: now, debtor: Some(debit_account), }, diff --git a/src/bin/testbench.rs b/src/bin/testbench.rs @@ -64,6 +64,7 @@ pub enum NexusCmd { List { raw_args: Vec<String>, }, + Wss, Exit, } @@ -92,7 +93,7 @@ fn check<R, E: Display>(res: Result<R, E>) -> bool { pub async fn nexus_cmd(cfg: &Config, cmd: &str) -> bool { let parts = shlex::split(cmd).unwrap(); - let args = std::iter::once("dummy_bin").chain(parts.iter().map(|it| it.as_str())); + let args = std::iter::once("libeufin_nexus").chain(parts.iter().map(|it| it.as_str())); match libeufin::Args::try_parse_from(args) { Ok(cmd) => check(run(cfg.clone(), cmd.cmd).await), @@ -284,6 +285,9 @@ async fn main() -> anyhow::Result<()> { } std::fs::remove_file(&ebics.bank_pub_keys_path)?; } + NexusCmd::Wss => { + nexus_cmd(&cfg.cfg, &format!("testing wss {log_flags}")).await; + } NexusCmd::Exit => return Ok(()), }, Err(e) => { diff --git a/src/config.rs b/src/config.rs @@ -56,6 +56,7 @@ impl EbicsKeysCfg { } } +#[derive(Clone)] pub struct EbicsHostCfg { pub base_url: url::Url, pub host_id: String, diff --git a/src/dialect.rs b/src/dialect.rs @@ -19,7 +19,7 @@ use taler_enum_meta::EnumMeta; -use crate::ebics::order::{Order, OrderDoc, Service}; +use crate::ebics::order::{BTF, Order, OrderDoc}; /** Supported EBICS standard */ #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -35,32 +35,32 @@ impl Standard { match self { Standard::SIX => match doc { OrderDoc::acknowledgement => vec![Order::HAC], - OrderDoc::status => vec![Order::BTD(Service { - name: "PSR".into(), + OrderDoc::status => vec![Order::BTD(BTF { + service: "PSR".into(), scope: Some("CH".into()), option: None, container: Some("ZIP".into()), msg: "pain.002".into(), version: Some("10".into()), })], - OrderDoc::report => vec![Order::BTD(Service { - name: "STM".into(), + OrderDoc::report => vec![Order::BTD(BTF { + service: "STM".into(), scope: Some("CH".into()), option: None, container: Some("ZIP".into()), msg: "camt.052".into(), version: Some("08".into()), })], - OrderDoc::statement => vec![Order::BTD(Service { - name: "EOP".into(), + OrderDoc::statement => vec![Order::BTD(BTF { + service: "EOP".into(), scope: Some("CH".into()), option: None, container: Some("ZIP".into()), msg: "camt.053".into(), version: Some("08".into()), })], - OrderDoc::notification => vec![Order::BTD(Service { - name: "REP".into(), + OrderDoc::notification => vec![Order::BTD(BTF { + service: "REP".into(), scope: Some("CH".into()), option: None, container: Some("ZIP".into()), @@ -71,16 +71,16 @@ impl Standard { Standard::GBIC => match doc { OrderDoc::acknowledgement => vec![Order::HAC], OrderDoc::status => vec![ - Order::BTD(Service { - name: "REP".into(), + Order::BTD(BTF { + service: "REP".into(), scope: Some("DE".into()), option: Some("SCI".into()), container: Some("ZIP".into()), msg: "pain.002".into(), version: None, }), - Order::BTD(Service { - name: "REP".into(), + Order::BTD(BTF { + service: "REP".into(), scope: Some("DE".into()), option: Some("SCT".into()), container: Some("ZIP".into()), @@ -88,16 +88,16 @@ impl Standard { version: None, }), ], - OrderDoc::report => vec![Order::BTD(Service { - name: "STM".into(), + OrderDoc::report => vec![Order::BTD(BTF { + service: "STM".into(), scope: Some("DE".into()), option: None, container: Some("ZIP".into()), msg: "camt.052".into(), version: None, })], - OrderDoc::statement => vec![Order::BTD(Service { - name: "EOP".into(), + OrderDoc::statement => vec![Order::BTD(BTF { + service: "EOP".into(), scope: Some("DE".into()), option: None, container: Some("ZIP".into()), @@ -105,16 +105,16 @@ impl Standard { version: None, })], OrderDoc::notification => vec![ - Order::BTD(Service { - name: "STM".into(), + Order::BTD(BTF { + service: "STM".into(), scope: Some("DE".into()), option: None, container: Some("ZIP".into()), msg: "camt.054".into(), version: None, }), - Order::BTD(Service { - name: "STM".into(), + Order::BTD(BTF { + service: "STM".into(), scope: Some("DE".into()), option: Some("SCI".into()), container: Some("ZIP".into()), @@ -128,16 +128,16 @@ impl Standard { pub fn direct_debit(&self) -> Order { match self { - Standard::SIX => Order::BTU(Service { - name: "MCT".into(), + Standard::SIX => Order::BTU(BTF { + service: "MCT".into(), scope: Some("CH".into()), option: None, container: None, msg: "pain.001".into(), version: Some("09".into()), }), - Standard::GBIC => Order::BTU(Service { - name: "SCT".into(), + Standard::GBIC => Order::BTU(BTF { + service: "SCT".into(), scope: None, option: None, container: None, @@ -150,8 +150,8 @@ impl Standard { pub fn instant_direct_debit(&self) -> Option<Order> { match self { Standard::SIX => None, - Standard::GBIC => Some(Order::BTU(Service { - name: "SCI".into(), + Standard::GBIC => Some(Order::BTU(BTF { + service: "SCI".into(), scope: Some("DE".into()), option: None, container: None, diff --git a/src/ebics/administrative.rs b/src/ebics/administrative.rs @@ -31,7 +31,7 @@ use crate::{ ebics::{ EbicsResponse, ebics_code::EbicsReturnCode, - order::{Order, Service}, + order::{BTF, Order}, }, xml, xml::{Xml, XmlAccess as _}, @@ -125,10 +125,10 @@ pub fn parse_hev(xml: &[u8]) -> xml::Result<EbicsResponse<Box<[VersionNumber]>>> }) } -fn service(n: Xml) -> xml::Result<Service> { +fn service(n: Xml) -> xml::Result<BTF> { let msg = n.one("MsgName")?; - Ok(Service { - name: n.one("ServiceName").parse()?, + Ok(BTF { + service: n.one("ServiceName").parse()?, scope: n.opt("Scope").parse()?, option: n.opt("ServiceOption").parse()?, container: n.opt("Container").parse_attr("containerType")?, diff --git a/src/ebics/bts.rs b/src/ebics/bts.rs @@ -27,7 +27,7 @@ use crate::{ crypto::ebics_pub_key_hash, ebics::{ EbicsResponse, PreparedUploadData, - order::{Order, Service}, + order::{BTF, Order}, }, keys::{BankPubKeysFile, ClientPriKeysFile}, utils::b64, @@ -71,9 +71,9 @@ fn bank_digest(w: &mut XmlWriter, bank: &BankPubKeysFile) { ) } -fn service(w: &mut XmlWriter, service: &Service) { - let Service { - name, +fn service(w: &mut XmlWriter, service: &BTF) { + let BTF { + service: name, scope, msg, version, diff --git a/src/ebics/key_management.rs b/src/ebics/key_management.rs @@ -23,16 +23,15 @@ use anyhow::bail; use aws_lc_rs::encoding::{AsDer, Pkcs8V1Der}; use base64::{Engine as _, prelude::BASE64_STANDARD}; use flate2::{Compression, write::ZlibEncoder}; -use reqwest::Client; use tracing::info; use crate::{ config::{EbicsHostCfg, EbicsKeysCfg}, crypto::{rsa_private_from_b64_x509_certificate, x509_certificate_from_rsa_private}, ebics::{ - EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper, EbicsResponse, + EbicsClient, EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper, EbicsResponse, bts::DataEncryptionInfo, decrypt_and_decompress_payload, ebics_code::EbicsReturnCode, - logger::EbicsLogger, order::Order, post_to_bank, + order::Order, }, keys::{self, BankPubKeysFile, ClientPriKeysFile, RsaPub}, xml, @@ -40,249 +39,243 @@ use crate::{ xml_sign::sign_ebics, }; -/** Perform an EBICS public key management [order] using [client] and update on disk state */ -pub async fn submit_client_keys( - keys_cfg: &EbicsKeysCfg, - host_cfg: &EbicsHostCfg, - client: &mut ClientPriKeysFile, - http: &Client, - ebics_logger: &EbicsLogger, - order: Order, -) -> Result<(), EbicsError> { - let ctx = EbicsCtx::new(&order); - if !matches!(order, Order::INI | Order::HIA) { - unreachable!("Only INI & HIA are supported for client keys"); - } - let res = key_management(host_cfg, client, http, ebics_logger, &order).await?; +impl EbicsClient { + /** Perform an EBICS public key management [order] using [client] and update on disk state */ + pub async fn submit_client_keys( + &self, + cfg: &EbicsKeysCfg, + client: &mut ClientPriKeysFile, + order: Order, + ) -> Result<(), EbicsError> { + let ctx = EbicsCtx::new(&order); + if !matches!(order, Order::INI | Order::HIA) { + unreachable!("Only INI & HIA are supported for client keys"); + } + let res = self.key_management(client, &order).await?; - if res.technical_code == EbicsReturnCode::EBICS_INVALID_USER_STATE - || res.technical_code == EbicsReturnCode::EBICS_INVALID_USER_OR_USER_STATE - { - return Err(EbicsErrKind::Custom(Cow::Owned(format!( + if res.technical_code == EbicsReturnCode::EBICS_INVALID_USER_STATE + || res.technical_code == EbicsReturnCode::EBICS_INVALID_USER_OR_USER_STATE + { + return Err(EbicsErrKind::Custom(Cow::Owned(format!( "status code {}: either your IDs are incorrect, or you already have keys registered with this bank", res.technical_code ))).ctx(&ctx)); + } + res.ok_or_fail().ctx(&ctx)?; + match order { + Order::INI => client.submitted_ini = true, + Order::HIA => client.submitted_hia = true, + _ => unreachable!("Only INI & HIA are supported for client keys"), + } + keys::persist_client_keys(client, cfg.client_priv_keys_path.as_ref()).ctx(&ctx)?; + // TODO better error: Could not update the $order state on disk + Ok(()) } - res.ok_or_fail().ctx(&ctx)?; - match order { - Order::INI => client.submitted_ini = true, - Order::HIA => client.submitted_hia = true, - _ => unreachable!("Only INI & HIA are supported for client keys"), - } - keys::persist_client_keys(client, keys_cfg.client_priv_keys_path.as_ref()).ctx(&ctx)?; - // TODO better error: Could not update the $order state on disk - Ok(()) -} - -/** Perform an EBICS private key management HPB using [client] */ -pub async fn hpb( - http: &Client, - cfg: &EbicsHostCfg, - logger: &EbicsLogger, - client: &ClientPriKeysFile, -) -> anyhow::Result<BankPubKeysFile> { - let order = Order::HPB; - let res = key_management(cfg, client, http, logger, &order).await?; - if res.technical_code == EbicsReturnCode::EBICS_AUTHENTICATION_FAILED { - bail!( - "{order} status code {}: could not download bank keys, send client keys (and/or related PDF document with --generate-registration-pdf) to the bank", - res.technical_code - ) - } - let order_data = res.ok_or_fail()?.expect("{order}: missing order data"); - fn rsa_pub_key(xml: Xml) -> xml::Result<RsaPub> { - xml.one("X509Data") - .one("X509Certificate") - .decode(rsa_private_from_b64_x509_certificate) - } + /** Perform an EBICS private key management HPB using [client] */ + pub async fn hpb(&self, client: &ClientPriKeysFile) -> anyhow::Result<BankPubKeysFile> { + let order = Order::HPB; + let res = self.key_management(client, &order).await?; + if res.technical_code == EbicsReturnCode::EBICS_AUTHENTICATION_FAILED { + bail!( + "{order} status code {}: could not download bank keys, send client keys (and/or related PDF document with --generate-registration-pdf) to the bank", + res.technical_code + ) + } + let order_data = res.ok_or_fail()?.expect("{order}: missing order data"); - Ok(Xml::parse(&order_data, "HPBResponseOrderData", |root| { - let auth_pub = root.one("AuthenticationPubKeyInfo")?; - let version = auth_pub.one("AuthenticationVersion")?.text(); - assert_eq!( - version, "X002", - "Expected authentication version X002 got unsupported {version}" - ); - let auth_pub = rsa_pub_key(auth_pub)?; + fn rsa_pub_key(xml: Xml) -> xml::Result<RsaPub> { + xml.one("X509Data") + .one("X509Certificate") + .decode(rsa_private_from_b64_x509_certificate) + } - let enc_pub = root.one("EncryptionPubKeyInfo")?; - let version = enc_pub.one("EncryptionVersion")?.text(); - assert_eq!( - version, "E002", - "Expected encryption version E002 got unsupported {version}" - ); - let enc_pub = rsa_pub_key(enc_pub)?; + Ok(Xml::parse(&order_data, "HPBResponseOrderData", |root| { + let auth_pub = root.one("AuthenticationPubKeyInfo")?; + let version = auth_pub.one("AuthenticationVersion")?.text(); + assert_eq!( + version, "X002", + "Expected authentication version X002 got unsupported {version}" + ); + let auth_pub = rsa_pub_key(auth_pub)?; - Ok(BankPubKeysFile { - auth: auth_pub, - enc: enc_pub, - accepted: false, - }) - })?) -} + let enc_pub = root.one("EncryptionPubKeyInfo")?; + let version = enc_pub.one("EncryptionVersion")?.text(); + assert_eq!( + version, "E002", + "Expected encryption version E002 got unsupported {version}" + ); + let enc_pub = rsa_pub_key(enc_pub)?; -async fn key_management( - cfg: &EbicsHostCfg, - client: &ClientPriKeysFile, - http: &Client, - ebics_logger: &EbicsLogger, - order: &Order, -) -> Result<EbicsResponse<Option<Vec<u8>>>, EbicsError> { - let EbicsHostCfg { - host_id, - user_id, - partner_id, - .. - } = cfg; - let ctx = EbicsCtx::new(order); - info!("Doing key request {order}"); + Ok(BankPubKeysFile { + auth: auth_pub, + enc: enc_pub, + accepted: false, + }) + })?) + } - let (name, security_medium) = match order { - Order::INI | Order::HIA => ("ebicsUnsecuredRequest", "0200"), - Order::HPB => ("ebicsNoPubKeyDigestsRequest", "0000"), - _ => unreachable!(), - }; + async fn key_management( + &self, + client: &ClientPriKeysFile, + order: &Order, + ) -> Result<EbicsResponse<Option<Vec<u8>>>, EbicsError> { + let EbicsHostCfg { + host_id, + user_id, + partner_id, + .. + } = &self.cfg; + let ctx = EbicsCtx::new(order); + info!("Doing key request {order}"); - fn xml_order_data( - cfg: &EbicsHostCfg, - name: &str, - schema: &str, - build: impl FnOnce(&mut XmlWriter), - ) -> String { - let xml = xml!(name ("xmlns":schema) ("xmlns:ds":"http://www.w3.org/2000/09/xmldsig#") { - @ build, - "PartnerID": &cfg.partner_id, - "UserID": &cfg.user_id - }); - // Deflate TODO write inside the compressor directly - let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default()); - encoder.write_all(xml.as_bytes()).unwrap(); - let compressed = encoder.finish().unwrap(); - BASE64_STANDARD.encode(&compressed) - } + let (name, security_medium) = match order { + Order::INI | Order::HIA => ("ebicsUnsecuredRequest", "0200"), + Order::HPB => ("ebicsNoPubKeyDigestsRequest", "0000"), + _ => unreachable!(), + }; - fn rsa_key_xml<K>(w: &mut XmlWriter, key: &K) - where - K: AsDer<Pkcs8V1Der<'static>>, - { - let der = key.as_der().unwrap(); - let b64 = BASE64_STANDARD.encode(der.as_ref()); - let lines = b64 - .as_bytes() - .chunks(64) - .map(|c| std::str::from_utf8(c).unwrap()) - .collect::<Vec<_>>() - .join("\n"); - let pem = - format!("-----BEGIN RSA PRIVATE KEY-----\n{lines}\n-----END RSA PRIVATE KEY-----\n"); - let cert = x509_certificate_from_rsa_private(&pem, "LibEuFin EBICS").unwrap(); - let der = cert.der(); - let b64 = BASE64_STANDARD.encode(der.as_ref()); + fn xml_order_data( + cfg: &EbicsHostCfg, + name: &str, + schema: &str, + build: impl FnOnce(&mut XmlWriter), + ) -> String { + let xml = xml!(name ("xmlns":schema) ("xmlns:ds":"http://www.w3.org/2000/09/xmldsig#") { + @ build, + "PartnerID": &cfg.partner_id, + "UserID": &cfg.user_id + }); + // Deflate TODO write inside the compressor directly + let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default()); + encoder.write_all(xml.as_bytes()).unwrap(); + let compressed = encoder.finish().unwrap(); + BASE64_STANDARD.encode(&compressed) + } - xml!(w, "ds:X509Data" { - "ds:X509Certificate": b64 - }) - } - let data = match order { - Order::INI => Some(xml_order_data( - cfg, - "SignaturePubKeyOrderData", - "http://www.ebics.org/S002", - |w| { - xml!(w, "SignaturePubKeyInfo" { - @ |w| rsa_key_xml(w, &client.sign), - "SignatureVersion": "A006" - }) - }, - )), - Order::HIA => Some(xml_order_data( - cfg, - "HIARequestOrderData", - "urn:org:ebics:H005", - |w| { - xml!(w, - "AuthenticationPubKeyInfo" { - @ |w| rsa_key_xml(w, &client.auth), - "AuthenticationVersion": "X002" - }, - "EncryptionPubKeyInfo" { - @ |w| rsa_key_xml(w, &client.enc), - "EncryptionVersion": "E002" - } - ) - }, - )), - Order::HPB => None, - _ => unreachable!(), - }; - let sign = matches!(order, Order::HPB); - let msg = xml!( - name - ("xmlns": "urn:org:ebics:H005") - ("xmlns:ds": "http://www.w3.org/2000/09/xmldsig#") - ("Version": "H005") - ("Revision": "1") + fn rsa_key_xml<K>(w: &mut XmlWriter, key: &K) + where + K: AsDer<Pkcs8V1Der<'static>>, { - "header" ("authenticate": "true") { - "static" { - "HostID": host_id, - @ |w: &mut XmlWriter| if *order == Order::HPB { - let nonce: u128 = rand::random(); - xml!(w, - "Nonce": format_args!("{:032x}", nonce), - "Timestamp": jiff::Timestamp::now() - ) - }, - "PartnerID": partner_id, - "UserID": user_id, - "OrderDetails" { - "AdminOrderType": order + let der = key.as_der().unwrap(); + let b64 = BASE64_STANDARD.encode(der.as_ref()); + let lines = b64 + .as_bytes() + .chunks(64) + .map(|c| std::str::from_utf8(c).unwrap()) + .collect::<Vec<_>>() + .join("\n"); + let pem = format!( + "-----BEGIN RSA PRIVATE KEY-----\n{lines}\n-----END RSA PRIVATE KEY-----\n" + ); + let cert = x509_certificate_from_rsa_private(&pem, "LibEuFin EBICS").unwrap(); + let der = cert.der(); + let b64 = BASE64_STANDARD.encode(der.as_ref()); + + xml!(w, "ds:X509Data" { + "ds:X509Certificate": b64 + }) + } + let data = match order { + Order::INI => Some(xml_order_data( + &self.cfg, + "SignaturePubKeyOrderData", + "http://www.ebics.org/S002", + |w| { + xml!(w, "SignaturePubKeyInfo" { + @ |w| rsa_key_xml(w, &client.sign), + "SignatureVersion": "A006" + }) + }, + )), + Order::HIA => Some(xml_order_data( + &self.cfg, + "HIARequestOrderData", + "urn:org:ebics:H005", + |w| { + xml!(w, + "AuthenticationPubKeyInfo" { + @ |w| rsa_key_xml(w, &client.auth), + "AuthenticationVersion": "X002" + }, + "EncryptionPubKeyInfo" { + @ |w| rsa_key_xml(w, &client.enc), + "EncryptionVersion": "E002" + } + ) + }, + )), + Order::HPB => None, + _ => unreachable!(), + }; + let sign = matches!(order, Order::HPB); + let msg = xml!( + name + ("xmlns": "urn:org:ebics:H005") + ("xmlns:ds": "http://www.w3.org/2000/09/xmldsig#") + ("Version": "H005") + ("Revision": "1") + { + "header" ("authenticate": "true") { + "static" { + "HostID": host_id, + @ |w: &mut XmlWriter| if *order == Order::HPB { + let nonce: u128 = rand::random(); + xml!(w, + "Nonce": format_args!("{:032x}", nonce), + "Timestamp": jiff::Timestamp::now() + ) + }, + "PartnerID": partner_id, + "UserID": user_id, + "OrderDetails" { + "AdminOrderType": order + }, + "SecurityMedium": security_medium }, - "SecurityMedium": security_medium + "mutable" }, - "mutable" - }, - @ |w: &mut XmlWriter| if sign { - xml!(w, "AuthSignature") - }, - "body" { - @ |w: &mut XmlWriter| if let Some(data) = data { - xml!(w, "DataTransfer" { - "OrderData": data - }) + @ |w: &mut XmlWriter| if sign { + xml!(w, "AuthSignature") + }, + "body" { + @ |w: &mut XmlWriter| if let Some(data) = data { + xml!(w, "DataTransfer" { + "OrderData": data + }) + } } } - } - ); - let signed = if sign { - sign_ebics(msg, &client.auth) - } else { - msg - }; - let res = post_to_bank(http, cfg.base_url.as_str(), signed, &ctx, ebics_logger).await?; - Xml::parse(&res, "ebicsKeyManagementResponse", |root| { - let body = root.one("body")?; - Ok(EbicsResponse { - technical_code: root - .one_signed("header") - .one("mutable") - .one("ReturnCode") - .parse()?, - bank_code: body.one_signed("ReturnCode").parse()?, - content: if let Some(data) = body.opt("DataTransfer")? { - let info = data.one_signed("DataEncryptionInfo")?; - let info = DataEncryptionInfo { - tx_key: info.one("TransactionKey").b64()?, - bank_pub_digest: info.one("EncryptionPubKeyDigest").b64()?, - }; - let chunk = data.one("OrderData").b64()?; - let decoded = decrypt_and_decompress_payload(&client.enc, info, vec![chunk]); - Some(decoded) - } else { - None - }, + ); + let signed = if sign { + sign_ebics(msg, &client.auth) + } else { + msg + }; + let res = self.post_to_bank(signed, &ctx).await?; + Xml::parse(&res, "ebicsKeyManagementResponse", |root| { + let body = root.one("body")?; + Ok(EbicsResponse { + technical_code: root + .one_signed("header") + .one("mutable") + .one("ReturnCode") + .parse()?, + bank_code: body.one_signed("ReturnCode").parse()?, + content: if let Some(data) = body.opt("DataTransfer")? { + let info = data.one_signed("DataEncryptionInfo")?; + let info = DataEncryptionInfo { + tx_key: info.one("TransactionKey").b64()?, + bank_pub_digest: info.one("EncryptionPubKeyDigest").b64()?, + }; + let chunk = data.one("OrderData").b64()?; + let decoded = decrypt_and_decompress_payload(&client.enc, info, vec![chunk]); + Some(decoded) + } else { + None + }, + }) }) - }) - .ctx(&ctx) + .ctx(&ctx) + } } diff --git a/src/ebics/mod.rs b/src/ebics/mod.rs @@ -32,7 +32,8 @@ use sqlx::PgPool; use tracing::{debug, info, trace}; use crate::{ - config::EbicsHostCfg, + EbicsLogs, + config::{EbicsHostCfg, NexusCfg}, crypto::{ decrypt_ebics_e002, decrypt_ebics_e002_key, digest_ebics_order_a006, encrypt_ebics_e002, gen_ebics_e002_key, sign_ebics_a006, @@ -193,32 +194,6 @@ impl EbicsErrKind { } } } - -async fn post_to_bank( - http: &Client, - url: &str, - xml: String, - ctx: &EbicsCtx<'_>, - logger: &EbicsLogger, -) -> Result<Vec<u8>, EbicsError> { - logger.log_request(ctx, &xml)?; - let res = http - .post(url) - .header(CONTENT_TYPE, HeaderValue::from_static("application/xml")) - .body(xml) - .send() - .await - .ctx(ctx)?; - let status = res.status(); - if status != StatusCode::OK { - logger.log_failure(ctx, res).await?; - return Err(EbicsErrKind::HTTP(status).ctx(ctx)); - } - let xml = res.bytes().await.ctx(ctx)?; - logger.log_response(ctx, &xml)?; - Ok(xml.into()) -} - pub struct EbicsResponse<T> { pub technical_code: EbicsReturnCode, pub bank_code: EbicsReturnCode, @@ -238,51 +213,243 @@ impl<T> EbicsResponse<T> { } } -/** POST an EBICS BTS request [xmlReq] using [client] returning a validated and parsed XML response */ -async fn post_bts( - cfg: &EbicsHostCfg, - http: &Client, - xml: String, - ctx: &EbicsCtx<'_>, - logger: &EbicsLogger, -) -> Result<BTSResponse, EbicsError> { - let xml = post_to_bank(http, cfg.base_url.as_str(), xml, ctx, logger).await?; - // TODO verify ebics - let res = parse_bts(&xml).ctx(ctx)?; - // TODO phase in logs ? - trace!(target: "ebics", - "{ctx}{}: {} {}", - std::fmt::from_fn(|f| { - if let Some(tx_id) = &res.content.tx_id { - write!(f, " {tx_id}") - } else { - Ok(()) - } - }), - res.technical_code, - res.bank_code - ); - res.ok_or_fail().ctx(ctx) +pub struct EbicsClient { + cfg: EbicsHostCfg, + pub http: Client, + logger: EbicsLogger, } -pub async fn hev( - http: &Client, - cfg: &EbicsHostCfg, - ebics_log: &EbicsLogger, -) -> Result<Box<[VersionNumber]>, EbicsError> { - let order = Order::HEV; - info!(target: "ebics", "Doing administrative request {order}"); - let msg = hev_msg(cfg); - let ctx = EbicsCtx::new(&order); - let res = post_to_bank(http, cfg.base_url.as_str(), msg, &ctx, ebics_log).await?; - parse_hev(&res).ctx(&ctx)?.ok_or_fail().ctx(&ctx) +impl EbicsClient { + pub fn new(cfg: &NexusCfg, log: EbicsLogs) -> anyhow::Result<Self> { + Ok(Self { + cfg: cfg.host()?.clone(), + http: Client::new(), + logger: EbicsLogger::new(log.dir)?, + }) + } + + async fn post_to_bank(&self, xml: String, ctx: &EbicsCtx<'_>) -> Result<Vec<u8>, EbicsError> { + self.logger.log_request(ctx, &xml)?; + let res = self + .http + .post(self.cfg.base_url.as_str()) + .header(CONTENT_TYPE, HeaderValue::from_static("application/xml")) + .body(xml) + .send() + .await + .ctx(ctx)?; + let status = res.status(); + if status != StatusCode::OK { + self.logger.log_failure(ctx, res).await?; + return Err(EbicsErrKind::HTTP(status).ctx(ctx)); + } + let xml = res.bytes().await.ctx(ctx)?; + self.logger.log_response(ctx, &xml)?; + Ok(xml.into()) + } + + /** POST an EBICS BTS request [xmlReq] using [client] returning a validated and parsed XML response */ + async fn post_bts(&self, xml: String, ctx: &EbicsCtx<'_>) -> Result<BTSResponse, EbicsError> { + let xml = self.post_to_bank(xml, ctx).await?; + // TODO verify ebics + let res = parse_bts(&xml).ctx(ctx)?; + // TODO phase in logs ? + trace!(target: "ebics", + "{ctx}{}: {} {}", + std::fmt::from_fn(|f| { + if let Some(tx_id) = &res.content.tx_id { + write!(f, " {tx_id}") + } else { + Ok(()) + } + }), + res.technical_code, + res.bank_code + ); + res.ok_or_fail().ctx(ctx) + } + + pub async fn hev(&self) -> Result<Box<[VersionNumber]>, EbicsError> { + let order = Order::HEV; + info!(target: "ebics", "Doing administrative request {order}"); + let msg = hev_msg(&self.cfg); + let ctx = EbicsCtx::new(&order); + let res = self.post_to_bank(msg, &ctx).await?; + parse_hev(&res).ctx(&ctx)?.ok_or_fail().ctx(&ctx) + } + + pub async fn haa( + &self, + db: &PgPool, + client: &ClientPriKeysFile, + bank: &BankPubKeysFile, + peek: bool, + ) -> Result<HAA, EbicsError> { + self.download( + db, + client, + bank, + &Order::HAA, + &None, + peek, + async |content| Ok(parse_haa(&content)?), + ) + .await + } + + /** + * Performs an EBICS download transaction of [order] between [startDate] and [endDate]. + * Download content is passed to [processing] + * + * It conducts init -> transfer -> processing -> receipt phases. + * + * Cancellations and failures are handled. + */ + pub async fn download<T>( + &self, + db: &PgPool, + client: &ClientPriKeysFile, + bank: &BankPubKeysFile, + order: &Order, + range: &Option<(Timestamp, Timestamp)>, + peek: bool, + processing: impl AsyncFnOnce(Vec<u8>) -> Result<T, EbicsErrKind>, + ) -> Result<T, EbicsError> { + let mut ctx = EbicsCtx::new(order); + debug!(target: "ebics", "Downloading order {order} {}", std::fmt::from_fn(|f| { + if let Some((start, end)) = range { + write!(f, " from {start} to {end}")? + } + Ok(()) + })); + + // Close interrupted + ctx = ctx.with_phase(Phase::Interrupt); + while let Some(tx_id) = ebics_first(db).await.ctx(&ctx)? { + let xml = download_receipt(&self.cfg, client, order, &tx_id, false); + if let Err(e) = self.post_bts(xml, &ctx).await { + if !matches!( + e.kind, + // Transaction already closed or expired - EBICS protocol error + EbicsErrKind::Code { + technical: EbicsReturnCode::EBICS_TX_UNKNOWN_TXID, + .. + } | + // Transaction already closed or expired - HTTP protocol error for non compliant banks + EbicsErrKind::HTTP(StatusCode::BAD_REQUEST) + ) { + return Err(e); + } else { + debug!(target: "ebics", "{e}") + } + } + ebics_remove(db, &tx_id).await.ctx(&ctx)?; + } + + // Init phase + ctx = ctx.with_phase(Phase::Init); + let xml = download_init(&self.cfg, bank, client, order, range); + let BTSResponse { + tx_id, + nb_segments, + segment, + data_encryption_info, + .. + } = self.post_bts(xml, &ctx).await?; + // TODO DAO add + let (tx_id, nb_segments, segment, encr_info) = ( + tx_id + .ok_or_else(|| EbicsErrKind::Custom("missing transaction ID".into())) + .ctx(&ctx)?, + nb_segments + .ok_or_else(|| EbicsErrKind::Custom("missing num segments".into())) + .ctx(&ctx)?, + segment + .ok_or_else(|| EbicsErrKind::Custom("missing OrderData".into())) + .ctx(&ctx)?, + data_encryption_info + .ok_or_else(|| EbicsErrKind::Custom("missing EncryptionInfo".into())) + .ctx(&ctx)?, + ); + ebics_register(db, &tx_id).await.ctx(&ctx)?; + + // Transfer phase + let mut segments = vec![segment]; + for segment_nb in 2..=nb_segments { + ctx = ctx.with_phase(Phase::Transfer(segment_nb)); + let xml = download_transfer(&self.cfg, client, order, nb_segments, segment_nb, &tx_id); + let BTSResponse { segment, .. } = self.post_bts(xml, &ctx).await?; + segments.push(segment.unwrap()); // TODO error + } + + // Processing phase + ctx = ctx.with_phase(Phase::Process); + let payload = decrypt_and_decompress_payload(&client.enc, encr_info, segments); + self.logger.log_payload(&ctx, &payload, order.file_type())?; + let res = processing(payload).await.ctx(&ctx); + + // Receipt phase + ctx = ctx.with_phase(Phase::Receipt); + let xml = download_receipt(&self.cfg, client, order, &tx_id, res.is_ok() && !peek); + if self.post_bts(xml, &ctx).await.is_ok() { + ebics_remove(db, &tx_id).await.ok(); + } + + res + } + + /** + * Performs an EBICS upload transaction of [order] using [payload]. + * + * It conducts init -> upload phases. + * + * Returns upload orderID + */ + pub async fn upload( + &self, + client: &ClientPriKeysFile, + bank: &BankPubKeysFile, + order: &Order, + payload: &str, + ) -> Result<CompactString, EbicsError> { + debug!(target: "ebics", "Uploading order {order}"); + let mut ctx = EbicsCtx::new(order); + + self.logger.log_payload(&ctx, payload.as_bytes(), "xml")?; + let payload = prepare_upload_payload(&self.cfg, client, bank, payload); + + // Init phase + ctx = ctx.with_phase(Phase::Init); + let xml = upload_init(&self.cfg, bank, client, order, &payload); + let BTSResponse { + tx_id, order_id, .. + } = self.post_bts(xml, &ctx).await?; + + let (tx_id, order_id) = ( + tx_id + .ok_or_else(|| EbicsErrKind::Custom("missing transaction ID".into())) + .ctx(&ctx)?, + order_id + .ok_or_else(|| EbicsErrKind::Custom("missing order ID".into())) + .ctx(&ctx)?, + ); + + // Transfer phase + for segment_nb in 1..=payload.nb_segments() { + ctx = ctx.with_phase(Phase::Transfer(segment_nb)); + let xml = upload_transfer(&self.cfg, client, order, &tx_id, &payload, segment_nb); + self.post_bts(xml, &ctx).await?; + } + + Ok(order_id) + } } pub struct PreparedUploadData { - pub encrypted_key: Vec<u8>, - pub signature_data: String, - pub digest: Digest, - pub payload: String, + encrypted_key: Vec<u8>, + signature_data: String, + digest: Digest, + payload: String, } impl PreparedUploadData { @@ -300,7 +467,7 @@ impl PreparedUploadData { } /** Decrypts and decompresses EBICS BTS payload */ -pub fn decrypt_and_decompress_payload( +fn decrypt_and_decompress_payload( client_encryption_key: &PrivateDecryptingKey, encryption_info: DataEncryptionInfo, segments: Vec<Vec<u8>>, @@ -315,109 +482,6 @@ pub fn decrypt_and_decompress_payload( decoder.finish().unwrap() } -/** - * Performs an EBICS download transaction of [order] between [startDate] and [endDate]. - * Download content is passed to [processing] - * - * It conducts init -> transfer -> processing -> receipt phases. - * - * Cancellations and failures are handled. - */ -pub async fn download<T>( - cfg: &EbicsHostCfg, - http: &Client, - db: &PgPool, - ebics_log: &EbicsLogger, - client: &ClientPriKeysFile, - bank: &BankPubKeysFile, - order: &Order, - range: &Option<(Timestamp, Timestamp)>, - peek: bool, - processing: impl AsyncFnOnce(Vec<u8>) -> Result<T, EbicsErrKind>, -) -> Result<T, EbicsError> { - let mut ctx = EbicsCtx::new(order); - debug!(target: "ebics", "Downloading order {order} {}", std::fmt::from_fn(|f| { - if let Some((start, end)) = range { - write!(f, " from {start} to {end}")? - } - Ok(()) - })); - - // Close interrupted - ctx = ctx.with_phase(Phase::Interrupt); - while let Some(tx_id) = ebics_first(db).await.ctx(&ctx)? { - let xml = download_receipt(cfg, client, order, &tx_id, false); - if let Err(e) = post_bts(cfg, http, xml, &ctx, ebics_log).await { - if !matches!( - e.kind, - // Transaction already closed or expired - EBICS protocol error - EbicsErrKind::Code { - technical: EbicsReturnCode::EBICS_TX_UNKNOWN_TXID, - .. - } | - // Transaction already closed or expired - HTTP protocol error for non compliant banks - EbicsErrKind::HTTP(StatusCode::BAD_REQUEST) - ) { - return Err(e); - } else { - debug!(target: "ebics", "{e}") - } - } - ebics_remove(db, &tx_id).await.ctx(&ctx)?; - } - - // Init phase - ctx = ctx.with_phase(Phase::Init); - let xml = download_init(cfg, bank, client, order, range); - let BTSResponse { - tx_id, - nb_segments, - segment, - data_encryption_info, - .. - } = post_bts(cfg, http, xml, &ctx, ebics_log).await?; - // TODO DAO add - let (tx_id, nb_segments, segment, encr_info) = ( - tx_id - .ok_or_else(|| EbicsErrKind::Custom("missing transaction ID".into())) - .ctx(&ctx)?, - nb_segments - .ok_or_else(|| EbicsErrKind::Custom("missing num segments".into())) - .ctx(&ctx)?, - segment - .ok_or_else(|| EbicsErrKind::Custom("missing OrderData".into())) - .ctx(&ctx)?, - data_encryption_info - .ok_or_else(|| EbicsErrKind::Custom("missing EncryptionInfo".into())) - .ctx(&ctx)?, - ); - ebics_register(db, &tx_id).await.ctx(&ctx)?; - - // Transfer phase - let mut segments = vec![segment]; - for segment_nb in 2..=nb_segments { - ctx = ctx.with_phase(Phase::Transfer(segment_nb)); - let xml = download_transfer(cfg, client, order, nb_segments, segment_nb, &tx_id); - let BTSResponse { segment, .. } = post_bts(cfg, http, xml, &ctx, ebics_log).await?; - segments.push(segment.unwrap()); // TODO error - } - - // Processing phase - ctx = ctx.with_phase(Phase::Process); - let payload = decrypt_and_decompress_payload(&client.enc, encr_info, segments); - ebics_log.log_payload(&ctx, &payload, order.file_type())?; - let res = processing(payload).await.ctx(&ctx); - - // Receipt phase - ctx = ctx.with_phase(Phase::Receipt); - let xml = download_receipt(cfg, client, order, &tx_id, res.is_ok() && !peek); - if post_bts(cfg, http, xml, &ctx, ebics_log).await.is_ok() { - ebics_remove(db, &tx_id).await.ok(); - } - - res -} - /** Signs, encrypts and format EBICS BTS payload */ fn prepare_upload_payload( cfg: &EbicsHostCfg, @@ -461,75 +525,3 @@ fn prepare_upload_payload( payload, } } - -/** - * Performs an EBICS upload transaction of [order] using [payload]. - * - * It conducts init -> upload phases. - * - * Returns upload orderID - */ -pub async fn upload( - cfg: &EbicsHostCfg, - http: &Client, - ebics_log: &EbicsLogger, - client: &ClientPriKeysFile, - bank: &BankPubKeysFile, - order: &Order, - payload: &str, -) -> Result<CompactString, EbicsError> { - debug!(target: "ebics", "Uploading order {order}"); - let mut ctx = EbicsCtx::new(order); - - ebics_log.log_payload(&ctx, payload.as_bytes(), "xml")?; - let payload = prepare_upload_payload(cfg, client, bank, payload); - - // Init phase - ctx = ctx.with_phase(Phase::Init); - let xml = upload_init(cfg, bank, client, order, &payload); - let BTSResponse { - tx_id, order_id, .. - } = post_bts(cfg, http, xml, &ctx, ebics_log).await?; - - let (tx_id, order_id) = ( - tx_id - .ok_or_else(|| EbicsErrKind::Custom("missing transaction ID".into())) - .ctx(&ctx)?, - order_id - .ok_or_else(|| EbicsErrKind::Custom("missing order ID".into())) - .ctx(&ctx)?, - ); - - // Transfer phase - for segment_nb in 1..=payload.nb_segments() { - ctx = ctx.with_phase(Phase::Transfer(segment_nb)); - let xml = upload_transfer(cfg, client, order, &tx_id, &payload, segment_nb); - post_bts(cfg, http, xml, &ctx, ebics_log).await?; - } - - Ok(order_id) -} - -pub async fn haa( - cfg: &EbicsHostCfg, - http: &Client, - db: &PgPool, - ebics_log: &EbicsLogger, - client: &ClientPriKeysFile, - bank: &BankPubKeysFile, - peek: bool, -) -> Result<HAA, EbicsError> { - download( - cfg, - http, - db, - ebics_log, - client, - bank, - &Order::HAA, - &None, - peek, - async |content| Ok(parse_haa(&content)?), - ) - .await -} diff --git a/src/ebics/order.rs b/src/ebics/order.rs @@ -27,8 +27,8 @@ pub enum Direction { } #[derive(Debug, Clone)] -pub struct Service { - pub name: CompactString, +pub struct BTF { + pub service: CompactString, pub scope: Option<CompactString>, pub option: Option<CompactString>, pub container: Option<CompactString>, @@ -36,9 +36,9 @@ pub struct Service { pub version: Option<CompactString>, } -impl PartialEq for Service { +impl PartialEq for BTF { fn eq(&self, other: &Self) -> bool { - self.name == other.name + self.service == other.service && self.scope == other.scope && self.option == other.option && self.container == other.container @@ -47,12 +47,34 @@ impl PartialEq for Service { } } +impl std::fmt::Display for BTF { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let BTF { + service: name, + scope, + option, + container, + msg, + version, + } = self; + write!(f, "{name}")?; + for part in [scope, container, option].into_iter().flatten() { + write!(f, "-{part}")?; + } + write!(f, "-{msg}")?; + if let Some(version) = version { + write!(f, ".{version}")?; + } + Ok(()) + } +} + #[derive(Debug, Clone, PartialEq)] pub enum Order { /// Download of a file identified by a BTF structure (Mandatory) - BTD(Service), + BTD(BTF), /// Upload of a file identified by a BTF structure (Mandatory) - BTU(Service), + BTU(BTF), /// Download retrievable order types (Optional) HAA, /// Download customer acknowledgment (Mandatory) @@ -98,8 +120,8 @@ pub enum Order { } impl Order { - pub const WSS_PARAMS: Self = Self::BTD(Service { - name: CompactString::const_new("OTH"), + pub const WSS_PARAMS: Self = Self::BTD(BTF { + service: CompactString::const_new("OTH"), scope: Some(CompactString::const_new("DE")), msg: CompactString::const_new("wssparam"), version: None, @@ -110,7 +132,7 @@ impl Order { pub fn doc(&self) -> Option<OrderDoc> { match self { Self::HAC => Some(OrderDoc::acknowledgement), - Self::BTD(Service { msg, .. }) => match msg.as_str() { + Self::BTD(BTF { msg, .. }) => match msg.as_str() { "pain.002" => Some(OrderDoc::status), "camt.052" => Some(OrderDoc::report), "camt.053" => Some(OrderDoc::statement), @@ -144,7 +166,7 @@ impl Order { pub fn file_type(&self) -> &str { match self { - Order::BTD(Service { container, .. }) | Order::BTU(Service { container, .. }) => { + Order::BTD(BTF { container, .. }) | Order::BTU(BTF { container, .. }) => { container.as_deref().unwrap_or("xml") } _ => "xml", @@ -179,10 +201,10 @@ impl Order { } } - pub fn from_parts(ty: &str, service: Option<Service>) -> Option<Self> { - match (ty, service) { - ("BTU", Some(service)) => Some(Self::BTU(service)), - ("BTD", Some(service)) => Some(Self::BTD(service)), + pub fn from_parts(ty: &str, btf: Option<BTF>) -> Option<Self> { + match (ty, btf) { + ("BTU", Some(btf)) => Some(Self::BTU(btf)), + ("BTD", Some(btf)) => Some(Self::BTD(btf)), ("HAA", None) => Some(Self::HAA), ("HAC", None) => Some(Self::HAC), ("HCA", None) => Some(Self::HCA), @@ -213,28 +235,9 @@ impl std::fmt::Display for Order { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.write_str(self.ty())?; match self { - Order::BTD(service) | Order::BTU(service) => { - let Service { - name, - scope, - option, - container, - msg, - version, - } = service; - write!(f, "-{name}")?; - for part in [scope, container, option].into_iter().flatten() { - write!(f, "-{part}")?; - } - write!(f, "-{msg}")?; - if let Some(version) = version { - write!(f, ".{version}")?; - } - } - _ => {} + Order::BTD(btf) | Order::BTU(btf) => write!(f, "-{btf}"), + _ => Ok(()), } - - Ok(()) } } diff --git a/src/iso20022/camt.rs b/src/iso20022/camt.rs @@ -203,17 +203,13 @@ fn incoming_id(n: Xml, sref: Option<&str>) -> xml::Result<InId> { } /** Parse transaction wire transfer subject */ -fn wire_transfer_subject(n: Xml) -> xml::Result<Option<Box<str>>> { - Ok(n.opt("RmtInf")?.map(|n| { - n.many("Ustrd") - .map(|n| n.text()) - .collect::<String>() - .into_boxed_str() - })) +fn wire_transfer_subject(n: Xml) -> xml::Result<Option<String>> { + Ok(n.opt("RmtInf")? + .map(|n| n.many("Ustrd").map(|n| n.text()).collect::<String>())) } /** Parse and format transaction return reasons */ -fn return_reason(n: Xml) -> xml::Result<Box<str>> { +fn return_reason(n: Xml) -> xml::Result<String> { let mut buf = String::new(); if let Some(n) = n.opt("RtrInf")? { let code: ReturnReason = n.one("Rsn").one("Cd").parse()?; @@ -231,7 +227,7 @@ fn return_reason(n: Xml) -> xml::Result<Box<str>> { } else if let Some(n) = wire_transfer_subject(n)? { return Ok(n); } - Ok(buf.into_boxed_str()) + Ok(buf) } /** Parse amount */ fn amount(n: Xml) -> xml::Result<Amount> { diff --git a/src/lib.rs b/src/lib.rs @@ -28,7 +28,6 @@ use anyhow::{anyhow, bail}; use compact_str::{CompactString, CompactStringExt}; use jiff::{Timestamp, civil::Date}; use rand::prelude::IndexedRandom; -use reqwest::Client; use sqlx::PgPool; use taler_build::long_version; use taler_common::{ @@ -43,7 +42,7 @@ use taler_common::{ use tracing::{debug, error, info, trace, warn}; use crate::{ - config::NexusCfg, + config::{EbicsKeysCfg, NexusCfg}, crypto::ebics_pub_key_hash, db::{ dbinit, @@ -54,15 +53,10 @@ use crate::{ pool, }, ebics::{ - EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper, + EbicsClient, EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper, administrative::VersionNumber, - download, ebics_code::EbicsReturnCode, - haa, hev, - key_management::{hpb, submit_client_keys}, - logger::EbicsLogger, order::{Order, OrderDoc}, - upload, }, iso20022::{ HacAction, @@ -78,6 +72,7 @@ use crate::{ }, list::ListCmd, model::{InTx, OutTx, PaymentBatch, SubmissionState, Tx}, + testing::TestingCmd, utils::hex_chunk_by_two, worker::register_tx, }; @@ -93,8 +88,10 @@ pub mod iso20022; pub mod keys; pub mod list; pub mod model; +pub mod testing; pub mod utils; pub mod worker; +pub mod ws; pub mod xml; pub mod xml_sign; @@ -198,8 +195,8 @@ pub enum Cmd { List(ListCmd), #[command(subcommand)] Config(ConfigCmd), - - Testing {}, + #[command(subcommand)] + Testing(TestingCmd), } #[derive(clap::Parser, Debug)] @@ -230,20 +227,16 @@ pub fn load_or_generate_client_keys(path: &Path) -> anyhow::Result<ClientPriKeys } pub async fn ebics_setup( - http: &Client, - cfg: &NexusCfg, - ebics_log: &EbicsLogger, + ebics: &EbicsClient, + cfg: &EbicsKeysCfg, force_keys_submissions: bool, auto_accept_keys: bool, ) -> anyhow::Result<()> { - let keys_cfg = cfg.keys()?; - let host_cfg = cfg.host()?; - - let mut client = load_or_generate_client_keys(keys_cfg.client_priv_keys_path.as_ref())?; - let bank = load_bank_keys(keys_cfg.bank_pub_keys_path.as_ref())?; + let mut client = load_or_generate_client_keys(cfg.client_priv_keys_path.as_ref())?; + let bank = load_bank_keys(cfg.bank_pub_keys_path.as_ref())?; // Check EBICS 3 support - let versions = hev(http, cfg.host()?, ebics_log).await?; + let versions = ebics.hev().await?; debug!(target: "setup", "HEV: {}", versions @@ -265,36 +258,40 @@ pub async fn ebics_setup( // Privs exist. Upload their pubs let keys_not_sub = !client.submitted_ini; if !client.submitted_ini || force_keys_submissions { - submit_client_keys(keys_cfg, host_cfg, &mut client, http, ebics_log, Order::INI).await?; + ebics + .submit_client_keys(cfg, &mut client, Order::INI) + .await?; } // Eject PDF if the keys were submitted for the first time, or the user asked. // TODO if (keysNotSub || generateRegistrationPdf) makePdf(clientKeys, hostCfg) if !client.submitted_hia || force_keys_submissions { - submit_client_keys(keys_cfg, host_cfg, &mut client, http, ebics_log, Order::HIA).await?; + ebics + .submit_client_keys(cfg, &mut client, Order::HIA) + .await?; } - let new = hpb(http, host_cfg, ebics_log, &client).await?; + let new = ebics.hpb(&client).await?; if let Some(current) = bank { // Check current bank keys if current.enc != new.enc { bail!( "On disk bank encryption key stored at {} doesn't match server key\nDisk: {}\nServer: {}", - keys_cfg.bank_pub_keys_path, + cfg.bank_pub_keys_path, hex_chunk_by_two(ebics_pub_key_hash(&current.enc.key)), hex_chunk_by_two(ebics_pub_key_hash(&new.enc.key)) ) } else if current.auth != new.auth { bail!( "On disk bank authentication key stored at {} doesn't match server key\nDisk: {}\nServer: {}", - keys_cfg.bank_pub_keys_path, + cfg.bank_pub_keys_path, hex_chunk_by_two(ebics_pub_key_hash(&current.auth.key)), hex_chunk_by_two(ebics_pub_key_hash(&new.auth.key)) ) } } else { // Accept bank keys - info!("Bank keys stored at {}", keys_cfg.bank_pub_keys_path); - persist_bank_keys(&new, keys_cfg.bank_pub_keys_path.as_ref())?; + info!("Bank keys stored at {}", cfg.bank_pub_keys_path); + persist_bank_keys(&new, cfg.bank_pub_keys_path.as_ref())?; }; let mut bank = new; if !bank.accepted { @@ -303,7 +300,7 @@ pub async fn ebics_setup( panic!("Cannot successfully finish the setup without accepting the bank keys"); } bank.accepted = true; - persist_bank_keys(&bank, keys_cfg.bank_pub_keys_path.as_ref())?; + persist_bank_keys(&bank, cfg.bank_pub_keys_path.as_ref())?; } // Check account information @@ -315,16 +312,14 @@ pub async fn ebics_setup( } pub async fn ebics_submit( + ebics: &EbicsClient, cfg: &NexusCfg, - http: &Client, client: &ClientPriKeysFile, bank: &BankPubKeysFile, - ebics_log: &EbicsLogger, db: &PgPool, transient: bool, ) -> anyhow::Result<()> { let ebics_cfg = cfg.ebics()?; - let host_cfg = cfg.host()?; let submit_cfg = cfg.submit()?; let submit_batch = async |order: &Order, @@ -352,7 +347,7 @@ pub async fn ebics_submit( .collect(), }; let xml = create_pain001(&msg, &ebics_cfg.dialect, instant).ctx(&ctx)?; - upload(host_cfg, http, ebics_log, client, bank, order, &xml).await + ebics.upload(client, bank, order, &xml).await }; let submit_all = async || -> anyhow::Result<()> { @@ -458,16 +453,14 @@ async fn register_camt(db: &PgPool, cfg: &NexusCfg, xml: &[u8]) -> anyhow::Resul } pub async fn ebics_fetch( + ebics: &EbicsClient, cfg: &NexusCfg, - http: &Client, client: &ClientPriKeysFile, bank: &BankPubKeysFile, - ebics_log: &EbicsLogger, db: &PgPool, documents: Option<&[OrderDoc]>, ) -> anyhow::Result<()> { let ebics_cfg = cfg.ebics()?; - let host_cfg = cfg.host()?; let register_file = async |doc: &OrderDoc, xml: Vec<u8>| -> anyhow::Result<()> { match doc { @@ -597,23 +590,13 @@ pub async fn ebics_fetch( for (doc, orders) in grouped_orders { if let Some(doc) = doc { for order in orders { - if let Err(e) = download( - host_cfg, - http, - db, - ebics_log, - client, - bank, - order, - &None, - false, - async |content| { + if let Err(e) = ebics + .download(db, client, bank, order, &None, false, async |content| { register_payload(&doc, content) .await .map_err(|e| EbicsErrKind::Custom(e.to_string().into())) - }, - ) - .await + }) + .await { if let EbicsErrKind::Code { bank, .. } = e.kind { match bank { @@ -651,7 +634,7 @@ pub async fn ebics_fetch( info!(target: "ebics-fetch", "Running at frequency"); - let mut haa = haa(host_cfg, http, db, ebics_log, client, bank, false).await?; + let mut haa = ebics.haa(db, client, bank, false).await?; debug!( "HAA: {}", std::fmt::from_fn(|f| f.write_str( @@ -682,10 +665,10 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { generate_registration_pdf, } => { let cfg = NexusCfg::parse(cfg)?; + let ebics = EbicsClient::new(&cfg, ebics_logs)?; ebics_setup( - &Client::new(), - &cfg, - &EbicsLogger::new(ebics_logs.dir)?, + &ebics, + cfg.keys()?, force_keys_resubmission, auto_accept_keys, ) @@ -698,37 +681,19 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { ebics, } => { let pool = pool(&cfg).await?; - let http = Client::new(); let cfg = NexusCfg::parse(cfg)?; let key_cfg = cfg.keys()?; + let ebics = EbicsClient::new(&cfg, ebics.logs)?; let (client, bank) = expect_full_keys(key_cfg)?; - ebics_fetch( - &cfg, - &http, - &client, - &bank, - &EbicsLogger::new(ebics.logs.dir)?, - &pool, - None, - ) - .await? + ebics_fetch(&ebics, &cfg, &client, &bank, &pool, None).await? } Cmd::EbicsSubmit { ebics } => { let pool = pool(&cfg).await?; - let http = Client::new(); let cfg = NexusCfg::parse(cfg)?; + let ebics = EbicsClient::new(&cfg, ebics.logs)?; let key_cfg = cfg.keys()?; let (client, bank) = expect_full_keys(key_cfg)?; - ebics_submit( - &cfg, - &http, - &client, - &bank, - &EbicsLogger::new(ebics.logs.dir)?, - &pool, - true, - ) - .await? + ebics_submit(&ebics, &cfg, &client, &bank, &pool, true).await? } Cmd::InitiatePayment { amount, @@ -777,8 +742,11 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> { let cfg = NexusCfg::parse(cfg)?; cmd.run(&pool, &cfg.currency).await?; } - Cmd::Config(config_cmd) => todo!(), - Cmd::Testing {} => todo!(), + Cmd::Config(cmd) => cmd.run(&cfg)?, + Cmd::Testing(cmd) => { + let pool = pool(&cfg).await?; + cmd.run(cfg, &pool).await?; + } } Ok(()) } @@ -818,7 +786,7 @@ pub mod test { LazyLock::new(|| payto("payto://iban/CH4189144589712575493?receiver-name=Test")); /** Generates an outgoing payment, given its subject */ - pub fn gen_out_pay(subject: impl Into<Box<str>>) -> OutTx { + pub fn gen_out_pay(subject: impl Into<String>) -> OutTx { OutTx { id: OutId { msg_id: None, @@ -855,7 +823,7 @@ pub mod test { } /** Generates an incoming payment, given its subject */ - pub fn gen_in_pay(subject: impl Into<Box<str>>) -> InTx { + pub fn gen_in_pay(subject: impl Into<String>) -> InTx { InTx { id: InId::new(None, Some(rand_ebics_id()), None), amount: Amount::new(&CURR, 44, 0), diff --git a/src/list.rs b/src/list.rs @@ -61,10 +61,10 @@ pub enum ListCmd { } impl ListCmd { - pub async fn run(&self, db: &PgPool, currency: &Currency) -> anyhow::Result<()> { + pub async fn run(self, db: &PgPool, currency: &Currency) -> anyhow::Result<()> { match self { ListCmd::Incoming { incomplete } => { - let txs = list::incoming(db, *incomplete, currency).await?; + let txs = list::incoming(db, incomplete, currency).await?; let out = &mut std::io::stdout().lock(); for InMetadata { id, @@ -130,7 +130,7 @@ impl ListCmd { } } ListCmd::Initiated { ack } => { - if *ack { + if ack { let txs = list::initiated_ack(db, currency).await?; let out = &mut std::io::stdout().lock(); for InitMetadataAck { diff --git a/src/model.rs b/src/model.rs @@ -253,7 +253,7 @@ pub struct InTx { pub id: InId, pub amount: Amount, pub credit_fee: Amount, - pub subject: Option<Box<str>>, + pub subject: Option<String>, pub execution_time: Timestamp, pub debtor: Option<PaytoURI>, } @@ -304,7 +304,7 @@ pub struct OutTx { pub id: OutId, pub amount: Amount, pub debit_fee: Amount, - pub subject: Option<Box<str>>, + pub subject: Option<String>, pub execution_time: Timestamp, pub creditor: Option<PaytoURI>, } @@ -401,7 +401,7 @@ pub struct OutReversal { pub e2e_id: CompactString, /** ISO20022 MessageId */ pub msg_id: Option<CompactString>, - pub reason: Box<str>, + pub reason: String, pub execution_time: Timestamp, } diff --git a/src/testing.rs b/src/testing.rs @@ -0,0 +1,236 @@ +/* +* 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 anyhow::{anyhow, bail}; +use compact_str::CompactString; +use jiff::{Timestamp, civil::Date}; +use sqlx::PgPool; +use taler_common::{ + config::Config, + types::{ + amount::Amount, + iban::{Country, IBAN}, + payto::TransferIbanPayto, + }, +}; +use tracing::debug; + +use crate::{ + EbicsClient, EbicsLogs, InTx, + config::NexusCfg, + ebics::{ + EbicsErrKind, + order::{BTF, Order}, + }, + keys::expect_full_keys, + list::ListCmd, + model::InId, + rand_ebics_id, + worker::register_incoming, + ws::listen_for_notification, +}; + +#[derive(clap::Subcommand, Debug)] +pub enum IbanCmd { + /// Generate fake IBANs for testing + Gen { country: Country }, +} + +impl IbanCmd { + pub fn run(self) -> anyhow::Result<()> { + match self { + IbanCmd::Gen { country } => { + println!("{}", IBAN::random(country)) + } + } + Ok(()) + } +} + +/// Testing helper commands +#[derive(clap::Subcommand, Debug)] +pub enum TestingCmd { + /// List incoming transactions + #[clap(subcommand)] + Iban(IbanCmd), + /// Genere a fake incoming payment + FakeIncoming { + /// The amount to transfer, payto 'amount' parameter takes the precedence + #[clap(long)] + amount: Option<Amount>, + + /// The payment credit fee + #[clap(long)] + credit_fee: Option<Amount>, + + /// The payment subject, payto 'message' parameter takes the precedence + #[clap(long)] + subject: Option<CompactString>, + + /// The debited account IBAN payto URI + payto: TransferIbanPayto, + }, + #[clap(subcommand)] + List(ListCmd), + /// Perform EBICS requests + EbicsBtd { + #[clap(long = "type", default_value_t = CompactString::const_new("BTD"))] + ty: CompactString, + #[clap(long)] + name: CompactString, + #[clap(long)] + scope: Option<CompactString>, + #[clap(long)] + message_name: CompactString, + #[clap(long)] + message_version: Option<CompactString>, + #[clap(long)] + container: Option<CompactString>, + #[clap(long)] + option: Option<CompactString>, + #[clap(flatten)] + logs: EbicsLogs, + /// Erliest timestamp of the downloaded documents + #[clap(long, value_name = "YYYY-MM-DD")] + pinned_start: Option<Date>, + /// Do not consume fetched documents + #[clap(long)] + peek: bool, + #[clap(long)] + dry_run: bool, + }, + /// Check transaction semantic + TxCheck, + /// "Listen to EBICS instant notification over websocket + Wss { + #[clap(flatten)] + logs: EbicsLogs, + }, +} + +impl TestingCmd { + pub async fn run(self, cfg: Config, db: &PgPool) -> anyhow::Result<()> { + match self { + TestingCmd::Iban(cmd) => cmd.run()?, + TestingCmd::FakeIncoming { + amount, + credit_fee, + subject, + payto, + } => { + let cfg = NexusCfg::parse(cfg)?; + let subject = payto + .subject + .as_ref() + .or(subject.as_ref()) + .ok_or(anyhow!("Mising subject"))?; + let amount = payto + .amount + .as_ref() + .or(amount.as_ref()) + .ok_or(anyhow!("Mising amount"))?; + + if cfg.currency != amount.currency { + bail!( + "Wrong currency: expected {} got {}", + cfg.currency, + amount.currency + ); + } + register_incoming( + db, + &cfg.ingest()?, + &InTx { + id: InId::new(None, Some(rand_ebics_id()), None), + amount: *amount, + credit_fee: credit_fee.unwrap_or(Amount::zero(&cfg.currency)), + subject: Some(subject.clone().into_string()), + execution_time: Timestamp::now(), + debtor: Some(payto.as_payto()), + }, + ) + .await?; + } + TestingCmd::List(list_cmd) => { + let cfg = NexusCfg::parse(cfg)?; + list_cmd.run(db, &cfg.currency).await?; + } + TestingCmd::EbicsBtd { + ty, + name, + scope, + message_name, + message_version, + container, + option, + logs, + pinned_start, + peek, + dry_run, + } => { + let cfg = NexusCfg::parse(cfg)?; + let order = Order::from_parts( + &ty, + Some(BTF { + service: name, + scope, + option, + container, + msg: message_name, + version: message_version, + }), + ) + .ok_or(anyhow!("Unknown ebics order"))?; + let (client, bank) = expect_full_keys(cfg.keys()?)?; + let ebics = EbicsClient::new(&cfg, logs)?; + ebics + .download( + db, + &client, + &bank, + &order, + &None, // TODO + peek, + async |_| { + if dry_run { + Err(EbicsErrKind::Custom("dry run".into())) + } else { + Ok(()) + } + }, + ) + .await?; + } + TestingCmd::TxCheck => todo!(), + TestingCmd::Wss { logs } => { + let cfg = NexusCfg::parse(cfg)?; + let ebics = EbicsClient::new(&cfg, logs)?; + let (client, bank) = expect_full_keys(cfg.keys()?)?; + let (sender, mut receiver) = tokio::sync::mpsc::channel(10); + tokio::spawn(async move { + while let Some(orders) = receiver.recv().await { + debug!(target: "testing", "{orders:?}") + } + }); + listen_for_notification(&ebics, db, &client, &bank, sender).await + } + } + Ok(()) + } +} diff --git a/src/ws.rs b/src/ws.rs @@ -0,0 +1,439 @@ +/* +* 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::time::Duration; + +use compact_str::CompactString; +use futures_util::TryStreamExt as _; +use reqwest::{Client, StatusCode}; +use reqwest_websocket::{Message, Upgrade}; +use serde::{Deserialize, Serialize}; +use sqlx::PgPool; +use taler_common::ExpoBackoffDecorr; +use thiserror::Error; +use tracing::{debug, error, info, trace}; + +use crate::{ + ebics::{ + EbicsClient, EbicsErrKind, + ebics_code::EbicsReturnCode, + order::{BTF, Order}, + }, + keys::{BankPubKeysFile, ClientPriKeysFile}, +}; + +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] +#[serde(rename_all = "UPPERCASE")] +pub struct WssParams { + pub url: String, + pub token: String, + pub ott: String, + pub validity: String, + pub partnerid: String, + pub userid: Option<String>, +} + +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] +#[serde(rename_all = "UPPERCASE")] +pub struct WssNotificationClass { + pub name: String, + pub vers: String, + pub timestamp: String, +} + +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] +#[serde(rename_all = "UPPERCASE")] +pub struct WssNotificationBTF { + pub service: CompactString, + pub scope: Option<CompactString>, + pub option: Option<CompactString>, + pub conttype: Option<CompactString>, + pub msgname: CompactString, + pub variant: Option<CompactString>, + pub version: Option<CompactString>, + pub format: Option<CompactString>, +} +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] +#[serde(rename_all = "UPPERCASE")] +pub struct WssInfo { + pub lang: String, + pub free: String, +} + +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] +#[serde(untagged)] +pub enum WssNotification { + // INFO + #[serde(rename_all = "UPPERCASE")] + GeneralInfo { + mclass: Vec<WssNotificationClass>, + info: Vec<WssInfo>, + }, + #[serde(rename_all = "UPPERCASE")] + NewData { + mclass: Vec<WssNotificationClass>, + partnerid: String, + userid: Option<String>, + btf: Vec<WssNotificationBTF>, + ordertype: Vec<String>, + }, +} + +impl WssParams { + async fn connect( + &self, + client: &Client, + mut lambda: impl AsyncFnMut(WssNotification), + ) -> Result<(), WssError> { + let Self { + url, + token, + partnerid, + userid, + .. + } = self; + let username = format!( + "{partnerid}{}", + std::fmt::from_fn(|f| if let Some(userid) = userid { + write!(f, "_{userid}") + } else { + Ok(()) + }) + ); + + let mut ws = client + .get( + url.replace("https://", "wss://") + .replace("http://", "ws://"), + ) + .basic_auth(username, Some(&token)) + .upgrade() + .send() + .await? + .into_websocket() + .await?; + trace!(target: "wss", "wait for ws msg"); + while let Some(msg) = ws.try_next().await? { + match msg { + Message::Text(str) => { + // TODO handle error + let msg: WssNotification = serde_json::from_str(&str)?; + trace!(target: "wss", "received: {msg:?}"); + lambda(msg).await; + } + Message::Binary(bytes) => { + // TODO what should we do ? + } + Message::Ping(_) | Message::Pong(_) => { + // Handled by tungstenite + } + Message::Close { code, reason } => { + debug!(target: "wss", "closed {code} {reason}"); + break; + } + } + trace!(target: "wss", "wait for ws msg"); + } + Ok(()) + } +} + +#[derive(Error, Debug)] +pub enum WssError { + #[error("ws: {0}")] + Ws(#[from] reqwest_websocket::Error), + #[error("ws JSON msg: {0}")] + ReqJson(#[from] serde_json::Error), +} + +pub async fn listen_for_notification( + ebics: &EbicsClient, + db: &PgPool, + client: &ClientPriKeysFile, + bank: &BankPubKeysFile, + sender: tokio::sync::mpsc::Sender<Vec<Order>>, +) { + let mut backoff = ExpoBackoffDecorr::new(Duration::from_secs(30), Duration::from_mins(30), 2.5); + loop { + let res: Result<(), anyhow::Error> = async { + let res = ebics + .download( + db, + client, + bank, + &Order::WSS_PARAMS, + &None, + false, + async |content| { + serde_json::from_slice::<WssParams>(&content) + .map_err(|e| EbicsErrKind::Custom(e.to_string().into())) + }, + ) + .await; + let params = match res { + Ok(params) => params, + Err(e) => { + if matches!( + e.kind, + // Expected EBICS error + EbicsErrKind::Code { + technical: EbicsReturnCode::EBICS_INVALID_ORDER_TYPE, + .. + } | + // Netzbon HTTP error + EbicsErrKind::HTTP(StatusCode::BAD_REQUEST) + ) { + // Failure is expected if this wss is not supported + info!(target: "ws", "Real-time EBICS notifications is not supported"); + return Ok(()); + } else { + return Err(e.into()); + } + } + }; + info!(target: "ws", "Listening to real-time EBICS notifications"); + trace!(target: "ws", "{params:?}"); + + params + .connect(&ebics.http, async |msg| { + backoff.reset(); + match msg { + WssNotification::GeneralInfo { info, .. } => { + for info in info { + info!(target: "ws", "info: {}", info.free); + } + } + WssNotification::NewData { btf, .. } => { + let orders = btf + .into_iter() + .map(|it| { + Order::BTD(BTF { + service: it.service, + scope: it.scope, + option: it.option, + container: it.conttype, + msg: it.msgname, + version: it.version, + }) + }) + .collect(); + sender.send(orders).await.ok(); + } + } + }) + .await?; + Ok(()) + } + .await; + if let Err(e) = res { + error!(target: "ws", "{e}"); + tokio::time::sleep(backoff.backoff()).await; + } else { + return; + } + } +} + +#[cfg(test)] +mod test { + use std::{fmt::Debug, fs::Permissions, os::unix::fs::PermissionsExt as _, time::Duration}; + + use axum::{ + extract::{ + WebSocketUpgrade, + ws::{CloseFrame, Message, Utf8Bytes}, + }, + http::HeaderMap, + routing::get, + }; + use reqwest::header::AUTHORIZATION; + use serde::{Serialize, de::DeserializeOwned}; + use taler_api::api::TalerRouter as _; + + use crate::ws::{WssNotification, WssParams}; + + // WSS params example from the spec + const PARAMS_EXAMPLE: &str = r#" + { + "URL": "http://bankmitwebsocket.de", + "TOKEN": "550e8400-e29b-11d4-a716-446655440000", + "OTT": "N", + "VALIDITY": "2019-03-21T10:35:22Z", + "PARTNERID": "K1234567", + "USERID": "USER4711" + } + "#; + // Authorization header example from the spec + const AUTH_EXAMPLE: &str = + "Basic SzEyMzQ1NjdfVVNFUjQ3MTE6NTUwZTg0MDAtZTI5Yi0xMWQ0LWE3MTYtNDQ2NjU1NDQwMDAw"; + // Notifications examples from the spec + const NOTIFICATION_EXAMPLES: [&str; 3] = [ + r#" + { + "MCLASS": [ + { + "NAME": "EBICS-HAA", + "VERS": "1.0", + "TIMESTAMP": "2019-05-13T12:21:50Z" + } + ], + "PARTNERID": "K1234567", + "USERID": "USER471", + "BTF": [ + { + "SERVICE": "REP", + "SCOPE": "DE", + "CONTTYPE": "ZIP", + "MSGNAME": "camt.054" + } + ], + "ORDERTYPE": [ + "C5N" + ] + } + "#, + r#" + { + "MCLASS": [ + { + "NAME": "EBICS-HAA", + "VERS": "1.0", + "TIMESTAMP": "2019-05-13T12:21:53Z" + } + ], + "PARTNERID": "K1234567", + "USERID": "USER471", + "BTF": [ + { + "SERVICE": "REP", + "SCOPE": "DE", + "CONTTYPE": "ZIP", + "MSGNAME": "camt.052" + }, + { + "SERVICE": "REP", + "SCOPE": "DE", + "OPTION": "SCI", + "CONTTYPE": "ZIP", + "MSGNAME": "pain.002" + } + ], + "ORDERTYPE": [ + "C52", + "CIZ" + ] + } + "#, + r#" + { + "MCLASS": [ + { + "NAME": "INFO", + "VERS": "1.0", + "TIMESTAMP": "2019-03-25T12:25:34Z" + } + ], + "INFO": [ + { + "LANG": "EN", + "FREE": " The EBICS-Service is limited on 30.03.2019 from 10:00 a.m. - 11:00a.m. due to maintenance work " + } + ] + } + "#, + ]; + + #[test] + pub fn serialization() { + fn roundrip<T: Serialize + DeserializeOwned + Eq + Debug>(src: &str) { + let it: T = serde_json::from_str(src).unwrap(); + let roundrip: T = serde_json::from_str(&serde_json::to_string(&it).unwrap()).unwrap(); + assert_eq!(it, roundrip); + } + roundrip::<WssParams>(PARAMS_EXAMPLE); + for ex in NOTIFICATION_EXAMPLES { + roundrip::<WssNotification>(ex); + } + } + + #[tokio::test] + pub async fn params() { + let path = "/tmp/libeufin_nexus_wss_test.sock"; + std::fs::remove_file(&path).ok(); + let server = axum::Router::new() + .route( + "/", + get(async |headers: HeaderMap, ws: WebSocketUpgrade| { + assert_eq!( + headers.get(AUTHORIZATION).map(|it| it.as_bytes()), + Some(AUTH_EXAMPLE.as_bytes()) + ); + ws.on_upgrade(async |mut it| { + for ex in NOTIFICATION_EXAMPLES { + it.send(Message::Text(Utf8Bytes::from_static(ex))) + .await + .unwrap(); + } + it.send(Message::Close(Some(CloseFrame { + code: 1000, + reason: Utf8Bytes::from_static("Test done"), + }))) + .await + .unwrap(); + }) + }), + ) + .serve( + taler_api::Serve::Unix { + path: path.into(), + permission: Permissions::from_mode(660), + }, + None, + ); + tokio::spawn(server); + for _ in 0..100 { + if std::fs::exists(path).unwrap() { + break; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + + let client = reqwest::ClientBuilder::new() + .unix_socket(path) + .build() + .unwrap(); + let params: WssParams = serde_json::from_str(PARAMS_EXAMPLE).unwrap(); + let mut count = 0; + params + .connect(&client, async |msg| { + count += 1; + // Check message number and type + assert!(count <= 3); + if count == 3 { + assert!(matches!(msg, WssNotification::GeneralInfo { .. })) + } else { + assert!(matches!(msg, WssNotification::NewData { .. })) + } + }) + .await + .unwrap(); + // Check receive all messages + assert_eq!(3, count); + } +}