commit b7fdacad8d43a68daabd6c05082939b57548e83d
parent 4c2ffb9cc4011bdabd62bb918268b8d6b88dcc05
Author: Antoine A <>
Date: Fri, 24 Apr 2026 10:38:01 +0200
nexus: add mock ebics bank and finish ebics-fetch
Diffstat:
18 files changed, 1591 insertions(+), 540 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -203,6 +203,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90"
dependencies = [
"axum-core",
+ "axum-macros",
"base64",
"bytes",
"form_urlencoded",
@@ -252,6 +253,17 @@ dependencies = [
]
[[package]]
+name = "axum-macros"
+version = "0.5.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "7aa268c23bfbbd2c4363b9cd302a4f504fb2a9dfe7e3451d66f35dd392e20aca"
+dependencies = [
+ "proc-macro2",
+ "quote",
+ "syn",
+]
+
+[[package]]
name = "base64"
version = "0.22.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1669,6 +1681,7 @@ dependencies = [
"taler-common",
"taler-enum-meta",
"taler-test-utils",
+ "tempfile",
"thiserror 2.0.18",
"tokio",
"tokio-tungstenite",
@@ -2471,9 +2484,9 @@ dependencies = [
[[package]]
name = "rustls"
-version = "0.23.38"
+version = "0.23.39"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "69f9466fb2c14ea04357e91413efb882e2a6d4a406e625449bc0a5d360d53a21"
+checksum = "7c2c118cb077cca2822033836dfb1b975355dfb784b5e8da48f7b6c5db74e60e"
dependencies = [
"aws-lc-rs",
"once_cell",
diff --git a/Cargo.toml b/Cargo.toml
@@ -12,7 +12,7 @@ roxmltree = "*"
base64 = "*"
pem = "*"
anyhow = "*"
-jiff = "*"
+jiff ={ version = "*", features = ["serde"]}
rand = "*"
getrandom = "*"
serde_json = "*"
@@ -43,6 +43,7 @@ sqlx = { version = "0.8", default-features = false, features = [
"runtime-tokio",
"tls-rustls-aws-lc-rs",
"uuid",
+ "json"
] }
compact_str = { version = "0.9.0", features = ["serde", "sqlx-postgres"] }
uuid = { version = "1.0", features = ["v4", "fast-rng"] }
@@ -56,7 +57,8 @@ indicatif = "0.18.0"
tracing-subscriber = "*"
owo-colors = "*"
shlex = "*"
-axum = { version = "*", features = ["ws"]}
+axum = { version = "*", features = ["ws", "macros"]}
reqwest-websocket = "*"
futures-util = "*"
-tokio-tungstenite = "*"
-\ No newline at end of file
+tokio-tungstenite = "*"
+tempfile = "*"
+\ No newline at end of file
diff --git a/src/config.rs b/src/config.rs
@@ -59,6 +59,7 @@ impl EbicsKeysCfg {
#[derive(Clone)]
pub struct EbicsHostCfg {
pub base_url: url::Url,
+ pub unix_path: Option<String>,
pub host_id: String,
pub user_id: String,
pub partner_id: String,
@@ -69,6 +70,7 @@ impl EbicsHostCfg {
let s = cfg.section("nexus-ebics");
Ok(Self {
base_url: s.url("host_base_url").require()?,
+ unix_path: s.path("UNIXPATH").opt()?,
host_id: s.str("host_id").require()?,
user_id: s.str("user_id").require()?,
partner_id: s.str("partner_id").require()?,
diff --git a/src/crypto.rs b/src/crypto.rs
@@ -131,15 +131,16 @@ pub fn gen_ebics_e002_key(pub_key: PublicEncryptingKey) -> ([u8; 16], Vec<u8>) {
(transaction_key, encrypted_key)
}
-pub fn encrypt_ebics_e002(transaction_key: &[u8; 16], data: &[u8]) -> Vec<u8> {
+pub fn encrypt_ebics_e002(transaction_key: &[u8; 16], mut data: Vec<u8>) -> Vec<u8> {
let block_size = 16;
let padding_len = block_size - (data.len() % block_size);
- let mut padded_data = data.to_vec();
+
+ // Add padding
for i in 0..padding_len {
if i == padding_len - 1 {
- padded_data.push(padding_len as u8);
+ data.push(padding_len as u8);
} else {
- padded_data.push(0);
+ data.push(0);
}
}
@@ -147,10 +148,10 @@ pub fn encrypt_ebics_e002(transaction_key: &[u8; 16], data: &[u8]) -> Vec<u8> {
let enc_key =
EncryptingKey::cbc(UnboundCipherKey::new(&AES_128, transaction_key).unwrap()).unwrap();
enc_key
- .less_safe_encrypt(&mut padded_data, EncryptionContext::Iv128(iv))
+ .less_safe_encrypt(&mut data, EncryptionContext::Iv128(iv))
.unwrap();
- padded_data
+ data
}
pub fn decrypt_ebics_e002(transaction_key: &DecryptingKey, mut encrypted_data: Vec<u8>) -> Vec<u8> {
@@ -224,7 +225,7 @@ mod test {
let key = PrivateDecryptingKey::generate(KeySize::Rsa2048).unwrap();
let (tx_key, encrypted_key) = gen_ebics_e002_key(key.public_key());
- let enc = encrypt_ebics_e002(&tx_key, data);
+ let enc = encrypt_ebics_e002(&tx_key, data.to_vec());
let key = decrypt_ebics_e002_key(key, &encrypted_key);
let dec = decrypt_ebics_e002(&key, enc);
assert_eq!(&data, &dec.as_slice());
diff --git a/src/db.rs b/src/db.rs
@@ -15,11 +15,13 @@
*/
use compact_str::CompactString;
-use sqlx::{PgPool, Row, postgres::PgRow};
+use jiff::Timestamp;
+use sqlx::{PgPool, Row, postgres::PgRow, types::Json};
+use taler_api::db::BindHelper;
use taler_common::config::Config;
use tokio::sync::watch::Sender;
-use crate::config::parse_db_cfg;
+use crate::{TaskStatus, config::parse_db_cfg};
pub mod exchange;
pub mod initiated;
@@ -41,7 +43,7 @@ pub async fn dbinit(cfg: &Config, reset: bool) -> anyhow::Result<PgPool> {
let db_cfg = parse_db_cfg(cfg)?;
let pool = taler_common::db::pool(db_cfg.cfg, SCHEMA).await?;
let mut db = pool.acquire().await?;
- taler_common::db::dbinit(&mut db, db_cfg.sql_dir.as_ref(), "magnet-bank", reset).await?;
+ taler_common::db::dbinit(&mut db, db_cfg.sql_dir.as_ref(), "libeufin-nexus", reset).await?;
Ok(pool)
}
@@ -92,6 +94,30 @@ pub async fn ebics_first(db: &PgPool) -> sqlx::Result<Option<CompactString>> {
.await
}
+/** Get current value for [key] */
+pub async fn get_task_status(db: &PgPool, key: &str) -> sqlx::Result<Option<TaskStatus>> {
+ sqlx::query_scalar::<_, Json<TaskStatus>>("SELECT value FROM kv WHERE key=$1")
+ .bind(key)
+ .fetch_optional(db)
+ .await
+ .map(|it| it.map(|it| it.0))
+}
+
+/** Update a TaskStatus timestamp */
+pub async fn update_task_status(
+ db: &PgPool,
+ key: &str,
+ timestamp: &Timestamp,
+ success: bool,
+) -> sqlx::Result<()> {
+ sqlx::query(if success {
+ "INSERT INTO kv (key, value) VALUES ($1, jsonb_build_object('last_successfull', $2, 'last_trial', $3)) ON CONFLICT (key) DO UPDATE SET value=EXCLUDED.value"
+ } else {
+ "INSERT INTO kv (key, value) VALUES ($1, jsonb_build_object('last_trial', $2)) ON CONFLICT (key) DO UPDATE SET value=jsonb_set(EXCLUDED.value, '{last_trial}'::text[], to_jsonb($3))"
+ }).bind(key).bind_timestamp(timestamp).bind_timestamp(timestamp).execute(db).await?;
+ Ok(())
+}
+
#[cfg(test)]
pub mod test {
use sqlx::{PgPool, Postgres, Row, pool::PoolConnection, postgres::PgRow};
diff --git a/src/dialect.rs b/src/dialect.rs
@@ -160,6 +160,16 @@ impl Standard {
})),
}
}
+
+ /*
+
+ /** All orders required for a dialect implementation to work */
+ fun downloadOrders(): Set<EbicsOrder> = (
+ // Administrative orders
+ sequenceOf(EbicsOrder.V3.HAA, EbicsOrder.V3.HKD)
+ // and documents orders
+ + OrderDoc.entries.flatMap { downloadDoc(it) }
+ ).toSet() */
}
/** Supported bank dialects */
diff --git a/src/ebics/administrative.rs b/src/ebics/administrative.rs
@@ -101,7 +101,7 @@ pub enum UserStatus {
pub fn hev_msg(cfg: &EbicsHostCfg) -> String {
xml!(
- "ebicsHEVRequest" ("xmlns": "http://www.ebics.org/H000") {
+ "ebicsHEVRequest" "xmlns"="http://www.ebics.org/H000" {
"HostID": &cfg.host_id
}
)
diff --git a/src/ebics/bts.rs b/src/ebics/bts.rs
@@ -30,7 +30,7 @@ use crate::{
ebics_code::EbicsReturnCode,
order::{BTF, Order},
},
- keys::{BankPubKeysFile, ClientPriKeysFile},
+ keys::{BankKeys, ClientKeys},
utils::b64,
xml,
xml::{Xml, XmlAccess, XmlWriter},
@@ -39,16 +39,16 @@ use crate::{
fn signed_request(
order: &Order,
- client: &ClientPriKeysFile,
+ client: &ClientKeys,
lambda: impl FnOnce(&mut XmlWriter),
) -> String {
let schema = order.schema();
let doc = xml!(
"ebicsRequest"
- ("xmlns": (format_args!("urn:org:ebics:{schema}")))
- ("xmlns:ds": "http://www.w3.org/2000/09/xmldsig#")
- ("Version": schema)
- ("Revision": "1")
+ "xmlns"=(format_args!("urn:org:ebics:{schema}"))
+ "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#"
+ "Version"=schema
+ "Revision"="1"
{
@ lambda
}
@@ -56,17 +56,11 @@ fn signed_request(
sign_ebics(doc, &client.auth)
}
-fn bank_digest(w: &mut XmlWriter, bank: &BankPubKeysFile) {
- xml!(w,
+fn bank_digest(w: &mut XmlWriter, bank: &BankKeys) {
+ xml!(w =>
"BankPubKeyDigests" {
- "Authentication"
- ("Version": "X002")
- ("Algorithm": "http://www.w3.org/2001/04/xmlenc#sha256")
- : b64(ebics_pub_key_hash(&bank.auth.key)),
- "Encryption"
- ("Version": "E002")
- ("Algorithm": "http://www.w3.org/2001/04/xmlenc#sha256")
- : b64(ebics_pub_key_hash(&bank.enc.key))
+ "Authentication" "Version"="X002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256" : b64(ebics_pub_key_hash(&bank.auth.key)),
+ "Encryption" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256" : b64(ebics_pub_key_hash(&bank.enc.key))
},
"SecurityMedium": "0000"
)
@@ -81,24 +75,24 @@ fn service(w: &mut XmlWriter, service: &BTF) {
container,
option,
} = service;
- xml!(w,
+ xml!(w =>
"Service" {
"ServiceName": name,
@ |w: &mut XmlWriter| {
if let Some(scope) = scope {
- xml!(w, "Scope": scope)
+ xml!(w => "Scope": scope)
}
if let Some(option) = option {
- xml!(w, "ServiceOption": option)
+ xml!(w => "ServiceOption": option)
}
if let Some(container) = container {
- xml!(w, "Container" ("containerType": container))
+ xml!(w => "Container" "containerType"=container)
}
if let Some(version) = version {
- xml!(w, "MsgName" ("version": version): msg)
+ xml!(w => "MsgName" "version"=version : msg)
} else {
- xml!(w, "MsgName": msg)
+ xml!(w => "MsgName": msg)
}
}
}
@@ -107,15 +101,15 @@ fn service(w: &mut XmlWriter, service: &BTF) {
pub fn d_init(
cfg: &EbicsHostCfg,
- bank: &BankPubKeysFile,
- client: &ClientPriKeysFile,
+ bank: &BankKeys,
+ client: &ClientKeys,
order: &Order,
range: &Option<(Timestamp, Timestamp)>,
) -> String {
let nonce: u128 = rand::random();
signed_request(order, client, |w| {
- xml!(w,
- "header" ("authenticate": "true") {
+ xml!(w =>
+ "header" "authenticate"="true" {
"static" {
"HostID": cfg.host_id,
"Nonce": format_args!("{:032x}", nonce),
@@ -125,11 +119,11 @@ pub fn d_init(
"OrderDetails" {
"AdminOrderType": order.ty(),
@ |w: &mut XmlWriter| if let Order::BTD(s) = order {
- xml!(w, "BTDOrderParams" {
+ xml!(w => "BTDOrderParams" {
@ |w: &mut XmlWriter| {
service(w, s);
if let Some((start, end)) = range {
- xml!(w,
+ xml!(w =>
"DateRange" {
"Start": Zoned::new(*start, TimeZone::UTC).date(),
"End": Zoned::new(*end, TimeZone::UTC).date()
@@ -139,7 +133,7 @@ pub fn d_init(
}
})
} else {
- xml!(w, "StandardOrderParams")
+ xml!(w => "StandardOrderParams")
}
},
@ |w: &mut XmlWriter| bank_digest(w, bank)
@@ -156,22 +150,22 @@ pub fn d_init(
pub fn d_transfer(
cfg: &EbicsHostCfg,
- client: &ClientPriKeysFile,
+ client: &ClientKeys,
order: &Order,
nb_segment: usize,
segment_nb: usize,
tx_id: &str,
) -> String {
signed_request(order, client, |w| {
- xml!(w,
- "header" ("authenticate": "true") {
+ xml!(w =>
+ "header" "authenticate"="true" {
"static" {
"HostID": cfg.host_id,
"TransactionID": tx_id
},
"mutable" {
"TransactionPhase": "Transfer",
- "SegmentNumber" ("lastSegment": (nb_segment == segment_nb)): segment_nb
+ "SegmentNumber" "lastSegment"=(nb_segment == segment_nb) : segment_nb
}
},
"AuthSignature",
@@ -182,14 +176,14 @@ pub fn d_transfer(
pub fn receipt(
cfg: &EbicsHostCfg,
- client: &ClientPriKeysFile,
+ client: &ClientKeys,
order: &Order,
tx_id: &str,
success: bool,
) -> String {
signed_request(order, client, |w| {
- xml!(w,
- "header" ("authenticate": "true") {
+ xml!(w =>
+ "header" "authenticate"="true" {
"static" {
"HostID": cfg.host_id,
"TransactionID": tx_id
@@ -200,7 +194,7 @@ pub fn receipt(
},
"AuthSignature",
"body" {
- "TransferReceipt" ("authenticate": "true") {
+ "TransferReceipt" "authenticate"="true" {
"ReceiptCode": (if success { "0" } else { "1"})
}
}
@@ -210,15 +204,15 @@ pub fn receipt(
pub fn u_init(
cfg: &EbicsHostCfg,
- bank: &BankPubKeysFile,
- client: &ClientPriKeysFile,
+ bank: &BankKeys,
+ client: &ClientKeys,
order: &Order,
data: &PreparedUploadData,
) -> String {
let nonce: u128 = rand::random();
signed_request(order, client, |w| {
- xml!(w,
- "header" ("authenticate": "true") {
+ xml!(w =>
+ "header" "authenticate"="true" {
"static" {
"HostID": cfg.host_id,
"Nonce": format_args!("{:032x}", nonce),
@@ -228,12 +222,12 @@ pub fn u_init(
"OrderDetails" {
"AdminOrderType": order.ty(),
@ |w: &mut XmlWriter| if let Order::BTU(s) = order {
- xml!(w, "BTUOrderParams" {
+ xml!(w => "BTUOrderParams" {
@ |w: &mut XmlWriter| service(w, s),
"SignatureFlag"
})
} else {
- xml!(w, "StandardOrderParams")
+ xml!(w => "StandardOrderParams")
}
},
@ |w: &mut XmlWriter| bank_digest(w, bank),
@@ -246,15 +240,12 @@ pub fn u_init(
"AuthSignature",
"body" {
"DataTransfer" {
- "DataEncryptionInfo" ("authenticate": "true") {
- "EncryptionPubKeyDigest"
- ("Version": "E002")
- ("Algorithm": "http://www.w3.org/2001/04/xmlenc#sha256")
- : b64(ebics_pub_key_hash(&bank.enc.key)),
+ "DataEncryptionInfo" "authenticate"="true" {
+ "EncryptionPubKeyDigest" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256": b64(ebics_pub_key_hash(&bank.enc.key)),
"TransactionKey": b64(&data.encrypted_key)
},
- "SignatureData" ("authenticate": "true"): data.signature_data,
- "DataDigest" ("SignatureVersion": "A006"): b64(data.digest)
+ "SignatureData" "authenticate"="true" : data.signature_data,
+ "DataDigest" "SignatureVersion"="A006" : b64(data.digest)
}
}
)
@@ -263,22 +254,22 @@ pub fn u_init(
pub fn u_transfer(
cfg: &EbicsHostCfg,
- client: &ClientPriKeysFile,
+ client: &ClientKeys,
order: &Order,
tx_id: &str,
data: &PreparedUploadData,
segment_nb: usize,
) -> String {
signed_request(order, client, |w| {
- xml!(w,
- "header" ("authenticate": "true") {
+ xml!(w =>
+ "header" "authenticate"="true" {
"static" {
"HostID": cfg.host_id,
"TransactionID": tx_id
},
"mutable" {
"TransactionPhase": "Transfer",
- "SegmentNumber" ("lastSegment": (data.nb_segments() == segment_nb)): segment_nb
+ "SegmentNumber" "lastSegment"=(data.nb_segments() == segment_nb) : segment_nb
}
},
"AuthSignature",
diff --git a/src/ebics/key_management.rs b/src/ebics/key_management.rs
@@ -33,7 +33,7 @@ use crate::{
bts::DataEncryptionInfo, decrypt_and_decompress_payload, ebics_code::EbicsReturnCode,
order::Order,
},
- keys::{self, BankPubKeysFile, ClientPriKeysFile, RsaPub},
+ keys::{self, BankKeys, ClientKeys, RsaPub},
xml,
xml::{Xml, XmlAccess as _, XmlWriter},
xml_sign::sign_ebics,
@@ -44,7 +44,7 @@ impl EbicsClient {
pub async fn submit_client_keys(
&self,
cfg: &EbicsKeysCfg,
- client: &mut ClientPriKeysFile,
+ client: &mut ClientKeys,
order: Order,
) -> Result<(), EbicsError> {
let ctx = EbicsCtx::new(&order);
@@ -73,7 +73,7 @@ impl EbicsClient {
}
/** Perform an EBICS private key management HPB using [client] */
- pub async fn hpb(&self, client: &ClientPriKeysFile) -> anyhow::Result<BankPubKeysFile> {
+ pub async fn hpb(&self, client: &ClientKeys) -> anyhow::Result<BankKeys> {
let order = Order::HPB;
let res = self.key_management(client, &order).await?;
if res.technical_code == EbicsReturnCode::EBICS_AUTHENTICATION_FAILED {
@@ -84,12 +84,6 @@ impl EbicsClient {
}
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)
- }
-
Ok(Xml::parse(&order_data, "HPBResponseOrderData", |root| {
let auth_pub = root.one("AuthenticationPubKeyInfo")?;
let version = auth_pub.one("AuthenticationVersion")?.text();
@@ -107,7 +101,7 @@ impl EbicsClient {
);
let enc_pub = rsa_pub_key(enc_pub)?;
- Ok(BankPubKeysFile {
+ Ok(BankKeys {
auth: auth_pub,
enc: enc_pub,
accepted: false,
@@ -117,7 +111,7 @@ impl EbicsClient {
async fn key_management(
&self,
- client: &ClientPriKeysFile,
+ client: &ClientKeys,
order: &Order,
) -> Result<EbicsResponse<Option<Vec<u8>>>, EbicsError> {
let EbicsHostCfg {
@@ -141,7 +135,7 @@ impl EbicsClient {
schema: &str,
build: impl FnOnce(&mut XmlWriter),
) -> String {
- let xml = xml!(name ("xmlns":schema) ("xmlns:ds":"http://www.w3.org/2000/09/xmldsig#") {
+ let xml = xml!(name "xmlns"=schema "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
@ build,
"PartnerID": &cfg.partner_id,
"UserID": &cfg.user_id
@@ -153,36 +147,13 @@ impl EbicsClient {
BASE64_STANDARD.encode(&compressed)
}
- 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());
-
- 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" {
+ xml!(w => "SignaturePubKeyInfo" {
@ |w| rsa_key_xml(w, &client.sign),
"SignatureVersion": "A006"
})
@@ -193,7 +164,7 @@ impl EbicsClient {
"HIARequestOrderData",
"urn:org:ebics:H005",
|w| {
- xml!(w,
+ xml!(w =>
"AuthenticationPubKeyInfo" {
@ |w| rsa_key_xml(w, &client.auth),
"AuthenticationVersion": "X002"
@@ -211,17 +182,17 @@ impl EbicsClient {
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")
+ "xmlns"="urn:org:ebics:H005"
+ "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#"
+ "Version"="H005"
+ "Revision"="1"
{
- "header" ("authenticate": "true") {
+ "header" "authenticate"="true" {
"static" {
"HostID": host_id,
@ |w: &mut XmlWriter| if *order == Order::HPB {
let nonce: u128 = rand::random();
- xml!(w,
+ xml!(w =>
"Nonce": format_args!("{:032x}", nonce),
"Timestamp": jiff::Timestamp::now()
)
@@ -236,11 +207,11 @@ impl EbicsClient {
"mutable"
},
@ |w: &mut XmlWriter| if sign {
- xml!(w, "AuthSignature")
+ xml!(w => "AuthSignature")
},
"body" {
@ |w: &mut XmlWriter| if let Some(data) = data {
- xml!(w, "DataTransfer" {
+ xml!(w => "DataTransfer" {
"OrderData": data
})
}
@@ -277,3 +248,31 @@ impl EbicsClient {
.ctx(&ctx)
}
}
+
+pub fn rsa_pub_key(xml: Xml) -> xml::Result<RsaPub> {
+ xml.one("X509Data")
+ .one("X509Certificate")
+ .decode(rsa_private_from_b64_x509_certificate)
+}
+
+pub 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());
+
+ xml!(w => "ds:X509Data" {
+ "ds:X509Certificate": b64
+ })
+}
diff --git a/src/ebics/mod.rs b/src/ebics/mod.rs
@@ -26,11 +26,11 @@ use flate2::write::ZlibDecoder;
use jiff::Timestamp;
use rand::{RngExt as _, distr::Alphanumeric};
use reqwest::{
- Client, StatusCode,
+ Client, ClientBuilder, StatusCode,
header::{CONTENT_TYPE, HeaderValue},
};
use sqlx::PgPool;
-use tracing::{debug, info, trace};
+use tracing::{debug, info, trace, warn};
use crate::{
EbicsLogs,
@@ -41,7 +41,7 @@ use crate::{
},
db::{ebics_first, ebics_register, ebics_remove},
ebics::{
- administrative::{HAA, VersionNumber, hev_msg, parse_haa, parse_hev},
+ administrative::{HAA, HKD, VersionNumber, hev_msg, parse_haa, parse_hev, parse_hkd},
bts::{
DInit, DTransfer, DataEncryptionInfo, U, d_init, d_transfer, parse_d_init,
parse_d_transfer, parse_receipt, parse_u_init, parse_u_transfer, receipt, u_init,
@@ -51,7 +51,7 @@ use crate::{
logger::EbicsLogger,
order::Order,
},
- keys::{BankPubKeysFile, ClientPriKeysFile},
+ keys::{BankKeys, ClientKeys},
utils::{b64, deflate},
xml,
};
@@ -287,9 +287,14 @@ pub struct EbicsClient {
impl EbicsClient {
pub fn new(cfg: &NexusCfg, log: EbicsLogs) -> anyhow::Result<Self> {
+ let cfg = cfg.host()?.clone();
+ let mut builder = ClientBuilder::new();
+ if let Some(unix_path) = &cfg.unix_path {
+ builder = builder.unix_socket(unix_path.as_str());
+ }
Ok(Self {
- cfg: cfg.host()?.clone(),
- http: Client::new(),
+ cfg,
+ http: builder.build()?,
logger: EbicsLogger::new(log.dir)?,
})
}
@@ -345,8 +350,8 @@ impl EbicsClient {
pub async fn haa(
&self,
db: &PgPool,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
peek: bool,
) -> Result<HAA, EbicsError> {
self.download(
@@ -361,6 +366,25 @@ impl EbicsClient {
.await
}
+ pub async fn hkd(
+ &self,
+ db: &PgPool,
+ client: &ClientKeys,
+ bank: &BankKeys,
+ peek: bool,
+ ) -> Result<HKD, EbicsError> {
+ self.download(
+ db,
+ client,
+ bank,
+ &Order::HKD,
+ &None,
+ peek,
+ async |content| Ok(parse_hkd(&content)?),
+ )
+ .await
+ }
+
/**
* Performs an EBICS download transaction of [order] between [startDate] and [endDate].
* Download content is passed to [processing]
@@ -372,8 +396,8 @@ impl EbicsClient {
pub async fn download<T>(
&self,
db: &PgPool,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
order: &Order,
range: &Option<(Timestamp, Timestamp)>,
peek: bool,
@@ -439,8 +463,13 @@ impl EbicsClient {
// Receipt phase
ctx = ctx.receipt(&tx_id);
let xml = receipt(&self.cfg, client, order, &tx_id, res.is_ok() && !peek);
- if self.post_bts(xml, &ctx, parse_receipt).await.is_ok() {
- ebics_remove(db, &tx_id).await.ok();
+ if let Err(e) = async {
+ self.post_bts(xml, &ctx, parse_receipt).await?;
+ ebics_remove(db, &tx_id).await.ctx(&ctx)
+ }
+ .await
+ {
+ warn!(target: "ebics", "{e}")
}
res
@@ -455,8 +484,8 @@ impl EbicsClient {
*/
pub async fn upload(
&self,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
order: &Order,
payload: &str,
) -> Result<CompactString, EbicsError> {
@@ -522,8 +551,8 @@ fn decrypt_and_decompress_payload(
/** Signs, encrypts and format EBICS BTS payload */
fn prepare_upload_payload(
cfg: &EbicsHostCfg,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
payload: &str,
) -> PreparedUploadData {
let digest = digest_ebics_order_a006(payload.as_bytes());
@@ -535,7 +564,7 @@ fn prepare_upload_payload(
let signature_data = {
let signed = sign_ebics_a006(digest.as_ref(), &client.sign);
let inner_signed_xml = xml!(
- "UserSignatureData" ("xmlns": "http://www.ebics.org/S002") {
+ "UserSignatureData" "xmlns"="http://www.ebics.org/S002" {
"OrderSignatureData" {
"SignatureVersion": "A006",
"SignatureValue": b64(&signed),
@@ -545,14 +574,14 @@ fn prepare_upload_payload(
}
);
let deflated = deflate(inner_signed_xml.as_bytes());
- let encrypted = encrypt_ebics_e002(&tx_key, &deflated);
+ let encrypted = encrypt_ebics_e002(&tx_key, deflated);
BASE64_STANDARD.encode(encrypted)
};
// Compress and encrypt payload
let payload = {
let deflated = deflate(payload.as_bytes());
- let encrypted = encrypt_ebics_e002(&tx_key, &deflated);
+ let encrypted = encrypt_ebics_e002(&tx_key, deflated);
BASE64_STANDARD.encode(encrypted)
};
PreparedUploadData {
@@ -581,8 +610,8 @@ pub struct TxCheckResult {
pub async fn tx_check(
ebics: &EbicsClient,
db: &PgPool,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
fetch: &Order,
submit: &Order,
) -> anyhow::Result<TxCheckResult> {
@@ -718,3 +747,481 @@ pub async fn tx_check(
Ok(result)
}
+
+#[cfg(test)]
+pub mod test {
+ use aws_lc_rs::{
+ encoding::AsDer,
+ rsa::{KeyPair, KeySize, PublicEncryptingKey, PublicKey},
+ };
+ use compact_str::CompactString;
+ use jiff::{
+ Timestamp, Zoned,
+ civil::{Date, date},
+ tz::TimeZone,
+ };
+
+ use crate::{
+ crypto::{ebics_pub_key_hash, encrypt_ebics_e002, gen_ebics_e002_key},
+ ebics::key_management::{rsa_key_xml, rsa_pub_key},
+ rand_ebics_id,
+ utils::{b64, deflate, inflate},
+ xml,
+ xml::{Xml, XmlAccess},
+ xml_sign::sign_ebics,
+ };
+
+ pub type Sequence = fn(&mut EbicsState, body: &[u8]) -> EbicsRes;
+
+ pub enum EbicsRes {
+ Ok(String),
+ BadRequest,
+ Failure,
+ }
+ pub struct EbicsState {
+ bank_sign: KeyPair,
+ bank_enc: KeyPair,
+ bank_auth: KeyPair,
+
+ client_sign: Option<PublicKey>,
+ client_enc: Option<PublicKey>,
+ client_auth: Option<PublicKey>,
+
+ tx_id: Option<CompactString>,
+ order_id: Option<CompactString>,
+ }
+
+ impl EbicsState {
+ pub fn new() -> Self {
+ Self {
+ bank_sign: KeyPair::generate(KeySize::Rsa2048).unwrap(),
+ bank_enc: KeyPair::generate(KeySize::Rsa2048).unwrap(),
+ bank_auth: KeyPair::generate(KeySize::Rsa2048).unwrap(),
+ client_sign: None,
+ client_enc: None,
+ client_auth: None,
+ tx_id: None,
+ order_id: None,
+ }
+ }
+
+ fn parse_unsecure_request(
+ body: &[u8],
+ order: &str,
+ root: &str,
+ parse: impl FnOnce(Xml) -> xml::Result<()>,
+ ) {
+ Xml::parse(body, "ebicsUnsecuredRequest", |n| {
+ let admin_order = n
+ .one("header")
+ .one("static")
+ .one("OrderDetails")
+ .one("AdminOrderType")?
+ .text();
+ assert_eq!(admin_order, order);
+ let chunk = n.one("body").one("DataTransfer").one("OrderData").b64()?;
+ let inflated = inflate(&chunk);
+ Xml::parse(&inflated, root, parse)
+ })
+ .unwrap()
+ }
+
+ fn parse_download_init(body: &[u8], order: &str) {
+ Xml::parse(body, "ebicsRequest", |root| {
+ let header = root.one("header")?;
+ let admin_order = header
+ .one("static")
+ .one("OrderDetails")
+ .one("AdminOrderType")?
+ .text();
+ assert_eq!(admin_order, order);
+ let phase = header.one("mutable").one("TransactionPhase")?.text();
+ assert_eq!(phase, "Initialisation");
+ Ok(())
+ })
+ .unwrap();
+ }
+
+ fn signed_response(&self, xml: String) -> EbicsRes {
+ EbicsRes::Ok(sign_ebics(xml, &self.bank_auth))
+ }
+
+ fn ebics_response_payload(&mut self, payload: &str, last: bool) -> EbicsRes {
+ let tx_id = self.tx_id.insert(rand_ebics_id());
+ let deflated = deflate(payload.as_bytes());
+ let client_enc = PublicEncryptingKey::from_der(
+ self.client_enc.as_ref().unwrap().as_der().unwrap().as_ref(),
+ )
+ .unwrap();
+ let (tx_key, encrypted_key) = gen_ebics_e002_key(client_enc);
+ let encrypted = encrypt_ebics_e002(&tx_key, deflated);
+ let xml = xml!("ebicsResponse" "xmlns"="http://www.ebics.org/H005" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
+ "header" "authenticate"="true" {
+ "static" {
+ "TransactionID": tx_id,
+ "NumSegments": "1"
+ },
+ "mutable" {
+ "TransactionPhase": "Initialisation",
+ "SegmentNumber" "lastSegment"=last : 1,
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK"
+ }
+ },
+ "AuthSignature",
+ "body" {
+ "DataTransfer" {
+ "DataEncryptionInfo" "authenticate"="true" {
+ "EncryptionPubKeyDigest" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256": b64(ebics_pub_key_hash(&self.client_enc.as_ref().unwrap())),
+ "TransactionKey": b64(&encrypted_key)
+ },
+ "OrderData": b64(&encrypted)
+ },
+ "ReturnCode" "authenticate"="true": "000000"
+ }
+ });
+ self.signed_response(xml)
+ }
+
+ fn ebics_response_no_data(&self) -> EbicsRes {
+ let xml = xml!("ebicsResponse" "xmlns"="http://www.ebics.org/H005" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
+ "header" "authenticate"="true" {
+ "static",
+ "mutable" {
+ "TransactionPhase": "Initialisation",
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK"
+ }
+ },
+ "AuthSignature",
+ "body" {
+ "ReturnCode" "authenticate"="true": "090005"
+ }
+ });
+ self.signed_response(xml)
+ }
+
+ pub fn hev(&mut self, body: &[u8]) -> EbicsRes {
+ Xml::parse(body, "ebicsHEVRequest", |root| {
+ root.one("HostID")?;
+ Ok(())
+ })
+ .unwrap();
+ EbicsRes::Ok(
+ xml!("ebicsHEVResponse" "xmlns"="http://www.ebics.org/H000" {
+ "SystemReturnCode" {
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK"
+ },
+ "VersionNumber" "ProtocolVersion"="H005" : "03.00"
+ }),
+ )
+ }
+
+ pub fn ini(&mut self, body: &[u8]) -> EbicsRes {
+ Self::parse_unsecure_request(body, "INI", "SignaturePubKeyOrderData", |root| {
+ let n = root.one("SignaturePubKeyInfo")?;
+ assert_eq!(n.one("SignatureVersion")?.text(), "A006");
+ self.client_sign = Some(rsa_pub_key(n)?.key);
+ Ok(())
+ });
+ EbicsRes::Ok(
+ xml!("ebicsKeyManagementResponse" "xmlns"="http://www.ebics.org/H000" {
+ "header" "authenticate"="true" {
+ "mutable" {
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK"
+ }
+ },
+ "body" {
+ "ReturnCode" "authenticate"="true" : "000000"
+ }
+ }),
+ )
+ }
+
+ pub fn hia(&mut self, body: &[u8]) -> EbicsRes {
+ Self::parse_unsecure_request(body, "HIA", "HIARequestOrderData", |root| {
+ let n = root.one("AuthenticationPubKeyInfo")?;
+ assert_eq!(n.one("AuthenticationVersion")?.text(), "X002");
+ self.client_auth = Some(rsa_pub_key(n)?.key);
+
+ let n = root.one("EncryptionPubKeyInfo")?;
+ assert_eq!(n.one("EncryptionVersion")?.text(), "E002");
+ self.client_enc = Some(rsa_pub_key(n)?.key);
+ Ok(())
+ });
+ EbicsRes::Ok(
+ xml!("ebicsKeyManagementResponse" "xmlns"="http://www.ebics.org/H000" {
+ "header" "authenticate"="true" {
+ "mutable" {
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK"
+ }
+ },
+ "body" {
+ "ReturnCode" "authenticate"="true" : "000000"
+ }
+ }),
+ )
+ }
+
+ pub fn hpb(&mut self, body: &[u8]) -> EbicsRes {
+ // Parse HPB request
+ Xml::parse(body, "ebicsNoPubKeyDigestsRequest", |root| {
+ let order = root
+ .one("header")
+ .one("static")
+ .one("OrderDetails")
+ .one("AdminOrderType")?
+ .text();
+ assert_eq!(order, "HPB");
+ Ok(())
+ })
+ .unwrap();
+
+ let payload = xml!("HPBResponseOrderData" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
+ "AuthenticationPubKeyInfo" {
+ @ |w| rsa_key_xml(w, &self.bank_auth),
+ "AuthenticationVersion": "X002"
+ },
+ "EncryptionPubKeyInfo" {
+ @ |w| rsa_key_xml(w, &self.bank_enc),
+ "EncryptionVersion": "E002"
+ }
+ });
+ let deflated = deflate(payload.as_bytes());
+ let client_enc = PublicEncryptingKey::from_der(
+ self.client_enc.as_ref().unwrap().as_der().unwrap().as_ref(),
+ )
+ .unwrap();
+ let (tx_key, encrypted_key) = gen_ebics_e002_key(client_enc);
+ let encrypted = encrypt_ebics_e002(&tx_key, deflated);
+ EbicsRes::Ok(
+ xml!("ebicsKeyManagementResponse" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" "xmlns"="http://www.ebics.org/H005" {
+ "header" "authenticate"="true"{
+ "mutable" {
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK"
+ }
+ },
+ "body" {
+ "DataTransfer" {
+ "DataEncryptionInfo" "authenticate"="true" {
+ "EncryptionPubKeyDigest" "Version"="E002" "Algorithm"="http://www.w3.org/2001/04/xmlenc#sha256": b64(ebics_pub_key_hash(&self.client_enc.as_ref().unwrap())),
+ "TransactionKey": b64(&encrypted_key)
+ },
+ "OrderData": b64(&encrypted)
+ },
+ "ReturnCode" "authenticate"="true": "000000"
+ }
+ }),
+ )
+ }
+
+ pub fn hkd(&mut self, body: &[u8]) -> EbicsRes {
+ Self::parse_download_init(body, "HKD");
+ self.ebics_response_payload(
+ &xml!("HKDResponseOrderData" {
+ "PartnerInfo" {
+ "AddressInfo",
+ "OrderInfo" {
+ "AdminOrderType": "BTD",
+ "Service" {
+ "ServiceName": "STM",
+ "Scope": "CH",
+ "Container" "containerType"="ZIP",
+ "MsgName" "version"="08": "camt.052"
+ },
+ "Description"
+ },
+ "OrderInfo" {
+ "AdminOrderType": "BTU",
+ "Service" {
+ "ServiceName": "SCT",
+ "MsgName": "pain.001"
+ },
+ "Description": "Direct Debit"
+ },
+ "OrderInfo" {
+ "AdminOrderType": "BTU",
+ "Service" {
+ "ServiceName": "SCI",
+ "Scope": "DE",
+ "MsgName": "pain.001"
+ },
+ "Description": "Instant Direct Debit"
+ }
+ }
+ }),
+ true,
+ )
+ }
+
+ pub fn haa(&mut self, body: &[u8]) -> EbicsRes {
+ Self::parse_download_init(body, "HAA");
+ self.ebics_response_payload(
+ &xml!("HAAResponseOrderData" {
+ "Service" {
+ "ServiceName": "STM",
+ "Scope": "CH",
+ "Container" "containerType"="ZIP",
+ "MsgName" "version"="08": "camt.052"
+ }
+ }),
+ true,
+ )
+ }
+
+ fn receipt(&mut self, body: &[u8], ok: bool) -> EbicsRes {
+ Xml::parse(body, "ebicsRequest", |root| {
+ let header = root.one("header")?;
+ let tx_id = header.one("static").one("TransactionID")?.text();
+ assert_eq!(tx_id, self.tx_id.as_deref().unwrap());
+ let phase = header.one("mutable").one("TransactionPhase")?.text();
+ assert_eq!(phase, "Receipt");
+ let code = root
+ .one("body")
+ .one("TransferReceipt")
+ .one("ReceiptCode")?
+ .text();
+ if ok {
+ assert_eq!(code, "0")
+ } else {
+ assert_eq!(code, "1")
+ }
+ Ok(())
+ })
+ .unwrap();
+ let tx_id = self.tx_id.take().unwrap();
+ self.signed_response(xml!("ebicsResponse" "xmlns"="http://www.ebics.org/H005" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
+ "header" "authenticate"="true" {
+ "static" {
+ "TransactionID": tx_id
+ },
+ "mutable" {
+ "TransactionPhase": "Receipt",
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK",
+ }
+ },
+ "AuthSignature",
+ "body" {
+ "ReturnCode" "authenticate"="true": "000000"
+ }
+ }))
+ }
+
+ pub fn receipt_ok(&mut self, body: &[u8]) -> EbicsRes {
+ self.receipt(body, true)
+ }
+
+ pub fn receipt_err(&mut self, body: &[u8]) -> EbicsRes {
+ self.receipt(body, false)
+ }
+
+ fn btd_date_check(&self, body: &[u8], pinned: Option<Date>) -> EbicsRes {
+ Xml::parse(body, "ebicsRequest", |root| {
+ let header = root.one("header")?;
+ let details = header.one("static").one("OrderDetails")?;
+ let admin_order = details.one("AdminOrderType")?.text();
+ assert_eq!(admin_order, "BTD");
+ let start = details
+ .one("BTDOrderParams")
+ .opt("DateRange")
+ .opt("Start")
+ .parse()?;
+ assert_eq!(start, pinned);
+ let phase = header.one("mutable").one("TransactionPhase")?.text();
+ assert_eq!(phase, "Initialisation");
+ Ok(())
+ })
+ .unwrap();
+ self.ebics_response_no_data()
+ }
+
+ pub fn btd_no_data(&mut self, body: &[u8]) -> EbicsRes {
+ self.btd_date_check(body, None)
+ }
+
+ pub fn btd_no_data_now(&mut self, body: &[u8]) -> EbicsRes {
+ self.btd_date_check(
+ body,
+ Some(Zoned::new(Timestamp::now(), TimeZone::UTC).date()),
+ )
+ }
+
+ pub fn btd_no_data_pinned(&mut self, body: &[u8]) -> EbicsRes {
+ self.btd_date_check(body, Some(date(2024, 06, 05)))
+ }
+
+ pub fn btu_init(&mut self, body: &[u8]) -> EbicsRes {
+ Self::parse_download_init(body, "BTU");
+ let tx_id = self.tx_id.insert(rand_ebics_id());
+ let order_id = self.order_id.insert(rand_ebics_id());
+ let xml = xml!("ebicsResponse" "xmlns"="http://www.ebics.org/H005" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
+ "header" "authenticate"="true" {
+ "static" {
+ "TransactionID": tx_id
+ },
+ "mutable" {
+ "TransactionPhase": "Initialisation",
+ "OrderID": order_id,
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK",
+ }
+ },
+ "AuthSignature",
+ "body" {
+ "ReturnCode" "authenticate"="true": "000000"
+ }
+ });
+ self.signed_response(xml)
+ }
+
+ pub fn btu_payload(&mut self, body: &[u8]) -> EbicsRes {
+ let tx_id = self.tx_id.as_ref().unwrap();
+ let order_id = self.order_id.as_ref().unwrap();
+ let segment_nb: CompactString = Xml::parse(body, "ebicsRequest", |root| {
+ let header = root.one("header")?;
+ let txid = header.one("static").one("TransactionID")?.text();
+ assert_eq!(txid, tx_id);
+ let mutable = header.one("mutable")?;
+ let phase = mutable.one("TransactionPhase")?.text();
+ assert_eq!(phase, "Transfer");
+ mutable.one("SegmentNumber").parse()
+ })
+ .unwrap();
+ self.signed_response(xml!("ebicsResponse" "xmlns"="http://www.ebics.org/H005" "xmlns:ds"="http://www.w3.org/2000/09/xmldsig#" {
+ "header" "authenticate"="true" {
+ "static" {
+ "TransactionID": tx_id
+ },
+ "mutable" {
+ "TransactionPhase": "Transfer",
+ "SegmentNumber": segment_nb,
+ "OrderID": order_id,
+ "ReturnCode": "000000",
+ "ReportText": "[EBICS_OK] OK",
+ }
+ },
+ "AuthSignature",
+ "body" {
+ "ReturnCode" "authenticate"="true": "000000"
+ }
+ }))
+ }
+
+ pub fn init_tx(&mut self, _: &[u8]) -> EbicsRes {
+ self.ebics_response_payload("", false)
+ }
+
+ pub fn failure(&mut self, _: &[u8]) -> EbicsRes {
+ EbicsRes::Failure
+ }
+
+ pub fn bad_request(&mut self, _: &[u8]) -> EbicsRes {
+ EbicsRes::BadRequest
+ }
+ }
+}
diff --git a/src/iso20022/pain001.rs b/src/iso20022/pain001.rs
@@ -71,9 +71,9 @@ pub fn create_pain001(
let total = ebics_amount(&msg.sum)?;
Ok(xml!(
"Document"
- ("xmlns": (format_args!("urn:iso:std:iso:20022:tech:xsd:pain.001.001.{version}")))
- ("xmlns:xsi": (format_args!("http://www.w3.org/2001/XMLSchema-instance")))
- ("xsi:schemaLocation": (format_args!("urn:iso:std:iso:20022:tech:xsd:pain.001.001.{version} pain.001.001.{version}{suffix}.xsd")))
+ "xmlns"=(format_args!("urn:iso:std:iso:20022:tech:xsd:pain.001.001.{version}"))
+ "xmlns:xsi"=(format_args!("http://www.w3.org/2001/XMLSchema-instance"))
+ "xsi:schemaLocation"=(format_args!("urn:iso:std:iso:20022:tech:xsd:pain.001.001.{version} pain.001.001.{version}{suffix}.xsd"))
{
"CstmrCdtTrfInitn" {
"GrpHdr" {
@@ -100,12 +100,12 @@ pub fn create_pain001(
"NbOfTxs": msg.txs.len(),
"CtrlSum": total,
@ |w: &mut XmlWriter| if dialect.standard() == Standard::GBIC {
- xml!(w, "PmtTpInf" {
+ xml!(w => "PmtTpInf" {
"SvcLvl" {
"Cd": "SEPA"
},
@ |w: &mut XmlWriter| if instant {
- xml!(w, "LclInstrm" {
+ xml!(w => "LclInstrm" {
"Cd": "INST"
})
}
@@ -125,9 +125,9 @@ pub fn create_pain001(
"DbtrAgt" {
"FinInstnId" {
@ |w: &mut XmlWriter| if let Some(bic) = &msg.debtor.bic {
- xml!(w, "BICFI": bic)
+ xml!(w => "BICFI": bic)
} else {
- xml!(w, "Othr" {
+ xml!(w => "Othr" {
"Id": "NOTPROVIDED"
})
}
@@ -136,17 +136,17 @@ pub fn create_pain001(
},
"ChrgBr": "SLEV",
@ |w: &mut XmlWriter| for tx in &msg.txs {
- xml!(w, "CdtTrfTxInf" {
+ xml!(w => "CdtTrfTxInf" {
"PmtId" {
"InstrId": tx.e2e_id,
// Used to uniquely identify transactions in other files
"EndToEndId": tx.e2e_id
},
"Amt" {
- "InstdAmt" ("Ccy": (tx.amount.currency)): ebics_amount(&tx.amount).unwrap()
+ "InstdAmt" "Ccy"=(tx.amount.currency) : ebics_amount(&tx.amount).unwrap()
},
@ |w: &mut XmlWriter| if let Some(bic) = &tx.creditor.bic {
- xml!(w, "CdtrAgt" {
+ xml!(w => "CdtrAgt" {
"FinInstnId" {
"BICFI": bic
}
diff --git a/src/keys.rs b/src/keys.rs
@@ -35,7 +35,7 @@ use taler_common::{
use crate::config::EbicsKeysCfg;
#[derive(Debug, serde::Serialize, serde::Deserialize)]
-pub struct ClientPriKeysFile {
+pub struct ClientKeys {
#[serde(
rename = "signature_private_key",
serialize_with = "ser_pkcs8",
@@ -58,7 +58,7 @@ pub struct ClientPriKeysFile {
pub submitted_hia: bool,
}
-impl ClientPriKeysFile {
+impl ClientKeys {
pub fn generate() -> anyhow::Result<Self> {
Ok(Self {
sign: RsaKeyPair::generate(KeySize::Rsa2048)?,
@@ -126,7 +126,7 @@ impl PartialEq for RsaPub {
impl Eq for RsaPub {}
#[derive(Debug, serde::Serialize, serde::Deserialize)]
-pub struct BankPubKeysFile {
+pub struct BankKeys {
#[serde(rename = "bank_encryption_public_key")]
pub enc: RsaPub,
#[serde(rename = "bank_authentication_public_key")]
@@ -170,20 +170,20 @@ where
}
/// Persist the bank keys file to disk
-pub fn persist_bank_keys(keys: &BankPubKeysFile, location: &Path) -> std::io::Result<()> {
+pub fn persist_bank_keys(keys: &BankKeys, location: &Path) -> std::io::Result<()> {
json_file::persist(location, keys)?;
// TODO better error message "bank public keys"
Ok(())
}
-pub fn persist_client_keys(keys: &ClientPriKeysFile, location: &Path) -> std::io::Result<()> {
+pub fn persist_client_keys(keys: &ClientKeys, location: &Path) -> std::io::Result<()> {
json_file::persist(location, keys)?;
// TODO better error message "client private keys"
Ok(())
}
/// Load the bank keys file from disk
-pub fn load_bank_keys(path: &Path) -> anyhow::Result<Option<BankPubKeysFile>> {
+pub fn load_bank_keys(path: &Path) -> anyhow::Result<Option<BankKeys>> {
match json_file::load(path) {
Ok(existing) => Ok(Some(existing)),
Err(e) if e.kind() == ErrorKind::NotFound => Ok(None),
@@ -196,7 +196,7 @@ pub fn load_bank_keys(path: &Path) -> anyhow::Result<Option<BankPubKeysFile>> {
}
/// Load the client keys file from disk
-pub fn load_client_keys(path: &Path) -> anyhow::Result<Option<ClientPriKeysFile>> {
+pub fn load_client_keys(path: &Path) -> anyhow::Result<Option<ClientKeys>> {
match json_file::load(path) {
Ok(existing) => Ok(Some(existing)),
Err(e) if e.kind() == ErrorKind::NotFound => Ok(None),
@@ -209,9 +209,7 @@ pub fn load_client_keys(path: &Path) -> anyhow::Result<Option<ClientPriKeysFile>
}
/// Load client and bank keys from disk and checks that the keying process has been fully completed
-pub fn expect_full_keys(
- cfg: &EbicsKeysCfg,
-) -> anyhow::Result<(ClientPriKeysFile, BankPubKeysFile)> {
+pub fn expect_full_keys(cfg: &EbicsKeysCfg) -> anyhow::Result<(ClientKeys, BankKeys)> {
let setup_cmd = "TODO";
let client_keys = load_client_keys(cfg.client_priv_keys_path.as_ref())?;
let Some(client_keys) = client_keys else {
diff --git a/src/lib.rs b/src/lib.rs
@@ -22,12 +22,14 @@ use std::{
io::{Cursor, Read},
path::{Path, PathBuf},
str::FromStr,
+ time::Duration,
};
use anyhow::{anyhow, bail};
use compact_str::{CompactString, CompactStringExt};
-use jiff::{Timestamp, civil::Date};
+use jiff::{Timestamp, Zoned, civil::Date, tz::TimeZone};
use rand::prelude::IndexedRandom;
+use serde::{Deserialize, Deserializer, Serialize, Serializer};
use sqlx::PgPool;
use taler_build::long_version;
use taler_common::{
@@ -37,20 +39,22 @@ use taler_common::{
types::{
amount::Amount,
payto::{FullIbanPayto, TransferIbanPayto},
+ utils::date_to_utc_ts,
},
};
+use tokio::{time::timeout, try_join};
use tracing::{debug, error, info, trace, warn};
use crate::{
config::{EbicsKeysCfg, NexusCfg},
crypto::ebics_pub_key_hash,
db::{
- dbinit,
+ dbinit, get_task_status,
initiated::{
batch_initiated, batch_status_update, batch_sub_failure, batch_sub_success, initiate,
initiated_submittable, order_failure, order_step, order_success, tx_status_update,
},
- pool,
+ pool, update_task_status,
},
ebics::{
EbicsClient, EbicsCtx, EbicsErrKind, EbicsError, EbicsErrorHelper,
@@ -67,7 +71,7 @@ use crate::{
status_code::{PaymentGroupStatus, PaymentTransactionStatus},
},
keys::{
- BankPubKeysFile, ClientPriKeysFile, expect_full_keys, load_bank_keys, load_client_keys,
+ BankKeys, ClientKeys, expect_full_keys, load_bank_keys, load_client_keys,
persist_bank_keys, persist_client_keys,
},
list::ListCmd,
@@ -75,6 +79,7 @@ use crate::{
testing::TestingCmd,
utils::hex_chunk_by_two,
worker::register_tx,
+ ws::listen_for_notification,
};
pub mod api;
@@ -88,6 +93,8 @@ pub mod iso20022;
pub mod keys;
pub mod list;
pub mod model;
+#[cfg(test)]
+pub mod test;
pub mod testing;
pub mod utils;
pub mod worker;
@@ -95,6 +102,11 @@ pub mod ws;
pub mod xml;
pub mod xml_sign;
+// KV
+const CHECKPOINT_KEY: &str = "checkpoint";
+const SUBMIT_TASK_KEY: &str = "submit_task";
+const FETCH_TASK_KEY: &str = "fetch_task";
+
pub const CONFIG_SOURCE: ConfigSource =
ConfigSource::new("libeufin", "libeufin-nexus", "libeufin-nexus");
@@ -210,14 +222,14 @@ pub struct Args {
}
/** Load client private keys at or create new ones if missing */
-pub fn load_or_generate_client_keys(path: &Path) -> anyhow::Result<ClientPriKeysFile> {
+pub fn load_or_generate_client_keys(path: &Path) -> anyhow::Result<ClientKeys> {
// If exists load from disk
let current = load_client_keys(path)?;
if let Some(current) = current {
return Ok(current);
}
// Else create new keys
- let new = ClientPriKeysFile::generate()?;
+ let new = ClientKeys::generate()?;
persist_client_keys(&new, path)?;
info!(
"New client private keys created at '{}'",
@@ -229,8 +241,9 @@ pub fn load_or_generate_client_keys(path: &Path) -> anyhow::Result<ClientPriKeys
pub async fn ebics_setup(
ebics: &EbicsClient,
cfg: &EbicsKeysCfg,
- force_keys_submissions: bool,
+ force_keys_submission: bool,
auto_accept_keys: bool,
+ generate_registration_pdf: bool,
) -> anyhow::Result<()> {
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())?;
@@ -257,14 +270,14 @@ pub async fn ebics_setup(
// Privs exist. Upload their pubs
let keys_not_sub = !client.submitted_ini;
- if !client.submitted_ini || force_keys_submissions {
+ if !client.submitted_ini || force_keys_submission {
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 {
+ if !client.submitted_hia || force_keys_submission {
ebics
.submit_client_keys(cfg, &mut client, Order::HIA)
.await?;
@@ -314,8 +327,8 @@ pub async fn ebics_setup(
pub async fn ebics_submit(
ebics: &EbicsClient,
cfg: &NexusCfg,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
db: &PgPool,
transient: bool,
) -> anyhow::Result<()> {
@@ -455,10 +468,14 @@ async fn register_camt(db: &PgPool, cfg: &NexusCfg, xml: &[u8]) -> anyhow::Resul
pub async fn ebics_fetch(
ebics: &EbicsClient,
cfg: &NexusCfg,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
db: &PgPool,
documents: Option<&[OrderDoc]>,
+ pinned_start: &Option<Timestamp>,
+ peek: bool,
+ transient: bool,
+ transient_checkpoint: bool,
) -> anyhow::Result<()> {
let ebics_cfg = cfg.ebics()?;
@@ -576,7 +593,7 @@ pub async fn ebics_fetch(
}
}
};
- let fetch = async |orders: &[Order]| -> anyhow::Result<bool> {
+ let fetch = async |orders: &[Order], since: Option<Timestamp>| -> anyhow::Result<bool> {
let mut grouped_orders = BTreeMap::new();
for order in orders {
@@ -591,11 +608,19 @@ pub async fn ebics_fetch(
if let Some(doc) = doc {
for order in orders {
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()))
- })
+ .download(
+ db,
+ client,
+ bank,
+ order,
+ &since.map(|it| (it, Timestamp::now())),
+ transient && peek,
+ async |content| {
+ register_payload(&doc, content)
+ .await
+ .map_err(|e| EbicsErrKind::Custom(e.to_string().into()))
+ },
+ )
.await
{
if let EbicsErrKind::Code { bank, .. } = e.kind {
@@ -626,30 +651,162 @@ pub async fn ebics_fetch(
.flat_map(|it| ebics_cfg.dialect.standard().downloads(it))
.collect();
- let last_fetch = Timestamp::UNIX_EPOCH;
+ let fetch_cfg = cfg.fetch()?;
- // TODO loop
+ let (sender, mut receiver) = tokio::sync::mpsc::channel::<Vec<Order>>(10);
- let now = Timestamp::now();
+ let fetch = async {
+ if transient {
+ info!(target: "ebics-fetch", "Transient mode: fetching once and returning");
+ } else {
+ info!(target: "ebics-fetch", "Running with a frequency of {}", fetch_cfg.frequency_raw);
+ }
- info!(target: "ebics-fetch", "Running at frequency");
+ // TODO loop
+
+ let mut last_fetch = Timestamp::UNIX_EPOCH;
+ loop {
+ let now = Timestamp::now();
+ let checkpoint = get_task_status(db, CHECKPOINT_KEY)
+ .await?
+ .unwrap_or_default();
+ let next_fetch = last_fetch + fetch_cfg.frequency;
+ let next_checkpoint = {
+ if let Some(last_trial) = checkpoint.last_trial {
+ // We run today at checkpointTime
+ let checkpoint_date = Zoned::new(now, TimeZone::UTC)
+ .with()
+ .time(fetch_cfg.checkpoint_time)
+ .build()
+ .unwrap();
+ // If we already ran today we ran tomorrow
+ if last_trial > checkpoint_date.timestamp() {
+ checkpoint_date.tomorrow().unwrap().timestamp()
+ } else {
+ checkpoint_date.timestamp()
+ }
+ } else {
+ // We never ran, we must checkpoint now
+ now
+ }
+ };
+
+ let mut success = true;
+ if
+ // Run transient checkpoint at request
+ (transient && transient_checkpoint)
+ // Or run recurrent checkpoint
+ || (!transient && now > next_checkpoint)
+ {
+ info!(target: "ebics-fetch", "Running checkpoint");
+
+ let since = if let Some(pinned_start) = pinned_start
+ && transient
+ && checkpoint
+ .last_successfull
+ .map(|it| *pinned_start <= it)
+ .unwrap_or(true)
+ {
+ Some(*pinned_start)
+ } else {
+ checkpoint.last_successfull
+ };
+ let res = async {
+ // We fetch HKD to only fetch supported EBICS orders and get the document versions
+ let hkd = ebics.hkd(db, client, bank, false).await?;
+ let mut supported_orders = hkd
+ .partner
+ .orders
+ .into_iter()
+ .map(|it| it.order)
+ .collect::<Vec<_>>();
+ debug!(
+ "HKD: {}",
+ std::fmt::from_fn(|f| f.write_str(
+ &supported_orders
+ .iter()
+ .map(|it| it.to_string())
+ .collect::<Vec<_>>()
+ .join(",")
+ ))
+ );
+ supported_orders
+ .retain(|order| orders.iter().find(|it| order.eq(it)).is_some());
+ fetch(&supported_orders, since).await
+ }
+ .await;
+ if let Err(e) = res {
+ success = false;
+ error!(target: "ebics-fetch", "{e}");
+ }
+ try_join!(
+ update_task_status(db, CHECKPOINT_KEY, &now, success),
+ update_task_status(db, FETCH_TASK_KEY, &now, success)
+ )?;
+ last_fetch = now;
+ } else if transient || now > next_fetch {
+ if !transient {
+ info!(target: "ebics-fetch", "Running at frequency");
+ }
+ let res = async {
+ // We fetch HAA to only fetch pending & supported EBICS orders and get the document versions
+ let mut haa = ebics.haa(db, client, bank, false).await?;
+ debug!(
+ "HAA: {}",
+ std::fmt::from_fn(|f| f.write_str(
+ &haa.orders
+ .iter()
+ .map(|it| it.to_string())
+ .collect::<Vec<_>>()
+ .join(",")
+ ))
+ );
+ haa.orders
+ .retain(|order| orders.iter().find(|it| order.eq(it)).is_some());
+ fetch(&haa.orders, *pinned_start).await
+ }
+ .await;
+ if let Err(e) = res {
+ success = false;
+ error!(target: "ebics-fetch", "{e}");
+ }
+ update_task_status(db, FETCH_TASK_KEY, &now, success).await?;
+ last_fetch = now;
+ }
- let mut haa = ebics.haa(db, client, bank, false).await?;
- debug!(
- "HAA: {}",
- std::fmt::from_fn(|f| f.write_str(
- &haa.orders
- .iter()
- .map(|it| it.to_string())
- .collect::<Vec<_>>()
- .join(",")
- ))
- );
+ if transient {
+ if success {
+ return anyhow::Ok(());
+ } else {
+ return Err(anyhow!("ebics-fetch failed"));
+ }
+ }
+
+ let delay = now.duration_until(next_fetch.min(next_checkpoint));
+ let tx = timeout(
+ Duration::from_millis(delay.abs().as_millis() as u64),
+ receiver.recv(),
+ )
+ .await;
+ if let Ok(Some(mut notification)) = tx {
+ notification.retain(|order| orders.iter().find(|it| order.eq(it)).is_some());
+ if !notification.is_empty() {
+ info!(target: "ebics-fetch", "Running at real-time notifications reception");
+ fetch(¬ification, None).await?;
+ }
+ }
+ }
+ };
+
+ if transient {
+ fetch.await?;
+ } else {
+ tokio::try_join!(fetch, async {
+ listen_for_notification(ebics, db, client, bank, sender).await;
+ Ok(())
+ })?;
+ }
- haa.orders
- .retain(|order| orders.iter().find(|it| order.eq(it)).is_some());
- fetch(&haa.orders).await?;
- // TODO notification
Ok(())
}
@@ -671,6 +828,7 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> {
cfg.keys()?,
force_keys_resubmission,
auto_accept_keys,
+ generate_registration_pdf,
)
.await?;
}
@@ -678,22 +836,36 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> {
pinned_start,
peek,
checkpoint,
- ebics,
+ ebics: EbicsArgs { logs, transient },
} => {
let pool = pool(&cfg).await?;
let cfg = NexusCfg::parse(cfg)?;
let key_cfg = cfg.keys()?;
- let ebics = EbicsClient::new(&cfg, ebics.logs)?;
+ let ebics = EbicsClient::new(&cfg, logs)?;
let (client, bank) = expect_full_keys(key_cfg)?;
- ebics_fetch(&ebics, &cfg, &client, &bank, &pool, None).await?
+ ebics_fetch(
+ &ebics,
+ &cfg,
+ &client,
+ &bank,
+ &pool,
+ None,
+ &pinned_start.map(|it| date_to_utc_ts(&it)),
+ peek,
+ transient,
+ transient && checkpoint,
+ )
+ .await?
}
- Cmd::EbicsSubmit { ebics } => {
+ Cmd::EbicsSubmit {
+ ebics: EbicsArgs { logs, transient },
+ } => {
let pool = pool(&cfg).await?;
let cfg = NexusCfg::parse(cfg)?;
- let ebics = EbicsClient::new(&cfg, ebics.logs)?;
+ let ebics = EbicsClient::new(&cfg, logs)?;
let key_cfg = cfg.keys()?;
let (client, bank) = expect_full_keys(key_cfg)?;
- ebics_submit(&ebics, &cfg, &client, &bank, &pool, true).await?
+ ebics_submit(&ebics, &cfg, &client, &bank, &pool, transient).await?
}
Cmd::InitiatePayment {
amount,
@@ -748,277 +920,21 @@ pub async fn run(cfg: Config, cmd: Cmd) -> anyhow::Result<()> {
Ok(())
}
-#[cfg(test)]
-pub mod test {
- use std::{str::FromStr as _, sync::LazyLock};
-
- use compact_str::CompactString;
- use jiff::Timestamp;
- use sqlx::PgPool;
- use taler_api::subject::{fmt_in_subject, fmt_out_subject, subject_fmt_qr_bill};
- use taler_common::{
- api_common::{EddsaPublicKey, EddsaSignature},
- db::IncomingType,
- types::{
- amount::{Amount, Currency},
- base32::Base32,
- payto::{IbanPayto, PaytoURI, payto},
- },
- };
- use url::Url;
-
- use crate::{
- config::{AccountType, NexusIngestCfg},
- db::{
- initiated::{PaymentInitiationResult, initiate},
- transfer::{RegistrationResult, transfer_register},
- },
- model::{InId, InTx, Initiated, OutId, OutTx},
- rand_ebics_id,
- worker::{register_incoming, register_outgoing},
- };
-
- pub const CURR: Currency = Currency::KUDOS;
- pub static ACCOUNT: LazyLock<PaytoURI> =
- LazyLock::new(|| payto("payto://iban/CH4189144589712575493?receiver-name=Test"));
-
- /** Generates an outgoing payment, given its subject */
- pub fn gen_out_pay(subject: impl Into<String>) -> OutTx {
- OutTx {
- id: OutId {
- msg_id: None,
- e2e_id: Some(rand_ebics_id()),
- sref: None,
- },
- amount: Amount::new(&CURR, 44, 0),
- debit_fee: Amount::zero(&CURR),
- creditor: Some(
- IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test")
- .unwrap()
- .as_payto(),
- ),
- subject: Some(subject.into()),
- execution_time: Timestamp::now(),
- }
- }
-
- /** Generates a payment initiation, given its subject and end-to-end ID */
- pub fn gen_init_pay(
- end_to_end_id: impl Into<CompactString>,
- subject: impl Into<String>,
- ) -> Initiated {
- Initiated {
- id: 0,
- amount: Amount::new(&CURR, 44, 0),
- creditor: IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test")
- .unwrap()
- .as_payto(),
- subject: subject.into(),
- initiation_time: Timestamp::now(),
- e2e_id: end_to_end_id.into(),
- }
- }
-
- /** Generates an incoming payment, given its subject */
- 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),
- credit_fee: Amount::zero(&CURR),
- debtor: Some(
- IbanPayto::from_str("payto://iban/DE84500105177118117964?receiver-name=John+Smith")
- .unwrap()
- .as_payto(),
- ),
- subject: Some(subject.into()),
- execution_time: Timestamp::now(),
- }
- }
-
- pub async fn gen_initiate(
- db: &PgPool,
- end_to_end_id: impl Into<CompactString>,
- subject: impl Into<String>,
- ) -> PaymentInitiationResult {
- let init = gen_init_pay(end_to_end_id, subject);
- initiate(
- &db,
- &init.amount,
- &init.subject,
- &init.creditor,
- &init.initiation_time,
- &init.e2e_id,
- )
- .await
- .unwrap()
- }
-
- const CFG: NexusIngestCfg = NexusIngestCfg::simple(AccountType::Exchange, &CURR);
-
- async fn prepare(db: &PgPool) -> String {
- let key = EddsaPublicKey::rand();
- let sig = EddsaSignature::rand();
- let reference_number = subject_fmt_qr_bill(key.as_ref());
- assert_eq!(
- RegistrationResult::Success,
- transfer_register(
- db,
- IncomingType::reserve,
- &key,
- &key,
- &sig,
- false,
- &reference_number,
- &Timestamp::now()
- )
- .await
- .unwrap()
- );
- return reference_number;
- }
-
- /// Register a talerable reserve prepared incoming transaction
- pub async fn prepared_in(db: &PgPool) {
- let ref_nb = prepare(db).await;
- register_incoming(db, &CFG, &gen_in_pay(ref_nb))
- .await
- .unwrap();
- }
-
- /// Register an incomplete talerable reserve prepared incoming transaction
- pub async fn prepared_incomplete_in(db: &PgPool) {
- let ref_nb = prepare(db).await;
- let incomplete = InTx {
- subject: None,
- debtor: None,
- ..gen_in_pay(ref_nb)
- };
- register_incoming(db, &CFG, &incomplete).await.unwrap();
- }
-
- /// Register a completed talerable reserve prepared incoming transaction
- pub async fn prepared_completeted_in(db: &PgPool) {
- let ref_nb = prepare(db).await;
- let original = gen_in_pay(ref_nb);
- let incomplete = InTx {
- subject: None,
- debtor: None,
- ..original.clone()
- };
- register_incoming(db, &CFG, &incomplete).await.unwrap();
- register_incoming(db, &CFG, &original).await.unwrap();
- }
-
- /// Register a talerable reserve incoming transaction
- pub async fn talerable_in(db: &PgPool) {
- register_incoming(
- db,
- &CFG,
- &gen_in_pay(fmt_in_subject(
- IncomingType::reserve,
- &EddsaPublicKey::rand(),
- )),
- )
- .await
- .unwrap();
- }
-
- /// Register a talerable kyc incoming transaction
- pub async fn talerable_kyc_in(db: &PgPool) {
- register_incoming(
- db,
- &CFG,
- &gen_in_pay(fmt_in_subject(IncomingType::kyc, &EddsaPublicKey::rand())),
- )
- .await
- .unwrap();
- }
-
- /// Register an incomplete talerable reserve incoming transaction
- pub async fn talerable_incomplete_in(db: &PgPool) {
- let incomplete = InTx {
- subject: None,
- debtor: None,
- ..gen_in_pay(fmt_in_subject(
- IncomingType::reserve,
- &EddsaPublicKey::rand(),
- ))
- };
- register_incoming(db, &CFG, &incomplete).await.unwrap();
- }
-
- /// Register a completed talerable reserve incoming transaction
- pub async fn talerable_completeted_in(db: &PgPool) {
- let original = gen_in_pay(fmt_in_subject(
- IncomingType::reserve,
- &EddsaPublicKey::rand(),
- ));
- let incomplete = InTx {
- subject: None,
- debtor: None,
- ..original.clone()
- };
- register_incoming(db, &CFG, &incomplete).await.unwrap();
- register_incoming(db, &CFG, &original).await.unwrap();
- }
-
- /// Register incoming malformed transaction
- pub async fn malformed_in(db: &PgPool) {
- register_incoming(db, &CFG, &gen_in_pay("ignored"))
- .await
- .unwrap();
- }
-
- /// Register incoming incomplete malformed incoming transaction
- pub async fn malformed_incomplete_in(db: &PgPool) {
- let incomplete = InTx {
- subject: None,
- debtor: None,
- ..gen_in_pay("ignored")
- };
- register_incoming(db, &CFG, &incomplete).await.unwrap();
- }
-
- /// Register incoming completed malformed transaction
- pub async fn malformed_completeted_in(db: &PgPool) {
- let original = gen_in_pay("ignored");
- let incomplete = InTx {
- subject: None,
- debtor: None,
- ..original.clone()
- };
- register_incoming(db, &CFG, &incomplete).await.unwrap();
- register_incoming(db, &CFG, &original).await.unwrap();
- }
-
- /** Register an outgoing transaction */
- pub async fn malformed_out(db: &PgPool) {
- register_outgoing(db, &gen_out_pay("ignored"))
- .await
- .unwrap();
- }
+#[derive(Debug, Serialize, Deserialize, Clone, Default)]
+pub struct TaskStatus {
+ #[serde(serialize_with = "ser_micros", deserialize_with = "de_micros", default)]
+ pub last_successfull: Option<Timestamp>,
+ #[serde(serialize_with = "ser_micros", deserialize_with = "de_micros", default)]
+ pub last_trial: Option<Timestamp>,
+}
- /** Register an incomplete outgoing transaction */
- pub async fn incomplete_out(db: &PgPool) {
- let incomplete = OutTx {
- subject: None,
- creditor: None,
- ..gen_out_pay("ignored")
- };
- register_outgoing(db, &incomplete).await.unwrap();
- }
+fn ser_micros<S: Serializer>(key: &Option<Timestamp>, serializer: S) -> Result<S::Ok, S::Error> {
+ key.map(|it| it.as_microsecond()).serialize(serializer)
+}
- /// Register outgoing talerable transaction
- pub async fn talerable_out(db: &PgPool) {
- register_outgoing(
- db,
- &gen_out_pay(fmt_out_subject(
- &Base32::rand(),
- &Url::from_str("https://exchange.test.com").unwrap(),
- None,
- )),
- )
- .await
- .unwrap();
- }
+fn de_micros<'de, D: Deserializer<'de>>(deserializer: D) -> Result<Option<Timestamp>, D::Error> {
+ Option::<i64>::deserialize(deserializer)?
+ .map(Timestamp::from_microsecond)
+ .transpose()
+ .map_err(|e| serde::de::Error::custom(e.to_string()))
}
diff --git a/src/test.rs b/src/test.rs
@@ -0,0 +1,568 @@
+/*
+* 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::{str::FromStr as _, sync::LazyLock};
+
+use compact_str::CompactString;
+use jiff::Timestamp;
+use sqlx::PgPool;
+use taler_api::subject::{fmt_in_subject, fmt_out_subject, subject_fmt_qr_bill};
+use taler_common::{
+ api_common::{EddsaPublicKey, EddsaSignature},
+ db::IncomingType,
+ types::{
+ amount::{Amount, Currency},
+ base32::Base32,
+ payto::{IbanPayto, PaytoURI, payto},
+ },
+};
+use url::Url;
+
+use crate::{
+ config::{AccountType, NexusIngestCfg},
+ db::{
+ initiated::{PaymentInitiationResult, initiate},
+ transfer::{RegistrationResult, transfer_register},
+ },
+ model::{InId, InTx, Initiated, OutId, OutTx},
+ rand_ebics_id,
+ worker::{register_incoming, register_outgoing},
+};
+
+pub const CURR: Currency = Currency::KUDOS;
+pub static ACCOUNT: LazyLock<PaytoURI> =
+ LazyLock::new(|| payto("payto://iban/CH4189144589712575493?receiver-name=Test"));
+
+/** Generates an outgoing payment, given its subject */
+pub fn gen_out_pay(subject: impl Into<String>) -> OutTx {
+ OutTx {
+ id: OutId {
+ msg_id: None,
+ e2e_id: Some(rand_ebics_id()),
+ sref: None,
+ },
+ amount: Amount::new(&CURR, 44, 0),
+ debit_fee: Amount::zero(&CURR),
+ creditor: Some(
+ IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test")
+ .unwrap()
+ .as_payto(),
+ ),
+ subject: Some(subject.into()),
+ execution_time: Timestamp::now(),
+ }
+}
+
+/** Generates a payment initiation, given its subject and end-to-end ID */
+pub fn gen_init_pay(
+ end_to_end_id: impl Into<CompactString>,
+ subject: impl Into<String>,
+) -> Initiated {
+ Initiated {
+ id: 0,
+ amount: Amount::new(&CURR, 44, 0),
+ creditor: IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test")
+ .unwrap()
+ .as_payto(),
+ subject: subject.into(),
+ initiation_time: Timestamp::now(),
+ e2e_id: end_to_end_id.into(),
+ }
+}
+
+/** Generates an incoming payment, given its subject */
+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),
+ credit_fee: Amount::zero(&CURR),
+ debtor: Some(
+ IbanPayto::from_str("payto://iban/DE84500105177118117964?receiver-name=John+Smith")
+ .unwrap()
+ .as_payto(),
+ ),
+ subject: Some(subject.into()),
+ execution_time: Timestamp::now(),
+ }
+}
+
+pub async fn gen_initiate(
+ db: &PgPool,
+ end_to_end_id: impl Into<CompactString>,
+ subject: impl Into<String>,
+) -> PaymentInitiationResult {
+ let init = gen_init_pay(end_to_end_id, subject);
+ initiate(
+ &db,
+ &init.amount,
+ &init.subject,
+ &init.creditor,
+ &init.initiation_time,
+ &init.e2e_id,
+ )
+ .await
+ .unwrap()
+}
+
+const CFG: NexusIngestCfg = NexusIngestCfg::simple(AccountType::Exchange, &CURR);
+
+async fn prepare(db: &PgPool) -> String {
+ let key = EddsaPublicKey::rand();
+ let sig = EddsaSignature::rand();
+ let reference_number = subject_fmt_qr_bill(key.as_ref());
+ assert_eq!(
+ RegistrationResult::Success,
+ transfer_register(
+ db,
+ IncomingType::reserve,
+ &key,
+ &key,
+ &sig,
+ false,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap()
+ );
+ return reference_number;
+}
+
+/// Register a talerable reserve prepared incoming transaction
+pub async fn prepared_in(db: &PgPool) {
+ let ref_nb = prepare(db).await;
+ register_incoming(db, &CFG, &gen_in_pay(ref_nb))
+ .await
+ .unwrap();
+}
+
+/// Register an incomplete talerable reserve prepared incoming transaction
+pub async fn prepared_incomplete_in(db: &PgPool) {
+ let ref_nb = prepare(db).await;
+ let incomplete = InTx {
+ subject: None,
+ debtor: None,
+ ..gen_in_pay(ref_nb)
+ };
+ register_incoming(db, &CFG, &incomplete).await.unwrap();
+}
+
+/// Register a completed talerable reserve prepared incoming transaction
+pub async fn prepared_completeted_in(db: &PgPool) {
+ let ref_nb = prepare(db).await;
+ let original = gen_in_pay(ref_nb);
+ let incomplete = InTx {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+ register_incoming(db, &CFG, &incomplete).await.unwrap();
+ register_incoming(db, &CFG, &original).await.unwrap();
+}
+
+/// Register a talerable reserve incoming transaction
+pub async fn talerable_in(db: &PgPool) {
+ register_incoming(
+ db,
+ &CFG,
+ &gen_in_pay(fmt_in_subject(
+ IncomingType::reserve,
+ &EddsaPublicKey::rand(),
+ )),
+ )
+ .await
+ .unwrap();
+}
+
+/// Register a talerable kyc incoming transaction
+pub async fn talerable_kyc_in(db: &PgPool) {
+ register_incoming(
+ db,
+ &CFG,
+ &gen_in_pay(fmt_in_subject(IncomingType::kyc, &EddsaPublicKey::rand())),
+ )
+ .await
+ .unwrap();
+}
+
+/// Register an incomplete talerable reserve incoming transaction
+pub async fn talerable_incomplete_in(db: &PgPool) {
+ let incomplete = InTx {
+ subject: None,
+ debtor: None,
+ ..gen_in_pay(fmt_in_subject(
+ IncomingType::reserve,
+ &EddsaPublicKey::rand(),
+ ))
+ };
+ register_incoming(db, &CFG, &incomplete).await.unwrap();
+}
+
+/// Register a completed talerable reserve incoming transaction
+pub async fn talerable_completeted_in(db: &PgPool) {
+ let original = gen_in_pay(fmt_in_subject(
+ IncomingType::reserve,
+ &EddsaPublicKey::rand(),
+ ));
+ let incomplete = InTx {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+ register_incoming(db, &CFG, &incomplete).await.unwrap();
+ register_incoming(db, &CFG, &original).await.unwrap();
+}
+
+/// Register incoming malformed transaction
+pub async fn malformed_in(db: &PgPool) {
+ register_incoming(db, &CFG, &gen_in_pay("ignored"))
+ .await
+ .unwrap();
+}
+
+/// Register incoming incomplete malformed incoming transaction
+pub async fn malformed_incomplete_in(db: &PgPool) {
+ let incomplete = InTx {
+ subject: None,
+ debtor: None,
+ ..gen_in_pay("ignored")
+ };
+ register_incoming(db, &CFG, &incomplete).await.unwrap();
+}
+
+/// Register incoming completed malformed transaction
+pub async fn malformed_completeted_in(db: &PgPool) {
+ let original = gen_in_pay("ignored");
+ let incomplete = InTx {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+ register_incoming(db, &CFG, &incomplete).await.unwrap();
+ register_incoming(db, &CFG, &original).await.unwrap();
+}
+
+/** Register an outgoing transaction */
+pub async fn malformed_out(db: &PgPool) {
+ register_outgoing(db, &gen_out_pay("ignored"))
+ .await
+ .unwrap();
+}
+
+/** Register an incomplete outgoing transaction */
+pub async fn incomplete_out(db: &PgPool) {
+ let incomplete = OutTx {
+ subject: None,
+ creditor: None,
+ ..gen_out_pay("ignored")
+ };
+ register_outgoing(db, &incomplete).await.unwrap();
+}
+
+/// Register outgoing talerable transaction
+pub async fn talerable_out(db: &PgPool) {
+ register_outgoing(
+ db,
+ &gen_out_pay(fmt_out_subject(
+ &Base32::rand(),
+ &Url::from_str("https://exchange.test.com").unwrap(),
+ None,
+ )),
+ )
+ .await
+ .unwrap();
+}
+
+mod ebics {
+ use std::{
+ fs::Permissions,
+ os::unix::fs::PermissionsExt as _,
+ sync::{Arc, Mutex},
+ time::Duration,
+ };
+
+ use axum::{body::Bytes, response::IntoResponse, routing::post};
+ use clap::Parser as _;
+ use reqwest::StatusCode;
+ use sqlx::{ConnectOptions, PgPool};
+ use taler_api::{Serve, api::TalerRouter as _};
+ use taler_common::config::Config;
+ use taler_test_utils::setup_tracing;
+ use tempfile::{TempDir, tempdir};
+ use tokio::net::UnixStream;
+
+ use crate::{
+ Args, CHECKPOINT_KEY, CONFIG_SOURCE,
+ db::test::db_setup,
+ ebics::test::{EbicsRes, EbicsState, Sequence},
+ run,
+ };
+
+ pub async fn nexus_cmd(cfg: &Config, cmd: &str) -> anyhow::Result<()> {
+ let parts = shlex::split(cmd).unwrap();
+ let args = std::iter::once("libeufin_nexus").chain(parts.iter().map(|it| it.as_str()));
+
+ let cmd = Args::try_parse_from(args).unwrap();
+ run(cfg.clone(), cmd.cmd).await
+ }
+
+ struct EbicsTestBank {
+ pub dir: TempDir,
+ pub sock_path: String,
+ pub sequence: Arc<Mutex<Vec<Sequence>>>,
+ }
+
+ impl EbicsTestBank {
+ pub async fn new() -> Self {
+ setup_tracing();
+ let dir = tempdir().unwrap();
+ let sock_path = dir.path().join("bank.sock").to_str().unwrap().to_string();
+ let sequence = Arc::new(Mutex::new(Vec::new()));
+ let server_sequence = sequence.clone();
+ let bank = Arc::new(Mutex::new(EbicsState::new()));
+ let server = axum::Router::new()
+ .route(
+ "/",
+ post(async move |body: Bytes| {
+ let sequence: Sequence = server_sequence.lock().unwrap().pop().unwrap();
+ let mut bank = bank.lock().unwrap();
+ let res = sequence(&mut *bank, &body);
+ match res {
+ EbicsRes::Ok(xml) => xml.into_response(),
+ EbicsRes::BadRequest => StatusCode::BAD_REQUEST.into_response(),
+ EbicsRes::Failure => StatusCode::SERVICE_UNAVAILABLE.into_response(),
+ }
+ }),
+ )
+ .serve(
+ Serve::Unix {
+ path: sock_path.clone(),
+ permission: Permissions::from_mode(660),
+ },
+ None,
+ );
+ tokio::spawn(server);
+ // Wait for server to start
+ for _ in 0..100 {
+ if UnixStream::connect(&sock_path).await.is_ok() {
+ break;
+ }
+ tokio::time::sleep(Duration::from_millis(10)).await;
+ }
+ Self {
+ dir,
+ sock_path,
+ sequence,
+ }
+ }
+
+ pub fn sequences(&self, sequences: &[Sequence]) {
+ let mut state = self.sequence.lock().unwrap();
+ assert_eq!(state.len(), 0);
+ state.extend(sequences.into_iter().rev());
+ }
+ }
+
+ impl Drop for EbicsTestBank {
+ fn drop(&mut self) {
+ assert_eq!(self.sequence.lock().unwrap().len(), 0);
+ }
+ }
+
+ async fn test_setup() -> (EbicsTestBank, Config, PgPool) {
+ let (_, db) = db_setup().await;
+ let test = EbicsTestBank::new().await;
+ let cfg = Config::from_mem_with_env(
+ CONFIG_SOURCE,
+ &format!(
+ "
+ [paths]
+ LIBEUFIN_NEXUS_HOME = {:?}
+
+ {}
+
+ [nexus-ebics]
+ UNIXPATH = {}
+
+ [libeufin-nexusdb-postgres]
+ CONFIG = postgresql:///{}
+ ",
+ test.dir.path(),
+ include_str!("../testbench/conf/mini.conf"),
+ test.sock_path,
+ db.connect_options().get_database().unwrap()
+ ),
+ )
+ .unwrap();
+ test.sequences(&[
+ EbicsState::hev,
+ EbicsState::ini,
+ EbicsState::hia,
+ EbicsState::hpb,
+ ]);
+ nexus_cmd(&cfg, "ebics-setup --auto-accept-keys")
+ .await
+ .unwrap();
+
+ (test, cfg, db)
+ }
+
+ #[tokio::test]
+ async fn setup() {
+ test_setup().await;
+ }
+
+ #[tokio::test]
+ async fn fetch_pinned_date() {
+ let (test, cfg, db) = test_setup().await;
+
+ let reset_checkpoint = async || {
+ let res = sqlx::query("DELETE FROM kv WHERE key=$1")
+ .bind(CHECKPOINT_KEY)
+ .execute(&db)
+ .await
+ .unwrap();
+ assert_eq!(res.rows_affected(), 1);
+ };
+
+ // Default transient
+ test.sequences(&[
+ EbicsState::haa,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient").await.unwrap();
+
+ // Pinned transient
+ test.sequences(&[
+ EbicsState::haa,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data_pinned,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient --pinned-start 2024-06-05")
+ .await
+ .unwrap();
+
+ // Init checkpoint
+ test.sequences(&[
+ EbicsState::hkd,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient --checkpoint")
+ .await
+ .unwrap();
+
+ // Default checkpoint
+ test.sequences(&[
+ EbicsState::hkd,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data_now,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient --checkpoint")
+ .await
+ .unwrap();
+
+ // Pinned checkpoint
+ test.sequences(&[
+ EbicsState::hkd,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data_pinned,
+ ]);
+ nexus_cmd(
+ &cfg,
+ "ebics-fetch --transient --checkpoint --pinned-start 2024-06-05",
+ )
+ .await
+ .unwrap();
+
+ // Reset checkpoint
+ reset_checkpoint().await;
+ test.sequences(&[
+ EbicsState::hkd,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient --checkpoint")
+ .await
+ .unwrap();
+
+ // Reset pinned checkpoint
+ reset_checkpoint().await;
+ test.sequences(&[
+ EbicsState::hkd,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data_pinned,
+ ]);
+ nexus_cmd(
+ &cfg,
+ "ebics-fetch --transient --checkpoint --pinned-start 2024-06-05",
+ )
+ .await
+ .unwrap();
+ }
+
+ #[tokio::test]
+ async fn close_pending_transaction() {
+ let (test, cfg, _) = test_setup().await;
+
+ // Failure before first segment
+ test.sequences(&[
+ // Failure to perform download
+ EbicsState::failure,
+ // Then continue
+ EbicsState::haa,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient")
+ .await
+ .unwrap_err();
+ nexus_cmd(&cfg, "ebics-fetch --transient").await.unwrap();
+
+ // Compliant server
+ test.sequences(&[
+ EbicsState::haa,
+ EbicsState::receipt_ok,
+ // Failure to perform download
+ EbicsState::init_tx,
+ EbicsState::failure,
+ // Retry fail once
+ EbicsState::failure,
+ // Retry fail twice
+ EbicsState::failure,
+ // Retry succeed
+ EbicsState::bad_request,
+ // Then continue
+ EbicsState::haa,
+ EbicsState::receipt_ok,
+ EbicsState::btd_no_data,
+ ]);
+ nexus_cmd(&cfg, "ebics-fetch --transient")
+ .await
+ .unwrap_err();
+ nexus_cmd(&cfg, "ebics-fetch --transient")
+ .await
+ .unwrap_err();
+ nexus_cmd(&cfg, "ebics-fetch --transient")
+ .await
+ .unwrap_err();
+ nexus_cmd(&cfg, "ebics-fetch --transient").await.unwrap();
+ }
+}
diff --git a/src/utils.rs b/src/utils.rs
@@ -20,7 +20,10 @@
use std::{fmt::Display, io::Write as _};
use base64::{display::Base64Display, prelude::BASE64_STANDARD};
-use flate2::{Compression, write::ZlibEncoder};
+use flate2::{
+ Compression,
+ write::{ZlibDecoder, ZlibEncoder},
+};
pub fn deflate(bytes: &[u8]) -> Vec<u8> {
let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
@@ -28,6 +31,12 @@ pub fn deflate(bytes: &[u8]) -> Vec<u8> {
encoder.finish().unwrap()
}
+pub fn inflate(bytes: &[u8]) -> Vec<u8> {
+ let mut encoder = ZlibDecoder::new(Vec::new());
+ encoder.write_all(bytes).unwrap();
+ encoder.finish().unwrap()
+}
+
pub fn hex_chunk_by_two<'a>(bytes: impl AsRef<[u8]> + 'a) -> impl Display + 'a {
std::fmt::from_fn(move |f| {
for b in bytes.as_ref() {
diff --git a/src/ws.rs b/src/ws.rs
@@ -35,7 +35,7 @@ use crate::{
ebics_code::EbicsReturnCode,
order::{BTF, Order},
},
- keys::{BankPubKeysFile, ClientPriKeysFile},
+ keys::{BankKeys, ClientKeys},
};
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)]
@@ -165,8 +165,8 @@ pub enum WssError {
pub async fn listen_for_notification(
ebics: &EbicsClient,
db: &PgPool,
- client: &ClientPriKeysFile,
- bank: &BankPubKeysFile,
+ client: &ClientKeys,
+ bank: &BankKeys,
sender: tokio::sync::mpsc::Sender<Vec<Order>>,
) {
let mut backoff = ExpoBackoffDecorr::new(Duration::from_secs(30), Duration::from_mins(30), 2.5);
diff --git a/src/xml.rs b/src/xml.rs
@@ -27,115 +27,125 @@ use roxmltree::{Document, Node};
#[macro_export]
macro_rules! xml {
+ // Trailing comma
+ ($w:ident => $(,)?) => {{}};
// Logic escape
- ($w:ident, @ $logic:expr$(, $($rest:tt)*)?) => {{
+ ($w:ident => @ $logic:expr$(, $($rest:tt)*)?) => {{
($logic)($w);
- $($crate::xml!($w, $($rest)*);)*
+ $($crate::xml!($w => $($rest)*);)*
}};
// Text element
- ($w:ident, $name:tt $(($k:tt: $v:tt))* : $content:expr $(, $($rest:tt)*)?) => {{
+ ($w:ident => $name:tt $($k:literal=$v:tt)* : $content:expr $(, $($rest:tt)*)?) => {{
$w.text(&$name, &[$((&$k, &$v)),*], &$content);
- $($crate::xml!($w, $($rest)*);)*
+ $($crate::xml!($w => $($rest)*);)*
}};
// Nested block
- ($w:ident, $name:tt $(($k:tt: $v:tt))* { $($body:tt)* }$(, $($rest:tt)*)?) => {{
+ ($w:ident => $name:tt $($k:literal=$v:tt)* { $($body:tt)* }$(, $($rest:tt)*)?) => {{
let name = &$name;
$w.open(&name, &[$((&$k, &$v)),*]);
- $crate::xml!($w, $($body)*);
+ $crate::xml!($w => $($body)*);
$w.close(&name);
- $($crate::xml!($w, $($rest)*);)*
+ $($crate::xml!($w => $($rest)*);)*
}};
// Empty element
- ($w:ident, $name:tt $(($k:tt: $v:tt))* $(, $($rest:tt)*)?) => {{
+ ($w:ident => $name:tt $($k:literal=$v:tt)* $(, $($rest:tt)*)?) => {{
$w.empty(&$name, &[$((&$k, &$v)),*]);
- $($crate::xml!($w, $($rest)*);)*
+ $($crate::xml!($w => $($rest)*);)*
+ }};
+ // Root builder
+ ($name:tt $($k:literal=$v:tt)* { $($body:tt)* }) => {{
+ let mut writer = $crate::xml::XmlWriter::init();
+ let w = &mut writer;
+ let name = &$name;
+ w.open(&name, &[$((&$k, &$v)),*]);
+ $crate::xml!(w => $($body)*);
+ w.close(&name);
+ writer.finish()
}};
- ($($xml:tt)*) => {
- {
- let mut writer = $crate::xml::XmlWriter::init();
- let w = &mut writer;
- $crate::xml!(w, $($xml)*);
- writer.finish()
- }
- };
}
pub struct XmlWriter {
- buff: String,
+ xml: String,
}
+
impl XmlWriter {
pub fn init() -> Self {
- let buff = r#"<?xml version="1.0" encoding="UTF-8" standalone="yes"?>"#.to_owned();
- Self { buff }
+ let mut xml = String::with_capacity(1024);
+ xml.push_str(r#"<?xml version="1.0" encoding="UTF-8" standalone="yes"?>"#);
+ Self { xml }
}
- pub fn open(&mut self, name: &dyn Display, attrs: &[(&dyn Display, &dyn Display)]) {
- self.buff.push('<');
+ pub fn open<N: Display>(&mut self, name: N, attrs: &[(&dyn Display, &dyn Display)]) {
+ self.xml.push('<');
self.write_tag_attrs(name, attrs);
- self.buff.push('>');
+ self.xml.push('>');
}
- pub fn close(&mut self, name: &dyn Display) {
- self.buff.push_str("</");
- self.buff.write_fmt(format_args!("{name}")).unwrap();
- self.buff.push('>');
+ pub fn close<N: Display>(&mut self, name: N) {
+ self.xml.push_str("</");
+ self.xml.write_fmt(format_args!("{name}")).unwrap();
+ self.xml.push('>');
}
- pub fn empty(&mut self, name: &dyn Display, attrs: &[(&dyn Display, &dyn Display)]) {
- self.buff.push('<');
+ pub fn empty<N: Display>(&mut self, name: N, attrs: &[(&dyn Display, &dyn Display)]) {
+ self.xml.push('<');
self.write_tag_attrs(name, attrs);
- self.buff.push_str("/>");
+ self.xml.push_str("/>");
}
- pub fn text(
+ pub fn text<N: Display, C: Display>(
&mut self,
- name: &dyn Display,
+ name: N,
attrs: &[(&dyn Display, &dyn Display)],
- content: &dyn Display,
+ content: C,
) {
- self.open(name, attrs);
+ self.open(&name, attrs);
self.write_escaped(content);
- self.close(name);
+ self.close(&name);
}
- fn write_tag_attrs(&mut self, name: &dyn Display, attrs: &[(&dyn Display, &dyn Display)]) {
- self.buff.write_fmt(format_args!("{name}")).unwrap();
+ fn write_tag_attrs<N: Display>(&mut self, name: N, attrs: &[(&dyn Display, &dyn Display)]) {
+ self.xml.write_fmt(format_args!("{name}")).unwrap();
for (key, value) in attrs {
- self.buff.push(' ');
- self.buff.write_fmt(format_args!("{key}")).unwrap();
- self.buff.push_str("=\"");
+ self.xml.push(' ');
+ self.xml.write_fmt(format_args!("{}", *key)).unwrap();
+ self.xml.push_str("=\"");
self.write_escaped(*value);
- self.buff.push('"');
+ self.xml.push('"');
}
}
- fn write_escaped(&mut self, content: &dyn Display) {
+ fn write_escaped<D: Display>(&mut self, content: D) {
std::fmt::write(self, format_args!("{content}")).unwrap();
}
pub fn finish(self) -> String {
- self.buff
+ self.xml
}
}
/// Write XML text content following XML escape rules
impl std::fmt::Write for XmlWriter {
fn write_str(&mut self, s: &str) -> std::fmt::Result {
- if s.contains(['<', '>', '&', '\'', '"']) {
- for c in s.chars() {
- match c {
- '<' => self.buff.push_str("<"),
- '>' => self.buff.push_str(">"),
- '&' => self.buff.push_str("&"),
- '\'' => self.buff.push_str("'"),
- '"' => self.buff.push_str("""),
- _ => self.buff.push(c),
- }
- }
- } else {
- self.buff.push_str(s);
+ // Single pass over bytes. For each special character, bulk-copy
+ // everything before it, then push the entity. No double-scan,
+ // no char-at-a-time pushing for clean runs.
+ let mut start = 0;
+ for (i, &b) in s.as_bytes().iter().enumerate() {
+ let entity = match b {
+ b'<' => "<",
+ b'>' => ">",
+ b'&' => "&",
+ b'\'' => "'",
+ b'"' => """,
+ _ => continue,
+ };
+ self.xml.push_str(&s[start..i]); // bulk copy of clean prefix
+ self.xml.push_str(entity);
+ start = i + 1;
}
+ self.xml.push_str(&s[start..]); // bulk copy of clean suffix
Ok(())
}
}
@@ -410,20 +420,19 @@ mod test {
#[test]
pub fn basic() {
assert_eq!(
- xml!("ebicsRequest" ("version": "H004") {
+ xml!("ebicsRequest" "version"="H004" {
"a" {
"b" {
- "c" ("attribute-of": "c") {
+ "c" "attribute-of"="c" {
"d" {
"e" {
- "f" ("nested": "true") {
+ "f" "nested"="true" {
"g" {
"h"
}
}
}
}
-
}
}
},
@@ -436,7 +445,7 @@ mod test {
#[test]
pub fn modularity() {
fn module(w: &mut XmlWriter) {
- xml!(w, "module");
+ xml!(w => "module");
}
assert_eq!(
xml!("root" { @ module }),
@@ -450,7 +459,7 @@ mod test {
xml!("iterable" {
"endOfDocument" {
@ |w: &mut XmlWriter| for i in 1..=10 {
- xml!(w, (format_args!("e{i}")) {
+ xml!(w => (format_args!("e{i}")) {
(format_args!("e{i}{i}")): (format_args!("{i}{i}{i}"))
})
}
diff --git a/testbench/conf/mini.conf b/testbench/conf/mini.conf
@@ -15,8 +15,8 @@ USER_ID = myuser
PARTNER_ID = myorg
# Account information
-IBAN = myiban
-BIC = mybic
+IBAN = CH7789144474425692816
+BIC = POFICHBEXXX
NAME = myname
[libeufin-bankdb-postgres]