commit f54a8fdf2147a89ff4c7c6f52e2edd9026e41708
parent c9b0056a4b800f3c39a8c4e5cc8d71b5ff0a216d
Author: Antoine A <>
Date: Fri, 24 Apr 2026 10:37:57 +0200
IncomingPaymentsTest
Diffstat:
9 files changed, 1977 insertions(+), 1928 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -34,9 +34,9 @@ dependencies = [
[[package]]
name = "anstream"
-version = "0.6.21"
+version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "43d5b281e737544384e969a5ccad3f1cdd24b48086a0fc1b2a5262a26b8f4f4a"
+checksum = "824a212faf96e9acacdbd09febd34438f8f711fb84e09a8916013cd7815ca28d"
dependencies = [
"anstyle",
"anstyle-parse",
@@ -49,15 +49,15 @@ dependencies = [
[[package]]
name = "anstyle"
-version = "1.0.13"
+version = "1.0.14"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "5192cca8006f1fd4f7237516f40fa183bb07f8fbdfedaa0036de5ea9b0b45e78"
+checksum = "940b3a0ca603d1eade50a4846a2afffd5ef57a9feac2c0e2ec2e14f9ead76000"
[[package]]
name = "anstyle-parse"
-version = "0.2.7"
+version = "1.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "4e7644824f0aa2c7b9384579234ef10eb7efb6a0deb83f9630a49594dd9c15c2"
+checksum = "52ce7f38b242319f7cabaa6813055467063ecdc9d355bbb4ce0c68908cd8130e"
dependencies = [
"utf8parse",
]
@@ -230,6 +230,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6"
[[package]]
+name = "base64ct"
+version = "1.8.3"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06"
+
+[[package]]
name = "bitflags"
version = "2.11.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -276,9 +282,9 @@ dependencies = [
[[package]]
name = "cc"
-version = "1.2.56"
+version = "1.2.57"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "aebf35691d1bfb0ac386a69bac2fde4dd276fb618cf8bf4f5318fe285e821bb2"
+checksum = "7a0dd1ca384932ff3641c8718a02769f1698e7563dc6974ffd03346116310423"
dependencies = [
"find-msvc-tools",
"jobserver",
@@ -329,9 +335,9 @@ dependencies = [
[[package]]
name = "clap"
-version = "4.5.60"
+version = "4.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "2797f34da339ce31042b27d23607e051786132987f595b02ba4f6a6dffb7030a"
+checksum = "b193af5b67834b676abd72466a96c1024e6a6ad978a1f484bd90b85c94041351"
dependencies = [
"clap_builder",
"clap_derive",
@@ -339,9 +345,9 @@ dependencies = [
[[package]]
name = "clap_builder"
-version = "4.5.60"
+version = "4.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "24a241312cea5059b13574bb9b3861cabf758b879c15190b37b6d6fd63ab6876"
+checksum = "714a53001bf66416adb0e2ef5ac857140e7dc3a0c48fb28b2f10762fc4b5069f"
dependencies = [
"anstream",
"anstyle",
@@ -351,9 +357,9 @@ dependencies = [
[[package]]
name = "clap_derive"
-version = "4.5.55"
+version = "4.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "a92793da1a46a5f2a02a6f4c46c6496b28c43638adea8306fcb0caa1634f24e5"
+checksum = "1110bd8a634a1ab8cb04345d8d878267d57c3cf1b38d91b71af6686408bbca6a"
dependencies = [
"heck",
"proc-macro2",
@@ -363,9 +369,9 @@ dependencies = [
[[package]]
name = "clap_lex"
-version = "1.0.0"
+version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "3a822ea5bc7590f9d40f1ba12c0dc3c2760f3482c6984db1573ad11031420831"
+checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9"
[[package]]
name = "cmake"
@@ -378,9 +384,9 @@ dependencies = [
[[package]]
name = "colorchoice"
-version = "1.0.4"
+version = "1.0.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "b05b61dc5112cbb17e4b6cd61790d9845d13888356391624cbe7e41efeac1e75"
+checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570"
[[package]]
name = "combine"
@@ -418,6 +424,12 @@ dependencies = [
]
[[package]]
+name = "const-oid"
+version = "0.9.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8"
+
+[[package]]
name = "convert_case"
version = "0.10.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -540,9 +552,9 @@ dependencies = [
[[package]]
name = "darling"
-version = "0.21.3"
+version = "0.23.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "9cdf337090841a411e2a7f3deb9187445851f91b309c0c0a29e05f74a00a48c0"
+checksum = "25ae13da2f202d56bd7f91c25fba009e7717a1e4a1cc98a76d844b65ae912e9d"
dependencies = [
"darling_core",
"darling_macro",
@@ -550,11 +562,10 @@ dependencies = [
[[package]]
name = "darling_core"
-version = "0.21.3"
+version = "0.23.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "1247195ecd7e3c85f83c8d2a366e4210d588e802133e1e355180a9870b517ea4"
+checksum = "9865a50f7c335f53564bb694ef660825eb8610e0a53d3e11bf1b0d3df31e03b0"
dependencies = [
- "fnv",
"ident_case",
"proc-macro2",
"quote",
@@ -564,9 +575,9 @@ dependencies = [
[[package]]
name = "darling_macro"
-version = "0.21.3"
+version = "0.23.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "d38308df82d1080de0afee5d069fa14b0326a88c14f15c5ccda35b4a6c414c81"
+checksum = "ac3984ec7bd6cfa798e62b4a642426a5be0e68f9401cfc2a01e3fa9ea2fcdb8d"
dependencies = [
"darling_core",
"quote",
@@ -594,6 +605,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d7a1e2f27636f116493b8b860f5546edb47c8d8f8ea73e1d2a20be88e28d1fea"
[[package]]
+name = "der"
+version = "0.7.10"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb"
+dependencies = [
+ "const-oid",
+ "pem-rfc7468",
+ "zeroize",
+]
+
+[[package]]
name = "der-parser"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -645,6 +667,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292"
dependencies = [
"block-buffer",
+ "const-oid",
"crypto-common",
"subtle",
]
@@ -771,6 +794,17 @@ dependencies = [
]
[[package]]
+name = "flume"
+version = "0.11.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095"
+dependencies = [
+ "futures-core",
+ "futures-sink",
+ "spin",
+]
+
+[[package]]
name = "fnv"
version = "1.0.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -814,6 +848,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7e3450815272ef58cec6d564423f6e755e25379b217b0bc688e295ba24df6b1d"
[[package]]
+name = "futures-executor"
+version = "0.3.32"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "baf29c38818342a3b26b5b923639e7b1f4a61fc5e76102d4b1981c6dc7a7579d"
+dependencies = [
+ "futures-core",
+ "futures-task",
+ "futures-util",
+]
+
+[[package]]
name = "futures-intrusive"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1388,6 +1433,9 @@ name = "lazy_static"
version = "1.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe"
+dependencies = [
+ "spin",
+]
[[package]]
name = "leb128fmt"
@@ -1417,6 +1465,7 @@ dependencies = [
"rand 0.10.0",
"rcgen",
"reedline",
+ "regex",
"reqwest",
"roxmltree",
"serde",
@@ -1438,6 +1487,12 @@ dependencies = [
]
[[package]]
+name = "libm"
+version = "0.2.16"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981"
+
+[[package]]
name = "libredox"
version = "0.1.14"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1450,6 +1505,16 @@ dependencies = [
]
[[package]]
+name = "libsqlite3-sys"
+version = "0.30.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "2e99fb7a497b1e3339bc746195567ed8d3e24945ecd636e3619d20b9de9e9149"
+dependencies = [
+ "pkg-config",
+ "vcpkg",
+]
+
+[[package]]
name = "linux-raw-sys"
version = "0.12.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1585,6 +1650,22 @@ dependencies = [
]
[[package]]
+name = "num-bigint-dig"
+version = "0.8.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "e661dda6640fad38e827a6d4a310ff4763082116fe217f279885c97f511bb0b7"
+dependencies = [
+ "lazy_static",
+ "libm",
+ "num-integer",
+ "num-iter",
+ "num-traits",
+ "rand 0.8.5",
+ "smallvec",
+ "zeroize",
+]
+
+[[package]]
name = "num-conv"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1600,12 +1681,24 @@ dependencies = [
]
[[package]]
+name = "num-iter"
+version = "0.1.45"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "1429034a0490724d0075ebb2bc9e875d6503c3cf69e235a8941aa757d83ef5bf"
+dependencies = [
+ "autocfg",
+ "num-integer",
+ "num-traits",
+]
+
+[[package]]
name = "num-traits"
version = "0.2.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841"
dependencies = [
"autocfg",
+ "libm",
]
[[package]]
@@ -1619,9 +1712,9 @@ dependencies = [
[[package]]
name = "once_cell"
-version = "1.21.3"
+version = "1.21.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "42f5e15c9953c5e4ccceeb2e7382a716482c34515315f7b03532b8b4e8393d2d"
+checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50"
[[package]]
name = "once_cell_polyfill"
@@ -1675,6 +1768,15 @@ dependencies = [
]
[[package]]
+name = "pem-rfc7468"
+version = "0.7.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "88b39c9bfcfc231068454382784bb460aae594343fb030d46e9f50a645418412"
+dependencies = [
+ "base64ct",
+]
+
+[[package]]
name = "percent-encoding"
version = "2.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -1736,6 +1838,33 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184"
[[package]]
+name = "pkcs1"
+version = "0.7.5"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "c8ffb9f10fa047879315e6625af03c164b16962a5368d724ed16323b68ace47f"
+dependencies = [
+ "der",
+ "pkcs8",
+ "spki",
+]
+
+[[package]]
+name = "pkcs8"
+version = "0.10.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "f950b2377845cebe5cf8b5165cb3cc1a5e0fa5cfa3e1f7f55707d8fd82e0a7b7"
+dependencies = [
+ "der",
+ "spki",
+]
+
+[[package]]
+name = "pkg-config"
+version = "0.3.32"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c"
+
+[[package]]
name = "plain"
version = "0.2.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2105,6 +2234,26 @@ dependencies = [
]
[[package]]
+name = "rsa"
+version = "0.9.10"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "b8573f03f5883dcaebdfcf4725caa1ecb9c15b2ef50c43a07b816e06799bb12d"
+dependencies = [
+ "const-oid",
+ "digest",
+ "num-bigint-dig",
+ "num-integer",
+ "num-traits",
+ "pkcs1",
+ "pkcs8",
+ "rand_core 0.6.4",
+ "signature",
+ "spki",
+ "subtle",
+ "zeroize",
+]
+
+[[package]]
name = "rustc-hash"
version = "2.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2349,9 +2498,9 @@ dependencies = [
[[package]]
name = "serde_with"
-version = "3.17.0"
+version = "3.18.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "381b283ce7bc6b476d903296fb59d0d36633652b633b27f64db4fb46dcbfc3b9"
+checksum = "dd5414fad8e6907dbdd5bc441a50ae8d6e26151a03b1de04d89a5576de61d01f"
dependencies = [
"serde_core",
"serde_with_macros",
@@ -2359,9 +2508,9 @@ dependencies = [
[[package]]
name = "serde_with_macros"
-version = "3.17.0"
+version = "3.18.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "a6d4e30573c8cb306ed6ab1dca8423eec9a463ea0e155f45399455e0368b27e0"
+checksum = "d3db8978e608f1fe7357e211969fd9abdcae80bac1ba7a3369bb7eb6b404eb65"
dependencies = [
"darling",
"proc-macro2",
@@ -2370,6 +2519,17 @@ dependencies = [
]
[[package]]
+name = "sha1"
+version = "0.10.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba"
+dependencies = [
+ "cfg-if",
+ "cpufeatures 0.2.17",
+ "digest",
+]
+
+[[package]]
name = "sha2"
version = "0.10.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2427,6 +2587,16 @@ dependencies = [
]
[[package]]
+name = "signature"
+version = "2.2.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de"
+dependencies = [
+ "digest",
+ "rand_core 0.6.4",
+]
+
+[[package]]
name = "simd-adler32"
version = "0.3.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2458,6 +2628,25 @@ dependencies = [
]
[[package]]
+name = "spin"
+version = "0.9.8"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67"
+dependencies = [
+ "lock_api",
+]
+
+[[package]]
+name = "spki"
+version = "0.7.3"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "d91ed6c858b01f942cd56b37a94b3e0a1798290327d1236e4d9cf4eaca44d29d"
+dependencies = [
+ "base64ct",
+ "der",
+]
+
+[[package]]
name = "sqlx"
version = "0.8.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2465,7 +2654,9 @@ checksum = "1fefb893899429669dcdd979aff487bd78f4064e5e7907e4269081e0ef7d97dc"
dependencies = [
"sqlx-core",
"sqlx-macros",
+ "sqlx-mysql",
"sqlx-postgres",
+ "sqlx-sqlite",
]
[[package]]
@@ -2501,6 +2692,7 @@ dependencies = [
"tokio-stream",
"tracing",
"url",
+ "uuid",
"webpki-roots 0.26.11",
]
@@ -2534,13 +2726,58 @@ dependencies = [
"serde_json",
"sha2",
"sqlx-core",
+ "sqlx-mysql",
"sqlx-postgres",
+ "sqlx-sqlite",
"syn",
"tokio",
"url",
]
[[package]]
+name = "sqlx-mysql"
+version = "0.8.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "aa003f0038df784eb8fecbbac13affe3da23b45194bd57dba231c8f48199c526"
+dependencies = [
+ "atoi",
+ "base64",
+ "bitflags",
+ "byteorder",
+ "bytes",
+ "crc",
+ "digest",
+ "dotenvy",
+ "either",
+ "futures-channel",
+ "futures-core",
+ "futures-io",
+ "futures-util",
+ "generic-array",
+ "hex",
+ "hkdf",
+ "hmac",
+ "itoa",
+ "log",
+ "md-5",
+ "memchr",
+ "once_cell",
+ "percent-encoding",
+ "rand 0.8.5",
+ "rsa",
+ "serde",
+ "sha1",
+ "sha2",
+ "smallvec",
+ "sqlx-core",
+ "stringprep",
+ "thiserror 2.0.18",
+ "tracing",
+ "uuid",
+ "whoami",
+]
+
+[[package]]
name = "sqlx-postgres"
version = "0.8.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2574,10 +2811,36 @@ dependencies = [
"stringprep",
"thiserror 2.0.18",
"tracing",
+ "uuid",
"whoami",
]
[[package]]
+name = "sqlx-sqlite"
+version = "0.8.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "c2d12fe70b2c1b4401038055f90f151b78208de1f9f89a7dbfd41587a10c3eea"
+dependencies = [
+ "atoi",
+ "flume",
+ "futures-channel",
+ "futures-core",
+ "futures-executor",
+ "futures-intrusive",
+ "futures-util",
+ "libsqlite3-sys",
+ "log",
+ "percent-encoding",
+ "serde",
+ "serde_urlencoded",
+ "sqlx-core",
+ "thiserror 2.0.18",
+ "tracing",
+ "url",
+ "uuid",
+]
+
+[[package]]
name = "stable_deref_trait"
version = "1.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2720,6 +2983,7 @@ dependencies = [
"compact_str",
"dashmap",
"http-body-util",
+ "itoa",
"jiff",
"listenfd",
"serde",
@@ -2746,10 +3010,10 @@ dependencies = [
"aws-lc-rs",
"clap",
"compact_str",
- "fastrand",
"glob",
"indexmap",
"jiff",
+ "rand 0.10.0",
"serde",
"serde_json",
"serde_path_to_error",
@@ -3055,9 +3319,9 @@ dependencies = [
[[package]]
name = "tracing-subscriber"
-version = "0.3.22"
+version = "0.3.23"
source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "2f30143827ddab0d256fd843b7a66d164e9f271cfa0dde49142c5ca0ca291f1e"
+checksum = "cb7f578e5945fb242538965c2d0b04418d38ec25c79d160cd279bf0731c8d319"
dependencies = [
"nu-ansi-term",
"sharded-slab",
@@ -3179,7 +3443,9 @@ version = "1.22.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a68d3c8f01c0cfa54a75291d83601161799e4a89a39e0929f4b0354d88757a37"
dependencies = [
+ "getrandom 0.4.2",
"js-sys",
+ "rand 0.10.0",
"wasm-bindgen",
]
@@ -3190,6 +3456,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65"
[[package]]
+name = "vcpkg"
+version = "0.2.15"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426"
+
+[[package]]
name = "version_check"
version = "0.9.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
diff --git a/Cargo.toml b/Cargo.toml
@@ -42,6 +42,8 @@ sqlx = { version = "0.8", default-features = false, features = [
"postgres",
"runtime-tokio",
"tls-rustls-aws-lc-rs",
+ "uuid",
] }
compact_str = { version = "0.9.0", features = ["serde", "sqlx-postgres"] }
-uuid = "*"
-\ No newline at end of file
+uuid = { version = "1.0", features = ["v4", "fast-rng"] }
+regex = "*"
diff --git a/database-versioning/libeufin-nexus-procedures.sql b/database-versioning/libeufin-nexus-procedures.sql
@@ -686,7 +686,7 @@ FROM prepared_transfers
WHERE authorization_pub = in_authorization_pub;
-- Check idempotency and delay garbage collection
-IF FOUND AND idempotent THEN
+IF idempotent THEN
UPDATE prepared_transfers
SET registered_at=in_timestamp
WHERE authorization_pub=in_authorization_pub;
diff --git a/src/config.rs b/src/config.rs
@@ -19,8 +19,13 @@
use std::cell::OnceCell;
+use jiff::Timestamp;
+use regex::Regex;
use taler_api::config::DbCfg;
-use taler_common::config::{Config, ValueErr};
+use taler_common::{
+ config::{Config, ValueErr},
+ types::amount::{Amount, Currency},
+};
pub fn parse_db_cfg(cfg: &Config) -> Result<DbCfg, ValueErr> {
DbCfg::parse(cfg.section("libeufin-nexusdb-postgres"))
@@ -60,6 +65,34 @@ impl EbicsHostCfg {
}
}
+#[derive(Debug, Clone, Copy)]
+pub enum AccountType {
+ Exchange,
+ Normal,
+}
+
+pub struct NexusIngestConfig {
+ pub account_type: AccountType,
+ pub ignore_transactions_before: Timestamp,
+ pub ignore_bounces_before: Timestamp,
+ pub restriction_payto_regex: Option<Regex>,
+ pub bounce_deduce_fee: bool,
+ pub bounce_fee: Amount,
+}
+
+impl NexusIngestConfig {
+ pub fn simple(account_type: AccountType, currency: &Currency) -> Self {
+ Self {
+ account_type,
+ ignore_transactions_before: Timestamp::UNIX_EPOCH,
+ ignore_bounces_before: Timestamp::UNIX_EPOCH,
+ restriction_payto_regex: None,
+ bounce_deduce_fee: false,
+ bounce_fee: Amount::zero(currency),
+ }
+ }
+}
+
pub struct NexusCfg {
pub cfg: Config,
pub keys: OnceCell<EbicsKeysCfg>,
diff --git a/src/db.rs b/src/db.rs
@@ -14,33 +14,23 @@
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 compact_str::CompactString;
+use jiff::Timestamp;
+use sqlx::{PgPool, QueryBuilder, Row, postgres::PgRow};
use taler_api::{
- db::{BindHelper, IncomingType, TypeHelper, history, page},
+ db::{BindHelper, IncomingType, TypeHelper, history},
subject::{IncomingSubject, OutgoingSubject},
};
use taler_common::{
- api_common::{HashCode, ShortHashCode},
- api_params::{History, Page},
- api_revenue::RevenueIncomingBankTransaction,
- api_wire::{
- IncomingBankTransaction, OutgoingBankTransaction, TransferListStatus, TransferState,
- TransferStatus,
- },
+ api_common::{EddsaPublicKey, EddsaSignature},
+ api_params::History,
+ api_wire::{IncomingBankTransaction, OutgoingBankTransaction},
config::Config,
- types::{
- amount::{Amount, Currency, Decimal},
- payto::PaytoImpl as _,
- },
+ types::amount::{Amount, Currency},
};
use tokio::sync::watch::{Receiver, Sender};
-use url::Url;
-use crate::{InitiatedPayment, OutgoingId, OutgoingPayment, config::parse_db_cfg};
+use crate::{IncomingPayment, InitiatedPayment, OutgoingId, OutgoingPayment, config::parse_db_cfg};
const SCHEMA: &str = "libeufin_nexus";
@@ -79,13 +69,13 @@ pub async fn notification_listener(
/// Outgoing payments initiation result
#[derive(Debug, PartialEq, Eq)]
-enum PaymentInitiationResult {
+pub enum PaymentInitiationResult {
Success(u64),
RequestUidReuse,
}
/// Initiate a new payment
-async fn initiate(
+pub async fn initiate(
pool: &PgPool,
payment: &InitiatedPayment,
) -> sqlx::Result<PaymentInitiationResult> {
@@ -268,6 +258,285 @@ pub async fn register_out_batch(
.await
}
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct InResult {
+ pub id: u64,
+ pub new: bool,
+ pub completed: bool,
+ pub pending: bool,
+ pub bounce_id: Option<CompactString>,
+}
+
+/** Incoming payments registration result */
+#[derive(Debug, PartialEq, Eq)]
+pub enum IncomingRegistrationResult {
+ Success(InResult),
+ ReservePubReuse,
+ MappingReuse,
+ UnknownMapping,
+}
+
+/** Register an incoming payment */
+pub async fn register_in(pool: &PgPool, payment: &IncomingPayment) -> sqlx::Result<InResult> {
+ sqlx::query(
+ "
+ SELECT out_found, out_completed, out_tx_id, out_bounce_id
+ FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,NULL)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .try_map(|r: PgRow| {
+ Ok(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: false,
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+/** Register an talerable incoming payment */
+pub async fn register_in_talerable(
+ pool: &PgPool,
+ payment: &IncomingPayment,
+ subject: &IncomingSubject,
+) -> sqlx::Result<IncomingRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT
+ out_reserve_pub_reuse,
+ out_mapping_reuse,
+ out_unknown_mapping,
+ out_found,
+ out_completed,
+ out_pending,
+ out_tx_id,
+ out_bounce_id
+ FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9::taler_incoming_type,$10,NULL)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(subject.ty())
+ .bind(subject.key())
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_reserve_pub_reuse")? {
+ IncomingRegistrationResult::ReservePubReuse
+ } else if r.try_get_flag("out_mapping_reuse")? {
+ IncomingRegistrationResult::MappingReuse
+ } else if r.try_get_flag("out_unknown_mapping")? {
+ IncomingRegistrationResult::UnknownMapping
+ } else {
+ IncomingRegistrationResult::Success(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: r.try_get("out_pending")?,
+ })
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+/** Register an talerable incoming payment */
+pub async fn register_in_qr_bill(
+ pool: &PgPool,
+ payment: &IncomingPayment,
+ reference: &str,
+) -> sqlx::Result<IncomingRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT
+ out_reserve_pub_reuse,
+ out_mapping_reuse,
+ out_unknown_mapping,
+ out_found,
+ out_completed,
+ out_pending,
+ out_tx_id,
+ out_bounce_id
+ FROM register_incoming($1,$2,$3,$4,$5,$6,$7,$8,NULL,NULL,$9)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(reference)
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_reserve_pub_reuse")? {
+ IncomingRegistrationResult::ReservePubReuse
+ } else if r.try_get_flag("out_mapping_reuse")? {
+ IncomingRegistrationResult::MappingReuse
+ } else if r.try_get_flag("out_unknown_mapping")? {
+ IncomingRegistrationResult::UnknownMapping
+ } else {
+ IncomingRegistrationResult::Success(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: r.try_get("out_pending")?,
+ })
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+/** Incoming payments bounce registration result */
+pub enum IncomingBounceRegistrationResult {
+ Success(InResult),
+ Talerable,
+}
+
+/** Register an incoming payment and bounce it */
+pub async fn register_in_malformed(
+ pool: &PgPool,
+ payment: &IncomingPayment,
+ bounce_amount: &Amount,
+ bounce_end_to_end_id: &str,
+ timestamp: &Timestamp,
+ cause: &str,
+) -> sqlx::Result<IncomingBounceRegistrationResult> {
+ sqlx::query(
+ "
+ SELECT out_found, out_tx_id, out_completed, out_bounce_id, out_talerable
+ FROM register_and_bounce_incoming($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12)
+ ",
+ )
+ .bind(&payment.amount)
+ .bind(
+ payment
+ .credit_fee
+ .as_ref()
+ .unwrap_or(&Amount::zero(&payment.amount.currency)),
+ )
+ .bind(&payment.subject)
+ .bind_timestamp(&payment.execution_time)
+ .bind(payment.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .bind(payment.id.uetr)
+ .bind(&payment.id.tx_id)
+ .bind(&payment.id.acct_svcr_ref)
+ .bind(bounce_amount)
+ .bind_timestamp(timestamp)
+ .bind(bounce_end_to_end_id)
+ .bind(cause)
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_talerable")? {
+ IncomingBounceRegistrationResult::Talerable
+ } else {
+ IncomingBounceRegistrationResult::Success(InResult {
+ id: r.try_get_u64("out_tx_id")?,
+ new: !r.try_get_flag("out_found")?,
+ completed: r.try_get_flag("out_completed")?,
+ bounce_id: r.try_get("out_bounce_id")?,
+ pending: false,
+ })
+ })
+ })
+ .fetch_one(pool)
+ .await
+}
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum RegistrationResult {
+ Success,
+ ReservePubReuse,
+ SubjectReuse,
+}
+
+pub async fn transfer_register(
+ db: &PgPool,
+ ty: IncomingType,
+ account_pub: &EddsaPublicKey,
+ auth_pub: &EddsaPublicKey,
+ auth_sig: &EddsaSignature,
+ recurrent: bool,
+ reference_number: &str,
+ timestamp: &Timestamp,
+) -> sqlx::Result<RegistrationResult> {
+ sqlx::query(
+ "
+ SELECT
+ out_subject_reuse,
+ out_reserve_pub_reuse
+ FROM register_prepared_transfers (
+ $1::taler_incoming_type,$2,$3,$4,$5,$6,$7
+ )
+ ",
+ )
+ .bind(ty)
+ .bind(account_pub)
+ .bind(auth_pub)
+ .bind(auth_sig)
+ .bind(recurrent)
+ .bind(reference_number)
+ .bind_timestamp(timestamp)
+ .try_map(|r: PgRow| {
+ Ok(if r.try_get_flag("out_subject_reuse")? {
+ RegistrationResult::SubjectReuse
+ } else if r.try_get_flag("out_reserve_pub_reuse")? {
+ RegistrationResult::ReservePubReuse
+ } else {
+ RegistrationResult::Success
+ })
+ })
+ .fetch_one(db)
+ .await
+}
+
+pub async fn transfer_unregister(
+ db: &PgPool,
+ auth_pub: &EddsaPublicKey,
+ timestamp: &Timestamp,
+) -> sqlx::Result<bool> {
+ sqlx::query("SELECT out_found FROM delete_prepared_transfers($1,$2)")
+ .bind(auth_pub)
+ .bind_timestamp(timestamp)
+ .try_map(|r: PgRow| r.try_get_flag(0))
+ .fetch_one(db)
+ .await
+}
+
pub async fn outgoing_history(
db: &PgPool,
currency: &Currency,
@@ -358,7 +627,7 @@ pub async fn incoming_history(
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(¤cy) {
+ credit_fee: if tmp == Amount::zero(currency) {
None
} else {
Some(tmp)
@@ -383,9 +652,7 @@ pub async fn incoming_history(
authorization_pub: r.try_get("authorization_pub")?,
authorization_sig: r.try_get("authorization_sig")?,
},
- IncomingType::wad => {
- unimplemented!("WAD is not yet supported")
- }
+ IncomingType::map => unimplemented!("MAP are never listed in the history"),
})
},
)
@@ -394,30 +661,40 @@ pub async fn incoming_history(
#[cfg(test)]
mod test {
+ use std::str::FromStr;
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_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,
- payto::{BankID, FullIbanPayto},
- },
+ api_common::{EddsaPublicKey, EddsaSignature},
+ types::amount::amount,
};
+ use uuid::Uuid;
- 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};
+ use crate::{
+ CONFIG_SOURCE, IncomingId, IncomingPayment, InitiatedPayment, OutgoingBatch,
+ config::{AccountType, NexusIngestConfig},
+ db::{
+ InResult, IncomingBounceRegistrationResult, OutgoingRegistrationResult,
+ RegistrationResult, register_in_malformed, transfer_register,
+ },
+ register_incoming,
+ };
+ use crate::{
+ OutgoingId, db::PaymentInitiationResult, db::batch_initiated, db::initiate,
+ db::initiated_ack,
+ };
+ use crate::{OutgoingPayment, rand_ebics_id, 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) {
@@ -426,58 +703,95 @@ mod test {
(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 {
+ /** Generates an outgoing payment, given its subject */
+ pub fn gen_out_pay(subject: impl Into<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(),
+ amount: Amount::new(&CURRENCY, 44, 0),
debit_fee: None,
creditor: Some(
IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test")
.unwrap()
.as_payto(),
),
- subject: Some(subject),
+ subject: Some(subject.into()),
execution_time: Timestamp::now(),
}
}
- /** Generates a payment initiation, given its subject */
- pub fn gen_init_pay(end_to_end_id: CompactString, subject: String) -> InitiatedPayment {
+ /** Generates a payment initiation, given its subject and end-to-end ID */
+ pub fn gen_init_pay(
+ end_to_end_id: CompactString,
+ subject: impl Into<String>,
+ ) -> InitiatedPayment {
InitiatedPayment {
id: 0,
- amount: Amount::from_str("KUDOS:44").unwrap(),
+ amount: Amount::new(&CURRENCY, 44, 0),
creditor: IbanPayto::from_str("payto://iban/CH4189144589712575493?receiver-name=Test")
.unwrap()
.as_payto(),
- subject,
+ subject: subject.into(),
initiation_time: Timestamp::now(),
end_to_end_id,
}
}
- async fn check_out_count(db: &PgPool, nb_incoming: u64, nb_talerable: u64) {
+ /** Generates an incoming payment, given its subject */
+ pub fn gen_in_pay(subject: impl Into<String>) -> IncomingPayment {
+ IncomingPayment {
+ id: IncomingId::new(None, Some(rand_ebics_id()), None),
+ amount: Amount::new(&CURRENCY, 44, 0),
+ credit_fee: None,
+ debtor: Some(
+ IbanPayto::from_str("payto://iban/DE84500105177118117964?receiver-name=John+Smith")
+ .unwrap()
+ .as_payto(),
+ ),
+ subject: Some(subject.into()),
+ execution_time: Timestamp::now(),
+ }
+ }
+
+ async fn check_in_count(
+ db: &PgPool,
+ nb_incoming: usize,
+ nb_bounce: usize,
+ nb_talerable: usize,
+ ) {
+ sqlx::query(
+ "
+ SELECT (SELECT count(*) FROM incoming_transactions) AS incoming,
+ (SELECT count(*) FROM bounced_transactions) AS bounce,
+ (SELECT count(*) FROM talerable_incoming_transactions) AS talerable
+ ",
+ )
+ .try_map(|r: PgRow| {
+ assert_eq!(
+ (r.try_get_u64(0)?, r.try_get_u64(1)?, r.try_get_u64(2)?),
+ (nb_incoming as u64, nb_bounce as u64, nb_talerable as u64)
+ );
+ Ok(())
+ })
+ .fetch_one(db)
+ .await
+ .unwrap();
+ }
+
+ async fn check_out_count(db: &PgPool, nb_outgoing: u64, nb_talerable: u64) {
sqlx::query(
"
- SELECT (SELECT count(*) FROM outgoing_transactions) AS incoming,
+ SELECT (SELECT count(*) FROM outgoing_transactions) AS outgoing,
(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)
+ (nb_outgoing, nb_talerable)
);
Ok(())
})
@@ -486,6 +800,53 @@ mod test {
.unwrap();
}
+ #[derive(Debug, PartialEq, Eq)]
+ enum Status {
+ Simple,
+ Pending,
+ Bounced,
+ Incomplete,
+ Reserve(EddsaPublicKey),
+ Kyc(EddsaPublicKey),
+ }
+
+ use Status::*;
+
+ async fn check_in(db: &PgPool, state: &[Status]) {
+ let current = sqlx::query(
+ "
+ SELECT pending_recurrent_incoming_transactions.authorization_pub IS NOT NULL, initiated_outgoing_transaction_id IS NOT NULL, debit_payto IS NULL OR subject IS NULL, type::text, metadata
+ FROM incoming_transactions
+ LEFT JOIN talerable_incoming_transactions USING (incoming_transaction_id)
+ LEFT JOIN pending_recurrent_incoming_transactions USING (incoming_transaction_id)
+ LEFT JOIN bounced_transactions USING (incoming_transaction_id)
+ ORDER BY incoming_transaction_id
+ ",
+ )
+ .try_map(|r: PgRow| {
+ Ok(
+ if r.try_get_flag(0)? {
+ Status::Pending
+ } else if r.try_get_flag(1)? {
+ Status::Bounced
+ } else if r.try_get_flag(2)? {
+ Status::Incomplete
+ } else {
+ match r.try_get(3)? {
+ None => Status::Simple,
+ Some(IncomingType::reserve) => Status::Reserve(r.try_get(4)?),
+ Some(IncomingType::kyc) => Status::Kyc(r.try_get(4)?),
+ Some(e) => unreachable!("{e:?}")
+ }
+ }
+ )
+ })
+ .fetch_all(db)
+ .await
+ .unwrap();
+ assert_eq!(state, current);
+ }
+
#[tokio::test]
async fn out_tx() {
let (_, db) = setup().await;
@@ -608,7 +969,7 @@ mod test {
// Init batch
let wtid = ShortHashCode::rand();
for subject in [
- format!("initiated by nexus"),
+ "initiated by nexus".to_string(),
format!("{} https://exchange.com/", ShortHashCode::rand()),
format!("{wtid} https://exchange.com/"),
format!("{wtid} https://exchange.com/"),
@@ -681,1441 +1042,610 @@ mod test {
.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,
-}
+ #[tokio::test]
+ async fn in_bounce() {
+ let (_, db) = setup().await;
-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
- )
- }
-}
+ // Creating and bouncing one incoming transaction
+ let payment = gen_in_pay("incoming and bounce");
+ let id = rand_ebics_id();
-#[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,
-}
+ let bounce_amount = amount("KUDOS:2.53");
+ let res = register_in_malformed(
+ &db,
+ &payment,
+ &bounce_amount,
+ &id,
+ &Timestamp::now(),
+ "manual bounce",
+ )
+ .await
+ .unwrap();
+ assert!(
+ matches!(
+ res,
+ IncomingBounceRegistrationResult::Success(InResult {
+ new: true,
+ id: _,
+ completed: false,
+ pending: false,
+ ref bounce_id
+ }) if bounce_id.as_ref() == Some(&id)
+ ),
+ "{res:?}"
+ );
+ // Idempotent
+ let res = register_in_malformed(
+ &db,
+ &payment,
+ &amount("KUDOS:2.5"),
+ &rand_ebics_id(),
+ &Timestamp::now(),
+ "other reason to bounce",
+ )
+ .await
+ .unwrap();
+ assert!(
+ matches!(
+ res,
+ IncomingBounceRegistrationResult::Success(InResult {
+ new: false,
+ id: _,
+ completed: false,
+ pending: false,
+ ref bounce_id
+ }) if bounce_id.as_ref() == Some(&id)
+ ),
+ "{res:?}"
+ );
-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
+ // Checking one incoming got created and bounced
+ sqlx::query(
+ "
+ SELECT
+ incoming_transactions.amount as in_amount,
+ initiated_outgoing_transactions.amount as bounce_amount
+ FROM incoming_transactions
+ JOIN bounced_transactions USING (incoming_transaction_id)
+ JOIN initiated_outgoing_transactions USING (initiated_outgoing_transaction_id)
+ ",
)
+ .try_map(|r: PgRow| {
+ assert_eq!(r.try_get_amount("in_amount", &CURRENCY)?, payment.amount);
+ assert_eq!(r.try_get_amount("bounce_amount", &CURRENCY)?, bounce_amount);
+ Ok(())
+ })
+ .fetch_one(&db)
+ .await
+ .unwrap();
}
-}
-#[derive(Debug, PartialEq, Eq)]
-pub struct Initiated {
- pub id: u64,
- pub amount: Amount,
- pub subject: String,
- pub creditor: FullHuPayto,
-}
+ #[tokio::test]
+ async fn in_simple() {
+ let (_, db) = setup().await;
-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
- )
- }
-}
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
-#[derive(Debug, Clone)]
-pub struct TxInAdmin {
- pub amount: Amount,
- pub subject: String,
- pub debtor: FullHuPayto,
- pub metadata: IncomingSubject,
-}
+ // Register
+ let incoming = gen_in_pay("test".to_owned());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Bounced]).await;
-/// 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
-}
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Bounced]).await;
-#[derive(Debug, PartialEq, Eq)]
-pub enum AddIncomingResult {
- Success {
- new: bool,
- row_id: u64,
- valued_at: Date,
- },
- ReservePubReuse,
-}
+ // Many
+ register_incoming(&db, &cfg, &gen_in_pay("another subject".to_owned()))
+ .await
+ .unwrap();
+ check_in(&db, &[Bounced, Bounced]).await;
-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
-}
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
+ .await
+ .unwrap();
-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,
-}
+ check_in(&db, &[Bounced, Bounced, Simple]).await;
-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
- }
+ let original = gen_in_pay("test 2".to_owned());
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
- 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)
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await;
+ // Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Bounced, Bounced, Simple, Incomplete]).await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(&db, &[Bounced, Bounced, Simple, Bounced]).await;
}
#[tokio::test]
- async fn kv() {
- let (mut db, _) = setup().await;
+ async fn in_talerable() {
+ let (_, db) = setup().await;
- let value = json!({
- "name": "Mr Smith",
- "no way": 32
- });
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+ let key = EddsaPublicKey::rand();
+ let subject = format!("test with {key} reserve pub");
- 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)
- );
- }
+ // Register
+ let incoming = gen_in_pay(subject.clone());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(key.clone())]).await;
- #[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
- }
- );
- }
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(key.clone())]).await;
- // 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()
- );
+ // Key reuse is bounced
+ register_incoming(&db, &cfg, &gen_in_pay(subject.clone()))
+ .await
+ .unwrap();
+ register_incoming(&db, &cfg, &gen_in_pay(format!("another {subject}")))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(key.clone()), Bounced, Bounced]).await;
- // Regular transaction
- routine(&mut db, &None, &None).await;
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(key.clone()), Bounced, Bounced, Simple]).await;
+
+ let new = EddsaPublicKey::rand();
+ let original = gen_in_pay(format!("test 2 with {new} reserve pub"));
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
- // Reserve transaction
- routine(
- &mut db,
- &Some(IncomingSubject::Reserve(EddsaPublicKey::rand())),
- &Some(IncomingSubject::Reserve(EddsaPublicKey::rand())),
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(
+ &db,
+ &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete],
)
.await;
-
- // Kyc transaction
- routine(
- &mut db,
- &Some(IncomingSubject::Kyc(EddsaPublicKey::rand())),
- &Some(IncomingSubject::Kyc(EddsaPublicKey::rand())),
+ // Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(
+ &db,
+ &[Reserve(key.clone()), Bounced, Bounced, Simple, Incomplete],
+ )
+ .await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(
+ &db,
+ &[Reserve(key.clone()), Bounced, Bounced, Simple, Reserve(new)],
)
.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;
+ async fn in_mapping() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+ 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 subject = format!("test with MAP:{auth_pub} auth pub");
- // Empty db
assert_eq!(
- db::incoming_history(&pool, &History::default(), fake_listen)
- .await
- .unwrap(),
- Vec::new()
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &first,
+ &auth_pub,
+ &auth_sig,
+ false,
+ &reference_number,
+ &Timestamp::now()
+ )
+ .await
+ .unwrap(),
+ RegistrationResult::Success
);
- 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()),
+ // Register
+ let incoming = gen_in_pay(subject.clone());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
+
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
+
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
+ .await
+ .unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple]).await;
+
+ let original = gen_in_pay(format!("test 2 for {subject}"));
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
};
- // 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
- }
- );
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
// Idempotent
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
+
+ let second = EddsaPublicKey::rand();
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
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &second,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
)
.await
- .expect("register tx in"),
- AddIncomingResult::Success {
- new: true,
- row_id: 2,
- valued_at: date
- }
+ .unwrap(),
+ RegistrationResult::Success
);
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
- // History
- assert_eq!(
- db::incoming_history(&pool, &History::default(), fake_listen)
+ // Key reuse is pending
+ for _ in 0..3 {
+ register_incoming(&db, &cfg, &gen_in_pay(subject.clone()))
.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,
- }
- );
+ .unwrap();
}
-
- // 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"),
- )),
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Pending,
+ Pending,
+ ],
)
.await;
- // Bounced transaction
- routine(&mut db, &TxOutKind::Bounce(21), &TxOutKind::Bounce(42)).await;
-
- // History
+ // Finish pending
+ let third = EddsaPublicKey::rand();
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)
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &third,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
+ )
.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
- }
+ .unwrap(),
+ RegistrationResult::Success
);
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Reserve(third.clone()),
+ Pending,
+ ],
+ )
+ .await;
}
#[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()
- );
+ async fn in_reference() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+ 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 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
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &first,
+ &auth_pub,
+ &auth_sig,
+ false,
+ &reference_number,
+ &Timestamp::now()
)
.await
- .expect("transfer"),
- TransferResult::RequestUidReuse
+ .unwrap(),
+ RegistrationResult::Success
);
- // 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();
+ // Register
+ let incoming = gen_in_pay(reference_number.clone());
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
- // Empty db
- assert!(db::pending_batch(&mut db, &now).await.unwrap().is_empty());
+ // Idempotent
+ register_incoming(&db, &cfg, &incoming).await.unwrap();
+ check_in(&db, &[Reserve(first.clone())]).await;
- // 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
- )
+ // Admin balance adjust is ignored
+ register_incoming(&db, &cfg, &gen_in_pay("ADMIN BALANCE ADJUST".to_owned()))
.await
- .expect("register tx in"),
- AddIncomingResult::Success {
- new: true,
- row_id: 1,
- valued_at: date
- }
- );
+ .unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple]).await;
- // 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
- }
- );
+ let original = gen_in_pay(reference_number.clone());
+ let incomplete = IncomingPayment {
+ subject: None,
+ debtor: None,
+ ..original.clone()
+ };
+ // Register incomplete transaction
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
// 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
- }
- );
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Incomplete]).await;
+ // Recover info when completed
+ register_incoming(&db, &cfg, &original).await.unwrap();
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
- // Bounce registered
+ let second = EddsaPublicKey::rand();
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
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &second,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
)
.await
- .expect("bounce"),
- BounceResult {
- tx_id: 1,
- tx_new: false,
- bounce_id: 2,
- bounce_new: true
- }
+ .unwrap(),
+ RegistrationResult::Success
);
- // Idempotent registered
+ check_in(&db, &[Reserve(first.clone()), Simple, Bounced]).await;
+
+ // Key reuse is pending
+ for _ in 0..3 {
+ register_incoming(&db, &cfg, &gen_in_pay(reference_number.clone()))
+ .await
+ .unwrap();
+ }
+ check_in(
+ &db,
+ &[
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Pending,
+ Pending,
+ ],
+ )
+ .await;
+
+ // Finish pending
+ let third = EddsaPublicKey::rand();
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
+ transfer_register(
+ &db,
+ IncomingType::reserve,
+ &third,
+ &auth_pub,
+ &auth_sig,
+ true,
+ &reference_number,
+ &Timestamp::now()
)
.await
- .expect("bounce"),
- BounceResult {
- tx_id: 1,
- tx_new: false,
- bounce_id: 2,
- bounce_new: false
- }
+ .unwrap(),
+ RegistrationResult::Success
);
-
- // Batch
- assert_eq!(
- db::pending_batch(&mut db, &now).await.unwrap(),
+ check_in(
+ &db,
&[
- 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
- }
- ]
- );
+ Reserve(first.clone()),
+ Simple,
+ Bounced,
+ Reserve(second.clone()),
+ Reserve(third.clone()),
+ Pending,
+ ],
+ )
+ .await;
}
#[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();
- }
+ async fn in_recover_info() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
- #[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(),
+ async fn check_content(db: &PgPool, p: &IncomingPayment) {
+ sqlx::query(
+ "
+ SELECT
+ uetr IS NOT DISTINCT FROM $1 AND
+ tx_id IS NOT DISTINCT FROM $2 AND
+ acct_svcr_ref IS NOT DISTINCT FROM $3 AND
+ subject IS NOT DISTINCT FROM $4 AND
+ debit_payto IS NOT DISTINCT FROM $5
+ FROM incoming_transactions ORDER BY incoming_transaction_id DESC LIMIT 1
+ ",
)
+ .bind(p.id.uetr)
+ .bind(&p.id.tx_id)
+ .bind(&p.id.acct_svcr_ref)
+ .bind(&p.subject)
+ .bind(p.debtor.as_ref().map(|it| it.as_ref().as_str()))
+ .try_map(|r: PgRow| {
+ assert!(r.try_get_flag(0)?);
+ Ok(())
+ })
+ .fetch_one(db)
.await
- .expect("transfer");
+ .unwrap();
}
- 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");
+ // Non talerable
+ for (i, id) in [
+ IncomingId::new(Some(Uuid::new_v4()), None, None),
+ IncomingId::new(None, Some(rand_ebics_id()), None),
+ IncomingId::new(None, None, Some(rand_ebics_id())),
+ ]
+ .iter()
+ .enumerate()
+ {
+ let payment = gen_in_pay("subject".to_owned());
+
+ // Register minimal
+ let partial = IncomingPayment {
+ id: id.clone(),
+ subject: None,
+ debtor: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &partial).await.unwrap();
+ check_content(&db, &partial).await;
+ check_in_count(&db, i + 1, i, 0).await;
+
+ // Recover ID
+ let full_id = IncomingId::new(
+ Some(id.uetr.unwrap_or_else(Uuid::new_v4)),
+ Some(id.tx_id.clone().unwrap_or_else(rand_ebics_id)),
+ Some(id.acct_svcr_ref.clone().unwrap_or_else(rand_ebics_id)),
+ );
+ let full = IncomingPayment {
+ id: full_id.clone(),
+ ..partial.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 1, i, 0).await;
+
+ // Recover subject & debtor
+ let full = IncomingPayment {
+ id: full_id,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 1, i + 1, 0).await;
}
- 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");
+ // Talerable
+ for (i, id) in [
+ IncomingId::new(Some(Uuid::new_v4()), None, None),
+ IncomingId::new(None, Some(rand_ebics_id()), None),
+ IncomingId::new(None, None, Some(rand_ebics_id())),
+ ]
+ .iter()
+ .enumerate()
+ {
+ let key = EddsaPublicKey::rand();
+ let payment = gen_in_pay(format!("test with {key} reserve pub"));
+
+ // Register minimal
+ let partial = IncomingPayment {
+ id: id.clone(),
+ subject: None,
+ debtor: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &partial).await.unwrap();
+ check_content(&db, &partial).await;
+ check_in_count(&db, i + 4, 3, i).await;
+
+ // Recover ID
+ let full_id = IncomingId::new(
+ Some(id.uetr.unwrap_or_else(Uuid::new_v4)),
+ Some(id.tx_id.clone().unwrap_or_else(rand_ebics_id)),
+ Some(id.acct_svcr_ref.clone().unwrap_or_else(rand_ebics_id)),
+ );
+ let full = IncomingPayment {
+ id: full_id.clone(),
+ ..partial.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 4, 3, i).await;
+
+ // Recover subject & debtor
+ let full = IncomingPayment {
+ id: full_id,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &full).await.unwrap();
+ check_content(&db, &full).await;
+ check_in_count(&db, i + 4, 3, i + 1).await;
}
- let pendings = db::pending_batch(&mut db, &start)
+ }
+
+ #[tokio::test]
+ pub async fn in_horror() {
+ let (_, db) = setup().await;
+ let cfg = NexusIngestConfig::simple(AccountType::Exchange, &CURRENCY);
+
+ // Check we do not bounce already registered talerable transaction
+ let key = EddsaPublicKey::rand();
+ let payment = gen_in_pay(format!("test with {key} reserve pub"));
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ assert_eq!(
+ register_in_malformed(
+ &db,
+ &payment,
+ &amount("KUDOS:2.53"),
+ &rand_ebics_id(),
+ &Timestamp::now(),
+ "manual bounce",
+ )
.await
- .expect("pending_batch");
- assert_eq!(pendings.len(), 83);
+ .unwrap(),
+ IncomingBounceRegistrationResult::Talerable
+ );
+ let incomplete = IncomingPayment {
+ subject: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ check_in(&db, &[Reserve(key.clone())]).await;
+
+ // Check we do not register as talerable bounced transaction
+ let new_key = EddsaPublicKey::rand();
+ let payment = gen_in_pay(format!("bounced {new_key}"));
+ let incomplete = IncomingPayment {
+ subject: None,
+ ..payment.clone()
+ };
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ register_incoming(&db, &cfg, &incomplete).await.unwrap();
+ register_incoming(&db, &cfg, &payment).await.unwrap();
+ check_in(&db, &[Reserve(key.clone()), Bounced]).await;
}
-}*/
+}
diff --git a/src/key_management.rs b/src/key_management.rs
@@ -43,7 +43,7 @@ use crate::{
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
-#[allow(non_camel_case_types)]
+#[allow(clippy::upper_case_acronyms)]
pub enum Order {
INI,
HIA,
@@ -153,7 +153,7 @@ pub async fn key_management(
cfg: &EbicsHostCfg,
client: &ClientPriKeysFile,
http: &Client,
- ebics_logger: &EbicsLogger,
+ _ebics_logger: &EbicsLogger,
order: Order,
) -> anyhow::Result<EbicsResponse<Option<String>>> {
info!("Doing key request {}", order.name());
diff --git a/src/lib.rs b/src/lib.rs
@@ -0,0 +1,688 @@
+/*
+* 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},
+ path::Path,
+};
+
+use anyhow::bail;
+use compact_str::{CompactString, ToCompactString};
+use jiff::{Timestamp, civil::Date};
+use rand::prelude::IndexedRandom;
+use reqwest::{
+ Client, StatusCode,
+ header::{CONTENT_TYPE, HeaderValue},
+};
+use sqlx::PgPool;
+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 tracing::{debug, info, warn};
+use uuid::Uuid;
+
+use crate::{
+ common::EbicsLogger,
+ config::{AccountType, EbicsHostCfg, NexusCfg, NexusIngestConfig},
+ db::{
+ InResult, IncomingBounceRegistrationResult, IncomingRegistrationResult,
+ OutgoingRegistrationResult, register_in, register_in_malformed, register_in_qr_bill,
+ register_in_talerable,
+ },
+ ebics_code::EbicsReturnCode,
+ key_management::{Order, hpb, submit_client_keys},
+ keys::{load_bank_keys, load_client_keys, persist_client_keys},
+};
+use crate::{keys::ClientPriKeysFile, xml::XmlReader};
+
+pub mod common;
+pub mod config;
+pub mod crypto;
+pub mod db;
+pub mod ebics_code;
+pub mod key_management;
+pub mod keys;
+pub mod testbench;
+pub mod xml;
+pub mod xml_sign;
+
+pub const CONFIG_SOURCE: ConfigSource =
+ ConfigSource::new("libeufin", "libeufin-nexus", "libeufin-nexus");
+
+const EBICS_ID_ALPHABET: &[u8] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789";
+
+pub fn rand_ebics_id() -> CompactString {
+ let mut rng = rand::rng();
+ (0..34)
+ .map(|_| *EBICS_ID_ALPHABET.choose(&mut rng).unwrap() as char)
+ .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 {
+ #[clap(flatten)]
+ pub common: CommonArgs,
+}
+
+pub async fn register_incoming(
+ db: &PgPool,
+ cfg: &NexusIngestConfig,
+ payment: &IncomingPayment,
+) -> sqlx::Result<()> {
+ let log_res = |res: InResult, kind: &str, suffix: &str| {
+ let fmt = std::fmt::from_fn(|f| {
+ write!(f, "{payment}")?;
+ if kind.is_empty() {
+ write!(f, " {kind}")?;
+ }
+ if res.new {
+ if let Some(id) = &res.bounce_id {
+ write!(f, " bounced in {id}")?;
+ }
+ } else {
+ if res.completed {
+ f.write_str(" completed")?;
+ if let Some(id) = &res.bounce_id {
+ write!(f, " bounced in {id}")?;
+ }
+ } else {
+ if let Some(id) = &res.bounce_id {
+ write!(f, " already bounced in {id}")?;
+ }
+ }
+ }
+ if suffix.is_empty() {
+ write!(f, " {suffix}")?;
+ }
+ Ok(())
+ });
+
+ if res.completed || res.new {
+ info!("{fmt}")
+ } else {
+ debug!("{fmt}")
+ }
+ };
+ let bounce = |cause: String| async move {
+ match cfg.account_type {
+ AccountType::Exchange => {
+ if payment.execution_time < cfg.ignore_bounces_before {
+ let res = register_in(db, payment).await?;
+ log_res(res, "", &format!("ignored bounce: {cause}"));
+ } else {
+ let mut bounce_amount = payment.amount.clone();
+ if let Some(credit_fee) = &payment.credit_fee
+ && cfg.bounce_deduce_fee
+ {
+ if let Some(res) = bounce_amount.try_sub(credit_fee) {
+ bounce_amount = res
+ } else {
+ let res = register_in(db, payment).await?;
+ log_res(
+ res,
+ "",
+ &format!("skip bounce (transfer fee higher than amount): {cause}"),
+ );
+ return Ok(());
+ }
+ }
+ if let Some(res) = bounce_amount.try_sub(&cfg.bounce_fee) {
+ bounce_amount = res
+ } else {
+ let res = register_in(db, payment).await?;
+ log_res(
+ res,
+ "",
+ &format!("skip bounce (bounce fee higher than amount): {cause}"),
+ );
+ return Ok(());
+ }
+ let res = register_in_malformed(
+ db,
+ payment,
+ &bounce_amount,
+ &rand_ebics_id(),
+ &Timestamp::now(),
+ &cause,
+ )
+ .await?;
+ match res {
+ IncomingBounceRegistrationResult::Talerable => {
+ warn!("{payment} tried to bounce a talerable transaction");
+ }
+ IncomingBounceRegistrationResult::Success(res) => {
+ log_res(res, "", &format!(": {cause}"));
+ }
+ }
+ }
+ }
+ AccountType::Normal => {
+ let res = register_in(db, payment).await?;
+ log_res(res, "", "");
+ }
+ }
+ sqlx::Result::<_, sqlx::Error>::Ok(())
+ };
+
+ // Check we have enough info to handle this transaction
+ if payment.debtor.is_none() {
+ // TODO payment.debtor.receiverName == null
+ let res = register_in(db, payment).await?;
+ log_res(res, "incomplete", "");
+ return Ok(());
+ }
+ // TODO if payment.debtor.is_none() && payment.debtor.map(|it| it.rec)
+ if let Some(regex) = &cfg.restriction_payto_regex
+ && let Some(debtor) = &payment.debtor
+ && !regex.is_match(debtor.as_ref().as_str())
+ {
+ bounce("restricted account".to_owned()).await?;
+ return Ok(());
+ }
+
+ if let Some(subject) = &payment.subject
+ && 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::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(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::Success(res) => {
+ log_res(res, "", "");
+ }
+ },
+ Err(e) => {
+ bounce(e.to_string()).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
+ let current = load_client_keys(path)?;
+ if let Some(current) = current {
+ return Ok(current);
+ }
+ // Else create new keys
+ let new = ClientPriKeysFile::generate()?;
+ persist_client_keys(&new, path)?;
+ info!(
+ "New client private keys created at '{}'",
+ path.to_string_lossy()
+ );
+ Ok(new)
+}
+
+pub async fn ebics_setup(
+ http: &Client,
+ cfg: &NexusCfg,
+ force_keys_submissions: bool,
+) -> anyhow::Result<()> {
+ let logger = EbicsLogger {};
+ let keys_cfg = cfg.keys()?;
+ let host_cfg = cfg.host()?;
+
+ let mut client = load_or_generate_client_keys(keys_cfg.client_priv_keys_path.as_ref())?;
+ let bank = load_bank_keys(keys_cfg.bank_pub_keys_path.as_ref())?;
+
+ // Check EBICS 3 support
+ let versions = hev(http, cfg.host()?).await?;
+ debug!(target: "setup",
+ "HEV: {}",
+ versions
+ .iter()
+ .map(|v| v.to_string())
+ .collect::<Vec<_>>()
+ .join(", ")
+ );
+ if !versions.contains(&VersionNumber {
+ number: "03.00".to_owned(),
+ schema: "H005".to_owned(),
+ }) && versions.contains(&VersionNumber {
+ number: "03.02".to_owned(),
+ schema: "H005".to_owned(),
+ }) {
+ bail!("EBICS 3 is not supported by your bank");
+ }
+
+ // Privs exist. Upload their pubs
+ let keys_not_sub = !client.submitted_ini;
+ if !client.submitted_ini || force_keys_submissions {
+ submit_client_keys(keys_cfg, host_cfg, &mut client, http, &logger, Order::INI).await?;
+ }
+ // Eject PDF if the keys were submitted for the first time, or the user asked.
+ // TODO if (keysNotSub || generateRegistrationPdf) makePdf(clientKeys, hostCfg)
+ if !client.submitted_hia || force_keys_submissions {
+ submit_client_keys(keys_cfg, host_cfg, &mut client, http, &logger, Order::HIA).await?;
+ }
+
+ let res = hpb(http, host_cfg, &logger, &client).await?;
+ dbg!(res);
+ // Fetch bank keys
+ Ok(())
+}
+
+pub async fn hev(http: &Client, cfg: &EbicsHostCfg) -> anyhow::Result<Vec<VersionNumber>> {
+ let phase = "HEV";
+ info!(target: "ebics", "Doing administrative request {phase}");
+ let msg = xml_build!(
+ "ebicsHEVRequest" ("xmlns": "http://www.ebics.org/H000") {
+ "HostID": &cfg.host_id
+ }
+ );
+ let res = post_to_bank(cfg.base_url.as_str(), http, msg).await?;
+ XmlReader::parse(&res, "ebicsHEVResponse", |root| {
+ let technical_code = root
+ .one("SystemReturnCode")
+ .one("ReturnCode")
+ .text()
+ .parse()
+ .unwrap();
+ let versions: Vec<_> = root
+ .each("VersionNumber")
+ .map(|n| VersionNumber {
+ number: n.text().parse().unwrap(),
+ schema: n.attr("ProtocolVersion").to_owned(),
+ })
+ .collect();
+ EbicsResponse {
+ technical_code,
+ bank_code: EbicsReturnCode::EBICS_OK,
+ content: versions,
+ }
+ })
+ .ok_or_fail(phase)
+}
+
+#[derive(Debug, thiserror::Error)]
+pub enum EbicsError {
+ #[error(transparent)]
+ Network(#[from] reqwest::Error),
+}
+
+async fn post_to_bank(url: &str, client: &Client, msg: String) -> anyhow::Result<String> {
+ let res = client
+ .post(url)
+ .header(CONTENT_TYPE, HeaderValue::from_static("application/xml"))
+ .body(msg)
+ .send()
+ .await?;
+ let status = res.status();
+ if status != StatusCode::OK {
+ bail!("bank http error {status}");
+ }
+ // Should parse xml here
+ let body = res.text().await?;
+ Ok(body)
+}
+
+pub struct EbicsResponse<T> {
+ pub technical_code: EbicsReturnCode,
+ pub bank_code: EbicsReturnCode,
+ pub content: T,
+}
+
+impl<T> EbicsResponse<T> {
+ fn ok_or_fail(self, phase: &str) -> anyhow::Result<T> {
+ if self.technical_code.is_error() {
+ bail!("{phase} has technical error: {:?}", self.technical_code)
+ } else if self.bank_code.is_error() {
+ bail!("{phase} has bank error: {:?}", self.bank_code)
+ } else {
+ Ok(self.content)
+ }
+ }
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub struct VersionNumber {
+ pub number: String,
+ pub schema: String,
+}
+
+impl Display for VersionNumber {
+ fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ let Self { number, schema } = self;
+ write!(f, "{number}:{schema}")
+ }
+}
diff --git a/src/main.rs b/src/main.rs
@@ -17,300 +17,10 @@
* <http://www.gnu.org/licenses/>
*/
-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,
- 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},
-};
-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 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)]
-struct Args {
- #[clap(flatten)]
- common: CommonArgs,
-}
+use clap::Parser as _;
+use libeufin::{Args, CONFIG_SOURCE, config::NexusCfg, ebics_setup};
+use reqwest::Client;
+use taler_common::taler_main;
fn main() {
let args = Args::parse();
@@ -320,188 +30,3 @@ fn main() {
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
- let current = load_client_keys(path)?;
- if let Some(current) = current {
- return Ok(current);
- }
- // Else create new keys
- let new = ClientPriKeysFile::generate()?;
- persist_client_keys(&new, path)?;
- info!(
- "New client private keys created at '{}'",
- path.to_string_lossy()
- );
- Ok(new)
-}
-
-pub async fn ebics_setup(
- http: &Client,
- cfg: &NexusCfg,
- force_keys_submissions: bool,
-) -> anyhow::Result<()> {
- let logger = EbicsLogger {};
- let keys_cfg = cfg.keys()?;
- let host_cfg = cfg.host()?;
-
- let mut client = load_or_generate_client_keys(keys_cfg.client_priv_keys_path.as_ref())?;
- let bank = load_bank_keys(keys_cfg.bank_pub_keys_path.as_ref())?;
-
- // Check EBICS 3 support
- let versions = hev(http, cfg.host()?).await?;
- debug!(target: "setup",
- "HEV: {}",
- versions
- .iter()
- .map(|v| v.to_string())
- .collect::<Vec<_>>()
- .join(", ")
- );
- if !versions.contains(&VersionNumber {
- number: "03.00".to_owned(),
- schema: "H005".to_owned(),
- }) && versions.contains(&VersionNumber {
- number: "03.02".to_owned(),
- schema: "H005".to_owned(),
- }) {
- bail!("EBICS 3 is not supported by your bank");
- }
-
- // Privs exist. Upload their pubs
- let keys_not_sub = !client.submitted_ini;
- if !client.submitted_ini || force_keys_submissions {
- submit_client_keys(keys_cfg, host_cfg, &mut client, http, &logger, Order::INI).await?;
- }
- // Eject PDF if the keys were submitted for the first time, or the user asked.
- // TODO if (keysNotSub || generateRegistrationPdf) makePdf(clientKeys, hostCfg)
- if !client.submitted_hia || force_keys_submissions {
- submit_client_keys(keys_cfg, host_cfg, &mut client, http, &logger, Order::HIA).await?;
- }
-
- let res = hpb(http, host_cfg, &logger, &client).await?;
- dbg!(res);
- // Fetch bank keys
- Ok(())
-}
-
-pub async fn hev(http: &Client, cfg: &EbicsHostCfg) -> anyhow::Result<Vec<VersionNumber>> {
- let phase = "HEV";
- info!(target: "ebics", "Doing administrative request {phase}");
- let msg = xml_build!(
- "ebicsHEVRequest" ("xmlns": "http://www.ebics.org/H000") {
- "HostID": &cfg.host_id
- }
- );
- let res = post_to_bank(cfg.base_url.as_str(), http, msg).await?;
- XmlReader::parse(&res, "ebicsHEVResponse", |root| {
- let technical_code = root
- .one("SystemReturnCode")
- .one("ReturnCode")
- .text()
- .parse()
- .unwrap();
- let versions: Vec<_> = root
- .each("VersionNumber")
- .map(|n| VersionNumber {
- number: n.text().parse().unwrap(),
- schema: n.attr("ProtocolVersion").to_owned(),
- })
- .collect();
- EbicsResponse {
- technical_code,
- bank_code: EbicsReturnCode::EBICS_OK,
- content: versions,
- }
- })
- .ok_or_fail(phase)
-}
-
-#[derive(Debug, thiserror::Error)]
-pub enum EbicsError {
- #[error(transparent)]
- Network(#[from] reqwest::Error),
-}
-
-async fn post_to_bank(url: &str, client: &Client, msg: String) -> anyhow::Result<String> {
- let res = client
- .post(url)
- .header(CONTENT_TYPE, HeaderValue::from_static("application/xml"))
- .body(msg)
- .send()
- .await?;
- let status = res.status();
- if status != StatusCode::OK {
- bail!("bank http error {status}");
- }
- // Should parse xml here
- let body = res.text().await?;
- Ok(body)
-}
-
-pub struct EbicsResponse<T> {
- pub technical_code: EbicsReturnCode,
- pub bank_code: EbicsReturnCode,
- pub content: T,
-}
-
-impl<T> EbicsResponse<T> {
- fn ok_or_fail(self, phase: &str) -> anyhow::Result<T> {
- if self.technical_code.is_error() {
- bail!("{phase} has technical error: {:?}", self.technical_code)
- } else if self.bank_code.is_error() {
- bail!("{phase} has bank error: {:?}", self.bank_code)
- } else {
- Ok(self.content)
- }
- }
-}
-
-#[derive(Debug, Clone, PartialEq, Eq)]
-pub struct VersionNumber {
- pub number: String,
- pub schema: String,
-}
-
-impl Display for VersionNumber {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
- let Self { number, schema } = self;
- write!(f, "{number}:{schema}")
- }
-}
diff --git a/src/xml.rs b/src/xml.rs
@@ -270,7 +270,7 @@ mod test {
xml_el!(w, "module");
}
assert_eq!(
- xml_build!("root" { @ |w| module(w) }),
+ xml_build!("root" { @ module }),
r#"<?xml version="1.0" encoding="UTF-8" standalone="yes"?><root><module/></root>"#
)
}