libeufin

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

commit c9b0056a4b800f3c39a8c4e5cc8d71b5ff0a216d
parent 646078079a51bdedc25dafdd0d24cb251f22840c
Author: Antoine A <>
Date:   Fri, 24 Apr 2026 10:37:57 +0200

OutgoingPaymentsTest

Diffstat:
MCargo.lock | 354++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
MCargo.toml | 9+++++++++
Msrc/config.rs | 5+++++
Asrc/db.rs | 2121+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/main.rs | 300+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Asrc/testbench.rs | 101+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
6 files changed, 2848 insertions(+), 42 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -24,6 +24,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" [[package]] +name = "android_system_properties" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "819e7219dbd41043ac279b19830f2efc897156490d7fd6ea916720117ee66311" +dependencies = [ + "libc", +] + +[[package]] name = "anstream" version = "0.6.21" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -225,6 +234,9 @@ name = "bitflags" version = "2.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "843867be96c8daad0d758b57df9392b6d8d271134fce549de6ce169ff98a92af" +dependencies = [ + "serde_core", +] [[package]] name = "block-buffer" @@ -304,6 +316,18 @@ dependencies = [ ] [[package]] +name = "chrono" +version = "0.4.44" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c673075a2e0e5f4a1dde27ce9dee1ea4558c7ffe648f576438a20ca1d2acc4b0" +dependencies = [ + "iana-time-zone", + "num-traits", + "serde", + "windows-link", +] + +[[package]] name = "clap" version = "4.5.60" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -394,6 +418,15 @@ dependencies = [ ] [[package]] +name = "convert_case" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "633458d4ef8c78b72454de2d54fd6ab2e60f9e02be22f3c6104cdc8a4e0fceb9" +dependencies = [ + "unicode-segmentation", +] + +[[package]] name = "core-foundation" version = "0.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -468,6 +501,34 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28" [[package]] +name = "crossterm" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d8b9f2e4c67f833b660cdb0a3523065869fb35570177239812ed4c905aeff87b" +dependencies = [ + "bitflags", + "crossterm_winapi", + "derive_more", + "document-features", + "mio", + "parking_lot", + "rustix", + "serde", + "signal-hook", + "signal-hook-mio", + "winapi", +] + +[[package]] +name = "crossterm_winapi" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "acdd7c62a3665c7f6830a51635d9ac9b23ed385797f70a83bb8bafe9c572ab2b" +dependencies = [ + "winapi", +] + +[[package]] name = "crypto-common" version = "0.1.7" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -556,6 +617,28 @@ dependencies = [ ] [[package]] +name = "derive_more" +version = "2.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d751e9e49156b02b44f9c1815bcb94b984cdcc4396ecc32521c739452808b134" +dependencies = [ + "derive_more-impl", +] + +[[package]] +name = "derive_more-impl" +version = "2.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "799a97264921d8623a957f6c3b9011f3b5492f557bbb7a5a19b7fa6d06ba8dcb" +dependencies = [ + "convert_case", + "proc-macro2", + "quote", + "rustc_version", + "syn", +] + +[[package]] name = "digest" version = "0.10.7" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -578,6 +661,15 @@ dependencies = [ ] [[package]] +name = "document-features" +version = "0.2.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d4b8a88685455ed29a21542a33abd9cb6510b6b129abadabdcef0f4c55bc8f61" +dependencies = [ + "litrs", +] + +[[package]] name = "dotenvy" version = "0.15.7" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -652,6 +744,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" [[package]] +name = "fd-lock" +version = "4.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ce92ff622d6dadf7349484f42c93271a0d49b7cc4d466a936405bacbe10aa78" +dependencies = [ + "cfg-if", + "rustix", + "windows-sys 0.59.0", +] + +[[package]] name = "find-msvc-tools" version = "0.1.9" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1011,6 +1114,30 @@ dependencies = [ ] [[package]] +name = "iana-time-zone" +version = "0.1.65" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e31bc9ad994ba00e440a8aa5c9ef0ec67d5cb5e5cb0cc7f8b744a35b389cc470" +dependencies = [ + "android_system_properties", + "core-foundation-sys", + "iana-time-zone-haiku", + "js-sys", + "log", + "wasm-bindgen", + "windows-core", +] + +[[package]] +name = "iana-time-zone-haiku" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f31827a206f56af32e590ba56d5d2d085f558508192593743f16b2306495269f" +dependencies = [ + "cc", +] + +[[package]] name = "icu_collections" version = "2.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1159,6 +1286,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" [[package]] +name = "itertools" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "413ee7dfc52ee1a4949ceeb7dbc8a33f2d6c088194d9f922fb8318faf1f01186" +dependencies = [ + "either", +] + +[[package]] name = "itoa" version = "1.0.17" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1261,27 +1397,9 @@ checksum = "09edd9e8b54e49e587e4f6295a7d29c3ea94d469cb40ab8ca70b288248a81db2" [[package]] name = "libc" -version = "0.2.182" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6800badb6cb2082ffd7b6a67e6125bb39f18782f793520caee8cb8846be06112" - -[[package]] -name = "libdeflate-sys" -version = "1.25.2" +version = "0.2.183" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "72753e0008ea87963d2f0770042d0df7abe51fafbb8dcaf618ac440f2f1fec0a" -dependencies = [ - "cc", -] - -[[package]] -name = "libdeflater" -version = "1.25.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d1ee41cf6fb1bb6030dfb59ffb7bc01ab26aade44142084c87f0fc7a1658fe71" -dependencies = [ - "libdeflate-sys", -] +checksum = "b5b646652bf6661599e1da8901b3b9522896f01e736bad5f723fe7a3a27f899d" [[package]] name = "libeufin" @@ -1291,18 +1409,21 @@ dependencies = [ "aws-lc-rs", "base64", "clap", + "compact_str", "flate2", "getrandom 0.4.2", "jiff", "pem", "rand 0.10.0", "rcgen", + "reedline", "reqwest", "roxmltree", "serde", "serde_json", - "strum", - "strum_macros", + "sqlx", + "strum 0.28.0", + "strum_macros 0.28.0", "taler-api", "taler-build", "taler-common", @@ -1311,6 +1432,7 @@ dependencies = [ "tokio", "tracing", "url", + "uuid", "x509-parser", "xml-canonicalization", ] @@ -1351,6 +1473,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6373607a59f0be73a39b6fe456b8192fcc3585f602af20751600e974dd455e77" [[package]] +name = "litrs" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "11d3d7f243d5c5a8b9bb5d6dd2b1602c0cb0b9db1621bafc7ed66e35ff9fe092" + +[[package]] name = "lock_api" version = "0.4.14" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1422,6 +1550,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a69bcab0ad47271a0234d9422b131806bf3968021e5dc9328caf2d4cd58557fc" dependencies = [ "libc", + "log", "wasi", "windows-sys 0.61.2", ] @@ -1701,9 +1830,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.13" +version = "0.11.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f1906b49b0c3bc04b5fe5d86a77925ae6524a19b816ae38ce1e426255f1d8a31" +checksum = "434b42fec591c96ef50e21e886936e66d3cc3f737104fdb9b737c40ffb94c098" dependencies = [ "aws-lc-rs", "bytes", @@ -1865,6 +1994,27 @@ dependencies = [ ] [[package]] +name = "reedline" +version = "0.46.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fe9e7c532bfc2759bc8a28902c04e8b993fc13ebd085ee4292eb1b230fa9beef" +dependencies = [ + "chrono", + "crossterm", + "fd-lock", + "itertools", + "nu-ansi-term", + "serde", + "strip-ansi-escapes", + "strum 0.26.3", + "strum_macros 0.26.4", + "thiserror 2.0.18", + "unicase", + "unicode-segmentation", + "unicode-width", +] + +[[package]] name = "regex" version = "1.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1961,6 +2111,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "357703d41365b4b27c590e3ed91eabb1b663f07c4c084095e60cbed4362dff0d" [[package]] +name = "rustc_version" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cfcb3a22ef46e85b45de6ee7e79d063319ebb6594faafcf1c225ea92ab6e9b92" +dependencies = [ + "semver", +] + +[[package]] name = "rusticata-macros" version = "4.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2080,9 +2239,9 @@ dependencies = [ [[package]] name = "schannel" -version = "0.1.28" +version = "0.1.29" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "891d81b926048e76efe18581bf793546b4c0eaf8448d72be8de2bbee5fd166e1" +checksum = "91c1b7e4904c873ef0710c1f407dde2e6287de2bebc1bbbf7d430bb7cbffd939" dependencies = [ "windows-sys 0.61.2", ] @@ -2237,6 +2396,27 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" [[package]] +name = "signal-hook" +version = "0.3.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d881a16cf4426aa584979d30bd82cb33429027e42122b169753d6ef1085ed6e2" +dependencies = [ + "libc", + "signal-hook-registry", +] + +[[package]] +name = "signal-hook-mio" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b75a19a7a740b25bc7944bdee6172368f988763b744e3d4dfe753f6b4ece40cc" +dependencies = [ + "libc", + "mio", + "signal-hook", +] + +[[package]] name = "signal-hook-registry" version = "1.4.8" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2269,12 +2449,12 @@ dependencies = [ [[package]] name = "socket2" -version = "0.6.2" +version = "0.6.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "86f4aa3ad99f2088c990dfa82d367e19cb29268ed67c574d10d0a4bfe71f07e0" +checksum = "3a766e1110788c36f4fa1c2b71b387a7815aa65f88ce0229841826633d93723e" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -2421,6 +2601,15 @@ dependencies = [ ] [[package]] +name = "strip-ansi-escapes" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2a8f8038e7e7969abb3f1b7c2a811225e9296da208539e0f79c5251d6cac0025" +dependencies = [ + "vte", +] + +[[package]] name = "strsim" version = "0.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2428,12 +2617,31 @@ checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" [[package]] name = "strum" +version = "0.26.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8fec0f0aef304996cf250b31b5a10dee7980c85da9d759361292b8bca5a18f06" + +[[package]] +name = "strum" version = "0.28.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9628de9b8791db39ceda2b119bbe13134770b56c138ec1d3af810d045c04f9bd" [[package]] name = "strum_macros" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4c6bee85a5a24955dc440386795aa378cd9cf82acd5f764469152d2270e581be" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "rustversion", + "syn", +] + +[[package]] +name = "strum_macros" version = "0.28.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ab85eea0270ee17587ed4156089e10b9e6880ee688791d45a905f5b1ca36f664" @@ -2509,10 +2717,10 @@ dependencies = [ "aws-lc-rs", "axum", "base64", + "compact_str", "dashmap", "http-body-util", "jiff", - "libdeflater", "listenfd", "serde", "serde_json", @@ -2523,6 +2731,7 @@ dependencies = [ "tokio", "tracing", "url", + "zlib-rs", ] [[package]] @@ -2560,8 +2769,8 @@ name = "taler-test-utils" version = "1.3.0" dependencies = [ "axum", + "flate2", "http-body-util", - "libdeflater", "serde", "serde_json", "serde_urlencoded", @@ -2577,9 +2786,9 @@ dependencies = [ [[package]] name = "tempfile" -version = "3.26.0" +version = "3.27.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "82a72c767771b47409d2345987fda8628641887d5466101319899796367354a0" +checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", "getrandom 0.4.2", @@ -2877,6 +3086,12 @@ 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" +checksum = "dbc4bc3a9f746d862c45cb89d705aa10f187bb96c76001afab07a0d35ce60142" + +[[package]] name = "unicode-bidi" version = "0.3.18" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2904,6 +3119,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7df058c713841ad818f1dc5d3fd88063241cc61f49f5fbea4b951e8cf5a8d71d" [[package]] +name = "unicode-segmentation" +version = "1.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f6ccf251212114b54433ec949fd6a7841275f9ada20dddd2f29e9ceea4501493" + +[[package]] +name = "unicode-width" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b4ac048d71ede7ee76d585517add45da530660ef4390e49b098733c6e897f254" + +[[package]] name = "unicode-xid" version = "0.2.6" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2969,6 +3196,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" [[package]] +name = "vte" +version = "0.14.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "231fdcd7ef3037e8330d8e17e61011a2c244126acc0a982f4040ac3f9f0bc077" +dependencies = [ + "memchr", +] + +[[package]] name = "walkdir" version = "2.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3199,6 +3435,41 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" [[package]] +name = "windows-core" +version = "0.62.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb" +dependencies = [ + "windows-implement", + "windows-interface", + "windows-link", + "windows-result", + "windows-strings", +] + +[[package]] +name = "windows-implement" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "053e2e040ab57b9dc951b72c264860db7eb3b0200ba345b4e4c3b14f67855ddf" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "windows-interface" +version = "0.59.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f316c4a2570ba26bbec722032c4099d8c8bc095efccdc15688708623367e358" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] name = "windows-link" version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3262,6 +3533,15 @@ dependencies = [ [[package]] name = "windows-sys" +version = "0.59.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e38bc4d79ed67fd075bcc251a1c39b32a1776bbe92e5bef1f0bf1f8c531853b" +dependencies = [ + "windows-targets 0.52.6", +] + +[[package]] +name = "windows-sys" version = "0.60.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb" @@ -3680,18 +3960,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.40" +version = "0.8.42" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a789c6e490b576db9f7e6b6d661bcc9799f7c0ac8352f56ea20193b2681532e5" +checksum = "f2578b716f8a7a858b7f02d5bd870c14bf4ddbbcf3a4c05414ba6503640505e3" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.40" +version = "0.8.42" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f65c489a7071a749c849713807783f70672b28094011623e200cb86dcb835953" +checksum = "7e6cc098ea4d3bd6246687de65af3f920c430e236bee1e3bf2e441463f08a02f" dependencies = [ "proc-macro2", "quote", diff --git a/Cargo.toml b/Cargo.toml @@ -37,3 +37,11 @@ strum = "0.28" strum_macros = "0.28" aws-lc-rs = { version = "*" } serde = { version = "*", features = ["derive"] } +reedline = "*" +sqlx = { version = "0.8", default-features = false, features = [ + "postgres", + "runtime-tokio", + "tls-rustls-aws-lc-rs", +] } +compact_str = { version = "0.9.0", features = ["serde", "sqlx-postgres"] } +uuid = "*" +\ No newline at end of file diff --git a/src/config.rs b/src/config.rs @@ -19,8 +19,13 @@ use std::cell::OnceCell; +use taler_api::config::DbCfg; use taler_common::config::{Config, ValueErr}; +pub fn parse_db_cfg(cfg: &Config) -> Result<DbCfg, ValueErr> { + DbCfg::parse(cfg.section("libeufin-nexusdb-postgres")) +} + pub struct EbicsKeysCfg { pub bank_pub_keys_path: String, pub client_priv_keys_path: String, diff --git a/src/db.rs b/src/db.rs @@ -0,0 +1,2121 @@ +/* + This file is part of TALER + Copyright (C) 2026 Taler Systems SA + + TALER 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. + + TALER 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 + TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> +*/ + +use std::fmt::Display; + +use jiff::{Timestamp, civil::Date, tz::TimeZone}; +use serde::{Serialize, de::DeserializeOwned}; +use sqlx::{PgConnection, PgExecutor, PgPool, QueryBuilder, Row, postgres::PgRow}; +use taler_api::{ + db::{BindHelper, IncomingType, TypeHelper, history, page}, + subject::{IncomingSubject, OutgoingSubject}, +}; +use taler_common::{ + api_common::{HashCode, ShortHashCode}, + api_params::{History, Page}, + api_revenue::RevenueIncomingBankTransaction, + api_wire::{ + IncomingBankTransaction, OutgoingBankTransaction, TransferListStatus, TransferState, + TransferStatus, + }, + config::Config, + types::{ + amount::{Amount, Currency, Decimal}, + payto::PaytoImpl as _, + }, +}; +use tokio::sync::watch::{Receiver, Sender}; +use url::Url; + +use crate::{InitiatedPayment, OutgoingId, OutgoingPayment, config::parse_db_cfg}; + +const SCHEMA: &str = "libeufin_nexus"; + +pub async fn pool(cfg: &Config) -> anyhow::Result<PgPool> { + let db = parse_db_cfg(cfg)?; + let pool = taler_common::db::pool(db.cfg, SCHEMA).await?; + Ok(pool) +} + +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?; + Ok(pool) +} + +pub async fn notification_listener( + pool: PgPool, + in_channel: Sender<i64>, + taler_in_channel: Sender<i64>, + taler_out_channel: Sender<i64>, +) -> sqlx::Result<()> { + taler_api::notification::notification_listener!(&pool, + "nexus_revenue_tx" => (row_id: i64) { + in_channel.send_replace(row_id); + }, + "nexus_incoming_tx" => (row_id: i64) { + taler_in_channel.send_replace(row_id); + }, + "nexus_outgoing_tx" => (row_id: i64) { + taler_out_channel.send_replace(row_id); + } + ) +} + +/// Outgoing payments initiation result +#[derive(Debug, PartialEq, Eq)] +enum PaymentInitiationResult { + Success(u64), + RequestUidReuse, +} + +/// Initiate a new payment +async fn initiate( + pool: &PgPool, + payment: &InitiatedPayment, +) -> sqlx::Result<PaymentInitiationResult> { + let res = sqlx::query( + " + INSERT INTO initiated_outgoing_transactions ( + amount, + subject, + credit_payto, + initiation_time, + end_to_end_id + ) VALUES ($1,$2,$3,$4,$5) + RETURNING initiated_outgoing_transaction_id + ", + ) + .bind(&payment.amount) + .bind(&payment.subject) + .bind(payment.creditor.as_ref().as_str()) + .bind_timestamp(&payment.initiation_time) + .bind(&payment.end_to_end_id) + .try_map(|r: PgRow| Ok(PaymentInitiationResult::Success(r.try_get_u64(0)?))) + .fetch_one(pool) + .await; + if let Err(e) = &res + && let Some(db_err) = e.as_database_error() + && db_err.code() == Some(std::borrow::Cow::Borrowed("23505")) + { + Ok(PaymentInitiationResult::RequestUidReuse) + } else { + res + } +} + +/// Group unbatched transaction into a single batch +pub async fn batch_initiated( + pool: &PgPool, + timestamp: &Timestamp, + ebics_id: &str, + require_ack: bool, +) -> sqlx::Result<()> { + sqlx::query("SELECT batch_outgoing_transactions($1, $2, $3)") + .bind_timestamp(timestamp) + .bind(ebics_id) + .bind(require_ack) + .execute(pool) + .await?; + Ok(()) +} + +pub async fn initiated_ack(db: &PgPool, id: u64) -> sqlx::Result<()> { + sqlx::query("UPDATE initiated_outgoing_transactions SET awaiting_ack=false WHERE initiated_outgoing_transaction_id=$1") + .bind(id as i64) + .execute(db) + .await?; + Ok(()) +} + +pub async fn unsettled_tx_in_batch( + db: &PgPool, + currency: &Currency, + msg_id: &str, + execution_time: &Timestamp, +) -> sqlx::Result<Vec<OutgoingPayment>> { + sqlx::query( + " + SELECT + end_to_end_id, + amount, + subject, + credit_payto + 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') + " + ) + .bind(msg_id) + .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 + }, + 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")? + } + )) + .fetch_all(db) + .await +} + +#[derive(Debug, PartialEq, Eq)] +pub struct OutgoingRegistrationResult { + pub id: u64, + pub initiated: bool, + pub new: bool, +} + +/** Register an outgoing payment reconciling it with its initiated payment counterpart if present */ +pub async fn register_out_tx( + pool: &PgPool, + payment: &OutgoingPayment, + subject: Option<&OutgoingSubject>, +) -> sqlx::Result<OutgoingRegistrationResult> { + sqlx::query( + " + SELECT out_tx_id, out_initiated, out_found + FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) + ", + ) + .bind(&payment.amount) + .bind( + payment + .debit_fee + .as_ref() + .unwrap_or(&Amount::zero(&payment.amount.currency)), + ) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str())) + .bind(&payment.id.end_to_end_id) + .bind(&payment.id.msg_id) + .bind(&payment.id.acct_svcr_ref) + .bind(subject.as_ref().map(|s| &s.wtid)) + .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str())) + .bind(subject.as_ref().map(|s| &s.metadata)) + .try_map(|r: PgRow| { + Ok(OutgoingRegistrationResult { + id: r.try_get_u64(0)?, + initiated: r.try_get_flag(1)?, + new: !r.try_get_flag(2)?, + }) + }) + .fetch_one(pool) + .await +} + +/// Register an outgoing batch +pub async fn register_out_batch( + pool: &PgPool, + currency: &Currency, + payment: &OutgoingPayment, + subject: Option<&OutgoingSubject>, +) -> sqlx::Result<OutgoingRegistrationResult> { + // TODO + sqlx::query( + " + SELECT out_tx_id, out_initiated, out_found + FROM register_outgoing($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) + ", + ) + .bind(&payment.amount) + .bind( + payment + .debit_fee + .as_ref() + .unwrap_or(&Amount::zero(currency)), + ) + .bind(&payment.subject) + .bind_timestamp(&payment.execution_time) + .bind(payment.creditor.as_ref().map(|it| it.as_ref().as_str())) + .bind(&payment.id.end_to_end_id) + .bind(&payment.id.msg_id) + .bind(&payment.id.acct_svcr_ref) + .bind(subject.as_ref().map(|s| &s.wtid)) + .bind(subject.as_ref().map(|s| s.exchange_base_url.as_str())) + .bind(subject.as_ref().map(|s| &s.metadata)) + .try_map(|r: PgRow| { + Ok(OutgoingRegistrationResult { + id: r.try_get_u64(0)?, + initiated: r.try_get_flag(1)?, + new: !r.try_get_flag(2)?, + }) + }) + .fetch_one(pool) + .await +} + +pub async fn outgoing_history( + db: &PgPool, + currency: &Currency, + params: &History, + listen: impl FnOnce() -> Receiver<i64>, +) -> sqlx::Result<Vec<OutgoingBankTransaction>> { + history( + db, + "outgoing_transaction_id", + params, + listen, + || { + QueryBuilder::new( + " + SELECT + outgoing_transaction_id + ,execution_time + ,(amount).val AS amount_val + ,(amount).frac AS amount_frac + ,(debit_fee).val AS debit_fee_val + ,(debit_fee).frac AS debit_fee_frac + ,credit_payto + ,wtid + ,exchange_base_url + ,metadata + FROM talerable_outgoing_transactions + JOIN outgoing_transactions USING(outgoing_transaction_id) + WHERE + ", + ) + }, + |r: PgRow| { + let tmp: Amount = r.try_get_amount("debit_fee", currency)?; + + Ok(OutgoingBankTransaction { + row_id: r.try_get_safeu64("outgoing_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + debit_fee: if tmp == Amount::zero(currency) { + None + } else { + Some(tmp) + }, + credit_account: r.try_get_payto("credit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + exchange_base_url: r.try_get_url("exchange_base_url")?, + wtid: r.try_get("wtid")?, + }) + }, + ) + .await +} + +pub async fn incoming_history( + db: &PgPool, + currency: &Currency, + params: &History, + listen: impl FnOnce() -> Receiver<i64>, +) -> sqlx::Result<Vec<IncomingBankTransaction>> { + history( + db, + "incoming_transaction_id", + params, + listen, + || { + QueryBuilder::new( + " + SELECT + incoming_transaction_id + ,execution_time + ,(amount).val AS amount_val + ,(amount).frac AS amount_frac + ,(credit_fee).val AS credit_fee_val + ,(credit_fee).frac AS credit_fee_frac + ,debit_payto + ,type + ,metadata + ,authorization_pub + ,authorization_sig + FROM talerable_incoming_transactions + JOIN incoming_transactions USING(incoming_transaction_id) + WHERE + ", + ) + }, + |r: PgRow| { + let tmp: Amount = r.try_get_amount("credit_fee", currency)?; + Ok(match r.try_get("type")? { + IncomingType::reserve => IncomingBankTransaction::Reserve { + row_id: r.try_get_safeu64("incoming_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + credit_fee: if tmp == Amount::zero(&currency) { + None + } else { + Some(tmp) + }, + debit_account: r.try_get_payto("debit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + reserve_pub: r.try_get("metadata")?, + authorization_pub: r.try_get("authorization_pub")?, + authorization_sig: r.try_get("authorization_sig")?, + }, + IncomingType::kyc => IncomingBankTransaction::Kyc { + row_id: r.try_get_safeu64("incoming_transaction_id")?, + amount: r.try_get_amount("amount", currency)?, + credit_fee: if tmp == Amount::zero(currency) { + None + } else { + Some(tmp) + }, + debit_account: r.try_get_payto("debit_payto")?, + date: r.try_get_timestamp("execution_time")?.into(), + account_pub: r.try_get("metadata")?, + authorization_pub: r.try_get("authorization_pub")?, + authorization_sig: r.try_get("authorization_sig")?, + }, + IncomingType::wad => { + unimplemented!("WAD is not yet supported") + } + }) + }, + ) + .await +} + +#[cfg(test)] +mod test { + use std::sync::LazyLock; + use std::{array::repeat, str::FromStr}; + + use compact_str::CompactString; + use jiff::Timestamp; + use rand::prelude::IndexedRandom; + use rand::seq::SliceRandom; + use sqlx::PgPool; + use sqlx::{PgConnection, Row, postgres::PgRow}; + use taler_api::db::TypeHelper; + use taler_common::types::{amount::Currency, payto::IbanPayto}; + use taler_common::{ + api_common::ShortHashCode, + types::{ + amount::Amount, + payto::{BankID, FullIbanPayto}, + }, + }; + + use crate::{CONFIG_SOURCE, InitiatedPayment, OutgoingBatch, db::OutgoingRegistrationResult}; + use crate::{OutgoingId, db::PaymentInitiationResult, db::batch_initiated, db::initiate, db::initiated_ack}; + use crate::{OutgoingPayment, register_outgoing, register_outgoing_batch}; + + const EBICS_ID_ALPHABET: &[u8] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"; + pub static CURRENCY: LazyLock<Currency> = LazyLock::new(|| "KUDOS".parse().unwrap()); + + async fn setup() -> (PgConnection, PgPool) { + let pool = taler_test_utils::db::db_test_setup(CONFIG_SOURCE).await; + let conn = pool.acquire().await.unwrap().leak(); + (conn, pool) + } + + pub fn rand_ebics_id() -> CompactString { + let mut rng = rand::rng(); + (0..34) + .map(|_| *EBICS_ID_ALPHABET.choose(&mut rng).unwrap() as char) + .collect() + } + + /** Generates an outgoing payment, given its subject and end-to-end ID */ + pub fn gen_out_pay(subject: String) -> OutgoingPayment { + OutgoingPayment { + id: OutgoingId { + msg_id: None, + end_to_end_id: Some(rand_ebics_id()), + acct_svcr_ref: None, + }, + amount: Amount::from_str("KUDOS:44").unwrap(), + debit_fee: None, + creditor: Some( + IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test") + .unwrap() + .as_payto(), + ), + subject: Some(subject), + execution_time: Timestamp::now(), + } + } + + /** Generates a payment initiation, given its subject */ + pub fn gen_init_pay(end_to_end_id: CompactString, subject: String) -> InitiatedPayment { + InitiatedPayment { + id: 0, + amount: Amount::from_str("KUDOS:44").unwrap(), + creditor: IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test") + .unwrap() + .as_payto(), + subject, + initiation_time: Timestamp::now(), + end_to_end_id, + } + } + + async fn check_out_count(db: &PgPool, nb_incoming: u64, nb_talerable: u64) { + sqlx::query( + " + SELECT (SELECT count(*) FROM outgoing_transactions) AS incoming, + (SELECT count(*) FROM talerable_outgoing_transactions) AS talerable + ", + ) + .try_map(|r: PgRow| { + assert_eq!( + (r.try_get_u64(0)?, r.try_get_u64(1)?), + (nb_incoming, nb_talerable) + ); + Ok(()) + }) + .fetch_one(db) + .await + .unwrap(); + } + + #[tokio::test] + async fn out_tx() { + let (_, db) = setup().await; + // Register initiated transactions + for subject in [ + "initiated by nexus".to_owned(), + format!("{} https://exchange.com/", ShortHashCode::rand()), + ] { + let payment = gen_out_pay(subject.clone()); + assert!(matches!( + initiate( + &db, + &gen_init_pay(payment.id.end_to_end_id.clone().unwrap(), subject), + ) + .await, + Ok(PaymentInitiationResult::Success(_)) + )); + let first = register_outgoing(&db, &payment).await.unwrap(); + assert_eq!( + first, + OutgoingRegistrationResult { + id: first.id, + initiated: true, + new: true + } + ); + assert_eq!( + register_outgoing(&db, &payment).await.unwrap(), + OutgoingRegistrationResult { + id: first.id, + initiated: true, + new: false + } + ); + let payment = OutgoingPayment { + id: OutgoingId { + msg_id: None, + end_to_end_id: None, + acct_svcr_ref: payment.id.end_to_end_id, + }, + ..payment + }; + let second = register_outgoing(&db, &payment).await.unwrap(); + assert_eq!( + second, + OutgoingRegistrationResult { + id: first.id + 1, + initiated: false, + new: true + } + ); + assert_eq!( + register_outgoing(&db, &payment).await.unwrap(), + OutgoingRegistrationResult { + id: second.id, + initiated: false, + new: false + } + ); + } + check_out_count(&db, 4, 1).await; + + // Register unknown + for subject in [ + "initiated by nexus".to_owned(), + format!("{} https://exchange.com/", ShortHashCode::rand()), + ] { + let payment = gen_out_pay(subject.clone()); + let res = register_outgoing(&db, &payment).await.unwrap(); + assert_eq!( + res, + OutgoingRegistrationResult { + id: res.id, + initiated: false, + new: true + } + ); + assert_eq!( + register_outgoing(&db, &payment).await.unwrap(), + OutgoingRegistrationResult { + id: res.id, + initiated: false, + new: false + } + ); + } + check_out_count(&db, 6, 2).await; + + // Register wtid reuse + let wtid = ShortHashCode::rand(); + for subject in [ + format!("{wtid} https://exchange.com/"), + format!("{wtid} https://exchange.com/"), + ] { + let payment = gen_out_pay(subject.clone()); + let res = register_outgoing(&db, &payment).await.unwrap(); + assert_eq!( + res, + OutgoingRegistrationResult { + id: res.id, + initiated: false, + new: true + } + ); + assert_eq!( + register_outgoing(&db, &payment).await.unwrap(), + OutgoingRegistrationResult { + id: res.id, + initiated: false, + new: false + } + ); + } + check_out_count(&db, 8, 3).await + } + + #[tokio::test] + async fn out_batch() { + let (_, db) = setup().await; + // Init batch + let wtid = ShortHashCode::rand(); + for subject in [ + format!("initiated by nexus"), + format!("{} https://exchange.com/", ShortHashCode::rand()), + format!("{wtid} https://exchange.com/"), + format!("{wtid} https://exchange.com/"), + ] { + assert!(matches!( + initiate(&db, &gen_init_pay(rand_ebics_id(), subject),).await, + Ok(PaymentInitiationResult::Success(_)) + )); + } + batch_initiated(&db, &Timestamp::now(), "BATCH", false) + .await + .unwrap(); + + // Register batch + register_outgoing_batch( + &db, + &CURRENCY, + &OutgoingBatch { + msg_id: "BATCH".into(), + execution_time: Timestamp::now(), + }, + ) + .await + .unwrap(); + check_out_count(&db, 4, 2).await; + + // Test manual ack + let mut txs = Vec::new(); + for nb in 0..3 { + let res = initiate(&db, &gen_init_pay(rand_ebics_id(), format!("tx {nb}"))).await; + if let Ok(PaymentInitiationResult::Success(id)) = &res { + txs.push(*id); + } else { + panic!("Expected success got {res:?}"); + } + } + + // Check not sent without ack + batch_initiated(&db, &Timestamp::now(), "BATCH_MANUAL", true) + .await + .unwrap(); + register_outgoing_batch( + &db, + &CURRENCY, + &OutgoingBatch { + msg_id: "BATCH_MANUAL".into(), + execution_time: Timestamp::now(), + }, + ) + .await + .unwrap(); + check_out_count(&db, 4, 2).await; + + // Check sent with ack + for tx in txs { + initiated_ack(&db, tx).await.unwrap(); + } + batch_initiated(&db, &Timestamp::now(), "BATCH_MANUAL", true) + .await + .unwrap(); + register_outgoing_batch( + &db, + &CURRENCY, + &OutgoingBatch { + msg_id: "BATCH_MANUAL".into(), + execution_time: Timestamp::now(), + }, + ) + .await + .unwrap(); + check_out_count(&db, 7, 2).await; + } +} + +/* +#[derive(Debug, Clone)] +pub struct TxIn { + pub code: u64, + pub amount: Amount, + pub subject: String, + pub debtor: FullHuPayto, + pub value_date: Date, + pub status: TxStatus, +} + +impl Display for TxIn { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let Self { + code, + amount, + subject, + debtor, + value_date, + status, + } = self; + write!( + f, + "{value_date} {code} {amount} ({} {}) {status:?} '{subject}'", + debtor.bban(), + debtor.name + ) + } +} + +#[derive(Debug, Clone)] +pub struct TxOut { + pub code: u64, + pub amount: Amount, + pub subject: String, + pub creditor: FullHuPayto, + pub value_date: Date, + pub status: TxStatus, +} + +impl Display for TxOut { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let Self { + code, + amount, + subject, + creditor, + value_date, + status, + } = self; + write!( + f, + "{value_date} {code} {amount} ({} {}) {status:?} '{subject}'", + creditor.bban(), + &creditor.name + ) + } +} + +#[derive(Debug, PartialEq, Eq)] +pub struct Initiated { + pub id: u64, + pub amount: Amount, + pub subject: String, + pub creditor: FullHuPayto, +} + +impl Display for Initiated { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let Self { + id, + amount, + subject, + creditor, + } = self; + write!( + f, + "{id} {amount} ({} {}) '{subject}'", + creditor.bban(), + &creditor.name + ) + } +} + +#[derive(Debug, Clone)] +pub struct TxInAdmin { + pub amount: Amount, + pub subject: String, + pub debtor: FullHuPayto, + pub metadata: IncomingSubject, +} + +/// Lock the database for worker execution +pub async fn worker_lock(e: &mut PgConnection) -> sqlx::Result<bool> { + sqlx::query("SELECT pg_try_advisory_lock(42)") + .try_map(|r: PgRow| r.try_get(0)) + .fetch_one(e) + .await +} + +#[derive(Debug, PartialEq, Eq)] +pub enum AddIncomingResult { + Success { + new: bool, + row_id: u64, + valued_at: Date, + }, + ReservePubReuse, +} + +pub async fn register_tx_in_admin( + db: &PgPool, + tx: &TxInAdmin, + now: &Timestamp, +) -> sqlx::Result<AddIncomingResult> { + sqlx::query( + " + SELECT out_reserve_pub_reuse, out_tx_row_id, out_valued_at, out_new + FROM register_tx_in(NULL, ($1, $2)::taler_amount, $3, $4, $5, $6, $7, $8, $6) + ", + ) + .bind_amount(&tx.amount) + .bind(&tx.subject) + .bind(tx.debtor.iban()) + .bind(&tx.debtor.name) + .bind_date(&now.to_zoned(TimeZone::UTC).date()) + .bind(tx.metadata.ty()) + .bind(tx.metadata.key()) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag(0)? { + AddIncomingResult::ReservePubReuse + } else { + AddIncomingResult::Success { + row_id: r.try_get_u64(1)?, + valued_at: r.try_get_date(2)?, + new: r.try_get(3)?, + } + }) + }) + .fetch_one(db) + .await +} + +pub async fn register_tx_in( + db: &mut PgConnection, + tx: &TxIn, + subject: &Option<IncomingSubject>, + now: &Timestamp, +) -> sqlx::Result<AddIncomingResult> { + sqlx::query( + " + SELECT out_reserve_pub_reuse, out_tx_row_id, out_valued_at, out_new + FROM register_tx_in($1, ($2, $3)::taler_amount, $4, $5, $6, $7, $8, $9, $10) + ", + ) + .bind(tx.code as i64) + .bind_amount(&tx.amount) + .bind(&tx.subject) + .bind(tx.debtor.iban()) + .bind(&tx.debtor.name) + .bind_date(&tx.value_date) + .bind(subject.as_ref().map(|it| it.ty())) + .bind(subject.as_ref().map(|it| it.key())) + .bind_timestamp(now) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag(0)? { + AddIncomingResult::ReservePubReuse + } else { + AddIncomingResult::Success { + row_id: r.try_get_u64(1)?, + valued_at: r.try_get_date(2)?, + new: r.try_get(3)?, + } + }) + }) + .fetch_one(db) + .await +} + +#[derive(Debug)] +pub enum TxOutKind { + Simple, + Bounce(u32), + Talerable(OutgoingSubject), +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, sqlx::Type)] +#[allow(non_camel_case_types)] +#[sqlx(type_name = "register_result")] +pub enum RegisterResult { + /// Already registered + idempotent, + /// Initiated transaction + known, + /// Recovered unknown outgoing transaction + recovered, +} + +#[derive(Debug, PartialEq, Eq)] +pub struct AddOutgoingResult { + pub result: RegisterResult, + pub row_id: u64, +} + +pub async fn register_tx_out( + db: &mut PgConnection, + tx: &TxOut, + kind: &TxOutKind, + now: &Timestamp, +) -> sqlx::Result<AddOutgoingResult> { + let query = sqlx::query( + " + SELECT out_result, out_tx_row_id + FROM register_tx_out($1, ($2, $3)::taler_amount, $4, $5, $6, $7, $8, $9, $10, $11) + ", + ) + .bind(tx.code as i64) + .bind_amount(&tx.amount) + .bind(&tx.subject) + .bind(tx.creditor.iban()) + .bind(&tx.creditor.name) + .bind_date(&tx.value_date); + let query = match kind { + TxOutKind::Simple => query + .bind(None::<&[u8]>) + .bind(None::<&str>) + .bind(None::<i64>), + TxOutKind::Bounce(bounced) => query + .bind(None::<&[u8]>) + .bind(None::<&str>) + .bind(*bounced as i64), + TxOutKind::Talerable(subject) => query + .bind(subject.0.as_ref()) + .bind(subject.1.as_ref()) + .bind(None::<i64>), + }; + query + .bind_timestamp(now) + .try_map(|r: PgRow| { + Ok(AddOutgoingResult { + result: r.try_get(0)?, + row_id: r.try_get_u64(1)?, + }) + }) + .fetch_one(db) + .await +} + +#[derive(Debug, PartialEq, Eq)] +pub struct OutFailureResult { + pub initiated_id: Option<u64>, + pub new: bool, +} + +pub async fn register_tx_out_failure( + db: &mut PgConnection, + code: u64, + bounced: Option<u32>, + now: &Timestamp, +) -> sqlx::Result<OutFailureResult> { + sqlx::query( + " + SELECT out_new, out_initiated_id + FROM register_tx_out_failure($1, $2, $3) + ", + ) + .bind(code as i64) + .bind(bounced.map(|i| i as i32)) + .bind_timestamp(now) + .try_map(|r: PgRow| { + Ok(OutFailureResult { + new: r.try_get(0)?, + initiated_id: r.try_get::<Option<i64>, _>(1)?.map(|i| i as u64), + }) + }) + .fetch_one(db) + .await +} + +#[derive(Debug, PartialEq, Eq)] +pub enum TransferResult { + Success { id: u64, initiated_at: Timestamp }, + RequestUidReuse, + WtidReuse, +} + +#[derive(Debug, Clone)] +pub struct Transfer { + pub request_uid: HashCode, + pub amount: Decimal, + pub exchange_base_url: Url, + pub wtid: ShortHashCode, + pub creditor: FullHuPayto, +} + +pub async fn make_transfer<'a>( + db: impl PgExecutor<'a>, + tx: &Transfer, + now: &Timestamp, +) -> sqlx::Result<TransferResult> { + let subject = format!("{} {}", tx.wtid, tx.exchange_base_url); + sqlx::query( + " + SELECT out_request_uid_reuse, out_wtid_reuse, out_initiated_row_id, out_initiated_at + FROM taler_transfer($1, $2, $3, ($4, $5)::taler_amount, $6, $7, $8, $9) + ", + ) + .bind(tx.request_uid.as_ref()) + .bind(tx.wtid.as_ref()) + .bind(&subject) + .bind_decimal(&tx.amount) + .bind(tx.exchange_base_url.as_str()) + .bind(tx.creditor.iban()) + .bind(&tx.creditor.name) + .bind_timestamp(now) + .try_map(|r: PgRow| { + Ok(if r.try_get_flag(0)? { + TransferResult::RequestUidReuse + } else if r.try_get_flag(1)? { + TransferResult::WtidReuse + } else { + TransferResult::Success { + id: r.try_get_u64(2)?, + initiated_at: r.try_get_timestamp(3)?, + } + }) + }) + .fetch_one(db) + .await +} + +#[derive(Debug, PartialEq, Eq)] +pub struct BounceResult { + pub tx_id: u64, + pub tx_new: bool, + pub bounce_id: u64, + pub bounce_new: bool, +} + +pub async fn register_bounce_tx_in( + db: &mut PgConnection, + tx: &TxIn, + amount: &Amount, + reason: &str, + now: &Timestamp, +) -> sqlx::Result<BounceResult> { + sqlx::query( + " + SELECT out_tx_row_id, out_tx_new, out_bounce_row_id, out_bounce_new + FROM register_bounce_tx_in($1, ($2, $3)::taler_amount, $4, $5, $6, $7, ($8, $9)::taler_amount, $10, $11) + ", + ) + .bind(tx.code as i64) + .bind_amount(&tx.amount) + .bind(&tx.subject) + .bind(tx.debtor.iban()) + .bind(&tx.debtor.name) + .bind_date(&tx.value_date) + .bind_amount(amount) + .bind(reason) + .bind_timestamp(now) + .try_map(|r: PgRow| { + Ok(BounceResult { + tx_id: r.try_get_u64(0)?, + tx_new: r.try_get(1)?, + bounce_id: r.try_get_u64(2)?, + bounce_new: r.try_get(3)?, + }) + }) + .fetch_one(db) + .await +} + +pub async fn transfer_page<'a>( + db: impl PgExecutor<'a>, + status: &Option<TransferState>, + params: &Page, +) -> sqlx::Result<Vec<TransferListStatus>> { + page( + db, + "initiated_id", + params, + || { + let mut builder = QueryBuilder::new( + " + SELECT + initiated_id, + status, + (amount).val as amount_val, + (amount).frac as amount_frac, + credit_account, + credit_name, + initiated_at + FROM transfer + JOIN initiated USING (initiated_id) + WHERE + ", + ); + if let Some(status) = status { + builder.push(" status = ").push_bind(status).push(" AND "); + } + builder + }, + |r: PgRow| { + Ok(TransferListStatus { + row_id: r.try_get_safeu64(0)?, + status: r.try_get(1)?, + amount: r.try_get_amount_i(2, &CURRENCY)?, + credit_account: r.try_get_iban(4)?.as_full_payto(r.try_get(5)?), + timestamp: r.try_get_timestamp(6)?.into(), + }) + }, + ) + .await +} +pub async fn transfer_by_id<'a>( + db: impl PgExecutor<'a>, + id: u64, +) -> sqlx::Result<Option<TransferStatus>> { + sqlx::query( + " + SELECT + status, + status_msg, + (amount).val as amount_val, + (amount).frac as amount_frac, + exchange_base_url, + wtid, + credit_account, + credit_name, + initiated_at + FROM transfer + JOIN initiated USING (initiated_id) + WHERE initiated_id = $1 + ", + ) + .bind(id as i64) + .try_map(|r: PgRow| { + Ok(TransferStatus { + status: r.try_get(0)?, + status_msg: r.try_get(1)?, + amount: r.try_get_amount_i(2, &CURRENCY)?, + origin_exchange_url: r.try_get(4)?, + wtid: r.try_get_base32(5)?, + credit_account: r.try_get_iban(6)?.as_full_payto(r.try_get(7)?), + timestamp: r.try_get_timestamp(8)?.into(), + }) + }) + .fetch_optional(db) + .await +} + +/** Get a batch of pending initiated transactions not attempted since [start] */ +pub async fn pending_batch<'a>( + db: impl PgExecutor<'a>, + start: &Timestamp, +) -> sqlx::Result<Vec<Initiated>> { + sqlx::query( + " + SELECT initiated_id, (amount).val, (amount).frac, subject, credit_account, credit_name + FROM initiated + WHERE magnet_code IS NULL + AND status='pending' + AND (last_submitted IS NULL OR last_submitted < $1) + LIMIT 100 + ", + ) + .bind_timestamp(start) + .try_map(|r: PgRow| { + Ok(Initiated { + id: r.try_get_u64(0)?, + amount: r.try_get_amount_i(1, &CURRENCY)?, + subject: r.try_get(3)?, + creditor: FullHuPayto::new(r.try_get_parse(4)?, r.try_get(5)?), + }) + }) + .fetch_all(db) + .await +} + +/** Get an initiated transaction matching the given magnet [code] */ +pub async fn initiated_by_code<'a>( + db: impl PgExecutor<'a>, + code: u64, +) -> sqlx::Result<Option<Initiated>> { + sqlx::query( + " + SELECT initiated_id, (amount).val, (amount).frac, subject, credit_account, credit_name + FROM initiated + WHERE magnet_code IS $1 + ", + ) + .bind(code as i64) + .try_map(|r: PgRow| { + Ok(Initiated { + id: r.try_get_u64(0)?, + amount: r.try_get_amount_i(1, &CURRENCY)?, + subject: r.try_get(3)?, + creditor: FullHuPayto::new(r.try_get_parse(4)?, r.try_get(5)?), + }) + }) + .fetch_optional(db) + .await +} + +/** Update status of a successful submitted initiated transaction */ +pub async fn initiated_submit_success<'a>( + db: impl PgExecutor<'a>, + id: u64, + timestamp: &Timestamp, + magnet_code: u64, +) -> sqlx::Result<()> { + sqlx::query( + " + UPDATE initiated + SET status='pending', submission_counter=submission_counter+1, last_submitted=$1, magnet_code=$2 + WHERE initiated_id=$3 + " + ).bind_timestamp(timestamp) + .bind(magnet_code as i64) + .bind(id as i64) + .execute(db).await?; + Ok(()) +} + +/** Update status of a permanently failed initiated transaction */ +pub async fn initiated_submit_permanent_failure<'a>( + db: impl PgExecutor<'a>, + id: u64, + timestamp: &Timestamp, + msg: &str, +) -> sqlx::Result<()> { + sqlx::query( + " + UPDATE initiated + SET status='permanent_failure', status_msg=$2 + WHERE initiated_id=$3 + ", + ) + .bind_timestamp(timestamp) + .bind(msg) + .bind(id as i64) + .execute(db) + .await?; + Ok(()) +} + +/** Check if an initiated transaction exist for a magnet code */ +pub async fn initiated_exists_for_code<'a>( + db: impl PgExecutor<'a>, + code: u64, +) -> sqlx::Result<Option<u64>> { + sqlx::query("SELECT initiated_id FROM initiated WHERE magnet_code=$1") + .bind(code as i64) + .try_map(|r| Ok(r.try_get::<i64, _>(0)? as u64)) + .fetch_optional(db) + .await +} + +/** Get JSON value from KV table */ +pub async fn kv_get<'a, T: DeserializeOwned + Unpin + Send>( + db: impl PgExecutor<'a>, + key: &str, +) -> sqlx::Result<Option<T>> { + sqlx::query("SELECT value FROM kv WHERE key=$1") + .bind(key) + .try_map(|r| Ok(r.try_get::<sqlx::types::Json<T>, _>(0)?.0)) + .fetch_optional(db) + .await +} + +/** Set JSON value in KV table */ +pub async fn kv_set<'a, T: Serialize>( + db: impl PgExecutor<'a>, + key: &str, + value: &T, +) -> sqlx::Result<()> { + sqlx::query("INSERT INTO kv (key, value) VALUES ($1, $2) ON CONFLICT (key) DO UPDATE SET value=EXCLUDED.value") + .bind(key) + .bind(sqlx::types::Json(value)) + .execute(db) + .await?; + Ok(()) +} + +#[cfg(test)] +mod test { + use jiff::{Span, Timestamp, Zoned}; + use serde_json::json; + use sqlx::{PgConnection, PgPool, postgres::PgRow}; + use taler_api::{ + db::TypeHelper, + subject::{IncomingSubject, OutgoingSubject}, + }; + use taler_common::{ + api_common::{EddsaPublicKey, HashCode, ShortHashCode}, + api_params::{History, Page}, + types::{ + amount::{amount, decimal}, + url, + utils::now_sql_stable_timestamp, + }, + }; + use tokio::sync::watch::Receiver; + + use crate::{ + constants::CONFIG_SOURCE, + db::{ + self, AddIncomingResult, AddOutgoingResult, BounceResult, Initiated, OutFailureResult, + TransferResult, TxIn, TxOut, TxOutKind, kv_get, kv_set, make_transfer, + register_bounce_tx_in, register_tx_in, register_tx_in_admin, register_tx_out, + }, + magnet_api::types::TxStatus, + magnet_payto, + }; + + use super::TxInAdmin; + + fn fake_listen<T: Default>() -> Receiver<T> { + tokio::sync::watch::channel(T::default()).1 + } + + async fn setup() -> (PgConnection, PgPool) { + let pool = taler_test_utils::db::db_test_setup(CONFIG_SOURCE).await; + let conn = pool.acquire().await.unwrap().leak(); + (conn, pool) + } + + #[tokio::test] + async fn kv() { + let (mut db, _) = setup().await; + + let value = json!({ + "name": "Mr Smith", + "no way": 32 + }); + + assert_eq!( + kv_get::<serde_json::Value>(&mut db, "value").await.unwrap(), + None + ); + kv_set(&mut db, "value", &value).await.unwrap(); + kv_set(&mut db, "value", &value).await.unwrap(); + assert_eq!( + kv_get::<serde_json::Value>(&mut db, "value").await.unwrap(), + Some(value) + ); + } + + #[tokio::test] + async fn tx_in() { + let (mut db, pool) = setup().await; + + async fn routine( + db: &mut PgConnection, + first: &Option<IncomingSubject>, + second: &Option<IncomingSubject>, + ) { + let (id, code) = + sqlx::query("SELECT count(*) + 1, COALESCE(max(magnet_code), 0) + 20 FROM tx_in") + .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get_u64(1)?))) + .fetch_one(&mut *db) + .await + .unwrap(); + let now = now_sql_stable_timestamp(); + let date = Zoned::now().date(); + let later = date.tomorrow().unwrap(); + let tx = TxIn { + code: code, + amount: amount("EUR:10"), + subject: "subject".to_owned(), + debtor: magnet_payto( + "payto://iban/HU30162000031000163100000000?receiver-name=name", + ), + value_date: date, + status: TxStatus::Completed, + }; + // Insert + assert_eq!( + register_tx_in(db, &tx, &first, &now) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: true, + row_id: id, + valued_at: date + } + ); + // Idempotent + assert_eq!( + register_tx_in( + db, + &TxIn { + value_date: later, + ..tx.clone() + }, + &first, + &now + ) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: false, + row_id: id, + valued_at: date + } + ); + // Many + assert_eq!( + register_tx_in( + db, + &TxIn { + code: code + 1, + value_date: later, + ..tx + }, + &second, + &now + ) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: true, + row_id: id + 1, + valued_at: later + } + ); + } + + // Empty db + assert_eq!( + db::revenue_history(&pool, &History::default(), fake_listen) + .await + .unwrap(), + Vec::new() + ); + assert_eq!( + db::incoming_history(&pool, &History::default(), fake_listen) + .await + .unwrap(), + Vec::new() + ); + + // Regular transaction + routine(&mut db, &None, &None).await; + + // Reserve transaction + routine( + &mut db, + &Some(IncomingSubject::Reserve(EddsaPublicKey::rand())), + &Some(IncomingSubject::Reserve(EddsaPublicKey::rand())), + ) + .await; + + // Kyc transaction + routine( + &mut db, + &Some(IncomingSubject::Kyc(EddsaPublicKey::rand())), + &Some(IncomingSubject::Kyc(EddsaPublicKey::rand())), + ) + .await; + + // History + assert_eq!( + db::revenue_history(&pool, &History::default(), fake_listen) + .await + .unwrap() + .len(), + 6 + ); + assert_eq!( + db::incoming_history(&pool, &History::default(), fake_listen) + .await + .unwrap() + .len(), + 4 + ); + } + + #[tokio::test] + async fn tx_in_admin() { + let (_, pool) = setup().await; + + // Empty db + assert_eq!( + db::incoming_history(&pool, &History::default(), fake_listen) + .await + .unwrap(), + Vec::new() + ); + + let now = now_sql_stable_timestamp(); + let later = now + Span::new().hours(2); + let date = Zoned::now().date(); + let tx = TxInAdmin { + amount: amount("EUR:10"), + subject: "subject".to_owned(), + debtor: magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"), + metadata: IncomingSubject::Reserve(EddsaPublicKey::rand()), + }; + // Insert + assert_eq!( + register_tx_in_admin(&pool, &tx, &now) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: true, + row_id: 1, + valued_at: date + } + ); + // Idempotent + assert_eq!( + register_tx_in_admin(&pool, &tx, &later) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: false, + row_id: 1, + valued_at: date + } + ); + // Many + assert_eq!( + register_tx_in_admin( + &pool, + &TxInAdmin { + subject: "Other".to_owned(), + metadata: IncomingSubject::Reserve(EddsaPublicKey::rand()), + ..tx.clone() + }, + &later + ) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: true, + row_id: 2, + valued_at: date + } + ); + + // History + assert_eq!( + db::incoming_history(&pool, &History::default(), fake_listen) + .await + .unwrap() + .len(), + 2 + ); + } + + #[tokio::test] + async fn tx_out() { + let (mut db, pool) = setup().await; + + async fn routine(db: &mut PgConnection, first: &TxOutKind, second: &TxOutKind) { + let (id, code) = + sqlx::query("SELECT count(*) + 1, COALESCE(max(magnet_code), 0) + 20 FROM tx_out") + .try_map(|r: PgRow| Ok((r.try_get_u64(0)?, r.try_get_u64(1)?))) + .fetch_one(&mut *db) + .await + .unwrap(); + let now = now_sql_stable_timestamp(); + let date = Zoned::now().date(); + let later = date.tomorrow().unwrap(); + let tx = TxOut { + code, + amount: amount("HUF:10"), + subject: "subject".to_owned(), + creditor: magnet_payto( + "payto://iban/HU30162000031000163100000000?receiver-name=name", + ), + value_date: date, + status: TxStatus::Completed, + }; + assert!(matches!( + make_transfer( + &mut *db, + &db::Transfer { + request_uid: HashCode::rand(), + amount: decimal("10"), + exchange_base_url: url("https://exchange.test.com/"), + wtid: ShortHashCode::rand(), + creditor: tx.creditor.clone() + }, + &now + ) + .await + .unwrap(), + TransferResult::Success { .. } + )); + db::initiated_submit_success(&mut *db, 1, &Timestamp::now(), tx.code) + .await + .expect("status success"); + + // Insert + assert_eq!( + register_tx_out(&mut *db, &tx, first, &now) + .await + .expect("register tx out"), + AddOutgoingResult { + result: db::RegisterResult::known, + row_id: id, + } + ); + // Idempotent + assert_eq!( + register_tx_out( + &mut *db, + &TxOut { + value_date: later, + ..tx.clone() + }, + first, + &now + ) + .await + .expect("register tx out"), + AddOutgoingResult { + result: db::RegisterResult::idempotent, + row_id: id, + } + ); + // Recovered + assert_eq!( + register_tx_out( + &mut *db, + &TxOut { + code: code + 1, + value_date: later, + ..tx.clone() + }, + second, + &now + ) + .await + .expect("register tx out"), + AddOutgoingResult { + result: db::RegisterResult::recovered, + row_id: id + 1, + } + ); + } + + // Empty db + assert_eq!( + db::outgoing_history(&pool, &History::default(), fake_listen) + .await + .unwrap(), + Vec::new() + ); + + // Regular transaction + routine(&mut db, &TxOutKind::Simple, &TxOutKind::Simple).await; + + // Talerable transaction + routine( + &mut db, + &TxOutKind::Talerable(OutgoingSubject( + ShortHashCode::rand(), + url("https://exchange.com"), + )), + &TxOutKind::Talerable(OutgoingSubject( + ShortHashCode::rand(), + url("https://exchange.com"), + )), + ) + .await; + + // Bounced transaction + routine(&mut db, &TxOutKind::Bounce(21), &TxOutKind::Bounce(42)).await; + + // History + assert_eq!( + db::outgoing_history(&pool, &History::default(), fake_listen) + .await + .unwrap() + .len(), + 2 + ); + } + + #[tokio::test] + async fn tx_out_failure() { + let (mut db, _) = setup().await; + + let now = now_sql_stable_timestamp(); + + // Unknown + assert_eq!( + db::register_tx_out_failure(&mut db, 42, None, &now) + .await + .unwrap(), + OutFailureResult { + initiated_id: None, + new: false + } + ); + assert_eq!( + db::register_tx_out_failure(&mut db, 42, Some(12), &now) + .await + .unwrap(), + OutFailureResult { + initiated_id: None, + new: false + } + ); + + // Initiated + let req = db::Transfer { + request_uid: HashCode::rand(), + amount: decimal("10"), + exchange_base_url: url("https://exchange.test.com/"), + wtid: ShortHashCode::rand(), + creditor: magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"), + }; + let payto = magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"); + assert_eq!( + make_transfer(&mut db, &req, &now).await.unwrap(), + TransferResult::Success { + id: 1, + initiated_at: now + } + ); + db::initiated_submit_success(&mut db, 1, &Timestamp::now(), 34) + .await + .expect("status success"); + assert_eq!( + db::register_tx_out_failure(&mut db, 34, None, &now) + .await + .unwrap(), + OutFailureResult { + initiated_id: Some(1), + new: true + } + ); + assert_eq!( + db::register_tx_out_failure(&mut db, 34, None, &now) + .await + .unwrap(), + OutFailureResult { + initiated_id: Some(1), + new: false + } + ); + + // Recovered bounce + let tx = TxIn { + code: 12, + amount: amount("HUF:11"), + subject: "malformed transaction".to_owned(), + debtor: payto, + value_date: Zoned::now().date(), + status: TxStatus::Completed, + }; + assert_eq!( + db::register_bounce_tx_in(&mut db, &tx, &tx.amount, "no reason", &now) + .await + .unwrap(), + BounceResult { + tx_id: 1, + tx_new: true, + bounce_id: 2, + bounce_new: true + } + ); + assert_eq!( + db::register_tx_out_failure(&mut db, 10, Some(12), &now) + .await + .unwrap(), + OutFailureResult { + initiated_id: Some(2), + new: true + } + ); + assert_eq!( + db::register_tx_out_failure(&mut db, 10, Some(12), &now) + .await + .unwrap(), + OutFailureResult { + initiated_id: Some(2), + new: false + } + ); + } + + #[tokio::test] + async fn transfer() { + let (mut db, _) = setup().await; + + // Empty db + assert_eq!(db::transfer_by_id(&mut db, 0).await.unwrap(), None); + assert_eq!( + db::transfer_page(&mut db, &None, &Page::default()) + .await + .unwrap(), + Vec::new() + ); + + let req = db::Transfer { + request_uid: HashCode::rand(), + amount: decimal("10"), + exchange_base_url: url("https://exchange.test.com/"), + wtid: ShortHashCode::rand(), + creditor: magnet_payto("payto://iban/HU02162000031000164800000000?receiver-name=name"), + }; + let now = now_sql_stable_timestamp(); + let later = now + Span::new().hours(2); + // Insert + assert_eq!( + make_transfer(&mut db, &req, &now).await.expect("transfer"), + TransferResult::Success { + id: 1, + initiated_at: now.into() + } + ); + // Idempotent + assert_eq!( + make_transfer(&mut db, &req, &later) + .await + .expect("transfer"), + TransferResult::Success { + id: 1, + initiated_at: now.into() + } + ); + // Request UID reuse + assert_eq!( + make_transfer( + &mut db, + &db::Transfer { + wtid: ShortHashCode::rand(), + ..req.clone() + }, + &now + ) + .await + .expect("transfer"), + TransferResult::RequestUidReuse + ); + // wtid reuse + assert_eq!( + make_transfer( + &mut db, + &db::Transfer { + request_uid: HashCode::rand(), + ..req.clone() + }, + &now + ) + .await + .expect("transfer"), + TransferResult::WtidReuse + ); + // Many + assert_eq!( + make_transfer( + &mut db, + &db::Transfer { + request_uid: HashCode::rand(), + wtid: ShortHashCode::rand(), + ..req + }, + &later + ) + .await + .expect("transfer"), + TransferResult::Success { + id: 2, + initiated_at: later.into() + } + ); + + // Get + assert!(db::transfer_by_id(&mut db, 1).await.unwrap().is_some()); + assert!(db::transfer_by_id(&mut db, 2).await.unwrap().is_some()); + assert!(db::transfer_by_id(&mut db, 3).await.unwrap().is_none()); + assert_eq!( + db::transfer_page(&mut db, &None, &Page::default()) + .await + .unwrap() + .len(), + 2 + ); + } + + #[tokio::test] + async fn bounce() { + let (mut db, _) = setup().await; + + let amount = amount("HUF:10"); + let payto = magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"); + let now = now_sql_stable_timestamp(); + let date = Zoned::now().date(); + + // Empty db + assert!(db::pending_batch(&mut db, &now).await.unwrap().is_empty()); + + // Insert + assert_eq!( + register_tx_in( + &mut db, + &TxIn { + code: 13, + amount: amount.clone(), + subject: "subject".to_owned(), + debtor: payto.clone(), + value_date: date, + status: TxStatus::Completed + }, + &None, + &now + ) + .await + .expect("register tx in"), + AddIncomingResult::Success { + new: true, + row_id: 1, + valued_at: date + } + ); + + // Bounce + assert_eq!( + register_bounce_tx_in( + &mut db, + &TxIn { + code: 12, + amount: amount.clone(), + subject: "subject".to_owned(), + debtor: payto.clone(), + value_date: date, + status: TxStatus::Completed + }, + &amount, + "good reason", + &now + ) + .await + .expect("bounce"), + BounceResult { + tx_id: 2, + tx_new: true, + bounce_id: 1, + bounce_new: true + } + ); + // Idempotent + assert_eq!( + register_bounce_tx_in( + &mut db, + &TxIn { + code: 12, + amount: amount.clone(), + subject: "subject".to_owned(), + debtor: payto.clone(), + value_date: date, + status: TxStatus::Completed + }, + &amount, + "good reason", + &now + ) + .await + .expect("bounce"), + BounceResult { + tx_id: 2, + tx_new: false, + bounce_id: 1, + bounce_new: false + } + ); + + // Bounce registered + assert_eq!( + register_bounce_tx_in( + &mut db, + &TxIn { + code: 13, + amount: amount.clone(), + subject: "subject".to_owned(), + debtor: payto.clone(), + value_date: date, + status: TxStatus::Completed + }, + &amount, + "good reason", + &now + ) + .await + .expect("bounce"), + BounceResult { + tx_id: 1, + tx_new: false, + bounce_id: 2, + bounce_new: true + } + ); + // Idempotent registered + assert_eq!( + register_bounce_tx_in( + &mut db, + &TxIn { + code: 13, + amount: amount.clone(), + subject: "subject".to_owned(), + debtor: payto.clone(), + value_date: date, + status: TxStatus::Completed + }, + &amount, + "good reason", + &now + ) + .await + .expect("bounce"), + BounceResult { + tx_id: 1, + tx_new: false, + bounce_id: 2, + bounce_new: false + } + ); + + // Batch + assert_eq!( + db::pending_batch(&mut db, &now).await.unwrap(), + &[ + Initiated { + id: 1, + amount: amount.clone(), + subject: "bounce: 12".to_owned(), + creditor: payto.clone() + }, + Initiated { + id: 2, + amount, + subject: "bounce: 13".to_owned(), + creditor: payto + } + ] + ); + } + + #[tokio::test] + async fn status() { + let (mut db, _) = setup().await; + + // Unknown transfer + db::initiated_submit_permanent_failure(&mut db, 1, &Timestamp::now(), "msg") + .await + .unwrap(); + db::initiated_submit_success(&mut db, 1, &Timestamp::now(), 12) + .await + .unwrap(); + } + + #[tokio::test] + async fn batch() { + let (mut db, _) = setup().await; + let start = Timestamp::now(); + let magnet_payto = + magnet_payto("payto://iban/HU30162000031000163100000000?receiver-name=name"); + + // Empty db + let pendings = db::pending_batch(&mut db, &start) + .await + .expect("pending_batch"); + assert_eq!(pendings.len(), 0); + + // Some transfers + for i in 0..3 { + make_transfer( + &mut db, + &db::Transfer { + request_uid: HashCode::rand(), + amount: decimal(format!("{}", i + 1)), + exchange_base_url: url("https://exchange.test.com/"), + wtid: ShortHashCode::rand(), + creditor: magnet_payto.clone(), + }, + &Timestamp::now(), + ) + .await + .expect("transfer"); + } + let pendings = db::pending_batch(&mut db, &start) + .await + .expect("pending_batch"); + assert_eq!(pendings.len(), 3); + + // Max 100 txs in batch + for i in 0..100 { + make_transfer( + &mut db, + &db::Transfer { + request_uid: HashCode::rand(), + amount: decimal(format!("{}", i + 1)), + exchange_base_url: url("https://exchange.test.com/"), + wtid: ShortHashCode::rand(), + creditor: magnet_payto.clone(), + }, + &Timestamp::now(), + ) + .await + .expect("transfer"); + } + let pendings = db::pending_batch(&mut db, &start) + .await + .expect("pending_batch"); + assert_eq!(pendings.len(), 100); + + // Skip uploaded + for i in 0..=10 { + db::initiated_submit_success(&mut db, i, &Timestamp::now(), i) + .await + .expect("status success"); + } + let pendings = db::pending_batch(&mut db, &start) + .await + .expect("pending_batch"); + assert_eq!(pendings.len(), 93); + + // Skip failed + for i in 0..=10 { + db::initiated_submit_permanent_failure(&mut db, 10 + i, &Timestamp::now(), "failure") + .await + .expect("status failure"); + } + let pendings = db::pending_batch(&mut db, &start) + .await + .expect("pending_batch"); + assert_eq!(pendings.len(), 83); + } +}*/ diff --git a/src/main.rs b/src/main.rs @@ -17,21 +17,38 @@ * <http://www.gnu.org/licenses/> */ -use std::{fmt::Display, path::Path}; +use std::{ + fmt::{Display, Write}, + path::Path, +}; use anyhow::bail; use clap::Parser; +use compact_str::{CompactString, ToCompactString}; +use jiff::{Timestamp, civil::Date}; use reqwest::{ Client, StatusCode, header::{CONTENT_TYPE, HeaderValue}, }; +use sqlx::PgPool; +use taler_api::subject::parse_outgoing; use taler_build::long_version; -use taler_common::{CommonArgs, config::parser::ConfigSource, taler_main}; -use tracing::{debug, info}; +use taler_common::{ + CommonArgs, + config::parser::ConfigSource, + taler_main, + types::{ + amount::{Amount, Currency}, + payto::{IbanPayto, PaytoURI}, + }, +}; +use tracing::{debug, info, warn}; +use uuid::Uuid; use crate::{ common::EbicsLogger, config::{EbicsHostCfg, NexusCfg}, + db::OutgoingRegistrationResult, ebics_code::EbicsReturnCode, key_management::{Order, hpb, submit_client_keys}, keys::{load_bank_keys, load_client_keys, persist_client_keys}, @@ -41,13 +58,252 @@ use crate::{keys::ClientPriKeysFile, xml::XmlReader}; mod common; pub mod config; mod crypto; +mod db; pub mod ebics_code; mod key_management; pub mod keys; +mod testbench; pub mod xml; pub mod xml_sign; -const SOURCE: ConfigSource = ConfigSource::new("libeufin", "libeufin-nexus", "libeufin-nexus"); +const CONFIG_SOURCE: ConfigSource = + ConfigSource::new("libeufin", "libeufin-nexus", "libeufin-nexus"); + +/// ID for incoming transactions +struct IncomingId { + /** ISO20022 UETR */ + uetr: Option<Uuid>, + /// ISO20022 TxID + tx_id: Option<CompactString>, + /// ISO20022 AcctSvcrRef + acct_svcr_ref: Option<CompactString>, +} +impl IncomingId { + 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 +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 +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 +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 */ +struct PaymentBatch { + pub id: u64, + pub msg_id: CompactString, + pub creation_date: Date, + pub sum: Amount, + pub payments: Vec<InitiatedPayment>, +} + +/** Initiated outgoing transaction */ +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)] @@ -58,13 +314,47 @@ struct Args { fn main() { let args = Args::parse(); - taler_main(SOURCE, args.common, |cfg| async move { + taler_main(CONFIG_SOURCE, args.common, |cfg| async move { let cfg = NexusCfg::parse(cfg)?; ebics_setup(&Client::new(), &cfg, false).await?; Ok(()) }) } +pub async fn register_outgoing( + db: &PgPool, + payment: &OutgoingPayment, +) -> sqlx::Result<OutgoingRegistrationResult> { + let metadata = payment + .subject + .as_ref() + .and_then(|s| parse_outgoing(s).ok()); + let res = db::register_out_tx(db, payment, metadata.as_ref()).await?; + if res.new { + if res.initiated { + info!("{payment}"); + } else { + warn!("{payment} recovered"); + } + } else { + debug!("{payment} already seen"); + } + Ok(res) +} + +pub async fn register_outgoing_batch( + db: &PgPool, + currency: &Currency, + batch: &OutgoingBatch, +) -> sqlx::Result<()> { + info!("{batch}"); + let txs = db::unsettled_tx_in_batch(db, currency, &batch.msg_id, &batch.execution_time).await?; + for tx in txs { + register_outgoing(db, &tx).await?; + } + 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/testbench.rs b/src/testbench.rs @@ -0,0 +1,101 @@ +/* +* This file is part of LibEuFin. +* Copyright (C) 2026 Taler Systems S.A. + +* LibEuFin is free software; you can redistribute it and/or modify +* it under the terms of the GNU Affero General Public License as +* published by the Free Software Foundation; either version 3, or +* (at your option) any later version. + +* LibEuFin is distributed in the hope that it will be useful, but +* WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY +* or FITNESS FOR A PARTICULAR PURPOSE. See the GNU Affero General +* Public License for more details. + +* You should have received a copy of the GNU Affero General Public +* License along with LibEuFin; see the file COPYING. If not, see +* <http://www.gnu.org/licenses/> +*/ + +use anyhow::bail; +use clap::{Parser, ValueEnum}; + +#[derive(Copy, Clone, PartialEq, Eq, PartialOrd, Ord, ValueEnum)] +enum Component { + Nexus, + Ebisync, +} + +#[derive(Parser)] +/// Run integration tests on banks provider +pub struct TestbenchCmd { + #[arg(value_enum)] + component: Component, + platform: String, +} + +pub async fn testbench(cmd: &TestbenchCmd) -> anyhow::Result<()> { + // List available platform + let platforms: Vec<_> = std::fs::read_dir("testbench/test/platform") + .unwrap() + .filter_map(|entry| { + let e = entry.unwrap(); + let filename = e.file_name(); + if filename == "config.json" { + None + } else { + Some( + filename + .to_string_lossy() + .strip_suffix(".conf") + .unwrap() + .to_owned(), + ) + } + }) + .collect(); + if !platforms.contains(&cmd.platform) { + bail!( + "Unknown platform '{}', expected one of {}", + cmd.platform, + platforms.join(", ") + ); + } + + // Augment config + let simple_cfg = + std::fs::read_to_string(format!("testbench/test/platform/{}.conf", cmd.platform)).unwrap(); + let conf = format!("test/{}/ebics.conf", cmd.platform); + std::fs::write(&conf, format!(r#" + {simple_cfg} + {} + [paths] + LIBEUFIN_NEXUS_HOME = test/{} + EBISYNC_HOME = test/{} + + [nexus-fetch] + FREQUENCY = 1h + CHECKPOINT_TIME_OF_DAY = 16:52 + + [ebisync-fetch] + FREQUENCY = 1h + CHECKPOINT_TIME_OF_DAY = 16:52 + DESTINATION = azure-blob-storage + AZURE_API_URL = http://localhost:10000/devstoreaccount1/ + AZURE_ACCOUNT_NAME = devstoreaccount1 + AZURE_ACCOUNT_KEY = Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw== + AZURE_CONTAINER = test + + [ebisync-submit] + SOURCE = ebisync-api + AUTH_METHOD = none + + [libeufin-nexusdb-postgres] + CONFIG = postgres:///libeufintestbench + + [ebisyncdb-postgres] + CONFIG = postgres:///libeufintestbench + "#, simple_cfg.replace("[nexus-ebics]", "[ebisync]").replace("[nexus-setup]", "[ebisync-setup]"), cmd.platform, cmd.platform)).unwrap(); + + Ok(()) +}