libeufin

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

commit d6d7828fd9e9117796f62e772d8c3a917a394f7c
parent f54a8fdf2147a89ff4c7c6f52e2edd9026e41708
Author: Antoine A <>
Date:   Fri, 24 Apr 2026 10:37:57 +0200

PaymentInitiationsTest

Diffstat:
MCargo.lock | 127+++++++++++++++++++++++++++++--------------------------------------------------
MCargo.toml | 2+-
Msrc/config.rs | 101++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Msrc/db.rs | 813++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Msrc/lib.rs | 353+++++++++++++++----------------------------------------------------------------
Msrc/main.rs | 2+-
Asrc/model.rs | 349+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/xml.rs | 1-
8 files changed, 1338 insertions(+), 410 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -430,6 +430,27 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" [[package]] +name = "const_format" +version = "0.2.35" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7faa7469a93a566e9ccc1c73fe783b4a65c274c5ace346038dca9c39fe0030ad" +dependencies = [ + "const_format_proc_macros", + "konst", +] + +[[package]] +name = "const_format_proc_macros" +version = "0.2.34" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d57c2eccfb16dbac1f4e61e206105db5820c9d26c3c472bc17c774259ef7744" +dependencies = [ + "proc-macro2", + "quote", + "unicode-xid", +] + +[[package]] name = "convert_case" version = "0.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1429,6 +1450,21 @@ dependencies = [ ] [[package]] +name = "konst" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "330f0e13e6483b8c34885f7e6c9f19b1a7bd449c673fbb948a51c99d66ef74f4" +dependencies = [ + "konst_macro_rules", +] + +[[package]] +name = "konst_macro_rules" +version = "0.2.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4933f3f57a8e9d9da04db23fb153356ecaf00cbd14aee46279c33dc80925c37" + +[[package]] name = "lazy_static" version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1458,6 +1494,7 @@ dependencies = [ "base64", "clap", "compact_str", + "const_format", "flate2", "getrandom 0.4.2", "jiff", @@ -1483,7 +1520,6 @@ dependencies = [ "url", "uuid", "x509-parser", - "xml-canonicalization", ] [[package]] @@ -1783,49 +1819,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] -name = "pest" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e0848c601009d37dfa3430c4666e147e49cdcf1b92ecd3e63657d8a5f19da662" -dependencies = [ - "memchr", - "ucd-trie", -] - -[[package]] -name = "pest_derive" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "11f486f1ea21e6c10ed15d5a7c77165d0ee443402f0780849d1768e7d9d6fe77" -dependencies = [ - "pest", - "pest_generator", -] - -[[package]] -name = "pest_generator" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8040c4647b13b210a963c1ed407c1ff4fdfa01c31d6d2a098218702e6664f94f" -dependencies = [ - "pest", - "pest_meta", - "proc-macro2", - "quote", - "syn", -] - -[[package]] -name = "pest_meta" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "89815c69d36021a140146f26659a81d6c2afa33d216d736dd4be5381a7362220" -dependencies = [ - "pest", - "sha2", -] - -[[package]] name = "pin-project-lite" version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1878,9 +1871,9 @@ checksum = "c33a9471896f1c69cecef8d20cbe2f7accd12527ce60845ff44c153bb2a21b49" [[package]] name = "portable-atomic-util" -version = "0.2.5" +version = "0.2.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7a9db96d7fa8782dd8c15ce32ffe8680bbd1e978a43bf51a34d39483540495f5" +checksum = "091397be61a01d4be58e7841595bd4bfedb15f1cd54977d79b8271e94ed799a3" dependencies = [ "portable-atomic", ] @@ -1929,15 +1922,6 @@ dependencies = [ ] [[package]] -name = "quick-xml" -version = "0.37.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "331e97a1af0bf59823e6eadffe373d7b27f485be8748f71471c662c1f269b7fb" -dependencies = [ - "memchr", -] - -[[package]] name = "quinn" version = "0.11.9" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2975,7 +2959,7 @@ dependencies = [ [[package]] name = "taler-api" -version = "1.3.0" +version = "1.4.0" dependencies = [ "aws-lc-rs", "axum", @@ -2983,7 +2967,6 @@ dependencies = [ "compact_str", "dashmap", "http-body-util", - "itoa", "jiff", "listenfd", "serde", @@ -3000,11 +2983,11 @@ dependencies = [ [[package]] name = "taler-build" -version = "1.3.0" +version = "1.4.0" [[package]] name = "taler-common" -version = "1.3.0" +version = "1.4.0" dependencies = [ "anyhow", "aws-lc-rs", @@ -3014,6 +2997,7 @@ dependencies = [ "indexmap", "jiff", "rand 0.10.0", + "regex", "serde", "serde_json", "serde_path_to_error", @@ -3030,7 +3014,7 @@ dependencies = [ [[package]] name = "taler-test-utils" -version = "1.3.0" +version = "1.4.0" dependencies = [ "axum", "flate2", @@ -3153,9 +3137,9 @@ dependencies = [ [[package]] name = "tinyvec" -version = "1.10.0" +version = "1.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bfa5fdc3bce6191a1dbc8c02d5c8bffcf557bafa17c124c5264a458f1b0613fa" +checksum = "3e61e67053d25a4e82c844e8424039d9745781b3fc4f32b8d55ed50f5f667ef3" dependencies = [ "tinyvec_macros", ] @@ -3344,12 +3328,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "562d481066bde0658276a35467c4af00bdc6ee726305698a55b86e61d7ad82bb" [[package]] -name = "ucd-trie" -version = "0.1.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2896d95c02a80c6d6a5d6e953d479f5ddf2dfdb6a244441010e373ac0fb88971" - -[[package]] name = "unicase" version = "2.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -4186,19 +4164,6 @@ dependencies = [ ] [[package]] -name = "xml-canonicalization" -version = "0.1.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d9a3101284404fdfe80cc80be42e4c1b5642ac76aae9ef10dae2eaf7884f66e1" -dependencies = [ - "pest", - "pest_derive", - "quick-xml", - "regex", - "tracing", -] - -[[package]] name = "yasna" version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" diff --git a/Cargo.toml b/Cargo.toml @@ -9,7 +9,6 @@ tokio = { version = "*", features = ["macros", "rt-multi-thread"] } tracing = "*" thiserror = "*" roxmltree = "*" -xml-canonicalization = "*" base64 = "*" pem = "*" anyhow = "*" @@ -47,3 +46,4 @@ sqlx = { version = "0.8", default-features = false, features = [ compact_str = { version = "0.9.0", features = ["serde", "sqlx-postgres"] } uuid = { version = "1.0", features = ["v4", "fast-rng"] } regex = "*" +const_format = { version = "0.2", features = ["rust_1_83"] } diff --git a/src/config.rs b/src/config.rs @@ -17,14 +17,20 @@ * <http://www.gnu.org/licenses/> */ -use std::cell::OnceCell; +use std::{cell::OnceCell, str::FromStr, time::Duration}; -use jiff::Timestamp; +use jiff::{ + Timestamp, + civil::{Date, Time}, +}; use regex::Regex; use taler_api::config::DbCfg; use taler_common::{ config::{Config, ValueErr}, - types::amount::{Amount, Currency}, + types::{ + amount::{Amount, Currency}, + utils::date_to_utc_ts, + }, }; pub fn parse_db_cfg(cfg: &Config) -> Result<DbCfg, ValueErr> { @@ -38,10 +44,10 @@ pub struct EbicsKeysCfg { impl EbicsKeysCfg { pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { - let sect = cfg.section("nexus-ebics"); + let s = cfg.section("nexus-ebics"); Ok(Self { - bank_pub_keys_path: sect.path("bank_public_keys_file").require()?, - client_priv_keys_path: sect.path("client_private_keys_file").require()?, + bank_pub_keys_path: s.path("bank_public_keys_file").require()?, + client_priv_keys_path: s.path("client_private_keys_file").require()?, }) } } @@ -55,17 +61,18 @@ pub struct EbicsHostCfg { impl EbicsHostCfg { pub fn parse(cfg: &Config) -> Result<Self, ValueErr> { - let sect = cfg.section("nexus-ebics"); + let s = cfg.section("nexus-ebics"); Ok(Self { - base_url: sect.url("host_base_url").require()?, - host_id: sect.str("host_id").require()?, - user_id: sect.str("user_id").require()?, - partner_id: sect.str("partner_id").require()?, + base_url: s.url("host_base_url").require()?, + host_id: s.str("host_id").require()?, + user_id: s.str("user_id").require()?, + partner_id: s.str("partner_id").require()?, }) } } -#[derive(Debug, Clone, Copy)] +#[derive(Debug, Clone, Copy, strum_macros::EnumString)] +#[strum(serialize_all = "snake_case")] pub enum AccountType { Exchange, Normal, @@ -73,38 +80,81 @@ pub enum AccountType { pub struct NexusIngestConfig { pub account_type: AccountType, - pub ignore_transactions_before: Timestamp, + pub ignore_txs_before: Timestamp, pub ignore_bounces_before: Timestamp, pub restriction_payto_regex: Option<Regex>, pub bounce_deduce_fee: bool, pub bounce_fee: Amount, + pub currency: Currency, } impl NexusIngestConfig { pub fn simple(account_type: AccountType, currency: &Currency) -> Self { Self { account_type, - ignore_transactions_before: Timestamp::UNIX_EPOCH, + ignore_txs_before: Timestamp::UNIX_EPOCH, ignore_bounces_before: Timestamp::UNIX_EPOCH, restriction_payto_regex: None, bounce_deduce_fee: false, bounce_fee: Amount::zero(currency), + currency: Currency::from_str("KUDOS").unwrap(), } } } +pub struct NexusFetchConfig { + pub frequency: Duration, + pub frequency_raw: String, + pub checkpoint_time: Time, + pub ignore_txs_before: Timestamp, + pub ignore_bounces_before: Timestamp, + pub restriction_payto_regex: Option<Regex>, + pub bounce_deduce_fee: bool, + pub bounce_fee: Amount, +} + +impl NexusFetchConfig { + pub fn parse(cfg: &Config, currency: &Currency) -> Result<Self, ValueErr> { + let s = cfg.section("nexus-fetch"); + + Ok(Self { + frequency: s.duration("frequency").require()?, + frequency_raw: s.str("frequency").require()?, + checkpoint_time: s.time("checkpoint_time_of_day").require()?, + ignore_txs_before: date_to_utc_ts( + &s.date("ignore_transactions_before").default(Date::MIN)?, + ), + ignore_bounces_before: date_to_utc_ts( + &s.date("ignore_bounces_before").default(Date::MIN)?, + ), + restriction_payto_regex: s.regex("restriction_payto_regex").opt()?, + bounce_deduce_fee: s.boolean("bounce_deduce_fee").default(false)?, + bounce_fee: s + .amount("bounce_fee", currency) + .default(Amount::zero(currency))?, + }) + } +} + pub struct NexusCfg { pub cfg: Config, + pub currency: Currency, + pub account_type: AccountType, pub keys: OnceCell<EbicsKeysCfg>, pub host: OnceCell<EbicsHostCfg>, + pub fetch: OnceCell<NexusFetchConfig>, } impl NexusCfg { pub fn parse(cfg: Config) -> Result<Self, ValueErr> { + let s = cfg.section("nexus-ebics"); Ok(Self { + currency: s.currency("currency").require()?, + account_type: s.parse("account type", "ACCOUNT_TYPE").require()?, cfg, keys: OnceCell::new(), host: OnceCell::new(), + fetch: OnceCell::new(), }) } @@ -127,4 +177,27 @@ impl NexusCfg { self.host.set(host).ok(); Ok(self.host.get().unwrap()) } + + pub fn fetch(&self) -> Result<&NexusFetchConfig, ValueErr> { + // TODO use get_or_try_init when stable + if let Some(fetch) = self.fetch.get() { + return Ok(fetch); + } + let fetch = NexusFetchConfig::parse(&self.cfg, &self.currency)?; + self.fetch.set(fetch).ok(); + Ok(self.fetch.get().unwrap()) + } + + pub fn ingest(&self) -> Result<NexusIngestConfig, ValueErr> { + let fetch = self.fetch()?; + Ok(NexusIngestConfig { + account_type: self.account_type, + ignore_txs_before: fetch.ignore_txs_before, + ignore_bounces_before: fetch.ignore_bounces_before, + restriction_payto_regex: fetch.restriction_payto_regex.clone(), + bounce_deduce_fee: fetch.bounce_deduce_fee, + bounce_fee: fetch.bounce_fee.clone(), + currency: self.currency.clone(), + }) + } } diff --git a/src/db.rs b/src/db.rs @@ -14,7 +14,10 @@ TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> */ +use std::collections::BTreeMap; + use compact_str::CompactString; +use const_format::formatcp; use jiff::Timestamp; use sqlx::{PgPool, QueryBuilder, Row, postgres::PgRow}; use taler_api::{ @@ -30,9 +33,15 @@ use taler_common::{ }; use tokio::sync::watch::{Receiver, Sender}; -use crate::{IncomingPayment, InitiatedPayment, OutgoingId, OutgoingPayment, config::parse_db_cfg}; +use crate::{ + IncomingPayment, OutgoingPayment, + config::parse_db_cfg, + model::{InitiatedPayment, OutgoingId, PaymentBatch, SubmissionState}, +}; const SCHEMA: &str = "libeufin_nexus"; +const UNSETTLED: &str = "status NOT IN ('success', 'permanent_failure', 'late_failure')"; +const PENDING: &str = "status IN ('unsubmitted', 'pending')"; pub async fn pool(cfg: &Config) -> anyhow::Result<PgPool> { let db = parse_db_cfg(cfg)?; @@ -133,13 +142,84 @@ pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<()> { Ok(()) } +pub async fn initiated_submittable( + db: &PgPool, + currency: &Currency, +) -> sqlx::Result<Vec<PaymentBatch>> { + const SELECT_PART: &str = " + SELECT initiated_outgoing_batch_id, message_id, creation_date, sum + FROM initiated_outgoing_batches + "; + let mut tx = db.begin().await?; + // We want to maximize the number of successfully submitted batches in the event + // of a malformed transaction or a persistent error classified as transient. We send + // the unsubmitted batches first, starting with the oldest by creation time. + // This is the happy path, giving every batch a chance while being fair on the + // basis of creation date. + // Then we retry the failed batches, starting with the oldest by submission time. + // This the bad path retrying each failed batch applying a rotation based on + // resubmission time. + let mut batches = sqlx::query(formatcp!( + " + ({SELECT_PART} WHERE status='unsubmitted' ORDER BY creation_date ASC) + UNION ALL + ({SELECT_PART} WHERE status='transient_failure' ORDER BY submission_date) + " + )) + .try_map(|r: PgRow| { + Ok(PaymentBatch { + id: r.try_get_u64("initiated_outgoing_batch_id")?, + msg_id: r.try_get("message_id")?, + creation_date: r.try_get_timestamp("creation_date")?, + sum: r.try_get_amount("sum", currency)?, + payments: Vec::new(), + }) + }) + .fetch_all(&mut *tx) + .await?; + let mut batch_map: BTreeMap<_, _> = batches.iter_mut().map(|it| (it.id, it)).collect(); + // Then load transactions + sqlx::query( + " + SELECT + initiated_outgoing_transaction_id + ,amount + ,subject + ,credit_payto + ,initiated_outgoing_transactions.initiation_time + ,end_to_end_id + ,initiated_outgoing_batch_id + FROM initiated_outgoing_transactions + JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id) + WHERE initiated_outgoing_batches.status IN ('unsubmitted', 'transient_failure') + ", + ) + .try_map(|r: PgRow| { + let payment = InitiatedPayment { + id: r.try_get_u64("initiated_outgoing_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + creditor: r.try_get_parse("credit_payto")?, + subject: r.try_get("subject")?, + initiation_time: r.try_get_timestamp("initiation_time")?, + end_to_end_id: r.try_get("end_to_end_id")?, + }; + let batch_id = r.try_get_u64("initiated_outgoing_batch_id")?; + batch_map.get_mut(&batch_id).unwrap().payments.push(payment); + Ok(()) + }) + .fetch_all(&mut *tx) + .await?; + tx.commit().await?; + Ok(batches) +} + pub async fn unsettled_tx_in_batch( db: &PgPool, currency: &Currency, msg_id: &str, execution_time: &Timestamp, ) -> sqlx::Result<Vec<OutgoingPayment>> { - sqlx::query( + sqlx::query(formatcp!( " SELECT end_to_end_id, @@ -149,28 +229,260 @@ pub async fn unsettled_tx_in_batch( FROM initiated_outgoing_transactions JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id) WHERE message_id = $1 - AND initiated_outgoing_transactions.status NOT IN ('success', 'permanent_failure', 'late_failure') + AND initiated_outgoing_transactions.{UNSETTLED} " - ) + )) .bind(msg_id) - .try_map(|r: PgRow| Ok( - OutgoingPayment { + .try_map(|r: PgRow| { + Ok(OutgoingPayment { id: OutgoingId { msg_id: Some(msg_id.into()), end_to_end_id: r.try_get("end_to_end_id")?, - acct_svcr_ref: None + acct_svcr_ref: None, }, amount: r.try_get_amount("amount", currency)?, debit_fee: None, subject: r.try_get("subject")?, execution_time: *execution_time, - creditor: r.try_get_opt_payto("credit_payto")? - } - )) + creditor: r.try_get_opt_payto("credit_payto")?, + }) + }) .fetch_all(db) .await } +/** Register submission success of order [orderId] for batch [id] at [timestamp] */ +pub async fn batch_sub_success( + db: &PgPool, + batch_id: u64, + timestamp: &Timestamp, + order_id: &str, +) -> sqlx::Result<()> { + let mut tx = db.begin().await?; + // Update batch status + let updated = sqlx::query( + " + UPDATE initiated_outgoing_batches + SET status = 'pending' + ,submission_date = $1 + ,status_msg = NULL + ,order_id = $2 + ,submission_counter = submission_counter + 1 + WHERE initiated_outgoing_batch_id = $3 AND order_id IS NULL + ", + ) + .bind_timestamp(timestamp) + .bind(order_id) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + if updated.rows_affected() > 0 { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = 'pending', status_msg = NULL + WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED} + " + )) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await +} + +/** Register submission failure with [msg] for batch [id] at [timestamp]*/ +pub async fn batch_sub_failure( + db: &PgPool, + batch_id: u64, + timestamp: &Timestamp, + msg: &str, +) -> sqlx::Result<()> { + let permanent = false; + let mut tx = db.begin().await?; + // Update batch status + sqlx::query( + " + UPDATE initiated_outgoing_batches + SET status = $1 + ,submission_date = $2 + ,status_msg = $3 + ,submission_counter = submission_counter + 1 + WHERE initiated_outgoing_batch_id = $4 + ", + ) + .bind(if permanent { + SubmissionState::permanent_failure + } else { + SubmissionState::transient_failure + }) + .bind_timestamp(timestamp) + .bind(msg) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = $1, status_msg = $2 + WHERE initiated_outgoing_batch_id = $3 AND {UNSETTLED} + " + )) + .bind(if permanent { + SubmissionState::permanent_failure + } else { + SubmissionState::transient_failure + }) + .bind(msg) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + tx.commit().await +} + +/** Register order step [msg] for [orderId] */ +pub async fn order_step(db: &PgPool, order_id: &str, msg: &str) -> sqlx::Result<()> { + let mut tx = db.begin().await?; + // Update batch status + let batch_id = sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_batches + SET status = 'pending', status_msg = $1 + WHERE order_id = $1 AND {PENDING} + RETURNING initiated_outgoing_batch_id + " + )) + .bind(msg) + .bind(order_id) + .try_map(|r: PgRow| r.try_get_u32(0)) + .fetch_optional(&mut *tx) + .await?; + if let Some(batch_id) = batch_id { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = 'pending', status_msg = $1 + WHERE initiated_outgoing_batch_id = $2 AND {PENDING} + " + )) + .bind(msg) + .bind(batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await +} + +/** Register order success for [orderId] and return message_id if found */ +pub async fn order_success(db: &PgPool, order_id: &str) -> sqlx::Result<Option<String>> { + let mut tx = db.begin().await?; + // Update batch status + let res = sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_batches + SET status = 'success' + WHERE order_id = $1 + RETURNING initiated_outgoing_batch_id, message_id + " + )) + .bind(order_id) + .try_map(|r: PgRow| Ok((r.try_get_u32(0)?, r.try_get(1)?))) + .fetch_optional(&mut *tx) + .await?; + if let Some((batch_id, _)) = &res { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = 'pending' + WHERE initiated_outgoing_batch_id = $1 AND {UNSETTLED} + " + )) + .bind(*batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await?; + Ok(res.map(|(_, msg_id)| msg_id)) +} + +/** Register order failure for [orderId] and return message_id and previous status_msg if found */ +pub async fn order_failure( + db: &PgPool, + order_id: &str, +) -> sqlx::Result<Option<(String, Option<String>)>> { + let mut tx = db.begin().await?; + // Update batch status + let res = sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_batches + SET status = 'permanent_failure' + WHERE order_id = $1 + RETURNING initiated_outgoing_batch_id, message_id, status_msg + " + )) + .bind(order_id) + .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get(1)?, r.try_get(2)?))) + .fetch_optional(&mut *tx) + .await?; + if let Some((batch_id, _, _)) = &res { + // Update unsettled batch's transaction status + sqlx::query(formatcp!( + " + UPDATE initiated_outgoing_transactions + SET status = 'permanent_failure' + WHERE initiated_outgoing_batch_id = $1 + " + )) + .bind(*batch_id as i64) + .execute(&mut *tx) + .await?; + } + tx.commit().await?; + Ok(res.map(|(_, msg_id, status_msg)| (msg_id, status_msg))) +} + +/** Register payment status [state] with [msg] for batch [msgId] */ +pub async fn batch_status_update( + db: &PgPool, + msg_id: &str, + state: SubmissionState, + msg: &str, +) -> sqlx::Result<bool> { + sqlx::query(formatcp!( + "SELECT out_ok FROM batch_status_update($1,$2,$3)" + )) + .bind(msg_id) + .bind(state) + .bind(msg) + .try_map(|r: PgRow| r.try_get(0)) + .fetch_one(db) + .await +} + +/** Register payment status [state] with [msg] for transaction [endToEndId] in batch [msgId] */ +pub async fn tx_status_update( + db: &PgPool, + end_to_end_id: &str, + msg_id: &str, + state: SubmissionState, + msg: &str, +) -> sqlx::Result<bool> { + sqlx::query(formatcp!( + "SELECT out_ok FROM tx_status_update($1,$2,$3,$4)" + )) + .bind(end_to_end_id) + .bind(msg_id) + .bind(state) + .bind(msg) + .try_map(|r: PgRow| r.try_get(0)) + .fetch_one(db) + .await +} + #[derive(Debug, PartialEq, Eq)] pub struct OutgoingRegistrationResult { pub id: u64, @@ -224,7 +536,6 @@ pub async fn register_out_batch( payment: &OutgoingPayment, subject: Option<&OutgoingSubject>, ) -> sqlx::Result<OutgoingRegistrationResult> { - // TODO sqlx::query( " SELECT out_tx_id, out_initiated, out_found @@ -665,35 +976,43 @@ mod test { use std::sync::LazyLock; use compact_str::CompactString; - use jiff::Timestamp; + use jiff::{Span, Timestamp, civil::Date}; use sqlx::PgPool; use sqlx::{PgConnection, Row, postgres::PgRow}; use taler_api::{ db::{IncomingType, TypeHelper}, subject::subject_fmt_qr_bill, }; - use taler_common::types::{amount::Currency, payto::IbanPayto}; - use taler_common::{api_common::ShortHashCode, types::amount::Amount}; + use taler_common::{ + api_common::ShortHashCode, + types::{amount::Amount, utils::date_to_utc_ts}, + }; use taler_common::{ api_common::{EddsaPublicKey, EddsaSignature}, types::amount::amount, }; + use taler_common::{ + config::Config, + types::{amount::Currency, payto::IbanPayto}, + }; use uuid::Uuid; use crate::{ - CONFIG_SOURCE, IncomingId, IncomingPayment, InitiatedPayment, OutgoingBatch, - config::{AccountType, NexusIngestConfig}, + CONFIG_SOURCE, IncomingPayment, OutgoingBatch, Tx, + config::{AccountType, NexusCfg, NexusIngestConfig}, db::{ InResult, IncomingBounceRegistrationResult, OutgoingRegistrationResult, - RegistrationResult, register_in_malformed, transfer_register, + RegistrationResult, batch_status_update, batch_sub_failure, batch_sub_success, + initiated_submittable, order_failure, order_step, order_success, register_in_malformed, + transfer_register, tx_status_update, }, - register_incoming, + model::{IncomingId, InitiatedPayment, OutgoingId, SubmissionState}, + register_incoming, register_tx, }; + use crate::{OutgoingPayment, rand_ebics_id, register_outgoing, register_outgoing_batch}; use crate::{ - OutgoingId, db::PaymentInitiationResult, db::batch_initiated, db::initiate, - db::initiated_ack, + db::PaymentInitiationResult, db::batch_initiated, db::initiate, db::initiated_ack, }; - use crate::{OutgoingPayment, rand_ebics_id, register_outgoing, register_outgoing_batch}; pub static CURRENCY: LazyLock<Currency> = LazyLock::new(|| "KUDOS".parse().unwrap()); @@ -725,7 +1044,7 @@ mod test { /** Generates a payment initiation, given its subject and end-to-end ID */ pub fn gen_init_pay( - end_to_end_id: CompactString, + end_to_end_id: impl Into<CompactString>, subject: impl Into<String>, ) -> InitiatedPayment { InitiatedPayment { @@ -736,7 +1055,7 @@ mod test { .as_payto(), subject: subject.into(), initiation_time: Timestamp::now(), - end_to_end_id, + end_to_end_id: end_to_end_id.into(), } } @@ -756,6 +1075,25 @@ mod test { } } + async fn check_count(db: &PgPool, nb_tx: usize, nb_bounce: usize) { + sqlx::query( + " + SELECT (SELECT count(*) FROM incoming_transactions) + (SELECT count(*) FROM outgoing_transactions) AS transactions, + (SELECT count(*) FROM bounced_transactions) AS bounce + ", + ) + .try_map(|r: PgRow| { + assert_eq!( + (r.try_get_u64(0)?, r.try_get_u64(1)?), + (nb_tx as u64, nb_bounce as u64) + ); + Ok(()) + }) + .fetch_one(db) + .await + .unwrap(); + } + async fn check_in_count( db: &PgPool, nb_incoming: usize, @@ -1237,7 +1575,7 @@ mod test { let first = EddsaPublicKey::rand(); let auth_pub = EddsaPublicKey::rand(); let auth_sig = EddsaSignature::rand(); - let reference_number = subject_fmt_qr_bill(auth_pub.slice()); + let reference_number = subject_fmt_qr_bill(auth_pub.as_slice()); let subject = format!("test with MAP:{auth_pub} auth pub"); assert_eq!( @@ -1362,7 +1700,7 @@ mod test { let first = EddsaPublicKey::rand(); let auth_pub = EddsaPublicKey::rand(); let auth_sig = EddsaSignature::rand(); - let reference_number = subject_fmt_qr_bill(auth_pub.slice()); + let reference_number = subject_fmt_qr_bill(auth_pub.as_slice()); assert_eq!( transfer_register( @@ -1648,4 +1986,429 @@ mod test { register_incoming(&db, &cfg, &payment).await.unwrap(); check_in(&db, &[Reserve(key.clone()), Bounced]).await; } + + #[tokio::test] + pub async fn initiated_skip() { + let (_, db) = setup().await; + let cfg = Config::from_file(CONFIG_SOURCE, Some("libeufin-nexus/conf/skip.conf")).unwrap(); + let cfg = NexusCfg::parse(cfg).unwrap(); + let cfg = cfg.ingest().unwrap(); + let millis = Span::new().milliseconds(10); + + async fn ingest(db: &PgPool, cfg: &NexusIngestConfig, execution_time: Timestamp) { + for tx in [ + Tx::In( + gen_in_pay(format!("test at {execution_time}")) + .with_execution_time(execution_time), + ), + Tx::Out( + gen_out_pay(format!("test at {execution_time}")) + .with_execution_time(execution_time), + ), + ] { + register_tx(db, cfg, &tx).await.unwrap() + } + } + + assert_eq!( + cfg.ignore_txs_before, + date_to_utc_ts(&Date::from_str("2024-04-04").unwrap()) + ); + assert_eq!( + cfg.ignore_bounces_before, + date_to_utc_ts(&Date::from_str("2024-06-12").unwrap()) + ); + + // No transaction at the beginning + check_count(&db, 0, 0).await; + + // Skipped transactions + ingest(&db, &cfg, cfg.ignore_txs_before - millis).await; + check_count(&db, 0, 0).await; + + // Skipped bounces + ingest(&db, &cfg, cfg.ignore_txs_before).await; + ingest(&db, &cfg, cfg.ignore_txs_before + millis).await; + ingest(&db, &cfg, cfg.ignore_bounces_before - millis).await; + check_count(&db, 6, 0).await; + + // Bounces + ingest(&db, &cfg, cfg.ignore_bounces_before).await; + ingest(&db, &cfg, cfg.ignore_bounces_before + millis).await; + check_count(&db, 10, 2).await; + } + + #[tokio::test] + pub async fn initiated_status() { + use SubmissionState::*; + + let (_, db) = setup().await; + + let check_parts = async |batch_id: u64, + batch_status: SubmissionState, + batch_msg: &str, + tx_status: SubmissionState, + tx_msg: &str, + settled_status: SubmissionState, + settled_msg: &str| { + // Check batch status + let msg_id: String = sqlx::query( + " + SELECT message_id, status, status_msg FROM initiated_outgoing_batches WHERE initiated_outgoing_batch_id=? + " + ).bind(batch_id as i64) + .try_map(|r: PgRow| { + let msg_id: String = r.try_get("message_id")?; + assert_eq!((batch_status, batch_msg), (r.try_get("status")?, r.try_get("status_msg")?), "{msg_id}"); + Ok(msg_id) + }).fetch_one(&db).await.unwrap(); + // Check tx status + sqlx::query( + " + SELECT end_to_end_id, status, status_msg FROM initiated_outgoing_transactions WHERE initiated_outgoing_batch_id=? + " + ).bind(batch_id as i64).try_map(|r: PgRow| { + let end_to_end_id: &str = r.try_get("end_to_end_id")?; + let expected = match end_to_end_id { + "TX" => (tx_status, tx_msg), + "TX_SETTLED" => (settled_status, settled_msg), + _ =>panic!("Unexpected tx $endToEndId") + }; + assert_eq!(expected, + (r.try_get("status")?, r.try_get("status_msg")?), + "{msg_id},{end_to_end_id}" + ); + Ok(()) + }).fetch_all(&db).await.unwrap(); + }; + + let check_batch_tx = async |batch_id: u64, + status: SubmissionState, + msg: &str, + tx_status: SubmissionState| { + check_parts(batch_id, status, msg, tx_status, msg, tx_status, msg).await; + }; + let check_batch = async |batch_id: u64, status: SubmissionState, msg: &str| { + check_batch_tx(batch_id, status, msg, status).await; + }; + let check_order_tx = async |order_id: &str, + status: SubmissionState, + msg: &str, + tx_status: SubmissionState| { + let batch_id = sqlx::query( + "SELECT initiated_outgoing_batch_id FROM initiated_outgoing_batches WHERE order_id=$1" + ).bind(order_id) + .try_map(|r: PgRow| { + r.try_get_u64(0) + }).fetch_one(&db).await.unwrap(); + check_batch_tx(batch_id, status, msg, tx_status).await; + }; + let check_order = async |order_id: &str, status: SubmissionState, msg: &str| { + check_order_tx(order_id, status, msg, status).await; + }; + + async fn test<F>(db: &PgPool, lambda: impl FnOnce(u64) -> F) + where + F: Future<Output = ()>, + { + // Reset DB + sqlx::query("DELETE FROM initiated_outgoing_transactions") + .execute(db) + .await + .unwrap(); + sqlx::query("DELETE FROM initiated_outgoing_batches") + .execute(db) + .await + .unwrap(); + // Create a test batch with three transactions + for id in ["TX", "TX_SETTLED"] { + assert!(matches!( + initiate(db, &gen_init_pay(id, "lol")).await.unwrap(), + PaymentInitiationResult::Success(_) + )); + } + batch_initiated(db, &Timestamp::now(), "BATCH", false) + .await + .unwrap(); + + // Create witness transactions and batch + for id in ["WITNESS_1", "WITNESS_2"] { + assert!(matches!( + initiate(db, &gen_init_pay(id, "lol")).await.unwrap(), + PaymentInitiationResult::Success(_) + )); + } + batch_initiated(db, &Timestamp::now(), "BATCH_WITNESS", false) + .await + .unwrap(); + for id in ["WITNESS_3", "WITNESS_4"] { + assert!(matches!( + initiate(db, &gen_init_pay(id, "lol")).await.unwrap(), + PaymentInitiationResult::Success(_) + )); + } + // Check everything is unsubmitted + sqlx::query( + " + SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches) + AND (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_transactions) + " + ).try_map(|r: PgRow| { + assert!(r.try_get_flag(0).unwrap()); + Ok(()) + }).fetch_one(db).await.unwrap(); + let submitibale = initiated_submittable(db, &CURRENCY).await.unwrap(); + lambda( + submitibale + .iter() + .find(|it| it.msg_id == "BATCH") + .unwrap() + .id, + ); + // Check witness status is unaltered + sqlx::query( + " + SELECT (SELECT bool_and(status = 'unsubmitted') FROM initiated_outgoing_batches WHERE message_id != 'BATCH') + AND (SELECT bool_and(initiated_outgoing_transactions.status = 'unsubmitted') + FROM initiated_outgoing_transactions JOIN initiated_outgoing_batches USING (initiated_outgoing_batch_id) + WHERE message_id != 'BATCH') + " + ).try_map(|r: PgRow| { + assert!(r.try_get(0)?); + Ok(()) + }).fetch_one(db).await.unwrap(); + } + + let now = Timestamp::now(); + + // Submission retry status + test(&db, async |batch_id| { + batch_sub_failure(&db, batch_id, &now, "First failure") + .await + .unwrap(); + check_batch(batch_id, transient_failure, "First failure").await; + batch_sub_failure(&db, batch_id, &now, "Second failure") + .await + .unwrap(); + check_batch(batch_id, transient_failure, "Second failure").await; + batch_sub_success(&db, batch_id, &now, "ORDER") + .await + .unwrap(); + check_order("ORDER", pending, "").await; + batch_sub_success(&db, batch_id, &now, "ORDER") + .await + .unwrap(); + check_order("ORDER", pending, "").await; + order_step(&db, "ORDER", "step msg").await.unwrap(); + check_order("ORDER", pending, "step msg").await; + order_step(&db, "ORDER", "success msg").await.unwrap(); + check_order("ORDER", pending, "success msg").await; + order_success(&db, "ORDER").await.unwrap(); + check_order_tx("ORDER", success, "success msg", pending).await; + order_step(&db, "ORDER", "late msg").await.unwrap(); + check_order_tx("ORDER", success, "success msg", pending).await; + }) + .await; + + // Order step message on failure + test(&db, async |batch_id| { + batch_sub_success(&db, batch_id, &now, "ORDER") + .await + .unwrap(); + check_order("ORDER", pending, "").await; + order_step(&db, "ORDER", "step msg").await.unwrap(); + check_order("ORDER", pending, "step msg").await; + order_step(&db, "ORDER", "failure msg").await.unwrap(); + check_order("ORDER", pending, "failure msg").await; + assert_eq!( + Some("failure msg"), + order_failure(&db, "ORDER") + .await + .unwrap() + .unwrap() + .1 + .as_deref() + ); + check_order("ORDER", permanent_failure, "late msg").await; + order_step(&db, "ORDER", "late msg").await.unwrap(); + check_order("ORDER", permanent_failure, "failure msg").await; + }) + .await; + + // Payment & batch status + test(&db, async |batch_id| { + check_batch(batch_id, unsubmitted, "").await; + batch_status_update(&db, "BATCH", pending, "progress") + .await + .unwrap(); + check_batch(batch_id, pending, "progress").await; + tx_status_update(&db, "TX_SETTLED", "", success, "success") + .await + .unwrap(); + check_parts( + batch_id, pending, "progress", pending, "progress", success, "success", + ) + .await; + batch_status_update(&db, "BATCH", transient_failure, "waiting") + .await + .unwrap(); + check_parts( + batch_id, + transient_failure, + "waiting", + transient_failure, + "waiting", + success, + "success", + ) + .await; + tx_status_update(&db, "TX", "BATCH", permanent_failure, "failure") + .await + .unwrap(); + check_parts( + batch_id, + success, + "", + permanent_failure, + "failure", + success, + "success", + ) + .await; + tx_status_update(&db, "TX_SETTLED", "BATCH", permanent_failure, "late") + .await + .unwrap(); + check_parts( + batch_id, + success, + "", + permanent_failure, + "failure", + late_failure, + "late", + ) + .await; + }) + .await; + + // Registration + test(&db, async |batch_id| { + check_batch(batch_id, unsubmitted, "").await; + register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX_SETTLED")) + .await + .unwrap(); + check_parts(batch_id, unsubmitted, "", unsubmitted, "", success, "").await; + register_outgoing(&db, &gen_out_pay("").with_e2e_id("TX").with_msg_id("BATCH")) + .await + .unwrap(); + check_parts(batch_id, success, "", success, "", success, "").await; + }) + .await; + + // Transaction failure take over batch failures + test(&db, async |batch_id| { + check_batch(batch_id, unsubmitted, "").await; + batch_status_update(&db, "BATCH", permanent_failure, "batch") + .await + .unwrap(); + check_parts( + batch_id, + permanent_failure, + "batch", + permanent_failure, + "batch", + permanent_failure, + "batch", + ) + .await; + tx_status_update(&db, "TX", "BATCH", permanent_failure, "tx") + .await + .unwrap(); + batch_status_update(&db, "BATCH", permanent_failure, "batch2") + .await + .unwrap(); + check_parts( + batch_id, + permanent_failure, + "batch", + permanent_failure, + "tx", + permanent_failure, + "batch", + ) + .await; + }) + .await; + + // Unknown order and batch + batch_sub_success(&db, 42, &now, "ORDER_X").await.unwrap(); + batch_sub_failure(&db, 42, &now, "").await.unwrap(); + order_step(&db, "ORDER_X", "msg").await.unwrap(); + batch_status_update(&db, "BATCH_X", success, "") + .await + .unwrap(); + tx_status_update(&db, "TX_X", "BATCH_X", success, "msg") + .await + .unwrap(); + assert!(order_success(&db, "ORDER_X").await.unwrap().is_none()); + assert!(order_failure(&db, "ORDER_X").await.unwrap().is_none()); + } + + #[tokio::test] + pub async fn initiated_submittables() { + let (_, db) = setup().await; + let now = Timestamp::now(); + for i in 0..6 { + assert!(matches!( + initiate(&db, &gen_init_pay(format!("PAY{i}"), "")).await, + Ok(PaymentInitiationResult::Success(_)) + )); + batch_initiated(&db, &now, &rand_ebics_id(), false) + .await + .unwrap(); + } + + let check_ids = async |ids: &[&str]| { + assert_eq!( + ids, + initiated_submittable(&db, &CURRENCY) + .await + .unwrap() + .iter() + .flat_map(|it| it.payments.iter().map(|it| it.end_to_end_id.as_str())) + .collect::<Vec<_>>() + ); + }; + check_ids(&["PAY0", "PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await; + + // Check submitted not submitable + batch_sub_success(&db, 1, &now, "ORDER1").await.unwrap(); + check_ids(&["PAY1", "PAY2", "PAY3", "PAY4", "PAY5"]).await; + + // Check transient failure submitable last + batch_sub_failure(&db, 2, &now, "Failure").await.unwrap(); + check_ids(&["PAY2", "PAY3", "PAY4", "PAY5", "PAY1"]).await; + + // Check persistent failure not submitable + batch_sub_success(&db, 4, &now, "ORDER3").await.unwrap(); + order_failure(&db, "ORDER3").await.unwrap(); + check_ids(&["PAY2", "PAY4", "PAY5", "PAY1"]).await; + batch_sub_success(&db, 5, &now, "ORDER4").await.unwrap(); + order_failure(&db, "ORDER4").await.unwrap(); + check_ids(&["PAY2", "PAY5", "PAY1"]).await; + + // Check rotation + batch_sub_failure(&db, 3, &Timestamp::now(), "FAILURE") + .await + .unwrap(); + check_ids(&["PAY5", "PAY1", "PAY2"]).await; + batch_sub_failure(&db, 6, &Timestamp::now(), "FAILURE") + .await + .unwrap(); + check_ids(&["PAY1", "PAY2", "PAY5"]).await; + batch_sub_failure(&db, 2, &Timestamp::now(), "FAILURE") + .await + .unwrap(); + check_ids(&["PAY2", "PAY5", "PAY1"]).await; + } } diff --git a/src/lib.rs b/src/lib.rs @@ -17,14 +17,11 @@ * <http://www.gnu.org/licenses/> */ -use std::{ - fmt::{Display, Write}, - path::Path, -}; +use std::{fmt::Display, path::Path}; use anyhow::bail; -use compact_str::{CompactString, ToCompactString}; -use jiff::{Timestamp, civil::Date}; +use compact_str::CompactString; +use jiff::Timestamp; use rand::prelude::IndexedRandom; use reqwest::{ Client, StatusCode, @@ -35,16 +32,8 @@ use taler_api::subject::{ IncomingSubject, parse_incoming_unstructured, parse_outgoing, subject_is_qr_bill, }; use taler_build::long_version; -use taler_common::{ - CommonArgs, - config::parser::ConfigSource, - types::{ - amount::{Amount, Currency}, - payto::PaytoURI, - }, -}; +use taler_common::{CommonArgs, config::parser::ConfigSource, types::amount::Currency}; use tracing::{debug, info, warn}; -use uuid::Uuid; use crate::{ common::EbicsLogger, @@ -57,6 +46,7 @@ use crate::{ ebics_code::EbicsReturnCode, key_management::{Order, hpb, submit_client_keys}, keys::{load_bank_keys, load_client_keys, persist_client_keys}, + model::{IncomingPayment, OutgoingBatch, OutgoingPayment}, }; use crate::{keys::ClientPriKeysFile, xml::XmlReader}; @@ -67,6 +57,7 @@ pub mod db; pub mod ebics_code; pub mod key_management; pub mod keys; +pub mod model; pub mod testbench; pub mod xml; pub mod xml_sign; @@ -83,257 +74,6 @@ pub fn rand_ebics_id() -> CompactString { .collect() } -/// ID for incoming transactions -#[derive(Debug, Clone)] -struct IncomingId { - /** ISO20022 UETR */ - uetr: Option<Uuid>, - /// ISO20022 TxID - tx_id: Option<CompactString>, - /// ISO20022 AcctSvcrRef - acct_svcr_ref: Option<CompactString>, -} -impl IncomingId { - pub fn new( - uetr: Option<Uuid>, - tx_id: Option<CompactString>, - acct_svcr_ref: Option<CompactString>, - ) -> Self { - assert!(uetr.is_some() || tx_id.is_some() || acct_svcr_ref.is_some()); - Self { - uetr, - tx_id, - acct_svcr_ref, - } - } - - pub fn r#ref(&self) -> CompactString { - self.uetr - .map(|e| e.to_compact_string()) - .or(self.tx_id.clone()) - .or(self.acct_svcr_ref.clone()) - .expect("must be at least one ref") - } -} - -impl std::fmt::Display for IncomingId { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.write_char('(')?; - let mut prepend = false; - if let Some(uetr) = &self.uetr { - write!(f, "uetr={uetr}")?; - prepend = true; - } - if let Some(tx_id) = &self.tx_id { - if prepend { - f.write_char(' ')?; - } - f.write_str("tx=")?; - f.write_str(tx_id)?; - prepend = true; - } - if let Some(acct_svcr_ref) = &self.acct_svcr_ref { - if prepend { - f.write_char(' ')?; - } - f.write_str("ref=")?; - f.write_str(acct_svcr_ref)?; - } - f.write_char(')')?; - Ok(()) - } -} - -/// ID for outgoing transactions -struct OutgoingId { - /// Unique msg ID generated by libeufin-nexus - /// ISO20022 MessageId - msg_id: Option<CompactString>, - /// Unique end-to-end ID generated by libeufin-nexus - /// ISO20022 EndToEndId or MessageId (retrocompatibility) - end_to_end_id: Option<CompactString>, - /// Unique end-to-end ID generated by the bank - /// ISO20022 AcctSvcrRef - acct_svcr_ref: Option<CompactString>, -} -impl OutgoingId { - pub fn r#ref(&self) -> CompactString { - self.end_to_end_id - .clone() - .or(self.acct_svcr_ref.clone()) - .or(self.acct_svcr_ref.clone()) - .expect("must be at least one ref") - } -} - -impl std::fmt::Display for OutgoingId { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.write_char('(')?; - let mut prepend = false; - if let Some(msg_id) = &self.msg_id - && self.msg_id != self.end_to_end_id - { - f.write_str("msg=")?; - f.write_str(msg_id)?; - prepend = true; - } - if let Some(end_to_end_id) = &self.end_to_end_id { - if prepend { - f.write_char(' ')?; - } - f.write_str("e2e=")?; - f.write_str(end_to_end_id)?; - prepend = true; - } - if let Some(acct_svcr_ref) = &self.acct_svcr_ref { - if prepend { - f.write_char(' ')?; - } - f.write_str("ref=")?; - f.write_str(acct_svcr_ref)?; - } - f.write_char(')')?; - Ok(()) - } -} - -/// ID for outgoing batches -pub struct BatchId { - /// Unique msg ID generated by libeufin-nexus - /// ISO20022 MessageId - pub msg_id: CompactString, - /// Unique end-to-end ID generated by the bank - /// ISO20022 AcctSvcrRef - pub acct_svcr_ref: Option<CompactString>, -} - -impl BatchId { - pub fn r#ref(&self) -> CompactString { - self.msg_id.clone() - } -} - -impl std::fmt::Display for BatchId { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.write_str("(msg=")?; - f.write_str(&self.msg_id)?; - if let Some(acct_svcr_ref) = &self.acct_svcr_ref { - f.write_str("ref=")?; - f.write_str(acct_svcr_ref)?; - } - f.write_char(')')?; - Ok(()) - } -} - -/// ISO20022 incoming payment -#[derive(Debug, Clone)] -pub struct IncomingPayment { - id: IncomingId, - amount: Amount, - credit_fee: Option<Amount>, - subject: Option<String>, - execution_time: Timestamp, - debtor: Option<PaytoURI>, -} - -impl Display for IncomingPayment { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let Self { - id, - amount, - credit_fee, - subject, - execution_time, - debtor, - } = self; - write!(f, "IN {execution_time} {amount}")?; - if let Some(credit_fee) = credit_fee { - write!(f, "-{credit_fee}")?; - } - write!(f, " {id}")?; - if let Some(creditor) = debtor { - write!(f, " creditor={creditor}")?; - } - if let Some(subject) = subject { - write!(f, " subject='{subject}'")?; - } - Ok(()) - } -} - -/// ISO20022 outgoing payment -pub struct OutgoingPayment { - id: OutgoingId, - amount: Amount, - debit_fee: Option<Amount>, - subject: Option<String>, - execution_time: Timestamp, - creditor: Option<PaytoURI>, -} - -impl Display for OutgoingPayment { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let Self { - id, - amount, - debit_fee, - subject, - execution_time, - creditor, - } = self; - write!(f, "OUT {execution_time} {amount}")?; - if let Some(debit_fee) = debit_fee { - write!(f, "-{debit_fee}")?; - } - write!(f, " {id}")?; - if let Some(creditor) = creditor { - write!(f, " creditor={creditor}")?; - } - if let Some(subject) = subject { - write!(f, " subject='{subject}'")?; - } - Ok(()) - } -} - -/** ISO20022 outgoing batch */ -pub struct OutgoingBatch { - /** ISO20022 MessageId */ - pub msg_id: CompactString, - pub execution_time: Timestamp, -} - -impl Display for OutgoingBatch { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - let Self { - msg_id, - execution_time, - } = self; - // TODO fmt date - write!(f, "BATCH {execution_time} {msg_id}") - } -} - -/** Batch of initiated outgoing payment to sent together */ -pub struct PaymentBatch { - pub id: u64, - pub msg_id: CompactString, - pub creation_date: Date, - pub sum: Amount, - pub payments: Vec<InitiatedPayment>, -} - -/** Initiated outgoing transaction */ -pub struct InitiatedPayment { - pub id: u64, - pub amount: Amount, - pub subject: String, - pub creditor: PaytoURI, - pub initiation_time: Timestamp, - pub end_to_end_id: CompactString, -} - #[derive(clap::Parser, Debug)] #[command(long_version = long_version(), about, long_about = None)] pub struct Args { @@ -380,7 +120,7 @@ pub async fn register_incoming( debug!("{fmt}") } }; - let bounce = |cause: String| async move { + let bounce = async |cause: &str| { match cfg.account_type { AccountType::Exchange => { if payment.execution_time < cfg.ignore_bounces_before { @@ -420,7 +160,7 @@ pub async fn register_incoming( &bounce_amount, &rand_ebics_id(), &Timestamp::now(), - &cause, + cause, ) .await?; match res { @@ -453,7 +193,7 @@ pub async fn register_incoming( && let Some(debtor) = &payment.debtor && !regex.is_match(debtor.as_ref().as_str()) { - bounce("restricted account".to_owned()).await?; + bounce("restricted account").await?; return Ok(()); } @@ -461,40 +201,30 @@ pub async fn register_incoming( && subject_is_qr_bill(subject) { match register_in_qr_bill(db, payment, subject).await? { - IncomingRegistrationResult::ReservePubReuse => { - bounce("reverse pub reuse".to_owned()).await? - } - IncomingRegistrationResult::MappingReuse => bounce("mapping reuse".to_owned()).await?, - IncomingRegistrationResult::UnknownMapping => { - bounce("unknown mapping".to_owned()).await? - } + IncomingRegistrationResult::ReservePubReuse => bounce("reverse pub reuse").await?, + IncomingRegistrationResult::MappingReuse => bounce("mapping reuse").await?, + IncomingRegistrationResult::UnknownMapping => bounce("unknown mapping").await?, IncomingRegistrationResult::Success(res) => { log_res(res, "", ""); } } } else { match parse_incoming_unstructured(payment.subject.as_deref().unwrap_or_default()) { - Ok(None) => bounce("missing public key".to_owned()).await?, + Ok(None) => bounce("missing public key").await?, Ok(Some(IncomingSubject::AdminBalanceAdjust)) => { let res = register_in(db, payment).await?; log_res(res, "admin balance adjust", ""); } Ok(Some(subject)) => match register_in_talerable(db, payment, &subject).await? { - IncomingRegistrationResult::ReservePubReuse => { - bounce("reverse pub reuse".to_owned()).await? - } - IncomingRegistrationResult::MappingReuse => { - bounce("mapping reuse".to_owned()).await? - } - IncomingRegistrationResult::UnknownMapping => { - bounce("unknown mapping".to_owned()).await? - } + IncomingRegistrationResult::ReservePubReuse => bounce("reverse pub reuse").await?, + IncomingRegistrationResult::MappingReuse => bounce("mapping reuse").await?, + IncomingRegistrationResult::UnknownMapping => bounce("unknown mapping").await?, IncomingRegistrationResult::Success(res) => { log_res(res, "", ""); } }, Err(e) => { - bounce(e.to_string()).await?; + bounce(&e.to_string()).await?; } } } @@ -536,6 +266,55 @@ pub async fn register_outgoing_batch( Ok(()) } +pub enum Tx { + In(IncomingPayment), + Out(OutgoingPayment), + Batch(OutgoingBatch), + Reversal, +} + +impl Tx { + pub fn execution_time(&self) -> &Timestamp { + match self { + Tx::In(IncomingPayment { execution_time, .. }) + | Tx::Out(OutgoingPayment { execution_time, .. }) + | Tx::Batch(OutgoingBatch { execution_time, .. }) => execution_time, + Tx::Reversal => todo!(), + } + } +} + +impl Display for Tx { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Tx::In(incoming_payment) => incoming_payment.fmt(f), + Tx::Out(outgoing_payment) => outgoing_payment.fmt(f), + Tx::Batch(outgoing_batch) => outgoing_batch.fmt(f), + Tx::Reversal => todo!(), + } + } +} + +pub async fn register_tx(db: &PgPool, cfg: &NexusIngestConfig, tx: &Tx) -> sqlx::Result<()> { + if tx.execution_time() < &cfg.ignore_txs_before { + debug!("IGNORE {tx}"); + } else { + match tx { + Tx::In(payment) => { + register_incoming(db, cfg, payment).await?; + } + Tx::Out(payment) => { + register_outgoing(db, payment).await?; + } + Tx::Batch(batch) => { + register_outgoing_batch(db, &cfg.currency, batch).await?; + } + Tx::Reversal => todo!(), + } + } + Ok(()) +} + /** Load client private keys at or create new ones if missing */ pub fn load_or_generate_client_keys(path: &Path) -> anyhow::Result<ClientPriKeysFile> { // If exists load from disk diff --git a/src/main.rs b/src/main.rs @@ -24,7 +24,7 @@ use taler_common::taler_main; fn main() { let args = Args::parse(); - taler_main(CONFIG_SOURCE, args.common, |cfg| async move { + taler_main(CONFIG_SOURCE, args.common, async |cfg| { let cfg = NexusCfg::parse(cfg)?; ebics_setup(&Client::new(), &cfg, false).await?; Ok(()) diff --git a/src/model.rs b/src/model.rs @@ -0,0 +1,349 @@ +/* +* This file is part of LibEuFin. +* Copyright (C) 2026 Taler Systems S.A. + +* LibEuFin is free software; you can redistribute it and/or modify +* it under the terms of the GNU Affero General Public License as +* published by the Free Software Foundation; either version 3, or +* (at your option) any later version. + +* LibEuFin is distributed in the hope that it will be useful, but +* WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY +* or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Affero General +* Public License for more details. + +* You should have received a copy of the GNU Affero General Public +* License along with LibEuFin; see the file COPYING. If not, see +* <http://www.gnu.org/licenses/> +*/ + +use std::fmt::{Display, Write as _}; + +use compact_str::{CompactString, ToCompactString as _}; +use jiff::Timestamp; +use taler_common::{ + api_wire::TransferState, + types::{amount::Amount, payto::PaytoURI}, +}; +use uuid::Uuid; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, sqlx::Type)] +#[allow(non_camel_case_types)] +#[sqlx(type_name = "submission_state")] +/** Outgoing transactions and batches submission status */ +pub enum SubmissionState { + // Initiated but not yet submitted + unsubmitted, + // Submission failed, retry possible + transient_failure, + // Submission succeed, pending settltment + pending, + // Definitive failure, will never succeed + permanent_failure, + // Definitive success, booked and settled + success, + // Late failure after a success, happens when a payment is returned + late_failure, +} + +impl SubmissionState { + fn to_transfer_status(self) -> TransferState { + match self { + SubmissionState::unsubmitted | SubmissionState::pending => TransferState::pending, + SubmissionState::transient_failure => TransferState::transient_failure, + SubmissionState::permanent_failure => TransferState::permanent_failure, + SubmissionState::success | SubmissionState::late_failure => TransferState::success, + } + } +} + +/// ID for incoming transactions +#[derive(Debug, Clone)] +pub struct IncomingId { + /** ISO20022 UETR */ + pub uetr: Option<Uuid>, + /// ISO20022 TxID + pub tx_id: Option<CompactString>, + /// ISO20022 AcctSvcrRef + pub acct_svcr_ref: Option<CompactString>, +} + +impl IncomingId { + pub fn new( + uetr: Option<Uuid>, + tx_id: Option<CompactString>, + acct_svcr_ref: Option<CompactString>, + ) -> Self { + assert!(uetr.is_some() || tx_id.is_some() || acct_svcr_ref.is_some()); + Self { + uetr, + tx_id, + acct_svcr_ref, + } + } + + pub fn r#ref(&self) -> CompactString { + self.uetr + .map(|e| e.to_compact_string()) + .or(self.tx_id.clone()) + .or(self.acct_svcr_ref.clone()) + .expect("must be at least one ref") + } +} + +impl std::fmt::Display for IncomingId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_char('(')?; + let mut prepend = false; + if let Some(uetr) = &self.uetr { + write!(f, "uetr={uetr}")?; + prepend = true; + } + if let Some(tx_id) = &self.tx_id { + if prepend { + f.write_char(' ')?; + } + f.write_str("tx=")?; + f.write_str(tx_id)?; + prepend = true; + } + if let Some(acct_svcr_ref) = &self.acct_svcr_ref { + if prepend { + f.write_char(' ')?; + } + f.write_str("ref=")?; + f.write_str(acct_svcr_ref)?; + } + f.write_char(')')?; + Ok(()) + } +} + +/// ID for outgoing transactions +pub struct OutgoingId { + /// Unique msg ID generated by libeufin-nexus + /// ISO20022 MessageId + pub msg_id: Option<CompactString>, + /// Unique end-to-end ID generated by libeufin-nexus + /// ISO20022 EndToEndId or MessageId (retrocompatibility) + pub end_to_end_id: Option<CompactString>, + /// Unique end-to-end ID generated by the bank + /// ISO20022 AcctSvcrRef + pub acct_svcr_ref: Option<CompactString>, +} + +impl OutgoingId { + pub fn r#ref(&self) -> CompactString { + self.end_to_end_id + .clone() + .or(self.acct_svcr_ref.clone()) + .or(self.acct_svcr_ref.clone()) + .expect("must be at least one ref") + } +} + +impl std::fmt::Display for OutgoingId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_char('(')?; + let mut prepend = false; + if let Some(msg_id) = &self.msg_id + && self.msg_id != self.end_to_end_id + { + f.write_str("msg=")?; + f.write_str(msg_id)?; + prepend = true; + } + if let Some(end_to_end_id) = &self.end_to_end_id { + if prepend { + f.write_char(' ')?; + } + f.write_str("e2e=")?; + f.write_str(end_to_end_id)?; + prepend = true; + } + if let Some(acct_svcr_ref) = &self.acct_svcr_ref { + if prepend { + f.write_char(' ')?; + } + f.write_str("ref=")?; + f.write_str(acct_svcr_ref)?; + } + f.write_char(')')?; + Ok(()) + } +} + +/// ID for outgoing batches +pub struct BatchId { + /// Unique msg ID generated by libeufin-nexus + /// ISO20022 MessageId + pub msg_id: CompactString, + /// Unique end-to-end ID generated by the bank + /// ISO20022 AcctSvcrRef + pub acct_svcr_ref: Option<CompactString>, +} + +impl BatchId { + pub fn r#ref(&self) -> CompactString { + self.msg_id.clone() + } +} + +impl std::fmt::Display for BatchId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("(msg=")?; + f.write_str(&self.msg_id)?; + if let Some(acct_svcr_ref) = &self.acct_svcr_ref { + f.write_str("ref=")?; + f.write_str(acct_svcr_ref)?; + } + f.write_char(')')?; + Ok(()) + } +} + +/// ISO20022 incoming payment +#[derive(Debug, Clone)] +pub struct IncomingPayment { + pub id: IncomingId, + pub amount: Amount, + pub credit_fee: Option<Amount>, + pub subject: Option<String>, + pub execution_time: Timestamp, + pub debtor: Option<PaytoURI>, +} + +impl IncomingPayment { + pub fn with_execution_time(self, execution_time: Timestamp) -> Self { + Self { + execution_time, + ..self + } + } +} + +impl Display for IncomingPayment { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let Self { + id, + amount, + credit_fee, + subject, + execution_time, + debtor, + } = self; + write!(f, "IN {execution_time} {amount}")?; + if let Some(credit_fee) = credit_fee { + write!(f, "-{credit_fee}")?; + } + write!(f, " {id}")?; + if let Some(creditor) = debtor { + write!(f, " creditor={creditor}")?; + } + if let Some(subject) = subject { + write!(f, " subject='{subject}'")?; + } + Ok(()) + } +} + +/// ISO20022 outgoing payment +pub struct OutgoingPayment { + pub id: OutgoingId, + pub amount: Amount, + pub debit_fee: Option<Amount>, + pub subject: Option<String>, + pub execution_time: Timestamp, + pub creditor: Option<PaytoURI>, +} + +impl OutgoingPayment { + pub fn with_execution_time(self, execution_time: Timestamp) -> Self { + Self { + execution_time, + ..self + } + } + + pub fn with_e2e_id(self, end_to_end_id: impl Into<CompactString>) -> Self { + Self { + id: OutgoingId { + end_to_end_id: Some(end_to_end_id.into()), + ..self.id + }, + ..self + } + } + + pub fn with_msg_id(self, msg_id: impl Into<CompactString>) -> Self { + Self { + id: OutgoingId { + msg_id: Some(msg_id.into()), + ..self.id + }, + ..self + } + } +} + +impl Display for OutgoingPayment { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let Self { + id, + amount, + debit_fee, + subject, + execution_time, + creditor, + } = self; + write!(f, "OUT {execution_time} {amount}")?; + if let Some(debit_fee) = debit_fee { + write!(f, "-{debit_fee}")?; + } + write!(f, " {id}")?; + if let Some(creditor) = creditor { + write!(f, " creditor={creditor}")?; + } + if let Some(subject) = subject { + write!(f, " subject='{subject}'")?; + } + Ok(()) + } +} + +/** ISO20022 outgoing batch */ +pub struct OutgoingBatch { + /** ISO20022 MessageId */ + pub msg_id: CompactString, + pub execution_time: Timestamp, +} + +impl Display for OutgoingBatch { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let Self { + msg_id, + execution_time, + } = self; + // TODO fmt date + write!(f, "BATCH {execution_time} {msg_id}") + } +} + +/** Batch of initiated outgoing payment to sent together */ +pub struct PaymentBatch { + pub id: u64, + pub msg_id: CompactString, + pub creation_date: Timestamp, + pub sum: Amount, + pub payments: Vec<InitiatedPayment>, +} + +/** Initiated outgoing transaction */ +pub struct InitiatedPayment { + pub id: u64, + pub amount: Amount, + pub subject: String, + pub creditor: PaytoURI, + pub initiation_time: Timestamp, + pub end_to_end_id: CompactString, +} diff --git a/src/xml.rs b/src/xml.rs @@ -234,7 +234,6 @@ impl<'node, 'input> XmlReader<'node, 'input> { } #[cfg(test)] - mod test { use crate::xml::XmlWriter;